refactor: move live sync orchestration

This commit is contained in:
Eric Wendland 2026-07-05 17:49:09 +02:00
commit 2b604dd184
4 changed files with 605 additions and 600 deletions

View file

@ -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");
}
}

View file

@ -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<SyncPeerRun, NodeError> {
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::<BTreeMap<_, _>>()
}) {
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<String>,
) -> Result<ControlResponse, NodeError> {
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<SyncStreamRun>,
) -> 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<SyncStreamRun>,
) -> 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<String, i64>>,
streams: &mut Vec<SyncStreamRun>,
) -> 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<String, i64>>,
streams: &mut Vec<SyncStreamRun>,
) -> 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<String, i64>>,
streams: &mut Vec<SyncStreamRun>,
) -> 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<String, i64>>,
streams: &mut Vec<SyncStreamRun>,
) -> 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<String, i64>>,
streams: &mut Vec<SyncStreamRun>,
) -> 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<String, i64>>,
streams: &mut Vec<SyncStreamRun>,
) -> 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<String, i64>>,
streams: &mut Vec<SyncStreamRun>,
) -> 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<String, i64>>,
streams: &mut Vec<SyncStreamRun>,
) -> 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 {

View file

@ -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<SyncPeerRun, NodeError> {
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::<BTreeMap<_, _>>()
}) {
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<String>,
) -> Result<ControlResponse, NodeError> {
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<SyncStreamRun>,
) -> 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<SyncStreamRun>,
) -> 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<String, i64>>,
streams: &mut Vec<SyncStreamRun>,
) -> 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<String, i64>>,
streams: &mut Vec<SyncStreamRun>,
) -> 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<String, i64>>,
streams: &mut Vec<SyncStreamRun>,
) -> 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<String, i64>>,
streams: &mut Vec<SyncStreamRun>,
) -> 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<String, i64>>,
streams: &mut Vec<SyncStreamRun>,
) -> 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<String, i64>>,
streams: &mut Vec<SyncStreamRun>,
) -> 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<String, i64>>,
streams: &mut Vec<SyncStreamRun>,
) -> 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<String, i64>>,
streams: &mut Vec<SyncStreamRun>,
) -> 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(())
}

View file

@ -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