From c661f760dc6246dcb212ef93f6e72cb820813aa7 Mon Sep 17 00:00:00 2001 From: mfw78 Date: Tue, 8 Sep 2026 12:34:27 +0000 Subject: [PATCH 1/7] test(cow): drive the engine against a forked mainnet Boots the shipped `shepherd` binary against an anvil fork of a real mainnet node, registers a conditional order on the deployed fork registry, and asserts the module indexes it and polls it to `Post` through the structured generator wire. Drives the binary rather than `BootScenario`: that harness wires a `FakeNode` and hardcodes its chain config, so it has no seam for a real endpoint, and only the binary reads `[chains.1] rpc_url`. Adding one would be a nexum-runtime change plus the pin cascade, to serve a test. Nothing is deployed. The registry, `OwnedTWAP` and `CowAccount7702` are all live on mainnet, so an anvil fork already has them. The owner is an EOA delegated to `CowAccount7702`, which is what #658 names for the acceptance smoke. It is load-bearing rather than incidental: the registry builds an ERC-1271 signature only for a POST, and that path asks the owner for `supportsInterface`. A bare EOA answers with empty returndata and the whole poll reverts with no payload; the delegate answers `FnSelectorNotRecognized`, a real revert, which `_buildSignature` catches as a non-Safe wallet. Skipped unless `SHEPHERD_FORK_RPC` names an endpoint to fork, since CI reaches no such node, and skipped when `anvil` or `cast` is absent. AI Assistance: Claude Code used for the harness and the environment survey. --- .../shepherd-engine/tests/anvil_fork_e2e.rs | 420 ++++++++++++++++++ 1 file changed, 420 insertions(+) create mode 100644 crates/shepherd-engine/tests/anvil_fork_e2e.rs diff --git a/crates/shepherd-engine/tests/anvil_fork_e2e.rs b/crates/shepherd-engine/tests/anvil_fork_e2e.rs new file mode 100644 index 00000000..6aa6e4bb --- /dev/null +++ b/crates/shepherd-engine/tests/anvil_fork_e2e.rs @@ -0,0 +1,420 @@ +//! End-to-end over a forked mainnet: the shipped `shepherd` binary +//! indexes a real `ConditionalOrderCreated` from the deployed fork +//! registry and polls it through the structured generator wire. +//! +//! This drives the binary rather than `BootScenario`, because that +//! harness wires a `FakeNode` and has no seam for a real endpoint. Only +//! the binary reads `[chains.1] rpc_url`, so only the binary can be +//! pointed at anvil. +//! +//! Skipped unless `SHEPHERD_FORK_RPC` names a mainnet endpoint to fork, +//! since CI reaches no such node. Needs `anvil` and `cast` on `PATH` and +//! `just build-modules` already run. + +use std::io::{BufRead, BufReader}; +use std::net::TcpListener; +use std::path::PathBuf; +use std::process::{Child, Command, Stdio}; +use std::sync::mpsc::{Receiver, RecvTimeoutError, channel}; +use std::time::Duration; + +/// The fork's mainnet registry, as `modules/ccow-monitor/component.toml` +/// pins it. Deployed at block 25674440 and never used, so every log the +/// run sees is one it created. +const REGISTRY: &str = "0xf9ba6F64c9b41Df1cEe76A50e2039D3847064232"; + +/// `OwnedTWAP`, deployed by the same CREATE2 broadcast as the registry. +/// `Owned` gates only `setDescriptor` and `setModule`, so registering +/// against it needs no owner. +const HANDLER: &str = "0x4e17a65d14e7f37d2a9f0389f17efd41aaa64c91"; + +/// `CowAccount7702`, the delegate #658 names for the acceptance smoke. +/// +/// The registry builds an ERC-1271 signature only for a POST, and that +/// path asks the owner for `supportsInterface`. A bare EOA answers with +/// empty returndata and the whole poll reverts with no payload; this +/// account answers `FnSelectorNotRecognized`, a real revert, which +/// `_buildSignature` catches and treats as a non-Safe wallet. +const ACCOUNT_7702: &str = "0x15236F06922A288B68e57Ab42e397920a3F3Bb99"; + +const WETH: &str = "0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2"; +const USDC: &str = "0xA0b86991c6218b36c1d19D4a2e9Eb0cE3606eB48"; + +/// How long to wait on any one engine log line. +const LOG_TIMEOUT: Duration = Duration::from_secs(90); + +/// A child killed when the test ends, however it ends. +struct Reaped(Child); + +impl Drop for Reaped { + fn drop(&mut self) { + let _ = self.0.kill(); + let _ = self.0.wait(); + } +} + +/// A port free at bind time. Racy by nature, which is why each spawn +/// waits for its own readiness rather than assuming the port took. +fn free_port() -> u16 { + TcpListener::bind("127.0.0.1:0") + .expect("bind an ephemeral port") + .local_addr() + .expect("the listener has an address") + .port() +} + +/// The orderbook mock, beside the binary under test. +/// +/// `CARGO_BIN_EXE_*` covers only this package's binaries, and the mock +/// belongs to `tools/orderbook-mock`, so it is found as a sibling of the +/// engine rather than named directly. +fn orderbook_mock_bin() -> PathBuf { + let engine = PathBuf::from(env!("CARGO_BIN_EXE_shepherd")); + let mock = engine + .parent() + .expect("the engine binary sits in a target directory") + .join("orderbook-mock"); + assert!( + mock.exists(), + "{} is missing; run `cargo build -p orderbook-mock` first", + mock.display(), + ); + mock +} + +fn workspace_root() -> PathBuf { + PathBuf::from(env!("CARGO_MANIFEST_DIR")) + .parent() + .and_then(std::path::Path::parent) + .expect("crates/ sits two levels under the workspace root") + .to_owned() +} + +/// The endpoint to fork, or `None` when this run cannot reach one. +fn fork_rpc_or_skip() -> Option { + let rpc = std::env::var("SHEPHERD_FORK_RPC") + .ok() + .filter(|s| !s.is_empty())?; + for tool in ["anvil", "cast"] { + if Command::new(tool).arg("--version").output().is_err() { + eprintln!("skipping: {tool} is not on PATH"); + return None; + } + } + Some(rpc) +} + +fn cast(args: &[&str]) -> String { + let out = Command::new("cast") + .args(args) + .output() + .unwrap_or_else(|e| panic!("cast {args:?}: {e}")); + assert!( + out.status.success(), + "cast {args:?} failed: {}", + String::from_utf8_lossy(&out.stderr), + ); + String::from_utf8_lossy(&out.stdout).trim().to_owned() +} + +/// Fork `rpc` at head and wait until the fork answers. +fn start_anvil(rpc: &str, port: u16) -> Reaped { + // Owned before the readiness loop, so the panic below still reaps it. + let child = Reaped( + Command::new("anvil") + .args(["--fork-url", rpc, "--port", &port.to_string(), "--silent"]) + .stdout(Stdio::null()) + .stderr(Stdio::null()) + .spawn() + .expect("spawn anvil"), + ); + let url = format!("http://127.0.0.1:{port}"); + for _ in 0..60 { + if Command::new("cast") + .args(["chain-id", "--rpc-url", &url]) + .output() + .is_ok_and(|o| o.status.success()) + { + return child; + } + std::thread::sleep(Duration::from_millis(500)); + } + panic!("anvil did not answer on {url}"); +} + +/// Register one TWAP conditional order against the forked registry. +/// +/// `t0` is derived from the fork's own clock and backdated, so part 0 is +/// tradeable at once. A pinned `t0` puts the part index far past `n` and +/// every poll then reports the series finished. +fn register_order(anvil: &str, owner: &str) -> String { + let now: u64 = cast(&[ + "block", + "latest", + "--field", + "timestamp", + "--rpc-url", + anvil, + ]) + .parse() + .expect("a numeric block timestamp"); + let t0 = now - 60; + let static_input = cast(&[ + "abi-encode", + "f((address,address,address,uint256,uint256,uint256,uint256,uint256,uint256,bytes32))", + &format!( + "({WETH},{USDC},{owner},1000000000000000,1000000,{t0},2,600,0,\ + 0x0000000000000000000000000000000000000000000000000000000000000000)" + ), + ]); + let calldata = cast(&[ + "calldata", + "create((address,bytes32,bytes),bool)", + &format!( + "({HANDLER},0x0000000000000000000000000000000000000000000000000000000000000001,{static_input})" + ), + "true", + ]); + cast(&[ + "send", + "--unlocked", + "--from", + owner, + REGISTRY, + &calldata, + "--rpc-url", + anvil, + "--json", + ]); + static_input +} + +/// An engine config pointing at this run's anvil and orderbook mock. +fn engine_config(dir: &std::path::Path, anvil: &str, orderbook: &str, manifest: &str) -> PathBuf { + let root = workspace_root(); + let wasm = root.join("target/wasm32-wasip2/release/ccow_monitor.wasm"); + assert!( + wasm.exists(), + "{} is missing; run `just build-modules` first", + wasm.display(), + ); + let state = dir.join("state"); + let config = format!( + r#" +[engine] +state_dir = "{state}" +log_level = "info" +# The wasm is rebuilt for every run, so a pin would have to be +# regenerated each time; the shipped dev configs make the same choice. +require_component_digest = false + +[chains.1] +rpc_url = "{anvil}" + +[[modules]] +id = "ccow-monitor" +path = "{wasm}" +manifest = "{manifest}" + +[extensions.videre.venues.cow] +chain = 1 +orderbook_url = "{orderbook}" +owner = "0x0000000000000000000000000000000000000001" +"#, + wasm = wasm.display(), + state = state.display(), + ); + let path = dir.join("engine.anvil.toml"); + std::fs::write(&path, config).expect("write the engine config"); + path +} + +/// The shipped manifest with `start_block` moved to the fork block. +/// +/// The shipped value is the registry's deployment block, roughly 258000 +/// behind head, and the registry holds no logs before the fork anyway. +fn manifest_at(dir: &std::path::Path, from_block: u64) -> PathBuf { + let shipped = workspace_root().join("modules/ccow-monitor/component.toml"); + let text = std::fs::read_to_string(&shipped).expect("read the shipped manifest"); + let rewritten = text + .lines() + .map(|line| { + if line.trim_start().starts_with("start_block") { + format!("start_block = {from_block}") + } else { + line.to_owned() + } + }) + .collect::>() + .join("\n"); + let path = dir.join("component.toml"); + std::fs::write(&path, rewritten).expect("write the test manifest"); + path +} + +/// Stream a child's stdout and stderr into one channel so the test can +/// wait on lines without blocking on a pipe that never closes. +/// +/// Both, because the engine's tracing subscriber picks its own stream +/// and a test that watched only one would hang for the full timeout with +/// nothing to show. +fn stream_lines(child: &mut Child) -> Receiver { + let (tx, rx) = channel(); + for pipe in [ + child + .stdout + .take() + .map(|p| Box::new(p) as Box), + child + .stderr + .take() + .map(|p| Box::new(p) as Box), + ] + .into_iter() + .flatten() + { + let tx = tx.clone(); + std::thread::spawn(move || { + for line in BufReader::new(pipe).lines().map_while(Result::ok) { + if tx.send(line).is_err() { + return; + } + } + }); + } + rx +} + +/// Wait for a line containing `needle` and return it, echoing what +/// arrived first when it never does. +fn wait_for(rx: &Receiver, needle: &str, seen: &mut Vec) -> String { + let deadline = std::time::Instant::now() + LOG_TIMEOUT; + loop { + let left = deadline.saturating_duration_since(std::time::Instant::now()); + match rx.recv_timeout(left) { + Ok(line) => { + if line.contains(needle) { + seen.push(line.clone()); + return line; + } + seen.push(line); + } + Err(RecvTimeoutError::Timeout) => { + panic!( + "no engine log line contained {needle:?}; saw:\n {}", + seen.join("\n ") + ) + } + Err(RecvTimeoutError::Disconnected) => { + panic!( + "the engine exited before logging {needle:?}; saw:\n {}", + seen.join("\n ") + ) + } + } + } +} + +/// The whole chain leg: a real registration on the forked registry is +/// indexed, then polled through `getTradeableOrderWithSignature` against +/// the deployed `OwnedTWAP`. +#[test] +fn ccow_monitor_indexes_and_polls_a_real_owned_twap() { + let Some(rpc) = fork_rpc_or_skip() else { + return; + }; + let dir = std::env::temp_dir().join(format!("shepherd-anvil-{}", std::process::id())); + std::fs::create_dir_all(&dir).expect("create the run directory"); + + let anvil_port = free_port(); + let _anvil = start_anvil(&rpc, anvil_port); + let anvil = format!("http://127.0.0.1:{anvil_port}"); + + let ob_port = free_port(); + let _ob = Reaped( + Command::new(orderbook_mock_bin()) + .args(["--port", &ob_port.to_string()]) + .stdout(Stdio::null()) + .stderr(Stdio::null()) + .spawn() + .expect("spawn the orderbook mock"), + ); + + let fork_block: u64 = cast(&["block-number", "--rpc-url", &anvil]) + .parse() + .expect("a numeric block number"); + let manifest = manifest_at(&dir, fork_block); + let config = engine_config( + &dir, + &anvil, + &format!("http://127.0.0.1:{ob_port}"), + &manifest.display().to_string(), + ); + + let mut engine = Command::new(env!("CARGO_BIN_EXE_shepherd")) + .args(["--engine-config", &config.display().to_string()]) + .current_dir(workspace_root()) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()) + .spawn() + .expect("spawn the shepherd binary"); + let lines = stream_lines(&mut engine); + let _engine = Reaped(engine); + + let mut seen = Vec::new(); + wait_for(&lines, "ccow-monitor", &mut seen); + + // anvil's first unlocked account owns the registration. + let owner = cast(&["rpc", "eth_accounts", "--rpc-url", &anvil]) + .trim_matches(|c| c == '[' || c == ']' || c == '"') + .split(',') + .next() + .expect("anvil funds at least one account") + .trim_matches('"') + .to_owned(); + // EIP-7702 delegation designator: 0xef0100 ++ implementation. Set + // directly rather than by authorisation tuple, because the test only + // needs the code to be there, not the signing ceremony that puts it + // there on a real chain. + cast(&[ + "rpc", + "anvil_setCode", + &owner, + &format!("0xef0100{}", ACCOUNT_7702.trim_start_matches("0x")), + "--rpc-url", + &anvil, + ]); + register_order(&anvil, &owner); + + wait_for(&lines, "indexed commitment:", &mut seen); + + // The poll fires on a block dispatch, so the registration needs one + // after it. Mining is cheap; the engine's own poller drives the rest. + for _ in 0..6 { + cast(&["rpc", "evm_mine", "--rpc-url", &anvil]); + std::thread::sleep(Duration::from_millis(500)); + } + + // `poll {key} -> {outcome}` is the module's own line, so this is the + // structured generator wire answering over a real eth_call against + // the deployed `OwnedTWAP`, not a fixture. + let polled = wait_for(&lines, "poll commitment:", &mut seen); + for bad in ["did not decode", "eth_call failed"] { + assert!( + !seen.iter().any(|l| l.contains(bad)), + "the poll leg reported {bad:?}:\n {}", + seen.join("\n "), + ); + } + assert!( + polled.contains("-> Post"), + "the generator posted, so the module must too; got: {polled}", + ); + + // The submit leg stops here. `orderbook-mock` answers with a + // synthetic uid rather than one derived from the order, so the + // venue's receipt check refuses it, correctly. Closing that loop + // needs a mock that echoes the real uid. + + let _ = std::fs::remove_dir_all(&dir); +} From 303f3aa0a50e48d6ae43e44dbad3c9db0294fa75 Mon Sep 17 00:00:00 2001 From: mfw78 Date: Tue, 8 Sep 2026 23:29:30 +0000 Subject: [PATCH 2/7] test(cow): derive the orderbook mock's UID from the order it was sent The mock answered with a counter in 56 bytes, ignoring the body. A submitting venue re-derives the UID locally and refuses a receipt that disagrees, so no submit through this mock could ever complete: the fork run reached `Post` and then dropped the commitment on `receipt mismatch`. It now parses the posted `OrderCreation` and returns the UID that order actually has, via `order_data()` and the chain's settlement domain, so the derivation is the one the venue checks with rather than a second implementation of it. A body it cannot read is refused rather than answered with something the venue would reject anyway. `--chain-id` selects the settlement domain, because a UID from the wrong domain is refused exactly like a synthetic one. This closes the venue leg of the fork run, which now asserts the submit lands. AI Assistance: Claude Code used for the mock and the tests. --- Cargo.lock | 1 + .../shepherd-engine/tests/anvil_fork_e2e.rs | 21 ++- tools/orderbook-mock/Cargo.toml | 3 + tools/orderbook-mock/src/main.rs | 121 ++++++++++++++---- 4 files changed, 115 insertions(+), 31 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index db64625b..fcd7bed1 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3993,6 +3993,7 @@ dependencies = [ "anyhow", "axum", "clap", + "cowprotocol", "rand 0.10.2", "reqwest", "serde", diff --git a/crates/shepherd-engine/tests/anvil_fork_e2e.rs b/crates/shepherd-engine/tests/anvil_fork_e2e.rs index 6aa6e4bb..a6decb84 100644 --- a/crates/shepherd-engine/tests/anvil_fork_e2e.rs +++ b/crates/shepherd-engine/tests/anvil_fork_e2e.rs @@ -79,6 +79,20 @@ fn orderbook_mock_bin() -> PathBuf { "{} is missing; run `cargo build -p orderbook-mock` first", mock.display(), ); + // Cargo does not rebuild another package's binary for this test, and + // a stale mock fails as a receipt mismatch rather than as a stale + // build, which is a long way from the cause. + let source = workspace_root().join("tools/orderbook-mock/src/main.rs"); + if let (Ok(built), Ok(src)) = ( + mock.metadata().and_then(|m| m.modified()), + source.metadata().and_then(|m| m.modified()), + ) { + assert!( + built >= src, + "{} is older than its source; run `cargo build -p orderbook-mock`", + mock.display(), + ); + } mock } @@ -411,10 +425,9 @@ fn ccow_monitor_indexes_and_polls_a_real_owned_twap() { "the generator posted, so the module must too; got: {polled}", ); - // The submit leg stops here. `orderbook-mock` answers with a - // synthetic uid rather than one derived from the order, so the - // venue's receipt check refuses it, correctly. Closing that loop - // needs a mock that echoes the real uid. + // And the post reaches the venue, closing the loop through the + // borsh body, the journal reservation and the orderbook adapter. + wait_for(&lines, "submitted", &mut seen); let _ = std::fs::remove_dir_all(&dir); } diff --git a/tools/orderbook-mock/Cargo.toml b/tools/orderbook-mock/Cargo.toml index af27e38c..50855076 100644 --- a/tools/orderbook-mock/Cargo.toml +++ b/tools/orderbook-mock/Cargo.toml @@ -12,6 +12,9 @@ path = "src/main.rs" [dependencies] anyhow.workspace = true +# The UID a real orderbook returns is derived from the posted order, so +# the mock reuses the same projection and hash the venue checks against. +cowprotocol = { version = "0.2.0", default-features = false } axum.workspace = true clap.workspace = true rand.workspace = true diff --git a/tools/orderbook-mock/src/main.rs b/tools/orderbook-mock/src/main.rs index 897e1b46..77d7a287 100644 --- a/tools/orderbook-mock/src/main.rs +++ b/tools/orderbook-mock/src/main.rs @@ -43,6 +43,11 @@ struct Cli { /// `InsufficientFee` (transient) and `InvalidSignature` (permanent). #[arg(long, default_value_t = 0.0)] error_rate: f64, + + /// Chain whose settlement domain the returned UID is derived under. + /// A UID from the wrong domain is refused by the submitting venue. + #[arg(long, default_value_t = 1)] + chain_id: u64, } #[derive(Debug, Default)] @@ -155,29 +160,44 @@ async fn post_orders(State(state): State>, body: String) -> impl I .into_response(); } - // Synthesise a deterministic-per-call OrderUid. The orderbook's - // real UID is `keccak(orderData) ++ owner ++ validTo`; for the - // load test the only requirement is that each response is a valid - // 56-byte hex (224 bits) so the host's cowprotocol decoder - // accepts it. - let n = state.counters.submits_ok.fetch_add(1, Ordering::Relaxed); - let _ = body; // intentionally ignored; load test does not validate the OrderCreation shape - let mut uid = [0u8; 56]; - uid[0..8].copy_from_slice(&n.to_be_bytes()); - let uid_hex = format!("\"0x{}\"", hex_encode_inline(&uid)); - (StatusCode::CREATED, uid_hex).into_response() + // The UID is derived from the posted order, not invented. A + // submitting venue re-derives it locally and refuses a receipt that + // disagrees, so a synthetic UID can never complete a submit. + let uid = match derive_uid(&body, state.cli.chain_id) { + Ok(uid) => uid, + Err(why) => { + state.counters.submits_err.fetch_add(1, Ordering::Relaxed); + tracing::warn!("refusing an order this mock cannot derive a UID for: {why}"); + return ( + StatusCode::BAD_REQUEST, + axum::Json(serde_json::json!({ + "errorType": "InvalidSignature", + "description": format!("mock could not derive a UID: {why}"), + })), + ) + .into_response(); + } + }; + state.counters.submits_ok.fetch_add(1, Ordering::Relaxed); + (StatusCode::CREATED, format!("\"{uid}\"")).into_response() } -/// Inline hex encoder; keeps the mock's dependency surface minimal. -fn hex_encode_inline(bytes: &[u8]) -> String { - use std::fmt::Write as _; - let mut s = String::with_capacity(bytes.len() * 2); - for b in bytes { - write!(s, "{b:02x}").expect("writing to String never fails"); - } - s +/// The UID a real orderbook would assign to this posted order. +/// +/// `OrderCreation::order_data` projects back the twelve signed fields +/// the UID was computed against, so the derivation here is the same one +/// the submitting venue checks with, not a second implementation of it. +fn derive_uid(body: &str, chain_id: u64) -> Result { + let creation: cowprotocol::OrderCreation = + serde_json::from_str(body).map_err(|e| format!("body is not an OrderCreation: {e}"))?; + let chain = + cowprotocol::Chain::try_from(chain_id).map_err(|_| format!("unknown chain {chain_id}"))?; + Ok(creation + .order_data() + .uid(&chain.settlement_domain(), creation.from)) } +/// Inline hex encoder; keeps the mock's dependency surface minimal. #[cfg(test)] mod tests { use super::*; @@ -197,17 +217,48 @@ mod tests { port: 0, latency_ms: 0, error_rate: 0.0, + chain_id: 1, } } + /// A minimal order in the shape the venue posts. + fn creation() -> cowprotocol::OrderCreation { + serde_json::from_str( + r#"{ + "sellToken": "0xc02aaa39b223fe8d0a0e5c4f27ead9083c756cc2", + "buyToken": "0xa0b86991c6218b36c1d19d4a2e9eb0ce3606eb48", + "receiver": "0x00112233445566778899aabbccddeeff00112233", + "sellAmount": "1000000000000000", + "buyAmount": "1000000", + "validTo": 2000000000, + "appData": "0x0000000000000000000000000000000000000000000000000000000000000000", + "feeAmount": "0", + "kind": "sell", + "partiallyFillable": false, + "sellTokenBalance": "erc20", + "buyTokenBalance": "erc20", + "signingScheme": "presign", + "signature": "0x", + "from": "0x00112233445566778899aabbccddeeff00112233" + }"#, + ) + .expect("the fixture is a valid OrderCreation") + } + + /// The venue re-derives the UID locally and refuses a receipt that + /// disagrees, so echoing the order's own UID is the whole point. #[tokio::test] - async fn post_orders_returns_56_byte_hex_uid() { + async fn post_orders_returns_the_uid_derived_from_the_order() { + let order = creation(); + let expected = order + .order_data() + .uid(&cowprotocol::Chain::Mainnet.settlement_domain(), order.from); let app = router_with(default_cli()); let resp = app .oneshot( Request::post("/api/v1/orders") .header("content-type", "application/json") - .body(Body::from(r#"{"any":"body"}"#)) + .body(Body::from(serde_json::to_string(&order).unwrap())) .unwrap(), ) .await @@ -216,18 +267,34 @@ mod tests { let body = axum::body::to_bytes(resp.into_body(), usize::MAX) .await .unwrap(); - let s = std::str::from_utf8(&body).unwrap(); - // JSON-encoded string: "0x..." (1 + 2 + 112 + 1 = 116 chars) - assert!(s.starts_with("\"0x")); - assert_eq!(s.len(), 116); + assert_eq!( + std::str::from_utf8(&body).unwrap(), + format!("\"{expected}\""), + ); + } + + /// A body this mock cannot read is refused rather than answered with + /// something the venue would reject anyway. + #[tokio::test] + async fn an_underivable_body_is_refused() { + let app = router_with(default_cli()); + let resp = app + .oneshot( + Request::post("/api/v1/orders") + .header("content-type", "application/json") + .body(Body::from(r#"{"any":"body"}"#)) + .unwrap(), + ) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::BAD_REQUEST); } #[tokio::test] async fn error_rate_one_always_returns_envelope() { let app = router_with(Cli { - port: 0, - latency_ms: 0, error_rate: 1.0, + ..default_cli() }); let resp = app .oneshot( From 7a0df9f015e987b702cbb14b1b0c2936ca7c8963 Mon Sep 17 00:00:00 2001 From: mfw78 Date: Tue, 8 Sep 2026 23:07:47 +0000 Subject: [PATCH 3/7] fix(cow): attribute an unreadable poll refusal before acting on it `classify_revert` fell through to `TryNextBlock` whenever it could not find a selector, so a refusal with no readable payload re-polled on every block indefinitely, silently, with nothing to distinguish it from a healthy scheduling gap. The poll interface is closed: a generator answers through `GeneratorResult` codes and the registry's own refusals carry selectors. An empty payload is outside that, so it names a defect. It does not say whose, and the candidates are not alike, so the classifier no longer guesses. `classify_revert` returns a `Refusal` and the keeper attributes it with one `eth_getCode`. A codeless owner cannot answer the registry's ERC-1271 probe, which makes that call revert in the caller's own frame with nothing attached. That is the registration answering for itself and it drops, loudly. Every other cause points outward: a gas cap too low for the handler, or a registry address that is not the fork. Those are the same for every commitment, so a drop would delete a whole watch set to report an operator's typo. They back off for an hour and name the likely cause. A failed `eth_getCode` reads as unattributable and takes that path too. Reachable in production by an ordinary mistake, and the codeless case is recoverable rather than permanent: an EOA can gain code through an EIP-7702 delegation, which is exactly the flow #658 describes. Closes #692. AI Assistance: Claude Code used for the fix and the tests. --- crates/composable-cow/src/fork.rs | 103 ++++++++++------- crates/composable-cow/src/lib.rs | 2 +- modules/ccow-monitor/src/keeper.rs | 175 +++++++++++++++++++++++++---- 3 files changed, 214 insertions(+), 66 deletions(-) diff --git a/crates/composable-cow/src/fork.rs b/crates/composable-cow/src/fork.rs index 026c89a0..ba3ab230 100644 --- a/crates/composable-cow/src/fork.rs +++ b/crates/composable-cow/src/fork.rs @@ -212,34 +212,47 @@ sol! { } } -/// Classify a failed poll `eth_call`. Every reachable revert is -/// deterministic on-chain state, so all are terminal; a re-`create` -/// re-indexes through its own event. A payload-free failure is the -/// transport, not the contract, so it stays retryable. +/// What a poll revert says about the commitment. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum Refusal { + /// The contract refused with a selector this build can read. + Named(Selector), + /// The node executed the call and refused with nothing readable. + /// + /// The poll interface is closed: a generator answers through + /// `GeneratorResult` codes and the registry's own refusals carry + /// selectors. An empty payload is outside that, so it names a + /// defect. It does not say whose, and the candidates are not alike: + /// a codeless owner is the registration's, a gas cap too low for the + /// handler or a wrong registry address is the operator's. + Unattributed, + /// The call never reached the node. + Transport, +} + +/// Classify a poll revert. Attribution of [`Refusal::Unattributed`] +/// needs a second chain read, so it happens at the call site. #[must_use] -pub fn classify_revert(err: &ChainError) -> Verdict { +pub fn classify_revert(err: &ChainError) -> Refusal { let ChainError::Rpc(rpc) = err else { - // `ChainError` is `#[non_exhaustive]`: transport faults and any - // future case are payload-free, so they stay retryable. - return Verdict::TryNextBlock { - reason: Selector::ZERO, - }; - }; - let Some(data) = rpc.data.as_deref() else { - return Verdict::TryNextBlock { - reason: Selector::ZERO, - }; + return Refusal::Transport; }; - let Some(reason) = data.get(..4).map(Selector::from_slice) else { - return Verdict::TryNextBlock { - reason: Selector::ZERO, - }; - }; - // Unrecognised still means the contract refused; a handler `Panic` - // lands here. - Verdict::Invalid { reason } + rpc.data + .as_deref() + .and_then(|data| data.get(..4)) + .map_or(Refusal::Unattributed, |s| { + Refusal::Named(Selector::from_slice(s)) + }) } +/// How long a payload-free revert takes the commitment out of the +/// rotation. +/// +/// Long enough that the poll budget is not spent on a deterministic +/// failure, short enough that a commitment recovers on its own if the +/// cause was the node rather than the registration. +pub const PAYLOAD_FREE_BACKOFF_S: u64 = 3_600; + /// Selectors the classifier recognises, for logging and tests. #[must_use] pub fn is_residual_selector(selector: Selector) -> bool { @@ -591,7 +604,7 @@ mod tests { mod residual_tests { use alloy_primitives::fixed_bytes; use alloy_sol_types::SolError; - use nexum_sdk::host::RpcError; + use nexum_sdk::host::{Fault, RpcError}; use super::*; @@ -614,11 +627,8 @@ mod residual_tests { ] .map(Selector::from) { - let verdict = classify_revert(&reverted(Some(selector.to_vec()))); - assert!( - matches!(verdict, Verdict::Invalid { reason } if reason == selector), - "{selector:?} produced {verdict:?}", - ); + let refusal = classify_revert(&reverted(Some(selector.to_vec()))); + assert_eq!(refusal, Refusal::Named(selector), "{selector:?}"); assert!(is_residual_selector(selector)); } } @@ -628,22 +638,33 @@ mod residual_tests { fn an_unrecognised_selector_is_terminal() { let panic_selector = Selector::new([0x4e, 0x48, 0x7b, 0x71]); assert!(!is_residual_selector(panic_selector)); - assert!(matches!( + assert_eq!( classify_revert(&reverted(Some(panic_selector.to_vec()))), - Verdict::Invalid { .. } - )); + Refusal::Named(panic_selector), + ); } + /// The closed poll interface specifies neither of these, so both + /// name a defect the caller has to attribute before acting. #[test] - fn a_payload_free_failure_stays_retryable() { - assert!(matches!( - classify_revert(&reverted(None)), - Verdict::TryNextBlock { .. } - )); - assert!(matches!( - classify_revert(&reverted(Some(vec![1, 2]))), - Verdict::TryNextBlock { .. } - )); + fn an_unreadable_payload_is_not_attributed_here() { + for data in [None, Some(vec![1, 2])] { + assert_eq!( + classify_revert(&reverted(data.clone())), + Refusal::Unattributed, + "{data:?}", + ); + } + } + + /// A transport fault never reached the node, so it says nothing + /// about the contract. + #[test] + fn a_transport_fault_is_its_own_class() { + assert_eq!( + classify_revert(&ChainError::Fault(Fault::Unavailable("node down".into()))), + Refusal::Transport, + ); } /// A rename upstream must fail here, not silently reclassify a diff --git a/crates/composable-cow/src/lib.rs b/crates/composable-cow/src/lib.rs index b56da0a9..5add6f41 100644 --- a/crates/composable-cow/src/lib.rs +++ b/crates/composable-cow/src/lib.rs @@ -16,7 +16,7 @@ pub mod poll; #[cfg(feature = "run")] pub mod run; -pub use fork::{Mapped, PollResult, Suppressed, classify_revert, map_verdict, to_verdict}; +pub use fork::{Mapped, PollResult, Refusal, Suppressed, classify_revert, map_verdict, to_verdict}; pub use poll::{NextPoll, Verdict}; #[cfg(feature = "run")] pub use run::run; diff --git a/modules/ccow-monitor/src/keeper.rs b/modules/ccow-monitor/src/keeper.rs index 9f1fcea4..ce0daf7e 100644 --- a/modules/ccow-monitor/src/keeper.rs +++ b/modules/ccow-monitor/src/keeper.rs @@ -11,7 +11,9 @@ use alloy_primitives::{Address, B256, Bytes, Selector, keccak256}; use alloy_sol_types::{SolCall, SolEvent, SolValue}; use composable::ConditionalOrderParams; -use composable_cow::fork::{classify_revert, decode_poll_return, map_verdict, to_verdict}; +use composable_cow::fork::{ + PAYLOAD_FREE_BACKOFF_S, Refusal, classify_revert, decode_poll_return, map_verdict, to_verdict, +}; use composable_cow::{Verdict, run}; use cow_venue::CowClient; // The poll path receives the order inside `PollResult`, so the bare type @@ -628,31 +630,72 @@ fn poll_one( to_verdict(map_verdict(&result, &signature), valid_to) }, ), - Err(err) => { - let outcome = classify_revert(&err); - match &err { - ChainError::Fault(fault) => { - tracing::warn!("eth_call failed ({fault}); retrying next block"); + Err(err) => match classify_revert(&err) { + Refusal::Transport => { + tracing::warn!("eth_call failed ({err}); retrying next block"); + Verdict::TryNextBlock { + reason: Selector::ZERO, } - // A permanent drop deserves its cause on the record: the - // selector and the node's message are unrecoverable once - // the commitment is gone. - ChainError::Rpc(rpc) if matches!(outcome, Verdict::Invalid { .. }) => { - let selector = rpc - .data - .as_deref() - .and_then(|data| data.get(..4)) - .map(|s| format!("{:#x}", Selector::from_slice(s))) - .unwrap_or_else(|| "none".to_string()); - tracing::warn!( - "eth_call reverted permanently (selector {selector}, {}); \ - dropping commitment", - rpc.message, - ); - } - _ => {} } - outcome + // A permanent drop deserves its cause on the record: the + // selector is unrecoverable once the commitment is gone. + Refusal::Named(reason) => { + tracing::warn!( + "eth_call reverted permanently (selector {reason:#x}, {err}); \ + dropping commitment" + ); + Verdict::Invalid { reason } + } + Refusal::Unattributed => attribute_unreadable(host, tick, owner, &err), + }, + } +} + +/// Decide who an unreadable refusal belongs to. +/// +/// The registry probes the owner for ERC-1271 support before it builds a +/// signature, and a codeless owner makes that probe revert in the +/// caller's own frame with nothing attached. So an empty payload plus an +/// owner with no code is the registration answering for itself, and no +/// later poll of this commitment changes it. +/// +/// Every other cause points outward, at a gas cap too low for the +/// handler or a registry address that is not the fork. Those are the +/// operator's, they are the same for every commitment, and removing a +/// watch set is not how to report them. +fn attribute_unreadable( + host: &H, + tick: &Tick, + owner: &Address, + err: &ChainError, +) -> Verdict { + let params = format!(r#"["{owner:#x}","latest"]"#); + let code = host + .request(tick.chain_id, "eth_getCode", ¶ms) + .ok() + .and_then(|json| parse_eth_call_result(&json)); + match code.as_deref() { + Some([]) => { + tracing::error!( + "dropping commitment: owner {owner:#x} has no code, so the registry cannot \ + build an ERC-1271 signature for it ({err})" + ); + Verdict::Invalid { + reason: Selector::ZERO, + } + } + // Includes the read failing: without an answer this is not + // attributable, and the safe reading of that is the loud one. + _ => { + tracing::error!( + "poll refused with no readable payload and owner {owner:#x} has code; \ + check the poll gas cap and the registry address ({err}); \ + backing off {PAYLOAD_FREE_BACKOFF_S}s" + ); + Verdict::WaitTimestamp { + wait_until: tick.epoch_s.saturating_add(PAYLOAD_FREE_BACKOFF_S), + reason: Selector::ZERO, + } } } } @@ -1873,6 +1916,90 @@ mod tests { assert!(!store.keys().any(|k| k.starts_with("submitted:"))); } + /// A codeless owner cannot answer the registry's ERC-1271 probe, so + /// the poll reverts in the caller's own frame with nothing attached. + /// That is the registration answering for itself. + #[test] + fn an_unreadable_refusal_with_a_codeless_owner_drops() { + use nexum_sdk::host::RpcError; + + let host = MockHost::new(); + let venue = MockVenue::default(); + let owner = address!("0011223344556677889900AABBCCDDEEFF001122"); + let params = sample_params(); + let key = seed_commitment(&host, owner, ¶ms); + + host.chain.respond_to( + "eth_call", + programmed_eth_call_params(owner, ¶ms), + Err(ChainError::Rpc(RpcError { + code: 3, + message: "execution reverted".into(), + data: None, + })), + ); + host.chain.respond_to( + "eth_getCode", + format!(r#"["{owner:#x}","latest"]"#), + Ok("\"0x\"".to_owned()), + ); + + let (result, logs) = capture_tracing(|| dispatch(&host, &venue, sample_block(1_000))); + result.unwrap(); + + assert!( + !host.store.snapshot().contains_key(&key), + "the commitment is removed", + ); + assert!( + logs.any(|e| e.message.contains("has no code")), + "and the drop says why", + ); + } + + /// Every other cause of an unreadable refusal points outward, at the + /// gas cap or the registry address. Those are the same for every + /// commitment, so removing one reports nothing and loses a watch. + #[test] + fn an_unreadable_refusal_with_a_contract_owner_backs_off() { + use nexum_sdk::host::RpcError; + + let host = MockHost::new(); + let venue = MockVenue::default(); + let owner = address!("0011223344556677889900AABBCCDDEEFF001122"); + let params = sample_params(); + let key = seed_commitment(&host, owner, ¶ms); + + host.chain.respond_to( + "eth_call", + programmed_eth_call_params(owner, ¶ms), + Err(ChainError::Rpc(RpcError { + code: 3, + message: "out of gas".into(), + data: None, + })), + ); + host.chain.respond_to( + "eth_getCode", + format!(r#"["{owner:#x}","latest"]"#), + Ok("\"0x60806040\"".to_owned()), + ); + + let (result, logs) = capture_tracing(|| dispatch(&host, &venue, sample_block(1_000))); + result.unwrap(); + + let store = host.store.snapshot(); + assert!(store.contains_key(&key), "the commitment survives"); + assert!( + store.keys().any(|k| k.starts_with("next_epoch:")), + "and is gated forward: {store:?}", + ); + assert!( + logs.any(|e| e.message.contains("gas cap")), + "and the operator is pointed at the likely cause", + ); + } + #[test] fn poll_invalid_drops_commitment_and_gates() { // A residual revert must delete the commitment and any stale From fcfab0557a76fcfa17e31b80393f577dd18c734c Mon Sep 17 00:00:00 2001 From: mfw78 Date: Wed, 9 Sep 2026 02:19:59 +0000 Subject: [PATCH 4/7] test(cow): share one harness, and register from a fresh owner Extracts the anvil fork, the orderbook mock and the engine into a `Harness` so a second case can reuse them. The run also stops registering from anvil's account zero. That is the well-known test key, and on mainnet it already carries an EIP-7702 delegation, which a fork inherits: an owner meant to have known code would silently have whatever mainnet gave it. The run funds and impersonates a fresh address instead, so the owner's code is exactly what the test put there. AI Assistance: Claude Code used for the refactor. --- .../shepherd-engine/tests/anvil_fork_e2e.rs | 228 ++++++++++++------ 1 file changed, 148 insertions(+), 80 deletions(-) diff --git a/crates/shepherd-engine/tests/anvil_fork_e2e.rs b/crates/shepherd-engine/tests/anvil_fork_e2e.rs index a6decb84..9218792f 100644 --- a/crates/shepherd-engine/tests/anvil_fork_e2e.rs +++ b/crates/shepherd-engine/tests/anvil_fork_e2e.rs @@ -37,6 +37,14 @@ const HANDLER: &str = "0x4e17a65d14e7f37d2a9f0389f17efd41aaa64c91"; /// `_buildSignature` catches and treats as a non-Safe wallet. const ACCOUNT_7702: &str = "0x15236F06922A288B68e57Ab42e397920a3F3Bb99"; +/// The owner the run registers from. +/// +/// A fresh address, not one of anvil's unlocked accounts. Account zero +/// is the well-known test key and already carries an EIP-7702 +/// delegation on mainnet, which a fork inherits, so its code would be +/// whatever mainnet says rather than what the test set. +const POST_OWNER: &str = "0x00000000000000000000000000000000CafeBabe"; + const WETH: &str = "0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2"; const USDC: &str = "0xA0b86991c6218b36c1d19D4a2e9Eb0cE3606eB48"; @@ -329,96 +337,158 @@ fn wait_for(rx: &Receiver, needle: &str, seen: &mut Vec) -> Stri } } -/// The whole chain leg: a real registration on the forked registry is -/// indexed, then polled through `getTradeableOrderWithSignature` against -/// the deployed `OwnedTWAP`. -#[test] -fn ccow_monitor_indexes_and_polls_a_real_owned_twap() { - let Some(rpc) = fork_rpc_or_skip() else { - return; - }; - let dir = std::env::temp_dir().join(format!("shepherd-anvil-{}", std::process::id())); - std::fs::create_dir_all(&dir).expect("create the run directory"); +/// One anvil fork, one orderbook mock and one engine, reaped together. +struct Harness { + anvil: String, + owner: String, + lines: Receiver, + seen: Vec, + dir: PathBuf, + _anvil: Reaped, + _ob: Reaped, + _engine: Reaped, +} - let anvil_port = free_port(); - let _anvil = start_anvil(&rpc, anvil_port); - let anvil = format!("http://127.0.0.1:{anvil_port}"); +impl Harness { + /// Fork `rpc`, stand the engine up against it, and wait until the + /// module is loaded. `tag` keeps concurrent tests off each other's + /// state directory, and `owner` is funded and impersonated so it can + /// register. + fn start(rpc: &str, tag: &str, owner: &str) -> Self { + let dir = std::env::temp_dir().join(format!("shepherd-anvil-{}-{tag}", std::process::id())); + let _ = std::fs::remove_dir_all(&dir); + std::fs::create_dir_all(&dir).expect("create the run directory"); + + let anvil_port = free_port(); + let _anvil = start_anvil(rpc, anvil_port); + let anvil = format!("http://127.0.0.1:{anvil_port}"); + + let ob_port = free_port(); + let _ob = Reaped( + Command::new(orderbook_mock_bin()) + .args(["--port", &ob_port.to_string()]) + .stdout(Stdio::null()) + .stderr(Stdio::null()) + .spawn() + .expect("spawn the orderbook mock"), + ); - let ob_port = free_port(); - let _ob = Reaped( - Command::new(orderbook_mock_bin()) - .args(["--port", &ob_port.to_string()]) - .stdout(Stdio::null()) - .stderr(Stdio::null()) + let fork_block: u64 = cast(&["block-number", "--rpc-url", &anvil]) + .parse() + .expect("a numeric block number"); + let manifest = manifest_at(&dir, fork_block); + let config = engine_config( + &dir, + &anvil, + &format!("http://127.0.0.1:{ob_port}"), + &manifest.display().to_string(), + ); + + let mut engine = Command::new(env!("CARGO_BIN_EXE_shepherd")) + .args(["--engine-config", &config.display().to_string()]) + .current_dir(workspace_root()) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()) .spawn() - .expect("spawn the orderbook mock"), - ); + .expect("spawn the shepherd binary"); + let lines = stream_lines(&mut engine); + let _engine = Reaped(engine); + + // A fresh address, not one of anvil's unlocked accounts. Account + // zero is the well-known test key, and on mainnet it already + // carries an EIP-7702 delegation, which a fork inherits: an + // owner meant to be codeless would silently have code. + let owner = owner.to_owned(); + cast(&[ + "rpc", + "anvil_setBalance", + &owner, + "0xde0b6b3a7640000", + "--rpc-url", + &anvil, + ]); + cast(&[ + "rpc", + "anvil_impersonateAccount", + &owner, + "--rpc-url", + &anvil, + ]); + + let mut harness = Self { + anvil, + owner, + lines, + seen: Vec::new(), + dir, + _anvil, + _ob, + _engine, + }; + harness.wait_for("ccow-monitor"); + harness + } - let fork_block: u64 = cast(&["block-number", "--rpc-url", &anvil]) - .parse() - .expect("a numeric block number"); - let manifest = manifest_at(&dir, fork_block); - let config = engine_config( - &dir, - &anvil, - &format!("http://127.0.0.1:{ob_port}"), - &manifest.display().to_string(), - ); + /// Delegate the owner EOA to `CowAccount7702`. + /// + /// The designator is written directly rather than by authorisation + /// tuple: the test needs the code to be there, not the signing + /// ceremony that puts it there on a real chain. + fn delegate_owner(&self) { + cast(&[ + "rpc", + "anvil_setCode", + &self.owner, + &format!("0xef0100{}", ACCOUNT_7702.trim_start_matches("0x")), + "--rpc-url", + &self.anvil, + ]); + } - let mut engine = Command::new(env!("CARGO_BIN_EXE_shepherd")) - .args(["--engine-config", &config.display().to_string()]) - .current_dir(workspace_root()) - .stdout(Stdio::piped()) - .stderr(Stdio::piped()) - .spawn() - .expect("spawn the shepherd binary"); - let lines = stream_lines(&mut engine); - let _engine = Reaped(engine); - - let mut seen = Vec::new(); - wait_for(&lines, "ccow-monitor", &mut seen); - - // anvil's first unlocked account owns the registration. - let owner = cast(&["rpc", "eth_accounts", "--rpc-url", &anvil]) - .trim_matches(|c| c == '[' || c == ']' || c == '"') - .split(',') - .next() - .expect("anvil funds at least one account") - .trim_matches('"') - .to_owned(); - // EIP-7702 delegation designator: 0xef0100 ++ implementation. Set - // directly rather than by authorisation tuple, because the test only - // needs the code to be there, not the signing ceremony that puts it - // there on a real chain. - cast(&[ - "rpc", - "anvil_setCode", - &owner, - &format!("0xef0100{}", ACCOUNT_7702.trim_start_matches("0x")), - "--rpc-url", - &anvil, - ]); - register_order(&anvil, &owner); + /// Mine, so a registration gets the block dispatch that polls it. + fn mine(&self, blocks: usize) { + for _ in 0..blocks { + cast(&["rpc", "evm_mine", "--rpc-url", &self.anvil]); + std::thread::sleep(Duration::from_millis(500)); + } + } - wait_for(&lines, "indexed commitment:", &mut seen); + fn wait_for(&mut self, needle: &str) -> String { + wait_for(&self.lines, needle, &mut self.seen) + } - // The poll fires on a block dispatch, so the registration needs one - // after it. Mining is cheap; the engine's own poller drives the rest. - for _ in 0..6 { - cast(&["rpc", "evm_mine", "--rpc-url", &anvil]); - std::thread::sleep(Duration::from_millis(500)); + fn saw(&self, needle: &str) -> bool { + self.seen.iter().any(|l| l.contains(needle)) } +} + +impl Drop for Harness { + fn drop(&mut self) { + let _ = std::fs::remove_dir_all(&self.dir); + } +} + +/// The whole path: a real registration on the forked registry is +/// indexed, polled to `Post` through the structured generator wire, and +/// submitted to the venue. +#[test] +fn ccow_monitor_polls_a_real_owned_twap_to_post() { + let Some(rpc) = fork_rpc_or_skip() else { + return; + }; + let mut h = Harness::start(&rpc, "post", POST_OWNER); + h.delegate_owner(); + register_order(&h.anvil, &h.owner.clone()); + + h.wait_for("indexed commitment:"); + h.mine(6); // `poll {key} -> {outcome}` is the module's own line, so this is the // structured generator wire answering over a real eth_call against // the deployed `OwnedTWAP`, not a fixture. - let polled = wait_for(&lines, "poll commitment:", &mut seen); + let polled = h.wait_for("poll commitment:"); for bad in ["did not decode", "eth_call failed"] { - assert!( - !seen.iter().any(|l| l.contains(bad)), - "the poll leg reported {bad:?}:\n {}", - seen.join("\n "), - ); + assert!(!h.saw(bad), "the poll leg reported {bad:?}"); } assert!( polled.contains("-> Post"), @@ -427,7 +497,5 @@ fn ccow_monitor_indexes_and_polls_a_real_owned_twap() { // And the post reaches the venue, closing the loop through the // borsh body, the journal reservation and the orderbook adapter. - wait_for(&lines, "submitted", &mut seen); - - let _ = std::fs::remove_dir_all(&dir); + h.wait_for("submitted"); } From b200f14fccc8762bff72ca5d5afb87d4307e3f9f Mon Sep 17 00:00:00 2001 From: mfw78 Date: Wed, 9 Sep 2026 05:40:04 +0000 Subject: [PATCH 5/7] fix(cow): end a submission the orderbook cannot be asked about Bumps the videre pin to 9ab1515 and supplies the CoW fault policy the seam there exists for. A receipt this keeper cannot correlate used to remove the commitment. The fork run reached `Post`, submitted, got a uid that did not match the one derived from the order it sent, and destroyed a commitment whose next poll would have minted a perfectly good order. The uid is `keccak(order) ++ owner ++ validTo`, a pure function of the body that was sent, so re-posting it yields the same disagreement every time, and the order cannot be asked after either: the only handle on it is a uid this keeper does not trust. So the submission ends and its reservation is released. Leaving it would have reconcile re-post the same body on every tick forever. The commitment is not the submission, and it survives: the next poll mints a later part with a later `validTo` and a different uid. A repeat at a later block says the disagreement is systematic rather than one order's, and then it drops. That parts from `is_terminal` deriving from `action`, deliberately and in the one place the two questions differ. Denials now route through the shipped table for the reconcile pass as well as the submit path. They did not before: reconcile took the platform default, which knows nothing of `classification.toml`, so a stranded reservation for a clearable refusal was released while the same refusal on the submit path backed off. Closes #693. AI Assistance: Claude Code used for the policy and the tests. --- Cargo.lock | 10 +-- crates/composable-cow/Cargo.toml | 2 +- crates/composable-cow/src/run.rs | 19 ++-- crates/composable-cow/tests/run.rs | 69 +++++++++++++++ crates/cow-venue/Cargo.toml | 6 +- crates/cow-venue/src/classification.rs | 115 +++++++++++++++++++++++++ crates/cow-venue/src/lib.rs | 4 +- crates/shepherd-engine/Cargo.toml | 4 +- modules/ccow-monitor/Cargo.toml | 2 +- modules/ethflow-watcher/Cargo.toml | 2 +- tools/orderbook-mock/src/main.rs | 21 ++++- 11 files changed, 226 insertions(+), 28 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index fcd7bed1..df0348e0 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -6036,7 +6036,7 @@ checksum = "0b928f33d975fc6ad9f86c8f283853ad26bdd5b10b7f1542aa2fa15e2289105a" [[package]] name = "videre-host" version = "0.1.0" -source = "git+https://github.com/nullislabs/videre-nexum-module?rev=fd8af027d9ae84823f8fcdcc6902b3af80cab451#fd8af027d9ae84823f8fcdcc6902b3af80cab451" +source = "git+https://github.com/nullislabs/videre-nexum-module?rev=9ab15152b9279974cb42c3ee3aff054344f6050d#9ab15152b9279974cb42c3ee3aff054344f6050d" dependencies = [ "anyhow", "derive_more", @@ -6054,7 +6054,7 @@ dependencies = [ [[package]] name = "videre-macros" version = "0.1.0" -source = "git+https://github.com/nullislabs/videre-nexum-module?rev=fd8af027d9ae84823f8fcdcc6902b3af80cab451#fd8af027d9ae84823f8fcdcc6902b3af80cab451" +source = "git+https://github.com/nullislabs/videre-nexum-module?rev=9ab15152b9279974cb42c3ee3aff054344f6050d#9ab15152b9279974cb42c3ee3aff054344f6050d" dependencies = [ "nexum-world", "proc-macro2", @@ -6066,7 +6066,7 @@ dependencies = [ [[package]] name = "videre-sdk" version = "0.1.0" -source = "git+https://github.com/nullislabs/videre-nexum-module?rev=fd8af027d9ae84823f8fcdcc6902b3af80cab451#fd8af027d9ae84823f8fcdcc6902b3af80cab451" +source = "git+https://github.com/nullislabs/videre-nexum-module?rev=9ab15152b9279974cb42c3ee3aff054344f6050d#9ab15152b9279974cb42c3ee3aff054344f6050d" dependencies = [ "borsh", "http", @@ -6082,7 +6082,7 @@ dependencies = [ [[package]] name = "videre-status-body" version = "0.1.0" -source = "git+https://github.com/nullislabs/videre-nexum-module?rev=fd8af027d9ae84823f8fcdcc6902b3af80cab451#fd8af027d9ae84823f8fcdcc6902b3af80cab451" +source = "git+https://github.com/nullislabs/videre-nexum-module?rev=9ab15152b9279974cb42c3ee3aff054344f6050d#9ab15152b9279974cb42c3ee3aff054344f6050d" dependencies = [ "borsh", "thiserror 2.0.18", @@ -6091,7 +6091,7 @@ dependencies = [ [[package]] name = "videre-test" version = "0.1.0" -source = "git+https://github.com/nullislabs/videre-nexum-module?rev=fd8af027d9ae84823f8fcdcc6902b3af80cab451#fd8af027d9ae84823f8fcdcc6902b3af80cab451" +source = "git+https://github.com/nullislabs/videre-nexum-module?rev=9ab15152b9279974cb42c3ee3aff054344f6050d#9ab15152b9279974cb42c3ee3aff054344f6050d" dependencies = [ "borsh", "hex", diff --git a/crates/composable-cow/Cargo.toml b/crates/composable-cow/Cargo.toml index 3809e6d3..f8ec91a2 100644 --- a/crates/composable-cow/Cargo.toml +++ b/crates/composable-cow/Cargo.toml @@ -22,7 +22,7 @@ nexum-sdk = { git = "https://github.com/nullislabs/nexum-runtime", rev = "2ed882 # `run` slice: the keeper run over the typed CoW client on the # `videre:venue/client` seam. cow-venue = { path = "../cow-venue", features = ["client", "assembly"], optional = true } -videre-sdk = { git = "https://github.com/nullislabs/videre-nexum-module", rev = "fd8af027d9ae84823f8fcdcc6902b3af80cab451", optional = true } +videre-sdk = { git = "https://github.com/nullislabs/videre-nexum-module", rev = "9ab15152b9279974cb42c3ee3aff054344f6050d", optional = true } tracing = { workspace = true, optional = true } [features] diff --git a/crates/composable-cow/src/run.rs b/crates/composable-cow/src/run.rs index 5376c6cd..02671725 100644 --- a/crates/composable-cow/src/run.rs +++ b/crates/composable-cow/src/run.rs @@ -11,13 +11,13 @@ //! never dropped. //! //! Store faults abort the run (the next tick replays it); a submission -//! failure folds into a [`RetryAction`], a `denied` refusal re-entering -//! the CoW classification by its errorType prefix ([`classify_denied`]). +//! failure folds into a [`RetryAction`] through [`CowFaults`], which the +//! reconcile pass reads too, so both paths classify a refusal alike. //! Diagnostics go through the guest `tracing` facade. use alloy_primitives::{Address, Bytes, hex}; use cow_venue::assembly::{gpv2_to_order_data, order_data_to_body}; -use cow_venue::{CowClient, CowIntent, CowIntentBody, CowVenue, SignedOrder, classify_denied}; +use cow_venue::{CowClient, CowFaults, CowIntent, CowIntentBody, CowVenue, SignedOrder}; use cowprotocol::GPv2OrderData; use nexum_sdk::host::{Fault, ListQuery, LocalStoreHost}; use nexum_sdk::keeper::{ @@ -26,11 +26,10 @@ use nexum_sdk::keeper::{ }; use std::task::Poll; +use videre_sdk::FaultPolicy as _; use videre_sdk::client::poll_once; -use videre_sdk::keeper::{retry_action, submission_key}; -use videre_sdk::{ - ClientError, IntentBody as _, SubmitOutcome, Venue as _, VenueFault, VenueTransport, -}; +use videre_sdk::keeper::submission_key; +use videre_sdk::{ClientError, IntentBody as _, SubmitOutcome, Venue as _, VenueTransport}; use crate::{NextPoll, Verdict}; @@ -57,6 +56,7 @@ where &journal, tick, videre_sdk::DEFAULT_RECONCILE_BUDGET, + &CowFaults, )) { Poll::Ready(res) => { res?; @@ -429,10 +429,7 @@ where tracing::error!("intent body encode failed: {err}"); } Err(ClientError::Venue(fault)) => { - let action = match &fault { - VenueFault::Denied(detail) => classify_denied(detail), - other => retry_action(other), - }; + let action = CowFaults.action(&fault); Retrier::new(host).apply(commitment, action, tick)?; match action { RetryAction::TryNextBlock => tracing::warn!("submit retry-next-block: {fault}"), diff --git a/crates/composable-cow/tests/run.rs b/crates/composable-cow/tests/run.rs index a80b5072..317be43f 100644 --- a/crates/composable-cow/tests/run.rs +++ b/crates/composable-cow/tests/run.rs @@ -1946,3 +1946,72 @@ fn an_insufficient_valid_to_keeps_the_commitment() { assert!(host.store.snapshot().contains_key(&key)); } + +/// The reconcile pass must read the shipped table, not the platform +/// default which knows nothing of it. +/// +/// A stranded reservation refused for a balance the owner can top up is +/// still owed, so it stays reserved. The default would release it, and +/// did before this policy reached reconcile. +#[test] +fn reconcile_keeps_a_reservation_the_table_says_is_clearable() { + let host = MockHost::new(); + seed_commitment(&host); + let key = "cow:0xdeadbeef"; + Journal::submitted(&host) + .reserve(key, b"body") + .expect("reserve"); + + let venue = MockVenue::default(); + venue.enqueue_submit(Err(VenueFault::Denied( + "InsufficientBalance: not enough sell token".into(), + ))); + + run( + &host, + &client(&venue), + &src(|_, _, _, _| Verdict::TryNextBlock { + reason: Selector::ZERO, + }), + &sample_tick(), + ) + .unwrap(); + + assert_eq!( + Journal::submitted(&host).mark(key).unwrap(), + Some(Mark::Reserved), + "a clearable refusal leaves the submit owed", + ); +} + +/// And a receipt it cannot correlate ends the submission, because the +/// uid is a pure function of the body: re-posting it fails identically, +/// so leaving the reservation would re-post it on every tick forever. +#[test] +fn reconcile_releases_a_reservation_whose_receipt_cannot_be_correlated() { + let host = MockHost::new(); + seed_commitment(&host); + let key = "cow:0xdeadbeef"; + Journal::submitted(&host) + .reserve(key, b"body") + .expect("reserve"); + + let venue = MockVenue::default(); + venue.enqueue_submit(Err(VenueFault::ReceiptMismatch)); + + run( + &host, + &client(&venue), + &src(|_, _, _, _| Verdict::TryNextBlock { + reason: Selector::ZERO, + }), + &sample_tick(), + ) + .unwrap(); + + assert_eq!( + Journal::submitted(&host).mark(key).unwrap(), + None, + "the reservation is released rather than re-posted forever", + ); +} diff --git a/crates/cow-venue/Cargo.toml b/crates/cow-venue/Cargo.toml index 5fcf92b0..cb64655b 100644 --- a/crates/cow-venue/Cargo.toml +++ b/crates/cow-venue/Cargo.toml @@ -16,7 +16,7 @@ workspace = true borsh = { workspace = true, optional = true } # Source of the `IntentBody` derive and trait the version enum implements, # and the typed intent client the `client` slice binds to the CoW venue. -videre-sdk = { git = "https://github.com/nullislabs/videre-nexum-module", rev = "fd8af027d9ae84823f8fcdcc6902b3af80cab451", optional = true } +videre-sdk = { git = "https://github.com/nullislabs/videre-nexum-module", rev = "9ab15152b9279974cb42c3ee3aff054344f6050d", optional = true } # `client` slice only: the keeper `RetryAction` the generated # classification table maps each errorType to. The TOML parse happens in # `build.rs`, so serde/toml/thiserror are build- and dev-only and never @@ -39,7 +39,7 @@ serde_json = { workspace = true, optional = true } http = { workspace = true, optional = true } url = { workspace = true, optional = true } # The registration seam: `VenueInvoker`, `VenueRegistry`, and `Liveness`. -videre-host = { git = "https://github.com/nullislabs/videre-nexum-module", rev = "fd8af027d9ae84823f8fcdcc6902b3af80cab451", optional = true } +videre-host = { git = "https://github.com/nullislabs/videre-nexum-module", rev = "9ab15152b9279974cb42c3ee3aff054344f6050d", optional = true } # `BoxFuture` plus `FutureExt::boxed`, the shape `VenueInvoker` returns. futures = { workspace = true, optional = true } # The venue owns its HTTP client; there is no scoped host import left to @@ -59,7 +59,7 @@ serde = { workspace = true } toml = { workspace = true } thiserror = { workspace = true } # The conformance kit: holds the body codec to its published vector set. -videre-test = { git = "https://github.com/nullislabs/videre-nexum-module", rev = "fd8af027d9ae84823f8fcdcc6902b3af80cab451" } +videre-test = { git = "https://github.com/nullislabs/videre-nexum-module", rev = "9ab15152b9279974cb42c3ee3aff054344f6050d" } # Parity tests: the upstream `retry_hint()` the shipped table is # reconciled against. cowprotocol = { version = "0.2.0", default-features = false } diff --git a/crates/cow-venue/src/classification.rs b/crates/cow-venue/src/classification.rs index 4e5a9c86..a88c4d26 100644 --- a/crates/cow-venue/src/classification.rs +++ b/crates/cow-venue/src/classification.rs @@ -112,6 +112,59 @@ pub fn is_already_submitted(error_type: OrderbookApiErrorType) -> bool { table().is_already_submitted(&error_type) } +/// The CoW orderbook's fault policy. +/// +/// Two things the platform default cannot know. +/// +/// A receipt this keeper cannot correlate ends the submission but not +/// the commitment. The default welds those together; this is the one +/// place they differ, and [`is_terminal`](videre_sdk::FaultPolicy::is_terminal) says why. +/// +/// Denials route through the shipped table, so the reconcile pass and +/// the submit path read the same policy. They did not before: reconcile +/// took the platform default while the submit path took the table, so a +/// stranded reservation for a clearable refusal was released while the +/// same refusal on the submit path backed off. +#[derive(Clone, Copy, Debug, Default)] +pub struct CowFaults; + +impl videre_sdk::FaultPolicy for CowFaults { + fn action(&self, fault: &videre_sdk::VenueFault) -> RetryAction { + use videre_sdk::VenueFault; + match fault { + VenueFault::Denied(detail) => classify_denied(detail), + // The commitment survives, because the next poll mints a + // different order: a later part, a later `validTo`, a + // different uid. A repeat at a later block says the + // disagreement is systematic rather than one order's, and + // then it drops. + _ if is_receipt_fault(fault) => RetryAction::DropOnRepeat, + other => videre_sdk::retry_action(other), + } + } + + /// A receipt fault ends the submission even though the commitment + /// lives on, so this deliberately parts from the derived default. + /// + /// The uid is `keccak(order) ++ owner ++ validTo`, a pure function + /// of the body that was sent, so re-posting it yields the same + /// disagreement every time. Leaving the reservation would have + /// reconcile re-post it on every tick forever, and the order cannot + /// be asked after either: the only handle on it is a uid this + /// keeper does not trust. + fn is_terminal(&self, fault: &videre_sdk::VenueFault) -> bool { + is_receipt_fault(fault) || matches!(self.action(fault), RetryAction::Drop) + } +} + +/// Whether `fault` says the receipt could not be correlated. +fn is_receipt_fault(fault: &videre_sdk::VenueFault) -> bool { + matches!( + fault, + videre_sdk::VenueFault::InvalidReceipt | videre_sdk::VenueFault::ReceiptMismatch + ) +} + /// Retry action for a coarse `denied` refusal: the `{errorType}:` /// prefix re-enters the table. /// @@ -376,4 +429,66 @@ mod tests { assert!(first.contains_key("error-type")); assert!(first.contains_key("action")); } + + /// The uid is a pure function of the body that was sent, so a + /// disagreement repeats identically. The submission ends, and the + /// commitment keeps its grace because the next poll mints a + /// different order. + #[test] + fn a_receipt_fault_ends_the_submission_but_not_the_commitment() { + use videre_sdk::FaultPolicy as _; + + for fault in [ + videre_sdk::VenueFault::ReceiptMismatch, + videre_sdk::VenueFault::InvalidReceipt, + ] { + assert_eq!( + CowFaults.action(&fault), + RetryAction::DropOnRepeat, + "{fault:?}" + ); + assert!( + CowFaults.is_terminal(&fault), + "{fault:?}: reconcile must not re-post a body that fails identically", + ); + } + } + + /// Denials route through the shipped table, so the reconcile pass + /// reads the same policy the submit path does. It took the platform + /// default before, which knows nothing of the table. + #[test] + fn a_denial_routes_through_the_table() { + use videre_sdk::FaultPolicy as _; + + let clearable = + videre_sdk::VenueFault::Denied("InsufficientBalance: not enough".to_owned()); + assert_eq!( + CowFaults.action(&clearable), + RetryAction::Backoff { seconds: 600 }, + ); + assert!(!CowFaults.is_terminal(&clearable)); + + let permanent = videre_sdk::VenueFault::Denied("WrongOwner: not yours".to_owned()); + assert_eq!(CowFaults.action(&permanent), RetryAction::Drop); + assert!(CowFaults.is_terminal(&permanent)); + } + + /// Everything the venue says nothing special about keeps the + /// platform default. + #[test] + fn other_faults_keep_the_platform_default() { + use videre_sdk::FaultPolicy as _; + + for fault in [ + videre_sdk::VenueFault::Timeout, + videre_sdk::VenueFault::InvalidBody("bad".to_owned()), + ] { + assert_eq!( + CowFaults.action(&fault), + videre_sdk::retry_action(&fault), + "{fault:?}", + ); + } + } } diff --git a/crates/cow-venue/src/lib.rs b/crates/cow-venue/src/lib.rs index f2829853..482926a1 100644 --- a/crates/cow-venue/src/lib.rs +++ b/crates/cow-venue/src/lib.rs @@ -60,6 +60,8 @@ pub use cowprotocol::Chain; pub use transport::{OrderbookHttp, Transport}; #[cfg(feature = "client")] -pub use classification::{ClassificationTable, classify, classify_denied, is_already_submitted}; +pub use classification::{ + ClassificationTable, CowFaults, classify, classify_denied, is_already_submitted, +}; #[cfg(feature = "client")] pub use client::{CowClient, CowVenue, VENUE_ID, intent_id}; diff --git a/crates/shepherd-engine/Cargo.toml b/crates/shepherd-engine/Cargo.toml index 8b5a5b79..9c4979b7 100644 --- a/crates/shepherd-engine/Cargo.toml +++ b/crates/shepherd-engine/Cargo.toml @@ -15,7 +15,7 @@ path = "src/main.rs" [dependencies] nexum-launch = { git = "https://github.com/nullislabs/nexum-runtime", rev = "2ed882f21e685921cd27d1edec870e8adbb312aa" } nexum-runtime = { git = "https://github.com/nullislabs/nexum-runtime", rev = "2ed882f21e685921cd27d1edec870e8adbb312aa" } -videre-host = { git = "https://github.com/nullislabs/videre-nexum-module", rev = "fd8af027d9ae84823f8fcdcc6902b3af80cab451" } +videre-host = { git = "https://github.com/nullislabs/videre-nexum-module", rev = "9ab15152b9279974cb42c3ee3aff054344f6050d" } # The CoW venue, a native `videre_host::VenueInvoker` linked into this # binary. It was a guest wasm component the operator wired by path until @@ -43,4 +43,4 @@ nexum-runtime = { git = "https://github.com/nullislabs/nexum-runtime", rev = "2e # The versioned status-body codec the intent-status envelope carries; the # tests build a `StatusBody` and encode it into an `IntentStatusUpdate`. # Pinned to the same videre rev as `videre-host`. -videre-status-body = { git = "https://github.com/nullislabs/videre-nexum-module", rev = "fd8af027d9ae84823f8fcdcc6902b3af80cab451" } +videre-status-body = { git = "https://github.com/nullislabs/videre-nexum-module", rev = "9ab15152b9279974cb42c3ee3aff054344f6050d" } diff --git a/modules/ccow-monitor/Cargo.toml b/modules/ccow-monitor/Cargo.toml index d5c6ce76..dbbac47f 100644 --- a/modules/ccow-monitor/Cargo.toml +++ b/modules/ccow-monitor/Cargo.toml @@ -12,7 +12,7 @@ crate-type = ["cdylib"] composable-cow = { path = "../../crates/composable-cow", features = ["run"] } cow-venue = { path = "../../crates/cow-venue", features = ["client"] } nexum-sdk = { git = "https://github.com/nullislabs/nexum-runtime", rev = "2ed882f21e685921cd27d1edec870e8adbb312aa" } -videre-sdk = { git = "https://github.com/nullislabs/videre-nexum-module", rev = "fd8af027d9ae84823f8fcdcc6902b3af80cab451" } +videre-sdk = { git = "https://github.com/nullislabs/videre-nexum-module", rev = "9ab15152b9279974cb42c3ee3aff054344f6050d" } alloy-primitives = { version = "1.6", default-features = false, features = ["std"] } alloy-sol-types = { version = "1.6", default-features = false, features = ["std"] } tracing = { version = "0.1", default-features = false } diff --git a/modules/ethflow-watcher/Cargo.toml b/modules/ethflow-watcher/Cargo.toml index 50322b00..bd94a3cc 100644 --- a/modules/ethflow-watcher/Cargo.toml +++ b/modules/ethflow-watcher/Cargo.toml @@ -11,7 +11,7 @@ crate-type = ["cdylib"] [dependencies] nexum-sdk = { git = "https://github.com/nullislabs/nexum-runtime", rev = "2ed882f21e685921cd27d1edec870e8adbb312aa" } -videre-sdk = { git = "https://github.com/nullislabs/videre-nexum-module", rev = "fd8af027d9ae84823f8fcdcc6902b3af80cab451" } +videre-sdk = { git = "https://github.com/nullislabs/videre-nexum-module", rev = "9ab15152b9279974cb42c3ee3aff054344f6050d" } cow-venue = { path = "../../crates/cow-venue", features = ["client", "assembly"] } cowprotocol = { version = "0.2.0", default-features = false } alloy-primitives = { version = "1.6", default-features = false, features = ["std"] } diff --git a/tools/orderbook-mock/src/main.rs b/tools/orderbook-mock/src/main.rs index 77d7a287..6dcdb172 100644 --- a/tools/orderbook-mock/src/main.rs +++ b/tools/orderbook-mock/src/main.rs @@ -168,10 +168,14 @@ async fn post_orders(State(state): State>, body: String) -> impl I Err(why) => { state.counters.submits_err.fetch_add(1, Ordering::Relaxed); tracing::warn!("refusing an order this mock cannot derive a UID for: {why}"); + // Deliberately not a real `errorType`: this is the mock + // failing to read the request, not the orderbook refusing + // the order, and borrowing a classified type would make a + // broken fixture look like a venue policy the table has an + // opinion about. return ( StatusCode::BAD_REQUEST, axum::Json(serde_json::json!({ - "errorType": "InvalidSignature", "description": format!("mock could not derive a UID: {why}"), })), ) @@ -274,9 +278,12 @@ mod tests { } /// A body this mock cannot read is refused rather than answered with - /// something the venue would reject anyway. + /// something the venue would reject anyway, and it carries no + /// `errorType`: the mock failed to read the request, so it must not + /// look like a venue policy the classification table has an opinion + /// about. #[tokio::test] - async fn an_underivable_body_is_refused() { + async fn an_underivable_body_is_refused_without_an_error_type() { let app = router_with(default_cli()); let resp = app .oneshot( @@ -288,6 +295,14 @@ mod tests { .await .unwrap(); assert_eq!(resp.status(), StatusCode::BAD_REQUEST); + let body = axum::body::to_bytes(resp.into_body(), usize::MAX) + .await + .unwrap(); + let parsed: serde_json::Value = serde_json::from_slice(&body).unwrap(); + assert!( + parsed.get("errorType").is_none(), + "a mock-side failure must not borrow a classified errorType: {parsed}", + ); } #[tokio::test] From b9d250ac39b14bbd8b0081e7ef87239f78a90fbd Mon Sep 17 00:00:00 2001 From: mfw78 Date: Wed, 9 Sep 2026 08:38:55 +0000 Subject: [PATCH 6/7] test(cow): drive the TWAP to retirement, not just its first part The run stopped after part 0 posted, so it proved the chain and venue legs connect and nothing else. A TWAP is a series, and the parts of it that had never been exercised end to end are the ones the last few changes were about. It now advances the fork's clock past the first part's window. That is load-bearing: part 0's post carries `nextPollTimestamp` at part 1's start, so the commitment is gated until then and nothing polls in between. Without the jump the second submission never arrives. Part 1 mints its own order, with a later `validTo` and so a different uid, and submits under its own journal key: the run asserts the two keys differ rather than that a second line appeared. After the final part the generator reports no successor, which is `NextPoll::Never`, and the commitment retires. That teardown was only ever unit-tested; this is the first time it runs against the deployed registry and handler. AI Assistance: Claude Code used for the lifecycle run. --- .../shepherd-engine/tests/anvil_fork_e2e.rs | 59 ++++++++++++++++--- 1 file changed, 50 insertions(+), 9 deletions(-) diff --git a/crates/shepherd-engine/tests/anvil_fork_e2e.rs b/crates/shepherd-engine/tests/anvil_fork_e2e.rs index 9218792f..458805c5 100644 --- a/crates/shepherd-engine/tests/anvil_fork_e2e.rs +++ b/crates/shepherd-engine/tests/anvil_fork_e2e.rs @@ -445,6 +445,20 @@ impl Harness { ]); } + /// Push the fork's clock forward, then mine so a block carries the + /// new timestamp. The parts of a TWAP are spaced in wall-clock + /// seconds, so this is the only way to reach the next one. + fn advance(&self, seconds: u64) { + cast(&[ + "rpc", + "evm_increaseTime", + &seconds.to_string(), + "--rpc-url", + &self.anvil, + ]); + self.mine(4); + } + /// Mine, so a registration gets the block dispatch that polls it. fn mine(&self, blocks: usize) { for _ in 0..blocks { @@ -468,15 +482,28 @@ impl Drop for Harness { } } -/// The whole path: a real registration on the forked registry is -/// indexed, polled to `Post` through the structured generator wire, and -/// submitted to the venue. +/// The journal key out of a `submitted {intent_id} (receipt ...)` line. +fn submitted_id(line: &str) -> &str { + line.split("submitted ") + .nth(1) + .and_then(|rest| rest.split_whitespace().next()) + .unwrap_or_else(|| panic!("no intent id in {line}")) +} + +/// A TWAP from registration to retirement, against the deployed +/// contracts. +/// +/// The fixture is two parts 600 seconds apart, so the run covers a +/// commitment being indexed, each part minting its own order and +/// submitting independently, and the series ending: after the final +/// part the generator reports no successor, which retires the +/// commitment. #[test] -fn ccow_monitor_polls_a_real_owned_twap_to_post() { +fn ccow_monitor_drives_a_real_owned_twap_to_completion() { let Some(rpc) = fork_rpc_or_skip() else { return; }; - let mut h = Harness::start(&rpc, "post", POST_OWNER); + let mut h = Harness::start(&rpc, "lifecycle", POST_OWNER); h.delegate_owner(); register_order(&h.anvil, &h.owner.clone()); @@ -492,10 +519,24 @@ fn ccow_monitor_polls_a_real_owned_twap_to_post() { } assert!( polled.contains("-> Post"), - "the generator posted, so the module must too; got: {polled}", + "the generator posted part 0, so the module must too; got: {polled}", + ); + let first = h.wait_for("submitted "); + + // Part 0's post carries `nextPollTimestamp` at part 1's start, so + // the commitment is gated until then and nothing polls in between. + // Only the clock moving reaches the next part. + h.advance(700); + let second = h.wait_for("submitted "); + assert_ne!( + submitted_id(&first), + submitted_id(&second), + "each part mints its own order, so each submits under its own key", ); - // And the post reaches the venue, closing the loop through the - // borsh body, the journal reservation and the orderbook adapter. - h.wait_for("submitted"); + // After the final part the generator reports no successor, which is + // `NextPoll::Never`, and the run retires the commitment: the row, + // both due-index entries, the submission index and the journal rows + // go together. + h.wait_for("completed commitment"); } From 71275e4226b7d8f2444fe8b1dd35854a5d951c9c Mon Sep 17 00:00:00 2001 From: mfw78 Date: Wed, 9 Sep 2026 09:02:03 +0000 Subject: [PATCH 7/7] test(cow): read the store back, and run five parts not two Two parts could pass for a fluke and the run only ever asserted on log lines, so a teardown that logged but left rows behind would have gone unnoticed. It is five parts now, and the run reads the engine's own redb store back at the end: after retirement no `commitment:`, `due-`, gate, refusal, `watch-sub:`, `exp-t:`, `submitted:`, `parked:`, `root:` or `context:` key survives. The engine is asked to exit rather than killed. redb loses whatever it had not flushed, and what is left is an older consistent snapshot, so a killed engine reads as a run that stopped mid-series: the first version of this failed that way, reporting a teardown that had in fact happened. Same backend the runtime writes, and the same keccak namespace, so the read cannot drift from the write. AI Assistance: Claude Code used for the store readback and the longer run. --- Cargo.lock | 1 + crates/shepherd-engine/Cargo.toml | 4 + .../shepherd-engine/tests/anvil_fork_e2e.rs | 131 +++++++++++++++--- 3 files changed, 116 insertions(+), 20 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index df0348e0..17c4662e 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -5271,6 +5271,7 @@ dependencies = [ "cow-venue", "nexum-launch", "nexum-runtime", + "redb", "serde", "tokio", "tracing", diff --git a/crates/shepherd-engine/Cargo.toml b/crates/shepherd-engine/Cargo.toml index 9c4979b7..11e7dbc5 100644 --- a/crates/shepherd-engine/Cargo.toml +++ b/crates/shepherd-engine/Cargo.toml @@ -44,3 +44,7 @@ nexum-runtime = { git = "https://github.com/nullislabs/nexum-runtime", rev = "2e # tests build a `StatusBody` and encode it into an `IntentStatusUpdate`. # Pinned to the same videre rev as `videre-host`. videre-status-body = { git = "https://github.com/nullislabs/videre-nexum-module", rev = "9ab15152b9279974cb42c3ee3aff054344f6050d" } +# The anvil-fork run reads the engine's own local store back once it has +# stopped, to prove a retired commitment left nothing behind. Same +# backend the runtime writes, so the read cannot drift from the write. +redb.workspace = true diff --git a/crates/shepherd-engine/tests/anvil_fork_e2e.rs b/crates/shepherd-engine/tests/anvil_fork_e2e.rs index 458805c5..88402e50 100644 --- a/crates/shepherd-engine/tests/anvil_fork_e2e.rs +++ b/crates/shepherd-engine/tests/anvil_fork_e2e.rs @@ -45,6 +45,13 @@ const ACCOUNT_7702: &str = "0x15236F06922A288B68e57Ab42e397920a3F3Bb99"; /// whatever mainnet says rather than what the test set. const POST_OWNER: &str = "0x00000000000000000000000000000000CafeBabe"; +/// Parts in the registered series, and the seconds between them. +/// +/// More than a couple, so a run that posted once and stopped cannot +/// pass for one that followed the series. +const PARTS: u64 = 5; +const PART_SECONDS: u64 = 600; + const WETH: &str = "0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2"; const USDC: &str = "0xA0b86991c6218b36c1d19D4a2e9Eb0cE3606eB48"; @@ -185,7 +192,7 @@ fn register_order(anvil: &str, owner: &str) -> String { "abi-encode", "f((address,address,address,uint256,uint256,uint256,uint256,uint256,uint256,bytes32))", &format!( - "({WETH},{USDC},{owner},1000000000000000,1000000,{t0},2,600,0,\ + "({WETH},{USDC},{owner},1000000000000000,1000000,{t0},{PARTS},{PART_SECONDS},0,\ 0x0000000000000000000000000000000000000000000000000000000000000000)" ), ]); @@ -346,7 +353,7 @@ struct Harness { dir: PathBuf, _anvil: Reaped, _ob: Reaped, - _engine: Reaped, + engine: Option, } impl Harness { @@ -392,7 +399,7 @@ impl Harness { .spawn() .expect("spawn the shepherd binary"); let lines = stream_lines(&mut engine); - let _engine = Reaped(engine); + let engine_child = Reaped(engine); // A fresh address, not one of anvil's unlocked accounts. Account // zero is the well-known test key, and on mainnet it already @@ -423,7 +430,7 @@ impl Harness { dir, _anvil, _ob, - _engine, + engine: Some(engine_child), }; harness.wait_for("ccow-monitor"); harness @@ -474,6 +481,61 @@ impl Harness { fn saw(&self, needle: &str) -> bool { self.seen.iter().any(|l| l.contains(needle)) } + + /// Ask the engine to exit and wait for it. + /// + /// Killing it instead loses whatever redb had not flushed, and what + /// is left is an older consistent snapshot: the run appears to have + /// stopped mid-series, which reads as a teardown that did not + /// happen rather than as a store that was never closed. + fn stop_engine_gracefully(&mut self) { + let Some(mut engine) = self.engine.take() else { + return; + }; + let pid = engine.0.id().to_string(); + let _ = Command::new("kill").args(["-TERM", &pid]).status(); + for _ in 0..300 { + match engine.0.try_wait() { + Ok(Some(_)) => return, + Ok(None) => std::thread::sleep(Duration::from_millis(100)), + Err(_) => return, + } + } + panic!("the engine did not exit within 30s of SIGTERM"); + } + + /// Every key the module can see, once the engine has stopped. + /// + /// redb holds the file exclusively while the engine runs, so this + /// stops it first. Anything already committed survives that, which + /// is the point: the assertions are about what the engine chose to + /// leave behind, not about how it exited. + fn stop_and_read_keys(&mut self) -> Vec { + use redb::{ReadableDatabase as _, ReadableTable as _}; + + self.stop_engine_gracefully(); + let path = self.dir.join("state/local-store.redb"); + let db = + redb::Database::open(&path).unwrap_or_else(|e| panic!("open {}: {e}", path.display())); + let read = db.begin_read().expect("begin read"); + let table = read + .open_table(redb::TableDefinition::<&[u8], &[u8]>::new( + "nexum:local-store", + )) + .expect("the runtime's local-store table"); + // Keys are namespaced with `keccak256(module_id)`. + let prefix = alloy_primitives::keccak256(b"ccow-monitor"); + table + .iter() + .expect("iterate") + .filter_map(Result::ok) + .filter_map(|(k, _): (redb::AccessGuard<'_, &[u8]>, _)| { + k.value() + .strip_prefix(prefix.as_slice()) + .map(|rest| String::from_utf8_lossy(rest).into_owned()) + }) + .collect() + } } impl Drop for Harness { @@ -493,11 +555,11 @@ fn submitted_id(line: &str) -> &str { /// A TWAP from registration to retirement, against the deployed /// contracts. /// -/// The fixture is two parts 600 seconds apart, so the run covers a +/// The fixture is five parts 600 seconds apart, so the run covers a /// commitment being indexed, each part minting its own order and -/// submitting independently, and the series ending: after the final -/// part the generator reports no successor, which retires the -/// commitment. +/// submitting independently, the gate between parts holding, and the +/// series ending: after the final part the generator reports no +/// successor, which retires the commitment and everything keyed to it. #[test] fn ccow_monitor_drives_a_real_owned_twap_to_completion() { let Some(rpc) = fork_rpc_or_skip() else { @@ -521,22 +583,51 @@ fn ccow_monitor_drives_a_real_owned_twap_to_completion() { polled.contains("-> Post"), "the generator posted part 0, so the module must too; got: {polled}", ); - let first = h.wait_for("submitted "); - // Part 0's post carries `nextPollTimestamp` at part 1's start, so + // Each post carries `nextPollTimestamp` at the next part's start, so // the commitment is gated until then and nothing polls in between. - // Only the clock moving reaches the next part. - h.advance(700); - let second = h.wait_for("submitted "); - assert_ne!( - submitted_id(&first), - submitted_id(&second), - "each part mints its own order, so each submits under its own key", + // Only the clock moving reaches the next part, which is what makes + // this a series rather than one poll repeated. + let mut ids = vec![submitted_id(&h.wait_for("submitted ")).to_owned()]; + for _ in 1..PARTS { + h.advance(PART_SECONDS + 100); + ids.push(submitted_id(&h.wait_for("submitted ")).to_owned()); + } + assert_eq!(ids.len(), PARTS as usize); + let distinct: std::collections::BTreeSet<_> = ids.iter().collect(); + assert_eq!( + distinct.len(), + ids.len(), + "each part mints its own order, so each submits under its own key: {ids:?}", ); // After the final part the generator reports no successor, which is - // `NextPoll::Never`, and the run retires the commitment: the row, - // both due-index entries, the submission index and the journal rows - // go together. + // `NextPoll::Never`, and the run retires the commitment. h.wait_for("completed commitment"); + + // And retiring means nothing keyed to it survives: the row itself, + // both due-index ranges and the pointer between them, the gates and + // the refusal marker, the submission index, its expiry entries, and + // every journal row the five parts wrote. + let left = h.stop_and_read_keys(); + for prefix in [ + "commitment:", + "due-b:", + "due-t:", + "due-at:", + "next_block:", + "next_epoch:", + "refused:", + "watch-sub:", + "exp-t:", + "submitted:", + "parked:", + "root:", + "context:", + ] { + assert!( + !left.iter().any(|k| k.starts_with(prefix)), + "{prefix} survived the teardown: {left:?}", + ); + } }