#![expect(clippy::unwrap_used, clippy::panic)]
use std::collections::BTreeMap;
use std::io::{Read, Write};
use std::net::TcpListener;
use std::net::TcpStream;
use std::path::Path;
use std::sync::atomic::AtomicBool;
use std::sync::atomic::AtomicUsize;
use std::sync::atomic::Ordering;
use std::sync::Arc;
use std::sync::Mutex;
use std::thread;
use std::time::Duration;
use alloy::primitives::{Address, B256};
use degenbot_cli_core::pool::{reset_verify_provider_build_count, verify_provider_build_count};
use degenbot_cli_core::{run, CliContext, Command, ExitCode, PoolCommand, PoolFamily, Prompter};
use degenbot_config::MapEnv;
use degenbot_db::{DegenbotDb, V3PoolRowInput, V4PoolRowInput};
use tempfile::TempDir;
const CHAIN: i64 = 8453;
static BUILD_COUNTER: Mutex<()> = Mutex::new(());
struct NoPrompt;
impl Prompter for NoPrompt {
fn confirm(&self, _message: &str, _default: bool) -> bool {
false
}
}
struct StubNode {
url: String,
chain_id_reads: Arc<AtomicUsize>,
shutdown: Arc<AtomicBool>,
}
impl StubNode {
fn spawn(chain_id: u64) -> Self {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let url = format!("http://{}", listener.local_addr().unwrap());
let chain_id_reads = Arc::new(AtomicUsize::new(0));
let shutdown = Arc::new(AtomicBool::new(false));
let reads = Arc::clone(&chain_id_reads);
let shutdown_flag = Arc::clone(&shutdown);
listener.set_nonblocking(true).unwrap();
thread::spawn(move || {
while !shutdown_flag.load(Ordering::Acquire) {
match listener.accept() {
Ok((mut stream, _addr)) => {
let _ = answer(&mut stream, chain_id, &reads);
}
Err(ref e) if e.kind() == std::io::ErrorKind::WouldBlock => {
thread::sleep(Duration::from_millis(5));
}
Err(_) => break,
}
}
});
Self {
url,
chain_id_reads,
shutdown,
}
}
fn chain_id_reads(&self) -> usize {
self.chain_id_reads.load(Ordering::Acquire)
}
}
impl Drop for StubNode {
fn drop(&mut self) {
self.shutdown.store(true, Ordering::Release);
}
}
fn answer(
stream: &mut TcpStream,
chain_id: u64,
chain_id_reads: &AtomicUsize,
) -> std::io::Result<()> {
stream.set_read_timeout(Some(Duration::from_secs(10)))?;
let mut buf = Vec::new();
let mut chunk = [0u8; 4096];
let header_end = loop {
let n = stream.read(&mut chunk)?;
if n == 0 {
return Ok(());
}
buf.extend_from_slice(&chunk[..n]);
if let Some(pos) = buf.windows(4).position(|w| w == b"\r\n\r\n") {
break pos + 4;
}
};
let headers = String::from_utf8_lossy(&buf[..header_end]);
let content_length = headers
.lines()
.find_map(|line| {
let (name, value) = line.split_once(':')?;
name.trim()
.eq_ignore_ascii_case("content-length")
.then(|| value.trim().parse::<usize>().ok())?
})
.unwrap_or(0);
while buf.len() < header_end + content_length {
let n = stream.read(&mut chunk)?;
if n == 0 {
break;
}
buf.extend_from_slice(&chunk[..n]);
}
let body = String::from_utf8_lossy(&buf[header_end..]);
let request: serde_json::Value =
serde_json::from_str(body.trim()).unwrap_or(serde_json::Value::Null);
let id = request["id"].clone();
let method = request["method"].as_str().unwrap_or("").to_string();
let payload = if method == "eth_chainId" {
chain_id_reads.fetch_add(1, Ordering::AcqRel);
serde_json::json!(format!("0x{chain_id:x}"))
} else {
let error = serde_json::json!({"code": -32601, "message": "method not found"});
let response = serde_json::json!({
"jsonrpc": "2.0",
"id": id,
"error": error,
});
return write_response(stream, &response);
};
let response = serde_json::json!({
"jsonrpc": "2.0",
"id": id,
"result": payload,
});
write_response(stream, &response)
}
fn write_response(stream: &mut TcpStream, response: &serde_json::Value) -> std::io::Result<()> {
let body = response.to_string();
let head = format!(
"HTTP/1.1 200 OK\r\ncontent-type: application/json\r\ncontent-length: {}\r\nconnection: close\r\n\r\n",
body.len(),
);
stream.write_all(head.as_bytes())?;
stream.write_all(body.as_bytes())?;
stream.flush()
}
fn env(rpc_url: &str) -> MapEnv {
let mut map = BTreeMap::new();
map.insert("DEGENBOT_DEFAULT_CHAIN_ID".to_string(), CHAIN.to_string());
map.insert(
"DEGENBOT_RPC_HTTP_CHAINID_8453".to_string(),
rpc_url.to_string(),
);
MapEnv::new(map)
}
fn ctx_for(rpc_url: &str, db: &Path) -> CliContext<'static> {
let e = Box::leak(Box::new(env(rpc_url)));
CliContext::new(e).with_database(db.display().to_string())
}
fn seed_v3_pool(db_path: &Path) -> Address {
let (db, _state) = DegenbotDb::open_for_writes(db_path).unwrap();
let factory = Address::repeat_byte(0xf1);
let exchange = db
.upsert_exchange(CHAIN, "uniswap_v3", factory, None)
.unwrap();
let pool = Address::repeat_byte(0x11);
db.upsert_v3_pools(
CHAIN,
"uniswap_v3",
exchange.id,
1_000_000,
&[V3PoolRowInput {
address: pool,
token0_address: Address::repeat_byte(0x22),
token1_address: Address::repeat_byte(0x33),
fee: 500,
tick_spacing: 10,
}],
)
.unwrap();
pool
}
fn seed_v4_pool(db_path: &Path) -> (String, String) {
let (db, _state) = DegenbotDb::open_for_writes(db_path).unwrap();
let factory = Address::repeat_byte(0xf1);
let exchange = db
.upsert_exchange(CHAIN, "uniswap_v4", factory, None)
.unwrap();
let manager = Address::repeat_byte(0x44);
db.upsert_pool_manager(manager, CHAIN, "uniswap_v4", None, exchange.id)
.unwrap();
let pool_hash = format!("0x{}", B256::from([0xAB; 32]));
db.upsert_v4_pools(
CHAIN,
&manager.to_checksum(None),
1_000_000,
&[V4PoolRowInput {
pool_hash: pool_hash.clone(),
hooks: Address::ZERO,
currency0_address: Address::repeat_byte(0x22),
currency1_address: Address::repeat_byte(0x33),
fee: 0,
tick_spacing: 60,
}],
)
.unwrap();
(pool_hash, manager.to_checksum(None))
}
fn verify_command(
rpc_url: &str,
pool: String,
family: PoolFamily,
pool_manager: Option<String>,
) -> Command {
Command::Pool(PoolCommand::Verify {
rpc_url: rpc_url.to_string(),
chain_id: CHAIN,
block_number: 42,
pool,
family,
pool_manager,
})
}
#[test]
fn pool_verify_builds_one_provider_and_runs_the_gate_green() {
let _guard = BUILD_COUNTER.lock().unwrap();
let _ = reset_verify_provider_build_count();
let stub = StubNode::spawn(CHAIN as u64);
let dir = TempDir::new().unwrap();
let db = dir.path().join("degenbot.db");
let pool = seed_v3_pool(&db);
let ctx = ctx_for(&stub.url, &db);
let outcome = run(
&verify_command(&stub.url, pool.to_checksum(None), PoolFamily::V3, None),
&ctx,
&NoPrompt,
);
assert_eq!(outcome.exit_code, ExitCode::Success, "{outcome:?}");
let Some(degenbot_cli_core::CommandReport::Pool(degenbot_cli_core::PoolReport::Verified {
divergences,
..
})) = outcome.report()
else {
panic!("expected Verified, got {:?}", outcome.report());
};
assert!(divergences.is_empty());
assert_eq!(
verify_provider_build_count(),
1,
"one provider build per verify run"
);
assert_eq!(
stub.chain_id_reads(),
1,
"the chain bind reads the endpoint chain exactly once"
);
}
#[test]
fn pool_verify_refuses_an_endpoint_serving_a_different_chain_v3() {
let _guard = BUILD_COUNTER.lock().unwrap();
let _ = reset_verify_provider_build_count();
let stub = StubNode::spawn(999);
let dir = TempDir::new().unwrap();
let db = dir.path().join("degenbot.db");
let pool = seed_v3_pool(&db);
let ctx = ctx_for(&stub.url, &db);
let outcome = run(
&verify_command(&stub.url, pool.to_checksum(None), PoolFamily::V3, None),
&ctx,
&NoPrompt,
);
assert_eq!(outcome.exit_code, ExitCode::Failure);
let message = outcome.error().unwrap().message();
assert!(
message.contains("reports chain id 999, not the expected 8453"),
"chain-binding refusal missing: {message}"
);
assert_eq!(
verify_provider_build_count(),
1,
"the refusal still rides the arm's single provider build"
);
assert_eq!(stub.chain_id_reads(), 1);
}
#[test]
fn pool_verify_refuses_an_endpoint_serving_a_different_chain_v4() {
let _guard = BUILD_COUNTER.lock().unwrap();
let _ = reset_verify_provider_build_count();
let stub = StubNode::spawn(999);
let dir = TempDir::new().unwrap();
let db = dir.path().join("degenbot.db");
let (pool_hash, manager) = seed_v4_pool(&db);
let ctx = ctx_for(&stub.url, &db);
let outcome = run(
&verify_command(&stub.url, pool_hash, PoolFamily::V4, Some(manager)),
&ctx,
&NoPrompt,
);
assert_eq!(outcome.exit_code, ExitCode::Failure);
let message = outcome.error().unwrap().message();
assert!(
message.contains("reports chain id 999, not the expected 8453"),
"chain-binding refusal missing: {message}"
);
assert_eq!(verify_provider_build_count(), 1);
}
#[test]
fn pool_verify_without_a_verified_pool_never_builds_a_provider() {
let _guard = BUILD_COUNTER.lock().unwrap();
let _ = reset_verify_provider_build_count();
let stub = StubNode::spawn(CHAIN as u64);
let dir = TempDir::new().unwrap();
let db = dir.path().join("degenbot.db");
degenbot_db::ops::create_new_database(&db).unwrap();
let ctx = ctx_for(&stub.url, &db);
let outcome = run(
&verify_command(
&stub.url,
Address::repeat_byte(0x77).to_checksum(None),
PoolFamily::V3,
None,
),
&ctx,
&NoPrompt,
);
assert_eq!(outcome.exit_code, ExitCode::Failure);
assert!(
matches!(
outcome.error(),
Some(degenbot_cli_core::CliError::InvalidArgument(_))
),
"expected the unknown-pool refusal, got {:?}",
outcome.error()
);
assert_eq!(
verify_provider_build_count(),
0,
"a run that verifies nothing never builds a provider"
);
assert_eq!(stub.chain_id_reads(), 0, "no dial without work");
}