Live sync KV stores over Iroh

This commit is contained in:
Eric Wendland 2026-05-18 18:36:03 +02:00
commit f85039367c
9 changed files with 473 additions and 13 deletions

View file

@ -20,7 +20,7 @@ use geth_discovery::{
use geth_document::{DocumentResource, DocumentState};
use geth_iroh::{EndpointStatus, GethIrohConfig, GethIrohEndpoint, GethRelayMode};
use geth_keychain::{KeychainOp, KeychainOpKind};
use geth_kv::{KvEntry, KvResource};
use geth_kv::{KvEntry, KvResource, KvSyncEntry};
use geth_pipe::{PipeConnection, PipeListener};
use geth_pubsub::PubsubMessage;
use geth_resource::ResourceDescriptor;
@ -259,6 +259,10 @@ pub async fn handle_request_async(
ControlRequest::SshRevocationSync { node: peer_node } => {
ssh_revocation_sync_from_peer(node, &peer_node).await
}
ControlRequest::KvSync {
node: peer_node,
name,
} => kv_sync_from_peer(node, &peer_node, &name).await,
other => handle_request(node, other),
}
}
@ -534,7 +538,8 @@ async fn peer_ping(node: &LocalNode, peer_node: &str) -> Result<ControlResponse,
)),
PeerControlResponse::CasFetched { .. }
| PeerControlResponse::SshCertSynced { .. }
| PeerControlResponse::SshRevocationSynced { .. } => Err(NodeError::IrohPeer(
| PeerControlResponse::SshRevocationSynced { .. }
| PeerControlResponse::KvSynced { .. } => Err(NodeError::IrohPeer(
"peer returned wrong response type to ping request".to_owned(),
)),
}
@ -641,7 +646,8 @@ async fn peer_auth_check(
)),
PeerControlResponse::CasFetched { .. }
| PeerControlResponse::SshCertSynced { .. }
| PeerControlResponse::SshRevocationSynced { .. } => Err(NodeError::IrohPeer(
| PeerControlResponse::SshRevocationSynced { .. }
| PeerControlResponse::KvSynced { .. } => Err(NodeError::IrohPeer(
"peer returned wrong response type to auth-check request".to_owned(),
)),
}
@ -777,7 +783,8 @@ async fn cas_fetch_from_peer(
PeerControlResponse::Pong { .. }
| PeerControlResponse::AuthChecked { .. }
| PeerControlResponse::SshCertSynced { .. }
| PeerControlResponse::SshRevocationSynced { .. } => Err(NodeError::IrohPeer(
| PeerControlResponse::SshRevocationSynced { .. }
| PeerControlResponse::KvSynced { .. } => Err(NodeError::IrohPeer(
"peer returned wrong response type to CAS fetch".to_owned(),
)),
}
@ -919,6 +926,90 @@ async fn ssh_revocation_sync_from_peer(
}
}
async fn kv_sync_from_peer(
node: &LocalNode,
peer_node: &str,
name: &str,
) -> Result<ControlResponse, NodeError> {
geth_kv::validate_kv_name(name).map_err(|_| NodeError::InvalidKvName(name.to_owned()))?;
let stream = format!("kv:{name}");
let since_ms =
load_live_sync_cursor(&Store::open(&node.paths.metadata_db())?, peer_node, &stream)?;
let response = request_peer_control(node, peer_node, "kv-sync", |peer_card, nonce| {
PeerControlRequest::KvSync {
peer_card,
name: name.to_owned(),
since_ms,
nonce,
}
})
.await?;
match response {
PeerControlResponse::KvSynced {
node_id,
agent_id,
endpoint_id,
name: response_name,
entries,
high_water_ms,
allowed,
reason,
note,
..
} if response_name == name => {
if !allowed {
return Ok(ControlResponse::KvSynced {
peer_node_id: node_id,
peer_agent_id: agent_id,
endpoint_id,
name: response_name,
entries_imported: 0,
allowed,
reason,
note,
});
}
let store = Store::open(&node.paths.metadata_db())?;
let kv = ensure_local_kv_store(&store, name)?;
let mut entries_imported = 0;
for entry in entries {
geth_kv::validate_kv_key(&entry.key)
.map_err(|_| NodeError::InvalidKvKey(entry.key.clone()))?;
let should_import = store
.get_kv_entry(&kv.kv_id, &entry.key)?
.is_none_or(|local| entry.updated_at_ms >= local.updated_at_ms);
if should_import {
store.set_kv_entry(&StoredKvEntry {
kv_id: kv.kv_id.clone(),
key: entry.key,
value: entry.value,
updated_at_ms: entry.updated_at_ms,
})?;
entries_imported += 1;
}
}
store_live_sync_cursor(&store, peer_node, &stream, high_water_ms)?;
Ok(ControlResponse::KvSynced {
peer_node_id: node_id,
peer_agent_id: agent_id,
endpoint_id,
name: response_name,
entries_imported,
allowed,
reason,
note,
})
}
PeerControlResponse::KvSynced { .. } => Err(NodeError::IrohPeer(
"peer KV sync response did not match request".to_owned(),
)),
PeerControlResponse::Error { message } => Err(NodeError::IrohPeer(message)),
_ => Err(NodeError::IrohPeer(
"peer returned wrong response type to KV sync".to_owned(),
)),
}
}
fn live_sync_cursor_key(peer_node: &str, stream: &str) -> String {
format!("live-sync:{peer_node}:{stream}")
}
@ -1011,6 +1102,10 @@ async fn request_peer_control(
| PeerControlResponse::SshRevocationSynced {
nonce: response_nonce,
..
}
| PeerControlResponse::KvSynced {
nonce: response_nonce,
..
} if response_nonce == &nonce => Ok(response),
PeerControlResponse::Error { .. } => Ok(response),
_ => Err(NodeError::IrohPeer(format!(
@ -1058,6 +1153,7 @@ 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()?;
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");
@ -1065,6 +1161,11 @@ async fn run_live_sync_once(node: &LocalNode) -> Result<(), NodeError> {
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");
}
}
}
Ok(())
}
@ -1334,6 +1435,68 @@ async fn handle_iroh_control_connection(
note: "SSH revocation sync authenticated endpoint/card binding and required ssh_revocation.sync on resource:ssh:revocations".to_owned(),
}
}
PeerControlRequest::KvSync {
peer_card,
name,
since_ms,
nonce,
} => {
geth_kv::validate_kv_name(&name).map_err(|_| NodeError::InvalidKvName(name.clone()))?;
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,
})?;
if let Some(kv) = store.get_kv_store_by_name(&name)? {
let capability = "kv.read".to_owned();
let high_water_ms = geth_store::now_ms();
let explanation = geth_auth::explain_auth_ops(
&load_auth_ops_for_resource(&store, &kv.resource_id)?,
PrincipalId::new(peer_card.node_id.to_string()),
ResourceId::new(kv.resource_id.clone()),
Capability::new(capability),
);
let entries = if explanation.allowed {
store
.list_kv_entries_since(&kv.kv_id, since_ms)?
.into_iter()
.map(|entry| KvSyncEntry {
key: entry.key,
value: entry.value,
updated_at_ms: entry.updated_at_ms,
})
.collect()
} else {
Vec::new()
};
PeerControlResponse::KvSynced {
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,
name,
entries,
high_water_ms,
allowed: explanation.allowed,
reason: explanation.reason,
evaluated_ops: explanation.evaluated_ops,
nonce,
note: "KV sync authenticated endpoint/card binding and required kv.read on the remote KV resource".to_owned(),
}
} else {
PeerControlResponse::Error {
message: format!("kv store not found: {name}"),
}
}
}
};
send.write_all(geth_control::encode_peer_response(&response)?.as_bytes())
.await
@ -1472,6 +1635,7 @@ pub fn handle_request(
ControlRequest::CasFetch { .. } => Err(NodeError::IrohEndpointUnavailable),
ControlRequest::SshCertSync { .. } => Err(NodeError::IrohEndpointUnavailable),
ControlRequest::SshRevocationSync { .. } => Err(NodeError::IrohEndpointUnavailable),
ControlRequest::KvSync { .. } => Err(NodeError::IrohEndpointUnavailable),
ControlRequest::ResourceList => Ok(ControlResponse::ResourceList {
resources: store
.list_resources()?
@ -2482,6 +2646,19 @@ fn kv_entry_from_stored(stored: StoredKvEntry) -> KvEntry {
}
}
fn ensure_local_kv_store(store: &Store, name: &str) -> Result<StoredKvStore, NodeError> {
if let Some(kv) = store.get_kv_store_by_name(name)? {
return Ok(kv);
}
let stored = StoredKvStore {
kv_id: format!("kv:{name}"),
resource_id: format!("resource:kv:{name}"),
name: name.to_owned(),
};
store.insert_kv_store(&stored)?;
Ok(stored)
}
fn document_resource_from_stored(stored: &StoredDocumentResource) -> DocumentResource {
DocumentResource {
id: stored.document_id.clone().into(),
@ -2985,6 +3162,23 @@ mod tests {
ControlResponse::SshRevocationAdded { revocation } => revocation.id,
other => panic!("unexpected SSH revocation response: {other:?}"),
};
handle_request(
&right,
ControlRequest::KvCreate {
name: "prefs".to_owned(),
},
)
.expect("right kv create");
handle_request(
&right,
ControlRequest::KvSet {
name: "prefs".to_owned(),
key: "apps/foo/theme".to_owned(),
value: "dark".to_owned(),
subject: None,
},
)
.expect("right kv set");
let ping = handle_request_async(
&left,
@ -3063,6 +3257,29 @@ mod tests {
other => panic!("unexpected denied CAS fetch response: {other:?}"),
}
let denied_kv_sync = handle_request_async(
&left,
ControlRequest::KvSync {
node: right_card.node_id.to_string(),
name: "prefs".to_owned(),
},
)
.await
.expect("denied kv sync");
match denied_kv_sync {
ControlResponse::KvSynced {
allowed,
entries_imported,
reason,
..
} => {
assert!(!allowed);
assert_eq!(entries_imported, 0);
assert!(reason.contains("no active direct or group grant"));
}
other => panic!("unexpected denied KV sync response: {other:?}"),
}
handle_request(
&right,
ControlRequest::AuthGrant {
@ -3073,6 +3290,16 @@ mod tests {
},
)
.expect("grant left peer");
handle_request(
&right,
ControlRequest::AuthGrant {
subject: left.node_id.clone(),
resource: "resource:kv:prefs".to_owned(),
capability: "kv.read".to_owned(),
grant_id: Some("grant:left-kv-read".to_owned()),
},
)
.expect("grant left kv read");
let allowed = handle_request_async(
&left,
@ -3131,6 +3358,45 @@ mod tests {
other => panic!("unexpected allowed CAS fetch response: {other:?}"),
}
let kv_synced = handle_request_async(
&left,
ControlRequest::KvSync {
node: right_card.node_id.to_string(),
name: "prefs".to_owned(),
},
)
.await
.expect("allowed kv sync");
match kv_synced {
ControlResponse::KvSynced {
allowed,
entries_imported,
reason,
note,
..
} => {
assert!(allowed);
assert_eq!(entries_imported, 1);
assert!(reason.contains("direct grant"));
assert!(note.contains("kv.read"));
}
other => panic!("unexpected allowed KV sync response: {other:?}"),
}
let synced_theme = handle_request(
&left,
ControlRequest::KvGet {
name: "prefs".to_owned(),
key: "apps/foo/theme".to_owned(),
},
)
.expect("left kv get synced");
match synced_theme {
ControlResponse::KvGet { entry: Some(entry) } => {
assert_eq!(entry.value, "dark");
}
other => panic!("unexpected synced KV get response: {other:?}"),
}
let denied_cert_sync = handle_request_async(
&left,
ControlRequest::SshCertSync {
@ -3280,6 +3546,16 @@ mod tests {
ControlResponse::SshRevocationAdded { revocation } => revocation.id,
other => panic!("unexpected second SSH revocation response: {other:?}"),
};
handle_request(
&right,
ControlRequest::KvSet {
name: "prefs".to_owned(),
key: "apps/foo/live".to_owned(),
value: "synced".to_owned(),
subject: None,
},
)
.expect("right second kv set");
run_live_sync_once(&left)
.await
@ -3299,6 +3575,13 @@ mod tests {
.iter()
.any(|revocation| revocation.revocation_id == second_revocation_id.as_str())
);
assert_eq!(
left_store
.get_kv_entry("kv:prefs", "apps/foo/live")
.expect("get live-synced kv")
.map(|entry| entry.value),
Some("synced".to_owned())
);
left_endpoint.shutdown().await;
right_endpoint.shutdown().await;