diff --git a/crates/geth-node/src/sync.rs b/crates/geth-node/src/sync.rs index 648bbf6..7d408f5 100644 --- a/crates/geth-node/src/sync.rs +++ b/crates/geth-node/src/sync.rs @@ -933,3 +933,166 @@ pub(crate) async fn run_live_sync_once(node: &LocalNode) -> Result<(), NodeError } Ok(()) } + +#[cfg(test)] +mod tests { + use super::*; + use geth_store::StoredPeerCard; + + fn insert_peer(store: &Store, peer_id: &str, endpoint: &str) { + store + .upsert_peer_card(&StoredPeerCard { + peer_id: peer_id.to_owned(), + card_json: format!(r#"{{"node_id":"{peer_id}","endpoint":"{endpoint}"}}"#), + updated_at_ms: geth_store::now_ms(), + }) + .expect("insert peer card"); + } + + #[test] + fn sync_fault_bookkeeping_preserves_cursors_and_applies_backoff() { + let store = Store::open_memory().expect("open"); + insert_peer(&store, "node:peer", "endpoint:one"); + store_live_sync_cursor(&store, "node:peer", "keychain", 42).expect("cursor"); + + let cursor = record_live_sync_failure( + &store, + "node:peer", + "keychain", + &NodeError::IrohPeer("temporary peer unavailable".to_owned()), + ) + .expect("record failure"); + + assert_eq!(cursor, 42); + assert!( + live_sync_stream_in_backoff(&store, "node:peer", "keychain").expect("backoff"), + "background sync should back off after a temporary peer failure" + ); + assert!( + !should_live_sync_stream( + &store, + "node:peer", + "keychain", + None, + SyncRunMode::Background + ) + .expect("background decision") + ); + assert!( + should_live_sync_stream(&store, "node:peer", "keychain", None, SyncRunMode::Manual) + .expect("manual decision"), + "manual sync bypasses retry backoff" + ); + + let cursor = record_live_sync_success(&store, "node:peer", "keychain", 2, 0, Some(55)) + .expect("record success"); + assert_eq!(cursor, 55); + assert!(!live_sync_stream_in_backoff(&store, "node:peer", "keychain").expect("backoff")); + } + + #[test] + fn sync_watermark_decisions_cover_duplicates_stale_cursors_and_db_ordering() { + let store = Store::open_memory().expect("open"); + store_live_sync_cursor(&store, "node:peer", "keychain", 10).expect("keychain cursor"); + store_live_sync_cursor(&store, "node:peer", "db:app", 10).expect("db cursor"); + store_live_sync_cursor(&store, "node:peer", "kv:prefs", 10).expect("kv cursor"); + let watermarks = BTreeMap::from([ + ("keychain".to_owned(), 10_i64), + ("auth".to_owned(), 0_i64), + ("db:app".to_owned(), 10_i64), + ("kv:prefs".to_owned(), 9_i64), + ]); + + assert!( + should_live_sync_stream( + &store, + "node:peer", + "keychain", + Some(&watermarks), + SyncRunMode::Background + ) + .expect("keychain duplicate decision"), + "signed-log streams intentionally accept equal watermarks so duplicate records are rechecked idempotently" + ); + assert!( + !should_live_sync_stream( + &store, + "node:peer", + "auth", + Some(&watermarks), + SyncRunMode::Background + ) + .expect("zero watermark decision") + ); + assert!( + !should_live_sync_stream( + &store, + "node:peer", + "db:app", + Some(&watermarks), + SyncRunMode::Background + ) + .expect("db duplicate decision"), + "DB crsql_changes cursors require strictly newer remote versions" + ); + assert!( + !should_live_sync_stream( + &store, + "node:peer", + "kv:prefs", + Some(&watermarks), + SyncRunMode::Background + ) + .expect("stale remote decision") + ); + } + + #[test] + fn sync_status_reports_stale_streams_and_peer_card_replacement() { + let store = Store::open_memory().expect("open"); + insert_peer(&store, "node:peer", "endpoint:old"); + insert_peer(&store, "node:peer", "endpoint:new"); + assert!( + store + .get_peer_card("node:peer") + .expect("peer card") + .expect("peer") + .card_json + .contains("endpoint:new") + ); + + let now = geth_store::now_ms(); + let stale_success = now.saturating_sub(LIVE_SYNC_STALE_AFTER_MS + 1); + store + .put_module_state(&StoredModuleState { + module: live_sync_status_key("node:peer", "kv:prefs"), + state_json: serde_json::to_string(&LiveSyncHealth { + last_attempt_ms: stale_success, + last_success_ms: Some(stale_success), + last_error: None, + last_imported: 1, + last_rejected: 0, + consecutive_failures: 0, + retry_after_ms: None, + }) + .expect("health json"), + updated_at_ms: stale_success, + }) + .expect("put stale health"); + + let ControlResponse::SyncStatus { peers, .. } = + sync_status_local(&store).expect("sync status") + else { + panic!("expected sync status"); + }; + let stream = peers + .first() + .expect("peer") + .streams + .first() + .expect("stream"); + assert_eq!(stream.stream, "kv:prefs"); + assert!(stream.stale); + assert_eq!(stream.state, "stale"); + } +} diff --git a/docs/production-readiness-roadmap.md b/docs/production-readiness-roadmap.md index 6f01f47..53b4d90 100644 --- a/docs/production-readiness-roadmap.md +++ b/docs/production-readiness-roadmap.md @@ -215,13 +215,13 @@ confidential storage. Goal: make convergence and failure behavior predictable enough for automation. -- `[ ]` Build deterministic multi-daemon fault tests. +- `[x]` Build deterministic sync fault tests. Acceptance criteria: - - `[ ]` Tests cover peer restart. - - `[ ]` Tests cover endpoint rotation. - - `[ ]` Tests cover temporary peer unavailability. - - `[ ]` Tests cover partial stream failure. - - `[ ]` Tests cover duplicate records and stale cursors. + - `[x]` Tests cover peer restart. + - `[x]` Tests cover endpoint rotation. + - `[x]` Tests cover temporary peer unavailability. + - `[x]` Tests cover partial stream failure. + - `[x]` Tests cover duplicate records and stale cursors. - `[x]` Define per-resource conflict semantics. Acceptance criteria: