From 2b604dd184ee01ba11bbef0d2b14a60aab73368d Mon Sep 17 00:00:00 2001 From: Eric Wendland Date: Sun, 5 Jul 2026 17:49:09 +0200 Subject: [PATCH] refactor: move live sync orchestration --- crates/geth-node/src/daemon.rs | 6 +- crates/geth-node/src/lib.rs | 596 +------------------------- crates/geth-node/src/sync.rs | 599 ++++++++++++++++++++++++++- docs/production-readiness-roadmap.md | 4 +- 4 files changed, 605 insertions(+), 600 deletions(-) diff --git a/crates/geth-node/src/daemon.rs b/crates/geth-node/src/daemon.rs index 129add0..a084a3c 100644 --- a/crates/geth-node/src/daemon.rs +++ b/crates/geth-node/src/daemon.rs @@ -1,6 +1,6 @@ //! Daemon task orchestration helpers. -use crate::LocalNode; +use crate::{LocalNode, sync::run_live_sync_once}; use geth_iroh::GethIrohEndpoint; use std::time::Duration; @@ -23,13 +23,13 @@ pub(crate) fn spawn_iroh_control_accept_loop(node: LocalNode, endpoint: GethIroh pub(crate) fn spawn_background_live_sync(node: LocalNode, interval_duration: Duration) { tokio::spawn(async move { - if let Err(error) = super::run_live_sync_once(&node).await { + if let Err(error) = run_live_sync_once(&node).await { tracing::debug!(%error, "initial live sync tick failed"); } let mut interval = tokio::time::interval(interval_duration); loop { interval.tick().await; - if let Err(error) = super::run_live_sync_once(&node).await { + if let Err(error) = run_live_sync_once(&node).await { tracing::debug!(%error, "live sync tick failed"); } } diff --git a/crates/geth-node/src/lib.rs b/crates/geth-node/src/lib.rs index 10c6e6b..4fd7215 100644 --- a/crates/geth-node/src/lib.rs +++ b/crates/geth-node/src/lib.rs @@ -18,7 +18,7 @@ use geth_control::{ CasBlob, CasProvider, ControlRequest, ControlResponse, KeychainStatusResponse, NativeBackendStatus, NodeIdResponse, OverlayPeer, OverlayWireRequest, OverlayWireResponse, PeerControlRequest, PeerControlResponse, PipeWireRequest, PipeWireResponse, StatusResponse, - SyncPeerRun, SyncStreamRun, SyncWatermark, + SyncWatermark, }; use geth_crypto::AgentKey; use geth_db::DbResource; @@ -70,8 +70,7 @@ use std::path::{Component, Path, PathBuf}; use std::sync::{Arc, Mutex}; use std::time::Duration; use sync::{ - KvDocsState, load_live_sync_cursor, record_live_sync_failure, record_live_sync_success, - should_live_sync_stream, store_live_sync_cursor, sync_status_local, + KvDocsState, load_live_sync_cursor, store_live_sync_cursor, sync_now, sync_status_local, }; use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; use tokio::net::{TcpListener, TcpStream, UnixListener, UnixStream}; @@ -4075,596 +4074,6 @@ fn explain_peer_or_bearer( }) } -async fn run_sync_for_peer(node: &LocalNode, peer_node: &str) -> Result { - let kv_stores = Store::open(&node.paths.metadata_db())?.list_kv_stores()?; - let documents = Store::open(&node.paths.metadata_db())?.list_document_resources()?; - let dbs = Store::open(&node.paths.metadata_db())?.list_db_resources()?; - let mut streams = Vec::new(); - let remote_watermarks = match sync_status_from_peer(node, peer_node) - .await - .map(|watermarks| { - watermarks - .into_iter() - .map(|watermark| (watermark.stream, watermark.high_water)) - .collect::>() - }) { - Ok(watermarks) => { - let store = Store::open(&node.paths.metadata_db())?; - let cursor_ms = record_live_sync_success(&store, peer_node, "sync-status", 0, 0, None)?; - streams.push(SyncStreamRun { - stream: "sync-status".to_owned(), - attempted: true, - success: true, - imported: 0, - rejected: 0, - cursor_ms, - error: None, - }); - Some(watermarks) - } - Err(error) => { - let store = Store::open(&node.paths.metadata_db())?; - let cursor_ms = record_live_sync_failure(&store, peer_node, "sync-status", &error)?; - streams.push(SyncStreamRun { - stream: "sync-status".to_owned(), - attempted: true, - success: false, - imported: 0, - rejected: 0, - cursor_ms, - error: Some(error.to_string()), - }); - None - } - }; - - sync_keychain_stream(node, peer_node, remote_watermarks.as_ref(), &mut streams).await?; - sync_auth_stream(node, peer_node, remote_watermarks.as_ref(), &mut streams).await?; - sync_ssh_cert_stream(node, peer_node, remote_watermarks.as_ref(), &mut streams).await?; - sync_ssh_revocation_stream(node, peer_node, remote_watermarks.as_ref(), &mut streams).await?; - - if let Some(remote_watermarks) = remote_watermarks.as_ref() { - for stream in remote_watermarks - .keys() - .filter(|stream| stream.starts_with("cas-tree:")) - { - sync_cas_tree_stream( - node, - peer_node, - stream, - Some(remote_watermarks), - &mut streams, - ) - .await?; - } - } - for kv in &kv_stores { - let stream = format!("kv:{}", kv.name); - sync_kv_stream( - node, - peer_node, - &kv.name, - &stream, - remote_watermarks.as_ref(), - &mut streams, - ) - .await?; - } - for document in &documents { - let stream = format!("document:{}", document.name); - sync_document_stream( - node, - peer_node, - &document.name, - &stream, - remote_watermarks.as_ref(), - &mut streams, - ) - .await?; - } - for db in &dbs { - let stream = format!("db:{}", db.name); - sync_db_stream( - node, - peer_node, - &db.name, - &stream, - remote_watermarks.as_ref(), - &mut streams, - ) - .await?; - } - Ok(SyncPeerRun { - peer_node_id: peer_node.to_owned(), - streams, - }) -} - -async fn sync_now( - node: &LocalNode, - requested_peer: Option, -) -> Result { - if node - .iroh_endpoint - .lock() - .map_err(|_| NodeError::RuntimeLockPoisoned)? - .is_none() - { - return Err(NodeError::IrohEndpointUnavailable); - } - let peers = if let Some(peer) = requested_peer { - vec![peer] - } else { - Store::open(&node.paths.metadata_db())? - .list_peer_cards()? - .into_iter() - .map(|peer| peer.peer_id) - .collect() - }; - let mut runs = Vec::new(); - for peer in peers { - runs.push(run_sync_for_peer(node, &peer).await?); - } - Ok(ControlResponse::SyncRan { - peers: runs, - note: "ran best-effort sync over Iroh; each imported keychain/auth operation was still verified from signed operation logs".to_owned(), - }) -} - -fn push_skipped_sync_stream( - store: &Store, - peer_node: &str, - stream: &str, - streams: &mut Vec, -) -> Result<(), NodeError> { - streams.push(SyncStreamRun { - stream: stream.to_owned(), - attempted: false, - success: true, - imported: 0, - rejected: 0, - cursor_ms: load_live_sync_cursor(store, peer_node, stream)?, - error: None, - }); - Ok(()) -} - -fn push_failed_sync_stream( - store: &Store, - peer_node: &str, - stream: &str, - error: NodeError, - streams: &mut Vec, -) -> Result<(), NodeError> { - let cursor_ms = record_live_sync_failure(store, peer_node, stream, &error)?; - streams.push(SyncStreamRun { - stream: stream.to_owned(), - attempted: true, - success: false, - imported: 0, - rejected: 0, - cursor_ms, - error: Some(error.to_string()), - }); - Ok(()) -} - -async fn sync_keychain_stream( - node: &LocalNode, - peer_node: &str, - remote_watermarks: Option<&BTreeMap>, - streams: &mut Vec, -) -> Result<(), NodeError> { - let stream = "keychain"; - let store = Store::open(&node.paths.metadata_db())?; - if !should_live_sync_stream(&store, peer_node, stream, remote_watermarks)? { - return push_skipped_sync_stream(&store, peer_node, stream, streams); - } - let since_ms = load_live_sync_cursor(&store, peer_node, stream)?; - match keychain_sync_from_peer_since(node, peer_node, since_ms).await { - Ok(ControlResponse::KeychainSynced { - ops_imported, - signatures_imported, - invalid_ops_rejected, - high_water_ms, - .. - }) => { - let store = Store::open(&node.paths.metadata_db())?; - let cursor_ms = record_live_sync_success( - &store, - peer_node, - stream, - ops_imported + signatures_imported, - invalid_ops_rejected, - Some(high_water_ms), - )?; - streams.push(SyncStreamRun { - stream: stream.to_owned(), - attempted: true, - success: true, - imported: ops_imported + signatures_imported, - rejected: invalid_ops_rejected, - cursor_ms, - error: None, - }); - Ok(()) - } - Ok(other) => push_failed_sync_stream( - &store, - peer_node, - stream, - NodeError::IrohPeer(format!("unexpected keychain sync response: {other:?}")), - streams, - ), - Err(error) => push_failed_sync_stream(&store, peer_node, stream, error, streams), - } -} - -async fn sync_auth_stream( - node: &LocalNode, - peer_node: &str, - remote_watermarks: Option<&BTreeMap>, - streams: &mut Vec, -) -> Result<(), NodeError> { - let stream = "auth"; - let store = Store::open(&node.paths.metadata_db())?; - if !should_live_sync_stream(&store, peer_node, stream, remote_watermarks)? { - return push_skipped_sync_stream(&store, peer_node, stream, streams); - } - let since_ms = load_live_sync_cursor(&store, peer_node, stream)?; - match auth_sync_from_peer_since(node, peer_node, since_ms).await { - Ok(ControlResponse::AuthSynced { - ops_imported, - signatures_imported, - invalid_ops_rejected, - high_water_ms, - .. - }) => { - let store = Store::open(&node.paths.metadata_db())?; - let cursor_ms = record_live_sync_success( - &store, - peer_node, - stream, - ops_imported + signatures_imported, - invalid_ops_rejected, - Some(high_water_ms), - )?; - streams.push(SyncStreamRun { - stream: stream.to_owned(), - attempted: true, - success: true, - imported: ops_imported + signatures_imported, - rejected: invalid_ops_rejected, - cursor_ms, - error: None, - }); - Ok(()) - } - Ok(other) => push_failed_sync_stream( - &store, - peer_node, - stream, - NodeError::IrohPeer(format!("unexpected auth sync response: {other:?}")), - streams, - ), - Err(error) => push_failed_sync_stream(&store, peer_node, stream, error, streams), - } -} - -async fn sync_ssh_cert_stream( - node: &LocalNode, - peer_node: &str, - remote_watermarks: Option<&BTreeMap>, - streams: &mut Vec, -) -> Result<(), NodeError> { - let stream = "ssh-certs"; - let store = Store::open(&node.paths.metadata_db())?; - if !should_live_sync_stream(&store, peer_node, stream, remote_watermarks)? { - return push_skipped_sync_stream(&store, peer_node, stream, streams); - } - match ssh_cert_sync_from_peer(node, peer_node, None).await { - Ok(ControlResponse::SshCertSynced { - requests_imported, - certificates_imported, - .. - }) => { - let store = Store::open(&node.paths.metadata_db())?; - let cursor_ms = record_live_sync_success( - &store, - peer_node, - stream, - requests_imported + certificates_imported, - 0, - None, - )?; - streams.push(SyncStreamRun { - stream: stream.to_owned(), - attempted: true, - success: true, - imported: requests_imported + certificates_imported, - rejected: 0, - cursor_ms, - error: None, - }); - Ok(()) - } - Ok(other) => push_failed_sync_stream( - &store, - peer_node, - stream, - NodeError::IrohPeer(format!("unexpected SSH cert sync response: {other:?}")), - streams, - ), - Err(error) => push_failed_sync_stream(&store, peer_node, stream, error, streams), - } -} - -async fn sync_ssh_revocation_stream( - node: &LocalNode, - peer_node: &str, - remote_watermarks: Option<&BTreeMap>, - streams: &mut Vec, -) -> Result<(), NodeError> { - let stream = "ssh-revocations"; - let store = Store::open(&node.paths.metadata_db())?; - if !should_live_sync_stream(&store, peer_node, stream, remote_watermarks)? { - return push_skipped_sync_stream(&store, peer_node, stream, streams); - } - match ssh_revocation_sync_from_peer(node, peer_node, None).await { - Ok(ControlResponse::SshRevocationSynced { - revocations_imported, - .. - }) => { - let store = Store::open(&node.paths.metadata_db())?; - let cursor_ms = - record_live_sync_success(&store, peer_node, stream, revocations_imported, 0, None)?; - streams.push(SyncStreamRun { - stream: stream.to_owned(), - attempted: true, - success: true, - imported: revocations_imported, - rejected: 0, - cursor_ms, - error: None, - }); - Ok(()) - } - Ok(other) => push_failed_sync_stream( - &store, - peer_node, - stream, - NodeError::IrohPeer(format!( - "unexpected SSH revocation sync response: {other:?}" - )), - streams, - ), - Err(error) => push_failed_sync_stream(&store, peer_node, stream, error, streams), - } -} - -async fn sync_cas_tree_stream( - node: &LocalNode, - peer_node: &str, - stream: &str, - remote_watermarks: Option<&BTreeMap>, - streams: &mut Vec, -) -> Result<(), NodeError> { - let store = Store::open(&node.paths.metadata_db())?; - if !should_live_sync_stream(&store, peer_node, stream, remote_watermarks)? { - return push_skipped_sync_stream(&store, peer_node, stream, streams); - } - let name = stream.trim_start_matches("cas-tree:"); - match cas_root_sync_from_peer(node, peer_node, name, None).await { - Ok(ControlResponse::CasRootSynced { - tree_bytes_imported, - sync_conflicts, - .. - }) => { - let store = Store::open(&node.paths.metadata_db())?; - let cursor = load_live_sync_cursor(&store, peer_node, stream)?; - let cursor_ms = record_live_sync_success( - &store, - peer_node, - stream, - usize::from(tree_bytes_imported), - sync_conflicts.len(), - Some(cursor), - )?; - streams.push(SyncStreamRun { - stream: stream.to_owned(), - attempted: true, - success: true, - imported: usize::from(tree_bytes_imported), - rejected: sync_conflicts.len(), - cursor_ms, - error: None, - }); - Ok(()) - } - Ok(other) => push_failed_sync_stream( - &store, - peer_node, - stream, - NodeError::IrohPeer(format!("unexpected file-root sync response: {other:?}")), - streams, - ), - Err(error) => push_failed_sync_stream(&store, peer_node, stream, error, streams), - } -} - -async fn sync_kv_stream( - node: &LocalNode, - peer_node: &str, - name: &str, - stream: &str, - remote_watermarks: Option<&BTreeMap>, - streams: &mut Vec, -) -> Result<(), NodeError> { - let store = Store::open(&node.paths.metadata_db())?; - if !should_live_sync_stream(&store, peer_node, stream, remote_watermarks)? { - return push_skipped_sync_stream(&store, peer_node, stream, streams); - } - match kv_sync_from_peer(node, peer_node, name, None).await { - Ok(ControlResponse::KvSynced { - entries_imported, .. - }) => { - let store = Store::open(&node.paths.metadata_db())?; - let cursor = load_live_sync_cursor(&store, peer_node, stream)?; - let cursor_ms = record_live_sync_success( - &store, - peer_node, - stream, - entries_imported, - 0, - Some(cursor), - )?; - streams.push(SyncStreamRun { - stream: stream.to_owned(), - attempted: true, - success: true, - imported: entries_imported, - rejected: 0, - cursor_ms, - error: None, - }); - Ok(()) - } - Ok(other) => push_failed_sync_stream( - &store, - peer_node, - stream, - NodeError::IrohPeer(format!("unexpected KV sync response: {other:?}")), - streams, - ), - Err(error) => push_failed_sync_stream(&store, peer_node, stream, error, streams), - } -} - -async fn sync_document_stream( - node: &LocalNode, - peer_node: &str, - name: &str, - stream: &str, - remote_watermarks: Option<&BTreeMap>, - streams: &mut Vec, -) -> Result<(), NodeError> { - let store = Store::open(&node.paths.metadata_db())?; - if !should_live_sync_stream(&store, peer_node, stream, remote_watermarks)? { - return push_skipped_sync_stream(&store, peer_node, stream, streams); - } - match document_sync_from_peer(node, peer_node, name, None).await { - Ok(ControlResponse::DocumentSynced { updated, .. }) => { - let store = Store::open(&node.paths.metadata_db())?; - let cursor = load_live_sync_cursor(&store, peer_node, stream)?; - let cursor_ms = record_live_sync_success( - &store, - peer_node, - stream, - usize::from(updated), - 0, - Some(cursor), - )?; - streams.push(SyncStreamRun { - stream: stream.to_owned(), - attempted: true, - success: true, - imported: usize::from(updated), - rejected: 0, - cursor_ms, - error: None, - }); - Ok(()) - } - Ok(other) => push_failed_sync_stream( - &store, - peer_node, - stream, - NodeError::IrohPeer(format!("unexpected document sync response: {other:?}")), - streams, - ), - Err(error) => push_failed_sync_stream(&store, peer_node, stream, error, streams), - } -} - -async fn sync_db_stream( - node: &LocalNode, - peer_node: &str, - name: &str, - stream: &str, - remote_watermarks: Option<&BTreeMap>, - streams: &mut Vec, -) -> Result<(), NodeError> { - let store = Store::open(&node.paths.metadata_db())?; - if !should_live_sync_stream(&store, peer_node, stream, remote_watermarks)? { - return push_skipped_sync_stream(&store, peer_node, stream, streams); - } - match db_sync_from_peer(node, peer_node, name, 100, None).await { - Ok(ControlResponse::DbSynced { - changes_applied, - max_db_version, - .. - }) => { - let store = Store::open(&node.paths.metadata_db())?; - let cursor = load_live_sync_cursor(&store, peer_node, stream)?; - let cursor_ms = record_live_sync_success( - &store, - peer_node, - stream, - changes_applied, - 0, - Some(max_db_version.unwrap_or(cursor)), - )?; - streams.push(SyncStreamRun { - stream: stream.to_owned(), - attempted: true, - success: true, - imported: changes_applied, - rejected: 0, - cursor_ms, - error: None, - }); - Ok(()) - } - Ok(other) => push_failed_sync_stream( - &store, - peer_node, - stream, - NodeError::IrohPeer(format!("unexpected DB sync response: {other:?}")), - streams, - ), - Err(error) => push_failed_sync_stream(&store, peer_node, stream, error, streams), - } -} - -async fn run_live_sync_once(node: &LocalNode) -> Result<(), NodeError> { - if node - .iroh_endpoint - .lock() - .map_err(|_| NodeError::RuntimeLockPoisoned)? - .is_none() - { - return Ok(()); - } - let peers = Store::open(&node.paths.metadata_db())?.list_peer_cards()?; - for peer in peers { - let run = run_sync_for_peer(node, &peer.peer_id).await?; - for stream in &run.streams { - tracing::debug!( - peer = %peer.peer_id, - stream = %stream.stream, - attempted = stream.attempted, - success = stream.success, - imported = stream.imported, - rejected = stream.rejected, - cursor_ms = stream.cursor_ms, - error = ?stream.error, - "live sync stream result" - ); - } - } - Ok(()) -} - async fn handle_iroh_control_connection( node: LocalNode, incoming: iroh::endpoint::Incoming, @@ -11953,6 +11362,7 @@ fn ssh_revocation_from_stored( #[cfg(test)] mod tests { use super::*; + use crate::sync::{record_live_sync_success, run_live_sync_once}; use tokio::io::AsyncReadExt; fn skip_iroh_integration_tests() -> bool { diff --git a/crates/geth-node/src/sync.rs b/crates/geth-node/src/sync.rs index 66c1643..82710fe 100644 --- a/crates/geth-node/src/sync.rs +++ b/crates/geth-node/src/sync.rs @@ -5,8 +5,8 @@ //! module handlers; this layer records whether streams are healthy and whether //! remote watermarks warrant another pull. -use crate::NodeError; -use geth_control::{ControlResponse, SyncPeerStatus, SyncStreamStatus}; +use crate::{LocalNode, NodeError}; +use geth_control::{ControlResponse, SyncPeerRun, SyncPeerStatus, SyncStreamRun, SyncStreamStatus}; use geth_store::{Store, StoredModuleState}; use std::collections::BTreeMap; @@ -228,3 +228,598 @@ pub(crate) fn should_live_sync_stream( Ok(*remote_high_water >= local_cursor) } } + +pub(crate) async fn run_sync_for_peer( + node: &LocalNode, + peer_node: &str, +) -> Result { + let kv_stores = Store::open(&node.paths.metadata_db())?.list_kv_stores()?; + let documents = Store::open(&node.paths.metadata_db())?.list_document_resources()?; + let dbs = Store::open(&node.paths.metadata_db())?.list_db_resources()?; + let mut streams = Vec::new(); + let remote_watermarks = + match super::sync_status_from_peer(node, peer_node) + .await + .map(|watermarks| { + watermarks + .into_iter() + .map(|watermark| (watermark.stream, watermark.high_water)) + .collect::>() + }) { + Ok(watermarks) => { + let store = Store::open(&node.paths.metadata_db())?; + let cursor_ms = + record_live_sync_success(&store, peer_node, "sync-status", 0, 0, None)?; + streams.push(SyncStreamRun { + stream: "sync-status".to_owned(), + attempted: true, + success: true, + imported: 0, + rejected: 0, + cursor_ms, + error: None, + }); + Some(watermarks) + } + Err(error) => { + let store = Store::open(&node.paths.metadata_db())?; + let cursor_ms = record_live_sync_failure(&store, peer_node, "sync-status", &error)?; + streams.push(SyncStreamRun { + stream: "sync-status".to_owned(), + attempted: true, + success: false, + imported: 0, + rejected: 0, + cursor_ms, + error: Some(error.to_string()), + }); + None + } + }; + + sync_keychain_stream(node, peer_node, remote_watermarks.as_ref(), &mut streams).await?; + sync_auth_stream(node, peer_node, remote_watermarks.as_ref(), &mut streams).await?; + sync_ssh_cert_stream(node, peer_node, remote_watermarks.as_ref(), &mut streams).await?; + sync_ssh_revocation_stream(node, peer_node, remote_watermarks.as_ref(), &mut streams).await?; + + if let Some(remote_watermarks) = remote_watermarks.as_ref() { + for stream in remote_watermarks + .keys() + .filter(|stream| stream.starts_with("cas-tree:")) + { + sync_cas_tree_stream( + node, + peer_node, + stream, + Some(remote_watermarks), + &mut streams, + ) + .await?; + } + } + for kv in &kv_stores { + let stream = format!("kv:{}", kv.name); + sync_kv_stream( + node, + peer_node, + &kv.name, + &stream, + remote_watermarks.as_ref(), + &mut streams, + ) + .await?; + } + for document in &documents { + let stream = format!("document:{}", document.name); + sync_document_stream( + node, + peer_node, + &document.name, + &stream, + remote_watermarks.as_ref(), + &mut streams, + ) + .await?; + } + for db in &dbs { + let stream = format!("db:{}", db.name); + sync_db_stream( + node, + peer_node, + &db.name, + &stream, + remote_watermarks.as_ref(), + &mut streams, + ) + .await?; + } + Ok(SyncPeerRun { + peer_node_id: peer_node.to_owned(), + streams, + }) +} + +pub(crate) async fn sync_now( + node: &LocalNode, + requested_peer: Option, +) -> Result { + if node + .iroh_endpoint + .lock() + .map_err(|_| NodeError::RuntimeLockPoisoned)? + .is_none() + { + return Err(NodeError::IrohEndpointUnavailable); + } + let peers = if let Some(peer) = requested_peer { + vec![peer] + } else { + Store::open(&node.paths.metadata_db())? + .list_peer_cards()? + .into_iter() + .map(|peer| peer.peer_id) + .collect() + }; + let mut runs = Vec::new(); + for peer in peers { + runs.push(run_sync_for_peer(node, &peer).await?); + } + Ok(ControlResponse::SyncRan { + peers: runs, + note: "ran best-effort sync over Iroh; each imported keychain/auth operation was still verified from signed operation logs".to_owned(), + }) +} + +fn push_skipped_sync_stream( + store: &Store, + peer_node: &str, + stream: &str, + streams: &mut Vec, +) -> Result<(), NodeError> { + streams.push(SyncStreamRun { + stream: stream.to_owned(), + attempted: false, + success: true, + imported: 0, + rejected: 0, + cursor_ms: load_live_sync_cursor(store, peer_node, stream)?, + error: None, + }); + Ok(()) +} + +fn push_failed_sync_stream( + store: &Store, + peer_node: &str, + stream: &str, + error: NodeError, + streams: &mut Vec, +) -> Result<(), NodeError> { + let cursor_ms = record_live_sync_failure(store, peer_node, stream, &error)?; + streams.push(SyncStreamRun { + stream: stream.to_owned(), + attempted: true, + success: false, + imported: 0, + rejected: 0, + cursor_ms, + error: Some(error.to_string()), + }); + Ok(()) +} + +async fn sync_keychain_stream( + node: &LocalNode, + peer_node: &str, + remote_watermarks: Option<&BTreeMap>, + streams: &mut Vec, +) -> Result<(), NodeError> { + let stream = "keychain"; + let store = Store::open(&node.paths.metadata_db())?; + if !should_live_sync_stream(&store, peer_node, stream, remote_watermarks)? { + return push_skipped_sync_stream(&store, peer_node, stream, streams); + } + let since_ms = load_live_sync_cursor(&store, peer_node, stream)?; + match super::keychain_sync_from_peer_since(node, peer_node, since_ms).await { + Ok(ControlResponse::KeychainSynced { + ops_imported, + signatures_imported, + invalid_ops_rejected, + high_water_ms, + .. + }) => { + let store = Store::open(&node.paths.metadata_db())?; + let cursor_ms = record_live_sync_success( + &store, + peer_node, + stream, + ops_imported + signatures_imported, + invalid_ops_rejected, + Some(high_water_ms), + )?; + streams.push(SyncStreamRun { + stream: stream.to_owned(), + attempted: true, + success: true, + imported: ops_imported + signatures_imported, + rejected: invalid_ops_rejected, + cursor_ms, + error: None, + }); + Ok(()) + } + Ok(other) => push_failed_sync_stream( + &store, + peer_node, + stream, + NodeError::IrohPeer(format!("unexpected keychain sync response: {other:?}")), + streams, + ), + Err(error) => push_failed_sync_stream(&store, peer_node, stream, error, streams), + } +} + +async fn sync_auth_stream( + node: &LocalNode, + peer_node: &str, + remote_watermarks: Option<&BTreeMap>, + streams: &mut Vec, +) -> Result<(), NodeError> { + let stream = "auth"; + let store = Store::open(&node.paths.metadata_db())?; + if !should_live_sync_stream(&store, peer_node, stream, remote_watermarks)? { + return push_skipped_sync_stream(&store, peer_node, stream, streams); + } + let since_ms = load_live_sync_cursor(&store, peer_node, stream)?; + match super::auth_sync_from_peer_since(node, peer_node, since_ms).await { + Ok(ControlResponse::AuthSynced { + ops_imported, + signatures_imported, + invalid_ops_rejected, + high_water_ms, + .. + }) => { + let store = Store::open(&node.paths.metadata_db())?; + let cursor_ms = record_live_sync_success( + &store, + peer_node, + stream, + ops_imported + signatures_imported, + invalid_ops_rejected, + Some(high_water_ms), + )?; + streams.push(SyncStreamRun { + stream: stream.to_owned(), + attempted: true, + success: true, + imported: ops_imported + signatures_imported, + rejected: invalid_ops_rejected, + cursor_ms, + error: None, + }); + Ok(()) + } + Ok(other) => push_failed_sync_stream( + &store, + peer_node, + stream, + NodeError::IrohPeer(format!("unexpected auth sync response: {other:?}")), + streams, + ), + Err(error) => push_failed_sync_stream(&store, peer_node, stream, error, streams), + } +} + +async fn sync_ssh_cert_stream( + node: &LocalNode, + peer_node: &str, + remote_watermarks: Option<&BTreeMap>, + streams: &mut Vec, +) -> Result<(), NodeError> { + let stream = "ssh-certs"; + let store = Store::open(&node.paths.metadata_db())?; + if !should_live_sync_stream(&store, peer_node, stream, remote_watermarks)? { + return push_skipped_sync_stream(&store, peer_node, stream, streams); + } + match super::ssh_cert_sync_from_peer(node, peer_node, None).await { + Ok(ControlResponse::SshCertSynced { + requests_imported, + certificates_imported, + .. + }) => { + let store = Store::open(&node.paths.metadata_db())?; + let cursor_ms = record_live_sync_success( + &store, + peer_node, + stream, + requests_imported + certificates_imported, + 0, + None, + )?; + streams.push(SyncStreamRun { + stream: stream.to_owned(), + attempted: true, + success: true, + imported: requests_imported + certificates_imported, + rejected: 0, + cursor_ms, + error: None, + }); + Ok(()) + } + Ok(other) => push_failed_sync_stream( + &store, + peer_node, + stream, + NodeError::IrohPeer(format!("unexpected SSH cert sync response: {other:?}")), + streams, + ), + Err(error) => push_failed_sync_stream(&store, peer_node, stream, error, streams), + } +} + +async fn sync_ssh_revocation_stream( + node: &LocalNode, + peer_node: &str, + remote_watermarks: Option<&BTreeMap>, + streams: &mut Vec, +) -> Result<(), NodeError> { + let stream = "ssh-revocations"; + let store = Store::open(&node.paths.metadata_db())?; + if !should_live_sync_stream(&store, peer_node, stream, remote_watermarks)? { + return push_skipped_sync_stream(&store, peer_node, stream, streams); + } + match super::ssh_revocation_sync_from_peer(node, peer_node, None).await { + Ok(ControlResponse::SshRevocationSynced { + revocations_imported, + .. + }) => { + let store = Store::open(&node.paths.metadata_db())?; + let cursor_ms = + record_live_sync_success(&store, peer_node, stream, revocations_imported, 0, None)?; + streams.push(SyncStreamRun { + stream: stream.to_owned(), + attempted: true, + success: true, + imported: revocations_imported, + rejected: 0, + cursor_ms, + error: None, + }); + Ok(()) + } + Ok(other) => push_failed_sync_stream( + &store, + peer_node, + stream, + NodeError::IrohPeer(format!( + "unexpected SSH revocation sync response: {other:?}" + )), + streams, + ), + Err(error) => push_failed_sync_stream(&store, peer_node, stream, error, streams), + } +} + +async fn sync_cas_tree_stream( + node: &LocalNode, + peer_node: &str, + stream: &str, + remote_watermarks: Option<&BTreeMap>, + streams: &mut Vec, +) -> Result<(), NodeError> { + let store = Store::open(&node.paths.metadata_db())?; + if !should_live_sync_stream(&store, peer_node, stream, remote_watermarks)? { + return push_skipped_sync_stream(&store, peer_node, stream, streams); + } + let name = stream.trim_start_matches("cas-tree:"); + match super::cas_root_sync_from_peer(node, peer_node, name, None).await { + Ok(ControlResponse::CasRootSynced { + tree_bytes_imported, + sync_conflicts, + .. + }) => { + let store = Store::open(&node.paths.metadata_db())?; + let cursor = load_live_sync_cursor(&store, peer_node, stream)?; + let cursor_ms = record_live_sync_success( + &store, + peer_node, + stream, + usize::from(tree_bytes_imported), + sync_conflicts.len(), + Some(cursor), + )?; + streams.push(SyncStreamRun { + stream: stream.to_owned(), + attempted: true, + success: true, + imported: usize::from(tree_bytes_imported), + rejected: sync_conflicts.len(), + cursor_ms, + error: None, + }); + Ok(()) + } + Ok(other) => push_failed_sync_stream( + &store, + peer_node, + stream, + NodeError::IrohPeer(format!("unexpected file-root sync response: {other:?}")), + streams, + ), + Err(error) => push_failed_sync_stream(&store, peer_node, stream, error, streams), + } +} + +async fn sync_kv_stream( + node: &LocalNode, + peer_node: &str, + name: &str, + stream: &str, + remote_watermarks: Option<&BTreeMap>, + streams: &mut Vec, +) -> Result<(), NodeError> { + let store = Store::open(&node.paths.metadata_db())?; + if !should_live_sync_stream(&store, peer_node, stream, remote_watermarks)? { + return push_skipped_sync_stream(&store, peer_node, stream, streams); + } + match super::kv_sync_from_peer(node, peer_node, name, None).await { + Ok(ControlResponse::KvSynced { + entries_imported, .. + }) => { + let store = Store::open(&node.paths.metadata_db())?; + let cursor = load_live_sync_cursor(&store, peer_node, stream)?; + let cursor_ms = record_live_sync_success( + &store, + peer_node, + stream, + entries_imported, + 0, + Some(cursor), + )?; + streams.push(SyncStreamRun { + stream: stream.to_owned(), + attempted: true, + success: true, + imported: entries_imported, + rejected: 0, + cursor_ms, + error: None, + }); + Ok(()) + } + Ok(other) => push_failed_sync_stream( + &store, + peer_node, + stream, + NodeError::IrohPeer(format!("unexpected KV sync response: {other:?}")), + streams, + ), + Err(error) => push_failed_sync_stream(&store, peer_node, stream, error, streams), + } +} + +async fn sync_document_stream( + node: &LocalNode, + peer_node: &str, + name: &str, + stream: &str, + remote_watermarks: Option<&BTreeMap>, + streams: &mut Vec, +) -> Result<(), NodeError> { + let store = Store::open(&node.paths.metadata_db())?; + if !should_live_sync_stream(&store, peer_node, stream, remote_watermarks)? { + return push_skipped_sync_stream(&store, peer_node, stream, streams); + } + match super::document_sync_from_peer(node, peer_node, name, None).await { + Ok(ControlResponse::DocumentSynced { updated, .. }) => { + let store = Store::open(&node.paths.metadata_db())?; + let cursor = load_live_sync_cursor(&store, peer_node, stream)?; + let cursor_ms = record_live_sync_success( + &store, + peer_node, + stream, + usize::from(updated), + 0, + Some(cursor), + )?; + streams.push(SyncStreamRun { + stream: stream.to_owned(), + attempted: true, + success: true, + imported: usize::from(updated), + rejected: 0, + cursor_ms, + error: None, + }); + Ok(()) + } + Ok(other) => push_failed_sync_stream( + &store, + peer_node, + stream, + NodeError::IrohPeer(format!("unexpected document sync response: {other:?}")), + streams, + ), + Err(error) => push_failed_sync_stream(&store, peer_node, stream, error, streams), + } +} + +async fn sync_db_stream( + node: &LocalNode, + peer_node: &str, + name: &str, + stream: &str, + remote_watermarks: Option<&BTreeMap>, + streams: &mut Vec, +) -> Result<(), NodeError> { + let store = Store::open(&node.paths.metadata_db())?; + if !should_live_sync_stream(&store, peer_node, stream, remote_watermarks)? { + return push_skipped_sync_stream(&store, peer_node, stream, streams); + } + match super::db_sync_from_peer(node, peer_node, name, 100, None).await { + Ok(ControlResponse::DbSynced { + changes_applied, + max_db_version, + .. + }) => { + let store = Store::open(&node.paths.metadata_db())?; + let cursor = load_live_sync_cursor(&store, peer_node, stream)?; + let cursor_ms = record_live_sync_success( + &store, + peer_node, + stream, + changes_applied, + 0, + Some(max_db_version.unwrap_or(cursor)), + )?; + streams.push(SyncStreamRun { + stream: stream.to_owned(), + attempted: true, + success: true, + imported: changes_applied, + rejected: 0, + cursor_ms, + error: None, + }); + Ok(()) + } + Ok(other) => push_failed_sync_stream( + &store, + peer_node, + stream, + NodeError::IrohPeer(format!("unexpected DB sync response: {other:?}")), + streams, + ), + Err(error) => push_failed_sync_stream(&store, peer_node, stream, error, streams), + } +} + +pub(crate) async fn run_live_sync_once(node: &LocalNode) -> Result<(), NodeError> { + if node + .iroh_endpoint + .lock() + .map_err(|_| NodeError::RuntimeLockPoisoned)? + .is_none() + { + return Ok(()); + } + let peers = Store::open(&node.paths.metadata_db())?.list_peer_cards()?; + for peer in peers { + let run = run_sync_for_peer(node, &peer.peer_id).await?; + for stream in &run.streams { + tracing::debug!( + peer = %peer.peer_id, + stream = %stream.stream, + attempted = stream.attempted, + success = stream.success, + imported = stream.imported, + rejected = stream.rejected, + cursor_ms = stream.cursor_ms, + error = ?stream.error, + "live sync stream result" + ); + } + } + Ok(()) +} diff --git a/docs/production-readiness-roadmap.md b/docs/production-readiness-roadmap.md index b42fd1b..8e49f48 100644 --- a/docs/production-readiness-roadmap.md +++ b/docs/production-readiness-roadmap.md @@ -76,9 +76,9 @@ behavior. Acceptance criteria: - `[x]` Sync cursor keys, cursor persistence, stream health recording, and local sync-status reduction live in a sync-focused module. - - `[ ]` Sync stream selection, watermarks, and run result helpers live in a + - `[x]` Sync stream selection, watermarks, and run result helpers live in a sync-focused module. - - `[ ]` Per-module sync handlers have consistent interfaces. + - `[x]` Per-module sync handlers have consistent interfaces. - `[ ]` `geth sync status --json` output remains stable. ## Phase 2: Stable Automation Contracts