diff --git a/AGENTS.md b/AGENTS.md index 605ad02..a25d42b 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -109,9 +109,10 @@ Roadmap items should be actionable and checkable: and `geth peer auth-check` over Iroh, signed peer-card LAN discovery payloads, authorized `geth cas fetch`, `geth ssh cert sync`, and `geth ssh revocation sync` over the Iroh control ALPN, a background SSH - metadata and KV live-sync loop with per-peer cursors in `module_state` and an - authorized sync-status summary that lets peers disclose only permitted stream - watermarks before module pulls, + metadata, keychain/auth, KV, DB, document, and file-root live-sync loop with + per-peer cursors in `module_state`, `geth sync status`/`geth sync now`, and + an authorized sync-status summary that lets peers disclose only permitted + stream watermarks before module pulls, untrusted discovery-backend trait, custom relay-map config, and Iroh local-network discovery toggle exist. - Canonical signed-operation envelopes exist for keychain/auth signature diff --git a/README.md b/README.md index ffa4414..aa25a92 100644 --- a/README.md +++ b/README.md @@ -116,6 +116,8 @@ The bootstrap implementation provides: - `geth keychain status` - `geth keychain sync ` - `geth auth sync ` +- `geth sync status` +- `geth sync now [node-id-or-name]` - `geth secret status` - `geth secret create ` - `geth secret rotate ` @@ -412,6 +414,13 @@ from a currently trusted admin key over the canonical keychain payload. This is the current replicated device-management substrate. It is still a pull-based operation log, not yet a CRDT or Keyhive-style convergent authority. +The daemon also runs best-effort live sync for imported peers. `geth sync now +[node]` triggers the same sync pass immediately, and `geth sync status` reports +the last local attempt, success, cursor, import count, rejection count, and +error per peer stream. Keychain and auth sync now use per-peer high-water +cursors, while receivers still verify every imported signed operation before it +can affect the reduced keychain or authorization views. + ## Authorization Direction The MVP defines the split between: diff --git a/crates/geth-cli/src/lib.rs b/crates/geth-cli/src/lib.rs index f6c7f57..5bd1680 100644 --- a/crates/geth-cli/src/lib.rs +++ b/crates/geth-cli/src/lib.rs @@ -37,6 +37,10 @@ pub enum Command { command: DaemonCommand, }, Status, + Sync { + #[command(subcommand)] + command: SyncCommand, + }, Node { #[command(subcommand)] command: NodeCommand, @@ -134,6 +138,12 @@ pub enum ServiceCommand { }, } +#[derive(Debug, Subcommand)] +pub enum SyncCommand { + Status, + Now { node: Option }, +} + #[derive(Debug, Subcommand)] pub enum NodeCommand { Id, @@ -825,6 +835,12 @@ pub async fn run() -> Result<()> { fn request_for_command(command: Command) -> Result { Ok(match command { Command::Status => ControlRequest::Status, + Command::Sync { + command: SyncCommand::Status, + } => ControlRequest::SyncStatus, + Command::Sync { + command: SyncCommand::Now { node }, + } => ControlRequest::SyncNow { node }, Command::Node { command: NodeCommand::Id, } => ControlRequest::NodeId, @@ -1851,6 +1867,7 @@ fn print_response(response: ControlResponse, json: bool) -> Result<()> { ops_imported, signatures_imported, invalid_ops_rejected, + high_water_ms, note, } => { println!("synced keychain from: {peer_node_id}"); @@ -1859,6 +1876,7 @@ fn print_response(response: ControlResponse, json: bool) -> Result<()> { println!("ops_imported: {ops_imported}"); println!("signatures_imported: {signatures_imported}"); println!("invalid_ops_rejected: {invalid_ops_rejected}"); + println!("high_water_ms: {high_water_ms}"); println!("note: {note}"); } ControlResponse::AuthSynced { @@ -1868,6 +1886,7 @@ fn print_response(response: ControlResponse, json: bool) -> Result<()> { ops_imported, signatures_imported, invalid_ops_rejected, + high_water_ms, note, } => { println!("synced auth from: {peer_node_id}"); @@ -1876,6 +1895,66 @@ fn print_response(response: ControlResponse, json: bool) -> Result<()> { println!("ops_imported: {ops_imported}"); println!("signatures_imported: {signatures_imported}"); println!("invalid_ops_rejected: {invalid_ops_rejected}"); + println!("high_water_ms: {high_water_ms}"); + println!("note: {note}"); + } + ControlResponse::SyncStatus { peers, note } => { + if peers.is_empty() { + println!("no sync peers"); + } else { + for peer in peers { + println!("peer: {}", peer.peer_node_id); + if peer.streams.is_empty() { + println!(" no sync attempts recorded"); + } + for stream in peer.streams { + println!( + " {}\tcursor={}\tlast_attempt={}\tlast_success={}\timported={}\trejected={}\terror={}", + stream.stream, + stream.cursor_ms, + stream + .last_attempt_ms + .map(|value| value.to_string()) + .unwrap_or_else(|| "never".to_owned()), + stream + .last_success_ms + .map(|value| value.to_string()) + .unwrap_or_else(|| "never".to_owned()), + stream.last_imported, + stream.last_rejected, + stream.last_error.unwrap_or_else(|| "-".to_owned()) + ); + } + } + } + println!("note: {note}"); + } + ControlResponse::SyncRan { peers, note } => { + if peers.is_empty() { + println!("no sync peers"); + } else { + for peer in peers { + println!("peer: {}", peer.peer_node_id); + for stream in peer.streams { + let state = if !stream.attempted { + "skipped" + } else if stream.success { + "ok" + } else { + "failed" + }; + println!( + " {}\t{}\tcursor={}\timported={}\trejected={}\terror={}", + stream.stream, + state, + stream.cursor_ms, + stream.imported, + stream.rejected, + stream.error.unwrap_or_else(|| "-".to_owned()) + ); + } + } + } println!("note: {note}"); } ControlResponse::SecretStatus { secrets } => { diff --git a/crates/geth-control/src/lib.rs b/crates/geth-control/src/lib.rs index 153ccd1..df041c3 100644 --- a/crates/geth-control/src/lib.rs +++ b/crates/geth-control/src/lib.rs @@ -189,6 +189,10 @@ pub enum ControlRequest { AuthSync { node: String, }, + SyncStatus, + SyncNow { + node: Option, + }, SecretStatus, SecretCreate { resource: String, @@ -576,6 +580,7 @@ pub enum ControlResponse { ops_imported: usize, signatures_imported: usize, invalid_ops_rejected: usize, + high_water_ms: i64, note: String, }, AuthSynced { @@ -585,6 +590,15 @@ pub enum ControlResponse { ops_imported: usize, signatures_imported: usize, invalid_ops_rejected: usize, + high_water_ms: i64, + note: String, + }, + SyncStatus { + peers: Vec, + note: String, + }, + SyncRan { + peers: Vec, note: String, }, SecretStatus { @@ -953,10 +967,12 @@ pub enum PeerControlRequest { }, KeychainSync { peer_card: PeerCard, + since_ms: i64, nonce: String, }, AuthSync { peer_card: PeerCard, + since_ms: i64, nonce: String, }, NodeEnrollmentSubmit { @@ -1089,6 +1105,7 @@ pub enum PeerControlResponse { remote_endpoint_id: String, ops: Vec, signatures: Vec, + high_water_ms: i64, nonce: String, note: String, }, @@ -1099,6 +1116,7 @@ pub enum PeerControlResponse { remote_endpoint_id: String, ops: Vec, signatures: Vec, + high_water_ms: i64, nonce: String, note: String, }, @@ -1351,6 +1369,40 @@ pub struct SyncWatermark { pub high_water: i64, } +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +pub struct SyncPeerStatus { + pub peer_node_id: String, + pub streams: Vec, +} + +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +pub struct SyncStreamStatus { + pub stream: String, + pub cursor_ms: i64, + pub last_attempt_ms: Option, + pub last_success_ms: Option, + pub last_error: Option, + pub last_imported: usize, + pub last_rejected: usize, +} + +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +pub struct SyncPeerRun { + pub peer_node_id: String, + pub streams: Vec, +} + +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +pub struct SyncStreamRun { + pub stream: String, + pub attempted: bool, + pub success: bool, + pub imported: usize, + pub rejected: usize, + pub cursor_ms: i64, + pub error: Option, +} + #[derive(Debug, thiserror::Error)] pub enum ControlError { #[error("json error: {0}")] @@ -1507,6 +1559,7 @@ mod tests { ops_imported: 2, signatures_imported: 2, invalid_ops_rejected: 1, + high_water_ms: 42, note: "trusted admin signatures only".to_owned(), }; assert_eq!( @@ -1535,6 +1588,7 @@ mod tests { ops_imported: 1, signatures_imported: 1, invalid_ops_rejected: 0, + high_water_ms: 42, note: "trusted admin signatures only".to_owned(), }; assert_eq!( @@ -1542,6 +1596,46 @@ mod tests { response ); + let response = ControlResponse::SyncStatus { + peers: vec![SyncPeerStatus { + peer_node_id: "node:peer".to_owned(), + streams: vec![SyncStreamStatus { + stream: "keychain".to_owned(), + cursor_ms: 42, + last_attempt_ms: Some(43), + last_success_ms: Some(43), + last_error: None, + last_imported: 2, + last_rejected: 0, + }], + }], + note: "local health".to_owned(), + }; + assert_eq!( + decode_response(&encode_response(&response).expect("encode")).expect("decode"), + response + ); + + let response = ControlResponse::SyncRan { + peers: vec![SyncPeerRun { + peer_node_id: "node:peer".to_owned(), + streams: vec![SyncStreamRun { + stream: "auth".to_owned(), + attempted: true, + success: true, + imported: 1, + rejected: 0, + cursor_ms: 44, + error: None, + }], + }], + note: "ran".to_owned(), + }; + assert_eq!( + decode_response(&encode_response(&response).expect("encode")).expect("decode"), + response + ); + let request = ControlRequest::SecretBearerVerify { secret: "bearer:test".to_owned(), resource: "resource:cas:local".to_owned(), diff --git a/crates/geth-node/src/lib.rs b/crates/geth-node/src/lib.rs index c57984f..1ade483 100644 --- a/crates/geth-node/src/lib.rs +++ b/crates/geth-node/src/lib.rs @@ -10,7 +10,7 @@ use geth_config::{GethConfig, GethPaths, RelayMode}; use geth_control::{ CasBlob, CasProvider, ControlRequest, ControlResponse, KeychainStatusResponse, NodeIdResponse, PeerControlRequest, PeerControlResponse, PipeWireRequest, PipeWireResponse, StatusResponse, - SyncWatermark, + SyncPeerRun, SyncPeerStatus, SyncStreamRun, SyncStreamStatus, SyncWatermark, }; use geth_crypto::AgentKey; use geth_db::DbResource; @@ -50,7 +50,7 @@ use geth_types::{ AuthOpId, BlobHash, Capability, DeviceId, KeyId, NodeId, PrincipalId, ResourceId, ResourceKind, ResourceName, SshCertId, SshCertRequestId, UnixMillis, UserId, }; -use std::collections::{BTreeMap, VecDeque}; +use std::collections::{BTreeMap, BTreeSet, VecDeque}; use std::path::{Component, Path, PathBuf}; use std::sync::{Arc, Mutex}; use std::time::Duration; @@ -186,6 +186,15 @@ struct LiveSyncCursor { cursor_ms: i64, } +#[derive(Debug, serde::Deserialize, serde::Serialize)] +struct LiveSyncHealth { + last_attempt_ms: i64, + last_success_ms: Option, + last_error: Option, + last_imported: usize, + last_rejected: usize, +} + struct PipeTcpConnectWire { peer_card: PeerCard, target_addr: String, @@ -652,6 +661,7 @@ pub async fn handle_request_async( keychain_sync_from_peer(node, &peer_node).await } ControlRequest::AuthSync { node: peer_node } => auth_sync_from_peer(node, &peer_node).await, + ControlRequest::SyncNow { node: peer_node } => sync_now(node, peer_node).await, ControlRequest::NodeEnrollSubmit { owner_node, request_id, @@ -2584,9 +2594,21 @@ async fn sync_status_from_peer( async fn keychain_sync_from_peer( node: &LocalNode, peer_node: &str, +) -> Result { + keychain_sync_from_peer_since(node, peer_node, 0).await +} + +async fn keychain_sync_from_peer_since( + node: &LocalNode, + peer_node: &str, + since_ms: i64, ) -> Result { let response = request_peer_control(node, peer_node, "keychain-sync", |peer_card, nonce| { - PeerControlRequest::KeychainSync { peer_card, nonce } + PeerControlRequest::KeychainSync { + peer_card, + since_ms, + nonce, + } }) .await?; match response { @@ -2596,12 +2618,18 @@ async fn keychain_sync_from_peer( endpoint_id, ops, signatures, + high_water_ms, note, .. } => { let store = Store::open(&node.paths.metadata_db())?; let mut local_ops = load_keychain_ops(&store)?; let mut trusted_admins = geth_keychain::reduce_keychain_ops(&local_ops).admin_keys; + let mut known_signatures = store + .list_keychain_signatures()? + .into_iter() + .map(|signature| (signature.op_id, signature.signer, signature.namespace)) + .collect::>(); let mut ops_imported = 0; let mut signatures_imported = 0; let mut invalid_ops_rejected = 0; @@ -2646,6 +2674,11 @@ async fn keychain_sync_from_peer( ops_imported += 1; } for signature in valid_signatures { + let signature_key = ( + signature.op_id.to_string(), + signature.signer.to_string(), + signature.namespace.clone(), + ); store.insert_keychain_signature(&StoredKeychainSignature { op_id: signature.op_id.to_string(), signer: signature.signer.to_string(), @@ -2654,9 +2687,10 @@ async fn keychain_sync_from_peer( signature: signature.signature.clone(), created_at_ms: signature.created_at.0, })?; - signatures_imported += 1; + signatures_imported += usize::from(known_signatures.insert(signature_key)); } } + store_live_sync_cursor(&store, peer_node, "keychain", high_water_ms)?; Ok(ControlResponse::KeychainSynced { peer_node_id: node_id, peer_agent_id: agent_id, @@ -2664,6 +2698,7 @@ async fn keychain_sync_from_peer( ops_imported, signatures_imported, invalid_ops_rejected, + high_water_ms, note: format!( "{note}; imported only keychain ops signed by currently trusted admin SSH keys" ), @@ -2679,9 +2714,21 @@ async fn keychain_sync_from_peer( async fn auth_sync_from_peer( node: &LocalNode, peer_node: &str, +) -> Result { + auth_sync_from_peer_since(node, peer_node, 0).await +} + +async fn auth_sync_from_peer_since( + node: &LocalNode, + peer_node: &str, + since_ms: i64, ) -> Result { let response = request_peer_control(node, peer_node, "auth-sync", |peer_card, nonce| { - PeerControlRequest::AuthSync { peer_card, nonce } + PeerControlRequest::AuthSync { + peer_card, + since_ms, + nonce, + } }) .await?; match response { @@ -2691,12 +2738,18 @@ async fn auth_sync_from_peer( endpoint_id, ops, signatures, + high_water_ms, note, .. } => { let store = Store::open(&node.paths.metadata_db())?; let trusted_admins = geth_keychain::reduce_keychain_ops(&load_keychain_ops(&store)?).admin_keys; + let mut known_signatures = store + .list_auth_signatures()? + .into_iter() + .map(|signature| (signature.op_id, signature.signer, signature.namespace)) + .collect::>(); let mut ops_imported = 0; let mut signatures_imported = 0; let mut invalid_ops_rejected = 0; @@ -2737,6 +2790,11 @@ async fn auth_sync_from_peer( store_auth_op(&store, &op)?; ops_imported += usize::from(existing.is_none()); for signature in valid_signatures { + let signature_key = ( + signature.op_id.to_string(), + signature.signer.to_string(), + signature.namespace.clone(), + ); store.insert_auth_signature(&StoredAuthSignature { op_id: signature.op_id.to_string(), signer: signature.signer.to_string(), @@ -2745,9 +2803,10 @@ async fn auth_sync_from_peer( signature: signature.signature.clone(), created_at_ms: signature.created_at.0, })?; - signatures_imported += 1; + signatures_imported += usize::from(known_signatures.insert(signature_key)); } } + store_live_sync_cursor(&store, peer_node, "auth", high_water_ms)?; Ok(ControlResponse::AuthSynced { peer_node_id: node_id, peer_agent_id: agent_id, @@ -2755,6 +2814,7 @@ async fn auth_sync_from_peer( ops_imported, signatures_imported, invalid_ops_rejected, + high_water_ms, note: format!("{note}; imported only auth ops signed by trusted admin SSH keys"), }) } @@ -2973,6 +3033,14 @@ fn live_sync_cursor_key(peer_node: &str, stream: &str) -> String { format!("live-sync:{peer_node}:{stream}") } +fn live_sync_status_prefix(peer_node: &str) -> String { + format!("live-sync-status:{peer_node}:") +} + +fn live_sync_status_key(peer_node: &str, stream: &str) -> String { + format!("{}{}", live_sync_status_prefix(peer_node), stream) +} + fn load_live_sync_cursor(store: &Store, peer_node: &str, stream: &str) -> Result { let key = live_sync_cursor_key(peer_node, stream); let Some(state) = store.get_module_state(&key)? else { @@ -2996,6 +3064,93 @@ fn store_live_sync_cursor( Ok(()) } +fn record_live_sync_success( + store: &Store, + peer_node: &str, + stream: &str, + imported: usize, + rejected: usize, + cursor_ms: Option, +) -> Result { + let now = geth_store::now_ms(); + if let Some(cursor_ms) = cursor_ms { + store_live_sync_cursor(store, peer_node, stream, cursor_ms)?; + } + store.put_module_state(&StoredModuleState { + module: live_sync_status_key(peer_node, stream), + state_json: serde_json::to_string(&LiveSyncHealth { + last_attempt_ms: now, + last_success_ms: Some(now), + last_error: None, + last_imported: imported, + last_rejected: rejected, + })?, + updated_at_ms: now, + })?; + load_live_sync_cursor(store, peer_node, stream) +} + +fn record_live_sync_failure( + store: &Store, + peer_node: &str, + stream: &str, + error: &NodeError, +) -> Result { + let now = geth_store::now_ms(); + let previous = store + .get_module_state(&live_sync_status_key(peer_node, stream))? + .and_then(|state| serde_json::from_str::(&state.state_json).ok()); + store.put_module_state(&StoredModuleState { + module: live_sync_status_key(peer_node, stream), + state_json: serde_json::to_string(&LiveSyncHealth { + last_attempt_ms: now, + last_success_ms: previous.and_then(|health| health.last_success_ms), + last_error: Some(error.to_string()), + last_imported: 0, + last_rejected: 0, + })?, + updated_at_ms: now, + })?; + load_live_sync_cursor(store, peer_node, stream) +} + +fn sync_status_local(store: &Store) -> Result { + let peers = store + .list_peer_cards()? + .into_iter() + .map(|peer| { + let prefix = live_sync_status_prefix(&peer.peer_id); + let mut streams = store + .list_module_states_with_prefix(&prefix)? + .into_iter() + .map(|state| { + let stream = state.module.trim_start_matches(&prefix).to_owned(); + let health: LiveSyncHealth = serde_json::from_str(&state.state_json)?; + let cursor_ms = load_live_sync_cursor(store, &peer.peer_id, &stream)?; + Ok(SyncStreamStatus { + stream, + cursor_ms, + last_attempt_ms: Some(health.last_attempt_ms), + last_success_ms: health.last_success_ms, + last_error: health.last_error, + last_imported: health.last_imported, + last_rejected: health.last_rejected, + }) + }) + .collect::, NodeError>>()?; + streams.sort_by(|left, right| left.stream.cmp(&right.stream)); + Ok(SyncPeerStatus { + peer_node_id: peer.peer_id, + streams, + }) + }) + .collect::, NodeError>>()?; + Ok(ControlResponse::SyncStatus { + peers, + note: "sync status is local daemon health for best-effort live sync; signed logs remain the durable source of truth".to_owned(), + }) +} + fn should_live_sync_stream( store: &Store, peer_node: &str, @@ -3024,6 +3179,34 @@ fn sync_watermarks_for_peer( peer_node: &str, ) -> Result, NodeError> { let mut watermarks = Vec::new(); + let keychain_high = load_keychain_ops(store)? + .into_iter() + .map(|op| op.created_at.0) + .chain( + load_keychain_signatures(store)? + .into_iter() + .map(|signature| signature.created_at.0), + ) + .max() + .unwrap_or(0); + watermarks.push(SyncWatermark { + stream: "keychain".to_owned(), + high_water: keychain_high, + }); + let auth_high = load_auth_ops(store)? + .into_iter() + .map(|op| op.created_at.0) + .chain( + load_auth_signatures(store)? + .into_iter() + .map(|signature| signature.created_at.0), + ) + .max() + .unwrap_or(0); + watermarks.push(SyncWatermark { + stream: "auth".to_owned(), + high_water: auth_high, + }); let peer = PrincipalId::new(peer_node.to_owned()); if can_sync_resource(store, &peer, "resource:ssh:certs", "ssh_cert.sync")? { let cert_high = store @@ -3469,6 +3652,567 @@ fn spawn_background_live_sync(node: LocalNode, interval_duration: Duration) { }); } +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 @@ -3479,96 +4223,11 @@ async fn run_live_sync_once(node: &LocalNode) -> Result<(), NodeError> { return Ok(()); } let peers = Store::open(&node.paths.metadata_db())?.list_peer_cards()?; - 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()?; for peer in peers { - let store = Store::open(&node.paths.metadata_db())?; - let remote_watermarks = match sync_status_from_peer(node, &peer.peer_id).await.map( - |watermarks| { - watermarks - .into_iter() - .map(|watermark| (watermark.stream, watermark.high_water)) - .collect::>() - }, - ) { - Ok(watermarks) => Some(watermarks), - Err(error) => { - tracing::debug!(peer = %peer.peer_id, %error, "sync status failed; falling back to direct live-sync probes"); - None - } - }; - if let Err(error) = keychain_sync_from_peer(node, &peer.peer_id).await { - tracing::debug!(peer = %peer.peer_id, %error, "live keychain sync failed"); - } - if let Err(error) = auth_sync_from_peer(node, &peer.peer_id).await { - tracing::debug!(peer = %peer.peer_id, %error, "live auth sync failed"); - } - if should_live_sync_stream( - &store, - &peer.peer_id, - "ssh-certs", - remote_watermarks.as_ref(), - )? { - if let Err(error) = ssh_cert_sync_from_peer(node, &peer.peer_id, None).await { - tracing::debug!(peer = %peer.peer_id, %error, "SSH cert live sync failed"); - } - } - if should_live_sync_stream( - &store, - &peer.peer_id, - "ssh-revocations", - remote_watermarks.as_ref(), - )? { - if let Err(error) = ssh_revocation_sync_from_peer(node, &peer.peer_id, None).await { - tracing::debug!(peer = %peer.peer_id, %error, "SSH revocation live sync failed"); - } - } - if let Some(remote_watermarks) = remote_watermarks.as_ref() { - for stream in remote_watermarks - .keys() - .filter(|stream| stream.starts_with("cas-tree:")) - { - if should_live_sync_stream(&store, &peer.peer_id, stream, Some(remote_watermarks))? - { - let name = stream.trim_start_matches("cas-tree:"); - if let Err(error) = - cas_root_sync_from_peer(node, &peer.peer_id, name, None).await - { - tracing::debug!(peer = %peer.peer_id, root = %name, %error, "file-root live sync failed"); - } - } - } - } - for kv in &kv_stores { - let stream = format!("kv:{}", kv.name); - if should_live_sync_stream(&store, &peer.peer_id, &stream, remote_watermarks.as_ref())? - { - if let Err(error) = kv_sync_from_peer(node, &peer.peer_id, &kv.name, None).await { - tracing::debug!(peer = %peer.peer_id, kv = %kv.name, %error, "KV live sync failed"); - } - } - } - for document in &documents { - let stream = format!("document:{}", document.name); - if should_live_sync_stream(&store, &peer.peer_id, &stream, remote_watermarks.as_ref())? - { - if let Err(error) = - document_sync_from_peer(node, &peer.peer_id, &document.name, None).await - { - tracing::debug!(peer = %peer.peer_id, document = %document.name, %error, "document live sync failed"); - } - } - } - for db in &dbs { - let stream = format!("db:{}", db.name); - if should_live_sync_stream(&store, &peer.peer_id, &stream, remote_watermarks.as_ref())? - { - if let Err(error) = - db_sync_from_peer(node, &peer.peer_id, &db.name, 100, None).await - { - tracing::debug!(peer = %peer.peer_id, db = %db.name, %error, "DB live sync failed"); - } + let run = run_sync_for_peer(node, &peer.peer_id).await?; + for stream in run.streams.iter().filter(|stream| !stream.success) { + if let Some(error) = &stream.error { + tracing::debug!(peer = %peer.peer_id, stream = %stream.stream, %error, "live sync stream failed"); } } } @@ -3656,7 +4315,11 @@ async fn handle_iroh_control_connection( note: "sync status authenticated endpoint/card binding and returns only streams for capabilities already granted to the caller".to_owned(), } } - PeerControlRequest::KeychainSync { peer_card, nonce } => { + PeerControlRequest::KeychainSync { + peer_card, + since_ms, + nonce, + } => { peer_card.validate_candidate()?; ensure_peer_card_matches_endpoint(&peer_card, &remote_endpoint_id)?; let discovered = DiscoveredPeer::candidate( @@ -3670,18 +4333,45 @@ async fn handle_iroh_control_connection( card_json: serde_json::to_string(&peer_card)?, updated_at_ms: discovered.discovered_at.0, })?; + let ops = load_keychain_ops(&store)?; + let signatures = load_keychain_signatures(&store)?; + let returned_op_ids = ops + .iter() + .filter(|op| op.created_at.0 >= since_ms) + .map(|op| op.id.to_string()) + .collect::>(); + let high_water_ms = ops + .iter() + .map(|op| op.created_at.0) + .chain(signatures.iter().map(|signature| signature.created_at.0)) + .max() + .unwrap_or(0); PeerControlResponse::KeychainSynced { node_id: node.node_id.clone(), agent_id: node.agent_id.clone(), endpoint_id: node.iroh_status.endpoint_id.clone().unwrap_or_default(), remote_endpoint_id, - ops: load_keychain_ops(&store)?, - signatures: load_keychain_signatures(&store)?, + ops: ops + .into_iter() + .filter(|op| op.created_at.0 >= since_ms) + .collect(), + signatures: signatures + .into_iter() + .filter(|signature| { + signature.created_at.0 >= since_ms + || returned_op_ids.contains(signature.op_id.as_str()) + }) + .collect(), + high_water_ms, nonce, note: "keychain sync returns signed operation-log data; receiver must verify OpenSSH signatures before import".to_owned(), } } - PeerControlRequest::AuthSync { peer_card, nonce } => { + PeerControlRequest::AuthSync { + peer_card, + since_ms, + nonce, + } => { peer_card.validate_candidate()?; ensure_peer_card_matches_endpoint(&peer_card, &remote_endpoint_id)?; let discovered = DiscoveredPeer::candidate( @@ -3695,13 +4385,36 @@ async fn handle_iroh_control_connection( card_json: serde_json::to_string(&peer_card)?, updated_at_ms: discovered.discovered_at.0, })?; + let ops = load_auth_ops(&store)?; + let signatures = load_auth_signatures(&store)?; + let returned_op_ids = ops + .iter() + .filter(|op| op.created_at.0 >= since_ms) + .map(|op| op.id.to_string()) + .collect::>(); + let high_water_ms = ops + .iter() + .map(|op| op.created_at.0) + .chain(signatures.iter().map(|signature| signature.created_at.0)) + .max() + .unwrap_or(0); PeerControlResponse::AuthSynced { node_id: node.node_id.clone(), agent_id: node.agent_id.clone(), endpoint_id: node.iroh_status.endpoint_id.clone().unwrap_or_default(), remote_endpoint_id, - ops: load_auth_ops(&store)?, - signatures: load_auth_signatures(&store)?, + ops: ops + .into_iter() + .filter(|op| op.created_at.0 >= since_ms) + .collect(), + signatures: signatures + .into_iter() + .filter(|signature| { + signature.created_at.0 >= since_ms + || returned_op_ids.contains(signature.op_id.as_str()) + }) + .collect(), + high_water_ms, nonce, note: "auth sync returns signed resource operation-log data; receiver must verify admin signatures before import".to_owned(), } @@ -5179,8 +5892,10 @@ pub fn handle_request( note: discovery_is_untrusted_note().to_owned(), }) } + ControlRequest::SyncStatus => sync_status_local(&store), ControlRequest::PeerPing { .. } => Err(NodeError::IrohEndpointUnavailable), ControlRequest::PeerAuthCheck { .. } => Err(NodeError::IrohEndpointUnavailable), + ControlRequest::SyncNow { .. } => Err(NodeError::IrohEndpointUnavailable), ControlRequest::CasFetch { .. } => Err(NodeError::IrohEndpointUnavailable), ControlRequest::CasRootSync { .. } => Err(NodeError::IrohEndpointUnavailable), ControlRequest::SshCertSync { .. } => Err(NodeError::IrohEndpointUnavailable), @@ -8905,6 +9620,16 @@ mod tests { let watermarks = sync_watermarks_for_peer(&store, "node:left").expect("watermarks"); + assert!( + watermarks + .iter() + .any(|watermark| watermark.stream == "auth") + ); + assert!( + watermarks + .iter() + .any(|watermark| watermark.stream == "keychain") + ); assert!(watermarks.contains(&SyncWatermark { stream: "kv:prefs".to_owned(), high_water: 42, @@ -8934,6 +9659,34 @@ mod tests { ); } + #[test] + fn sync_status_local_reports_recorded_stream_health() { + let store = Store::open_memory().expect("open"); + store + .upsert_peer_card(&StoredPeerCard { + peer_id: "node:left".to_owned(), + card_json: "{}".to_owned(), + updated_at_ms: 1, + }) + .expect("insert peer"); + store_live_sync_cursor(&store, "node:left", "keychain", 42).expect("store cursor"); + record_live_sync_success(&store, "node:left", "keychain", 3, 1, Some(42)) + .expect("record success"); + + let response = sync_status_local(&store).expect("sync status"); + let ControlResponse::SyncStatus { peers, .. } = response else { + panic!("unexpected sync status response"); + }; + assert_eq!(peers.len(), 1); + assert_eq!(peers[0].peer_node_id, "node:left"); + assert_eq!(peers[0].streams.len(), 1); + assert_eq!(peers[0].streams[0].stream, "keychain"); + assert_eq!(peers[0].streams[0].cursor_ms, 42); + assert_eq!(peers[0].streams[0].last_imported, 3); + assert_eq!(peers[0].streams[0].last_rejected, 1); + assert!(peers[0].streams[0].last_error.is_none()); + } + #[test] fn file_root_sync_records_automatic_three_tree_conflicts() { let store = Store::open_memory().expect("open"); diff --git a/crates/geth-store/src/lib.rs b/crates/geth-store/src/lib.rs index e66fab9..918fc7f 100644 --- a/crates/geth-store/src/lib.rs +++ b/crates/geth-store/src/lib.rs @@ -905,6 +905,28 @@ impl Store { } } + pub fn list_module_states_with_prefix( + &self, + prefix: &str, + ) -> Result, StoreError> { + let like = format!("{prefix}%"); + let mut stmt = self.conn.prepare( + r#"SELECT module, state_json, updated_at_ms + FROM module_state + WHERE module LIKE ?1 + ORDER BY module"#, + )?; + let rows = stmt.query_map(params![like], |row| { + Ok(StoredModuleState { + module: row.get(0)?, + state_json: row.get(1)?, + updated_at_ms: row.get(2)?, + }) + })?; + rows.collect::, _>>() + .map_err(StoreError::from) + } + pub fn insert_auth_op(&self, op: &StoredAuthOp) -> Result<(), StoreError> { self.conn.execute( r#"INSERT OR REPLACE INTO auth_ops(op_id, resource_id, op_json, created_at_ms) @@ -1672,12 +1694,26 @@ mod tests { updated_at_ms: 43, }; store.put_module_state(&state).expect("put state"); + store + .put_module_state(&StoredModuleState { + module: "live-sync-status:node:laptop:ssh-certs".to_owned(), + state_json: r#"{"last_attempt_ms":44}"#.to_owned(), + updated_at_ms: 44, + }) + .expect("put status state"); assert_eq!( store .get_module_state("live-sync:node:laptop:ssh-certs") .expect("get state"), Some(state) ); + assert_eq!( + store + .list_module_states_with_prefix("live-sync-status:node:laptop:") + .expect("list prefix") + .len(), + 1 + ); } #[test] diff --git a/docs/architecture.md b/docs/architecture.md index c1ec745..de487c5 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -106,6 +106,12 @@ status summary over the same protected Iroh control ALPN. The serving peer validates endpoint/card binding and returns only watermarks for streams where the caller already has the required resource capability, which reduces blind polling without letting discovery reveal private resource names. +Keychain and auth operation logs are also advertised through this watermark +path. Pulls are delta-style by per-peer cursor, but every received operation is +still verified against trusted-admin OpenSSH signatures before import. +Operators can run `geth sync now [node]` to trigger the same best-effort pass +immediately and `geth sync status` to inspect locally recorded last-attempt, +last-success, cursor, import/rejection counts, and errors for each peer stream. ## Resource Model diff --git a/docs/roadmap.md b/docs/roadmap.md index 03f2699..c412a02 100644 --- a/docs/roadmap.md +++ b/docs/roadmap.md @@ -150,7 +150,12 @@ geth-to-geth connections without granting trust from discovery alone. the relevant resource capability. - `[x]` Background live-sync skips per-module pulls when the authorized remote watermark has not advanced. + - `[x]` `geth sync status` reports last local attempt, success, cursor, + import/rejection counts, and error for each recorded peer stream. + - `[x]` `geth sync now [node]` triggers the same best-effort sync pass that + background live sync uses. - `[x]` Tests verify unauthorized streams are omitted from summary output. + - `[x]` Tests verify persisted stream health is exposed in local sync status. ## Phase 2: Trust And Authorization @@ -188,6 +193,8 @@ resource-scoped capability decisions. - `[x]` `geth node rename/revoke` require an admin signing key. - `[x]` `geth keychain sync ` verifies signatures from currently trusted admin keys before accepting keychain ops. + - `[x]` Keychain live sync advertises and consumes per-peer high-water + cursors instead of blindly re-requesting the full log on every tick. - `[x]` `geth node enroll request` creates an agent-key-signed enrollment request with requested node name and capabilities. - `[x]` `geth node enroll submit/import/list` moves pending enrollment @@ -217,6 +224,8 @@ resource-scoped capability decisions. - `[x]` Enrollment approval signs capability grants as auth ops. - `[x]` `geth auth sync ` imports only auth ops signed by currently trusted admin keys. + - `[x]` Auth live sync advertises and consumes per-peer high-water cursors + instead of blindly re-requesting the full log on every tick. - `[x]` `geth auth grant/revoke --signing-key` records signed auth ops. - `[x]` Resource auth operation reducer.