Skip unchanged live-sync streams

This commit is contained in:
Eric Wendland 2026-05-18 22:17:27 +02:00
commit bc0ba4d169
6 changed files with 416 additions and 15 deletions

View file

@ -571,6 +571,10 @@ pub enum PeerControlRequest {
capability: String,
nonce: String,
},
SyncStatus {
peer_card: PeerCard,
nonce: String,
},
CasFetch {
peer_card: PeerCard,
hash: BlobHash,
@ -643,6 +647,15 @@ pub enum PeerControlResponse {
nonce: String,
note: String,
},
SyncStatus {
node_id: String,
agent_id: String,
endpoint_id: String,
remote_endpoint_id: String,
watermarks: Vec<SyncWatermark>,
nonce: String,
note: String,
},
CasFetched {
node_id: String,
agent_id: String,
@ -755,6 +768,12 @@ pub enum PeerControlResponse {
},
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub struct SyncWatermark {
pub stream: String,
pub high_water: i64,
}
#[derive(Debug, thiserror::Error)]
pub enum ControlError {
#[error("json error: {0}")]
@ -1118,6 +1137,24 @@ mod tests {
response
);
let response = PeerControlResponse::SyncStatus {
node_id: "node:peer".to_owned(),
agent_id: "agent:peer".to_owned(),
endpoint_id: "endpoint:peer".to_owned(),
remote_endpoint_id: "endpoint:caller".to_owned(),
watermarks: vec![SyncWatermark {
stream: "kv:prefs".to_owned(),
high_water: 42,
}],
nonce: "nonce".to_owned(),
note: "status".to_owned(),
};
assert_eq!(
decode_peer_response(&encode_peer_response(&response).expect("encode"))
.expect("decode"),
response
);
let response = PeerControlResponse::CasFetched {
node_id: "node:peer".to_owned(),
agent_id: "agent:peer".to_owned(),
@ -1161,6 +1198,26 @@ mod tests {
request
);
let request = PeerControlRequest::SyncStatus {
peer_card: PeerCard {
node_id: "node:caller".into(),
agent_id: "agent:caller".into(),
endpoints: Vec::new(),
issued_at: geth_types::UnixMillis(1),
signature: geth_discovery::SignatureMetadata {
namespace: "geth.peer-card.v1@geth.local".to_owned(),
signer: "agent:caller".to_owned(),
public_key: "key".to_owned(),
signature: "sig".to_owned(),
},
},
nonce: "nonce".to_owned(),
};
assert_eq!(
decode_peer_request(&encode_peer_request(&request).expect("encode")).expect("decode"),
request
);
let response = PeerControlResponse::SshCertSynced {
node_id: "node:peer".to_owned(),
agent_id: "agent:peer".to_owned(),

View file

@ -9,7 +9,7 @@ use geth_cas::{
use geth_config::{GethConfig, GethPaths, RelayMode};
use geth_control::{
CasBlob, ControlRequest, ControlResponse, KeychainStatusResponse, NodeIdResponse,
PeerControlRequest, PeerControlResponse, StatusResponse,
PeerControlRequest, PeerControlResponse, StatusResponse, SyncWatermark,
};
use geth_crypto::AgentKey;
use geth_db::DbResource;
@ -554,7 +554,8 @@ async fn peer_ping(node: &LocalNode, peer_node: &str) -> Result<ControlResponse,
PeerControlResponse::AuthChecked { .. } => Err(NodeError::IrohPeer(
"peer returned auth-check response to ping request".to_owned(),
)),
PeerControlResponse::CasFetched { .. }
PeerControlResponse::SyncStatus { .. }
| PeerControlResponse::CasFetched { .. }
| PeerControlResponse::SshCertSynced { .. }
| PeerControlResponse::SshRevocationSynced { .. }
| PeerControlResponse::KvSynced { .. }
@ -666,7 +667,8 @@ async fn peer_auth_check(
PeerControlResponse::Pong { .. } => Err(NodeError::IrohPeer(
"peer returned pong to auth-check request".to_owned(),
)),
PeerControlResponse::CasFetched { .. }
PeerControlResponse::SyncStatus { .. }
| PeerControlResponse::CasFetched { .. }
| PeerControlResponse::SshCertSynced { .. }
| PeerControlResponse::SshRevocationSynced { .. }
| PeerControlResponse::KvSynced { .. }
@ -808,6 +810,7 @@ async fn cas_fetch_from_peer(
PeerControlResponse::Error { message } => Err(NodeError::IrohPeer(message)),
PeerControlResponse::Pong { .. }
| PeerControlResponse::AuthChecked { .. }
| PeerControlResponse::SyncStatus { .. }
| PeerControlResponse::SshCertSynced { .. }
| PeerControlResponse::SshRevocationSynced { .. }
| PeerControlResponse::KvSynced { .. }
@ -1220,6 +1223,28 @@ async fn document_sync_from_peer(
}
}
async fn sync_status_from_peer(
node: &LocalNode,
peer_node: &str,
) -> Result<Vec<SyncWatermark>, NodeError> {
let response = request_peer_control(node, peer_node, "sync-status", |peer_card, nonce| {
PeerControlRequest::SyncStatus { peer_card, nonce }
})
.await?;
match response {
PeerControlResponse::SyncStatus {
watermarks, note, ..
} => {
tracing::debug!(peer = %peer_node, %note, streams = watermarks.len(), "peer sync status received");
Ok(watermarks)
}
PeerControlResponse::Error { message } => Err(NodeError::IrohPeer(message)),
_ => Err(NodeError::IrohPeer(
"peer returned wrong response type to sync status".to_owned(),
)),
}
}
async fn db_sync_from_peer(
node: &LocalNode,
peer_node: &str,
@ -1335,6 +1360,121 @@ fn store_live_sync_cursor(
Ok(())
}
fn should_live_sync_stream(
store: &Store,
peer_node: &str,
stream: &str,
remote_watermarks: Option<&BTreeMap<String, i64>>,
) -> Result<bool, NodeError> {
let Some(remote_watermarks) = remote_watermarks else {
return Ok(true);
};
let Some(remote_high_water) = remote_watermarks.get(stream) else {
return Ok(false);
};
if *remote_high_water == 0 {
return Ok(false);
}
let local_cursor = load_live_sync_cursor(store, peer_node, stream)?;
if stream.starts_with("db:") {
Ok(*remote_high_water > local_cursor)
} else {
Ok(*remote_high_water >= local_cursor)
}
}
fn sync_watermarks_for_peer(
store: &Store,
peer_node: &str,
) -> Result<Vec<SyncWatermark>, NodeError> {
let mut watermarks = Vec::new();
let peer = PrincipalId::new(peer_node.to_owned());
if can_sync_resource(store, &peer, "resource:ssh:certs", "ssh_cert.sync")? {
let cert_high = store
.list_ssh_cert_requests()?
.into_iter()
.map(|request| request.created_at_ms)
.chain(
store
.list_ssh_certificates()?
.into_iter()
.map(|certificate| certificate.imported_at_ms),
)
.max()
.unwrap_or(0);
watermarks.push(SyncWatermark {
stream: "ssh-certs".to_owned(),
high_water: cert_high,
});
}
if can_sync_resource(
store,
&peer,
"resource:ssh:revocations",
"ssh_revocation.sync",
)? {
let revocation_high = store
.list_ssh_revocations()?
.into_iter()
.map(|revocation| revocation.created_at_ms)
.max()
.unwrap_or(0);
watermarks.push(SyncWatermark {
stream: "ssh-revocations".to_owned(),
high_water: revocation_high,
});
}
for kv in store.list_kv_stores()? {
if can_sync_resource(store, &peer, &kv.resource_id, "kv.read")? {
let high_water = store
.list_kv_entries_since(&kv.kv_id, 0)?
.into_iter()
.map(|entry| entry.updated_at_ms)
.max()
.unwrap_or(0);
watermarks.push(SyncWatermark {
stream: format!("kv:{}", kv.name),
high_water,
});
}
}
for document in store.list_document_resources()? {
if can_sync_resource(store, &peer, &document.resource_id, "document.read")? {
watermarks.push(SyncWatermark {
stream: format!("document:{}", document.name),
high_water: document.updated_at_ms,
});
}
}
for db in store.list_db_resources()? {
if can_sync_resource(store, &peer, &db.resource_id, "db.sync")? {
if let Ok(metadata) = geth_db::crsqlite_change_metadata(Path::new(&db.path)) {
watermarks.push(SyncWatermark {
stream: format!("db:{}", db.name),
high_water: metadata.max_db_version.unwrap_or(0),
});
}
}
}
watermarks.sort_by(|left, right| left.stream.cmp(&right.stream));
Ok(watermarks)
}
fn can_sync_resource(
store: &Store,
peer: &PrincipalId,
resource: &str,
capability: &str,
) -> Result<bool, NodeError> {
Ok(geth_auth::explain_auth_ops(
&load_auth_ops_for_resource(store, resource)?,
peer.clone(),
ResourceId::new(resource.to_owned()),
Capability::new(capability.to_owned()),
)
.allowed)
}
async fn request_peer_control(
node: &LocalNode,
peer_node: &str,
@ -1420,6 +1560,10 @@ async fn request_peer_control(
| PeerControlResponse::DbSynced {
nonce: response_nonce,
..
}
| PeerControlResponse::SyncStatus {
nonce: response_nonce,
..
} if response_nonce == &nonce => Ok(response),
PeerControlResponse::Error { .. } => Ok(response),
_ => Err(NodeError::IrohPeer(format!(
@ -1471,25 +1615,68 @@ async fn run_live_sync_once(node: &LocalNode) -> Result<(), NodeError> {
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 {
if let Err(error) = ssh_cert_sync_from_peer(node, &peer.peer_id).await {
tracing::debug!(peer = %peer.peer_id, %error, "SSH cert live sync failed");
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::<BTreeMap<_, _>>()
},
) {
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 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).await {
tracing::debug!(peer = %peer.peer_id, %error, "SSH cert live sync failed");
}
}
if let Err(error) = ssh_revocation_sync_from_peer(node, &peer.peer_id).await {
tracing::debug!(peer = %peer.peer_id, %error, "SSH revocation 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).await {
tracing::debug!(peer = %peer.peer_id, %error, "SSH revocation live sync failed");
}
}
for kv in &kv_stores {
if let Err(error) = kv_sync_from_peer(node, &peer.peer_id, &kv.name).await {
tracing::debug!(peer = %peer.peer_id, kv = %kv.name, %error, "KV live sync failed");
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).await {
tracing::debug!(peer = %peer.peer_id, kv = %kv.name, %error, "KV live sync failed");
}
}
}
for document in &documents {
if let Err(error) = document_sync_from_peer(node, &peer.peer_id, &document.name).await {
tracing::debug!(peer = %peer.peer_id, document = %document.name, %error, "document live sync failed");
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).await
{
tracing::debug!(peer = %peer.peer_id, document = %document.name, %error, "document live sync failed");
}
}
}
for db in &dbs {
if let Err(error) = db_sync_from_peer(node, &peer.peer_id, &db.name, 100).await {
tracing::debug!(peer = %peer.peer_id, db = %db.name, %error, "DB live sync failed");
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).await {
tracing::debug!(peer = %peer.peer_id, db = %db.name, %error, "DB live sync failed");
}
}
}
}
@ -1546,6 +1733,31 @@ async fn handle_iroh_control_connection(
note: "peer endpoint authenticated by Iroh and peer-card signature; candidate status does not grant resource capabilities".to_owned(),
}
}
PeerControlRequest::SyncStatus { peer_card, nonce } => {
peer_card.validate_candidate()?;
ensure_peer_card_matches_endpoint(&peer_card, &remote_endpoint_id)?;
let discovered = DiscoveredPeer::candidate(
peer_card.clone(),
UnixMillis(geth_store::now_ms()),
DiscoverySource::PeerExchange,
)?;
let store = Store::open(&node.paths.metadata_db())?;
store.upsert_peer_card(&StoredPeerCard {
peer_id: peer_card.node_id.to_string(),
card_json: serde_json::to_string(&peer_card)?,
updated_at_ms: discovered.discovered_at.0,
})?;
let watermarks = sync_watermarks_for_peer(&store, peer_card.node_id.as_str())?;
PeerControlResponse::SyncStatus {
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,
watermarks,
nonce,
note: "sync status authenticated endpoint/card binding and returns only streams for capabilities already granted to the caller".to_owned(),
}
}
PeerControlRequest::AuthCheck {
peer_card,
resource,
@ -3707,6 +3919,114 @@ mod tests {
.expect("insert mock crsqlite change");
}
fn grant_test_capability(store: &Store, peer: &str, resource: &str, capability: &str) {
let created_at = UnixMillis(geth_store::now_ms());
store_auth_op(
store,
&AuthOp {
id: generated_auth_op_id("grant-create", resource, capability, created_at),
resource: ResourceId::new(resource.to_owned()),
created_at,
kind: AuthOpKind::GrantCreate {
grant_id: format!("grant:{peer}:{resource}:{capability}"),
principal: PrincipalId::new(peer.to_owned()),
capabilities: vec![Capability::new(capability.to_owned())],
},
},
)
.expect("store auth grant");
}
#[test]
fn sync_watermarks_include_only_authorized_streams() {
let store = Store::open_memory().expect("open");
store
.insert_resource(&StoredResource {
resource_id: "resource:kv:prefs".to_owned(),
kind: "kv".to_owned(),
name: "prefs".to_owned(),
status: "active".to_owned(),
})
.expect("insert kv resource");
store
.insert_kv_store(&StoredKvStore {
kv_id: "kv:prefs".to_owned(),
resource_id: "resource:kv:prefs".to_owned(),
name: "prefs".to_owned(),
})
.expect("insert kv");
store
.set_kv_entry(&StoredKvEntry {
kv_id: "kv:prefs".to_owned(),
key: "theme".to_owned(),
value: "dark".to_owned(),
updated_at_ms: 42,
})
.expect("set kv");
store
.insert_resource(&StoredResource {
resource_id: "resource:document:notes".to_owned(),
kind: "document".to_owned(),
name: "notes".to_owned(),
status: "active".to_owned(),
})
.expect("insert document resource");
store
.insert_document_resource(&StoredDocumentResource {
document_id: "document:notes".to_owned(),
resource_id: "resource:document:notes".to_owned(),
name: "notes".to_owned(),
state_json: "{}".to_owned(),
updated_at_ms: 99,
})
.expect("insert document");
let dir = tempfile::tempdir().expect("tempdir");
let db_path = dir.path().join("notes.sqlite");
create_mock_crsqlite_db(&db_path, Some(7));
store
.insert_resource(&StoredResource {
resource_id: "resource:db:notes".to_owned(),
kind: "db".to_owned(),
name: "notes".to_owned(),
status: "active".to_owned(),
})
.expect("insert db resource");
store
.insert_db_resource(&StoredDbResource {
db_id: "db:notes".to_owned(),
resource_id: "resource:db:notes".to_owned(),
name: "notes".to_owned(),
path: db_path.display().to_string(),
})
.expect("insert db");
grant_test_capability(&store, "node:left", "resource:kv:prefs", "kv.read");
grant_test_capability(&store, "node:left", "resource:db:notes", "db.sync");
let watermarks = sync_watermarks_for_peer(&store, "node:left").expect("watermarks");
assert!(watermarks.contains(&SyncWatermark {
stream: "kv:prefs".to_owned(),
high_water: 42,
}));
assert!(watermarks.contains(&SyncWatermark {
stream: "db:notes".to_owned(),
high_water: 7,
}));
assert!(
!watermarks
.iter()
.any(|watermark| watermark.stream == "document:notes")
);
assert!(
!watermarks
.iter()
.any(|watermark| watermark.stream == "ssh-certs")
);
}
#[test]
fn lan_discovery_address_selection_uses_iroh_direct_addresses() {
let key = AgentKey::generate();