diff --git a/Cargo.lock b/Cargo.lock index db64625b..17c4662e 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3993,6 +3993,7 @@ dependencies = [ "anyhow", "axum", "clap", + "cowprotocol", "rand 0.10.2", "reqwest", "serde", @@ -5270,6 +5271,7 @@ dependencies = [ "cow-venue", "nexum-launch", "nexum-runtime", + "redb", "serde", "tokio", "tracing", @@ -6035,7 +6037,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", @@ -6053,7 +6055,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", @@ -6065,7 +6067,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", @@ -6081,7 +6083,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", @@ -6090,7 +6092,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/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/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..11e7dbc5 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,8 @@ 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" } +# 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 new file mode 100644 index 00000000..88402e50 --- /dev/null +++ b/crates/shepherd-engine/tests/anvil_fork_e2e.rs @@ -0,0 +1,633 @@ +//! 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"; + +/// 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"; + +/// 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"; + +/// 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(), + ); + // 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 +} + +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},{PARTS},{PART_SECONDS},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 ") + ) + } + } + } +} + +/// 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: Option, +} + +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 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_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 + // 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: Some(engine_child), + }; + harness.wait_for("ccow-monitor"); + harness + } + + /// 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, + ]); + } + + /// 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 { + cast(&["rpc", "evm_mine", "--rpc-url", &self.anvil]); + std::thread::sleep(Duration::from_millis(500)); + } + } + + fn wait_for(&mut self, needle: &str) -> String { + wait_for(&self.lines, needle, &mut self.seen) + } + + 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 { + fn drop(&mut self) { + let _ = std::fs::remove_dir_all(&self.dir); + } +} + +/// 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 five parts 600 seconds apart, so the run covers a +/// commitment being indexed, each part minting its own order and +/// 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 { + return; + }; + let mut h = Harness::start(&rpc, "lifecycle", 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 = h.wait_for("poll commitment:"); + for bad in ["did not decode", "eth_call failed"] { + assert!(!h.saw(bad), "the poll leg reported {bad:?}"); + } + assert!( + polled.contains("-> Post"), + "the generator posted part 0, so the module must too; got: {polled}", + ); + + // 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, 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. + 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:?}", + ); + } +} 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/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 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/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..6dcdb172 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,48 @@ 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}"); + // 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!({ + "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 +221,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 +271,45 @@ 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, 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_without_an_error_type() { + 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); + 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] 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(