pub mod service; use base64::Engine; use geth_auth::{AuthExplanation, AuthOp, AuthOpKind}; use geth_cas::{ BlobInfoSummary, FileConflict, FileConflictKind, FileConflictResolution, FileConflictStatus, FileRoot, FileRootScan, LocalCas, hash_path, }; use geth_config::{GethConfig, GethPaths, RelayMode}; use geth_control::{ CasBlob, CasProvider, ControlRequest, ControlResponse, KeychainStatusResponse, NodeIdResponse, PeerControlRequest, PeerControlResponse, StatusResponse, SyncWatermark, }; use geth_crypto::AgentKey; use geth_db::DbResource; use geth_discovery::{ DiscoveredPeer, DiscoverySource, EndpointCandidate, PEER_CARD_LAN_DISCOVERY_SERVICE, PeerCard, discovery_is_untrusted_note, peer_card_from_txt_attributes, peer_card_txt_attributes, }; use geth_document::{DocumentResource, DocumentState}; use geth_iroh::{EndpointStatus, GethIrohConfig, GethIrohEndpoint, GethRelayMode}; use geth_keychain::{KeychainOp, KeychainOpKind}; use geth_kv::{KvEntry, KvResource, KvSyncEntry}; use geth_pipe::{PipeConnection, PipeListener}; use geth_pubsub::PubsubMessage; use geth_resource::ResourceDescriptor; use geth_secrets::{BearerAccess, ResourceMasterSecret}; use geth_ssh_identity::{ SshCertApproval, SshCertKind, SshCertRequest, SshCertRequestStatus, SshCertificateRecord, SshRevocationEntry, SshRevocationExportFormat, SshRevocationKind, build_ssh_cert_sign_command, cert_request_id, certificate_id, openssh_krl_spec, parse_openssh_krl_spec, revocation_id, ssh_public_key_fingerprint, write_openssh_krl, }; use geth_ssh_proxy::SshProxyConnection; use geth_store::{ Store, StoredAuthOp, StoredDbResource, StoredDocumentResource, StoredFileConflict, StoredFileRoot, StoredKeychainOp, StoredKvEntry, StoredKvStore, StoredModuleState, StoredPeerCard, StoredResource, StoredResourceSecret, StoredSshCertRequest, StoredSshCertificate, StoredSshRevocation, }; use geth_types::{ AuthOpId, BlobHash, Capability, KeyId, NodeId, PrincipalId, ResourceId, ResourceKind, ResourceName, SshCertId, SshCertRequestId, UnixMillis, }; use std::collections::{BTreeMap, VecDeque}; use std::path::Path; use std::sync::{Arc, Mutex}; use std::time::Duration; use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; use tokio::net::{UnixListener, UnixStream}; #[derive(Debug, thiserror::Error)] pub enum NodeError { #[error("config error: {0}")] Config(#[from] geth_config::ConfigError), #[error("crypto error: {0}")] Crypto(#[from] geth_crypto::CryptoError), #[error("store error: {0}")] Store(#[from] geth_store::StoreError), #[error("cas error: {0}")] Cas(#[from] geth_cas::CasError), #[error("db error: {0}")] Db(#[from] geth_db::DbError), #[error("control error: {0}")] Control(#[from] geth_control::ControlError), #[error("json error: {0}")] Json(#[from] serde_json::Error), #[error("io error: {0}")] Io(#[from] std::io::Error), #[error("invalid resource kind: {0}")] InvalidResourceKind(String), #[error("invalid db resource name: {0}")] InvalidDbName(String), #[error("db path does not exist or is not a file: {0}")] InvalidDbPath(String), #[error("db resource not found: {0}")] DbNotFound(String), #[error("invalid kv store name: {0}")] InvalidKvName(String), #[error("invalid kv key: {0}")] InvalidKvKey(String), #[error("kv store not found: {0}")] KvNotFound(String), #[error("invalid document name: {0}")] InvalidDocumentName(String), #[error("document not found: {0}")] DocumentNotFound(String), #[error("document error: {0}")] Document(#[from] geth_document::DocumentError), #[error("resource not found: {0}")] ResourceNotFound(String), #[error("unauthorized: {0}")] Unauthorized(String), #[error("secrets error: {0}")] Secrets(#[from] geth_secrets::SecretsError), #[error("pubsub error: {0}")] Pubsub(#[from] geth_pubsub::PubsubError), #[error("pipe error: {0}")] Pipe(#[from] geth_pipe::PipeError), #[error("runtime state lock poisoned")] RuntimeLockPoisoned, #[error("invalid ssh certificate kind: {0}")] InvalidSshCertKind(String), #[error("invalid ssh certificate request status: {0}")] InvalidSshCertStatus(String), #[error("invalid ssh revocation kind: {0}")] InvalidSshRevocationKind(String), #[error("ssh certificate request not found: {0}")] SshCertRequestNotFound(String), #[error("ssh certificate request must include at least one principal")] MissingSshCertPrincipal, #[error("ssh certificate flow error: {0}")] SshCertFlow(#[from] geth_ssh_identity::SshCertFlowError), #[error("ssh identity error: {0}")] SshIdentity(#[from] geth_ssh_identity::SshIdentityError), #[error("discovery error: {0}")] Discovery(#[from] geth_discovery::DiscoveryError), #[error("cannot export peer card before the daemon has an Iroh EndpointID")] IrohEndpointUnavailable, #[error("peer candidate not found: {0}")] PeerNotFound(String), #[error("iroh peer error: {0}")] IrohPeer(String), } #[derive(Clone, Debug)] pub struct LocalNode { pub paths: GethPaths, pub agent_id: String, pub node_id: String, pub iroh_status: EndpointStatus, iroh_endpoint: Arc>>, runtime: Arc, } #[derive(Debug)] struct NodeRuntime { pubsub: Mutex, pipes: Mutex, } #[derive(Debug, Default)] struct PubsubRuntime { messages: VecDeque, } #[derive(Debug, Default)] struct PipeRuntime { listeners: BTreeMap, connections: VecDeque, } const PUBSUB_RING_LIMIT: usize = 256; const LIVE_SYNC_INTERVAL: Duration = Duration::from_secs(30); const PIPE_CONNECTION_RING_LIMIT: usize = 256; #[derive(Debug, serde::Deserialize, serde::Serialize)] struct LiveSyncCursor { cursor_ms: i64, } pub fn init_node(paths: &GethPaths) -> Result { paths.ensure_base_dirs()?; if !paths.config_file().exists() { std::fs::write(paths.config_file(), GethConfig::default_toml())?; } let key = AgentKey::load_or_create(&paths.agent_key())?; let agent_id = key.agent_id().to_string(); let node_id = stable_node_id(&agent_id); let store = Store::open(&paths.metadata_db())?; store.upsert_agent(&agent_id, &key.public_key_hex())?; store.upsert_node(&node_id, "local", &agent_id)?; store.insert_resource(&StoredResource { resource_id: "resource:cas:local".to_owned(), kind: ResourceKind::Cas.to_string(), name: "local-cas".to_owned(), status: "active".to_owned(), })?; Ok(LocalNode { paths: paths.clone(), agent_id, node_id, iroh_status: EndpointStatus::scaffolded(), iroh_endpoint: Arc::new(Mutex::new(None)), runtime: Arc::new(NodeRuntime { pubsub: Mutex::new(PubsubRuntime::default()), pipes: Mutex::new(PipeRuntime::default()), }), }) } pub fn open_node(paths: &GethPaths) -> Result { init_node(paths) } pub async fn run_daemon(paths: GethPaths) -> Result<(), NodeError> { let mut node = init_node(&paths)?; let _iroh_endpoint = start_daemon_iroh_endpoint(&mut node).await?; if let Some(endpoint) = node .iroh_endpoint .lock() .map_err(|_| NodeError::RuntimeLockPoisoned)? .clone() { spawn_iroh_control_accept_loop(node.clone(), endpoint); spawn_background_live_sync(node.clone()); } let _peer_card_lan_discovery = start_peer_card_lan_discovery(&node).await; if Path::new(&paths.socket_path()).exists() { std::fs::remove_file(paths.socket_path())?; } let listener = UnixListener::bind(paths.socket_path())?; tracing::info!(socket = %paths.socket_path().display(), "geth daemon listening"); loop { let (stream, _) = listener.accept().await?; let node = node.clone(); tokio::spawn(async move { if let Err(error) = handle_stream(node, stream).await { tracing::warn!(%error, "control request failed"); } }); } } pub async fn send_control( paths: &GethPaths, request: ControlRequest, ) -> Result { let mut stream = UnixStream::connect(paths.socket_path()).await?; stream .write_all(geth_control::encode_request(&request)?.as_bytes()) .await?; stream.shutdown().await?; let mut reader = BufReader::new(stream); let mut line = String::new(); reader.read_line(&mut line).await?; Ok(geth_control::decode_response(&line)?) } pub async fn handle_request_async( node: &LocalNode, request: ControlRequest, ) -> Result { match request { ControlRequest::PeerCardExport { out } => export_peer_card(node, out, true).await, ControlRequest::PeerPing { node: peer_node } => peer_ping(node, &peer_node).await, ControlRequest::PeerAuthCheck { node: peer_node, resource, capability, } => peer_auth_check(node, &peer_node, resource, capability).await, ControlRequest::CasFetch { node: peer_node, hash, } => cas_fetch_from_peer(node, &peer_node, hash).await, ControlRequest::SshCertSync { node: peer_node } => { ssh_cert_sync_from_peer(node, &peer_node).await } 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, ControlRequest::DbSync { node: peer_node, name, limit, } => db_sync_from_peer(node, &peer_node, &name, limit).await, ControlRequest::PubsubPub { topic, message, node: Some(peer_node), } => pubsub_publish_to_peer(node, &peer_node, topic, message).await, ControlRequest::PubsubSub { topic, node: Some(peer_node), } => pubsub_subscribe_from_peer(node, &peer_node, topic).await, ControlRequest::PipeConnect { target, node: Some(peer_node), } => pipe_connect_to_peer(node, &peer_node, target).await, ControlRequest::SshProxyConnect { node: peer_node } => { ssh_proxy_connect_to_peer(node, &peer_node).await } ControlRequest::DocumentSync { node: peer_node, name, } => document_sync_from_peer(node, &peer_node, &name).await, other => handle_request(node, other), } } async fn handle_stream(node: LocalNode, stream: UnixStream) -> Result<(), NodeError> { let mut reader = BufReader::new(stream); let mut line = String::new(); reader.read_line(&mut line).await?; let request = geth_control::decode_request(&line)?; let response = match handle_request_async(&node, request).await { Ok(response) => response, Err(error) => ControlResponse::Error { message: error.to_string(), }, }; let mut stream = reader.into_inner(); stream .write_all(geth_control::encode_response(&response)?.as_bytes()) .await?; Ok(()) } async fn start_peer_card_lan_discovery(node: &LocalNode) -> Option { if !node.iroh_status.local_discovery || !node.iroh_status.enabled { return None; } let card = match local_peer_card(node, DiscoverySource::Mdns, true).await { Ok(card) => card, Err(error) => { tracing::warn!(%error, "signed peer-card LAN discovery disabled"); return None; } }; let attributes = match peer_card_txt_attributes(&card) { Ok(attributes) => attributes, Err(error) => { tracing::warn!(%error, "could not encode peer-card LAN TXT payload"); return None; } }; let Some((port, addrs)) = lan_discovery_addrs(&card) else { tracing::warn!( "signed peer-card LAN discovery disabled because no direct Iroh address is known" ); return None; }; let paths = node.paths.clone(); let self_node_id = node.node_id.clone(); let discoverer = match swarm_discovery::Discoverer::new_interactive( PEER_CARD_LAN_DISCOVERY_SERVICE.to_owned(), lan_discovery_peer_id(&node.agent_id), ) .with_addrs(port, addrs) .with_txt_attributes(attributes) { Ok(discoverer) => discoverer, Err(error) => { tracing::warn!(%error, "could not build peer-card LAN discovery"); return None; } } .with_callback(move |_peer_id, peer| { if peer.is_expiry() { return; } let card = match peer_card_from_txt_attributes(peer.txt_attributes()) { Ok(card) => card, Err(error) => { tracing::debug!(%error, "ignoring invalid peer-card LAN discovery payload"); return; } }; if card.node_id.as_str() == self_node_id { return; } let discovered = match DiscoveredPeer::candidate( card.clone(), UnixMillis(geth_store::now_ms()), DiscoverySource::Mdns, ) { Ok(discovered) => discovered, Err(error) => { tracing::debug!(%error, "ignoring invalid LAN peer-card candidate"); return; } }; match Store::open(&paths.metadata_db()).and_then(|store| { store.upsert_peer_card(&StoredPeerCard { peer_id: card.node_id.to_string(), card_json: serde_json::to_string(&card).map_err(geth_store::StoreError::from)?, updated_at_ms: discovered.discovered_at.0, }) }) { Ok(()) => tracing::debug!(node = %card.node_id, "stored LAN peer-card candidate"), Err(error) => tracing::warn!(%error, "could not store LAN peer-card candidate"), } }); match discoverer.spawn(&tokio::runtime::Handle::current()) { Ok(guard) => { tracing::info!( service = PEER_CARD_LAN_DISCOVERY_SERVICE, "signed peer-card LAN discovery running" ); Some(guard) } Err(error) => { tracing::warn!(%error, "could not start signed peer-card LAN discovery"); None } } } fn lan_discovery_peer_id(agent_id: &str) -> String { format!("geth-{}", geth_crypto::blake3_hex(agent_id.as_bytes())) } fn lan_discovery_addrs(card: &PeerCard) -> Option<(u16, Vec)> { let mut parsed = card .endpoints .iter() .flat_map(|endpoint| endpoint.direct_addresses.iter()) .filter_map(|addr| addr.parse::().ok()) .collect::>(); parsed.sort_unstable(); parsed.dedup(); let port = parsed.first()?.port(); let addrs = parsed .into_iter() .filter(|addr| addr.port() == port) .map(|addr| addr.ip()) .collect::>(); (!addrs.is_empty()).then_some((port, addrs)) } async fn export_peer_card( node: &LocalNode, out: Option, include_node_addr: bool, ) -> Result { let card = local_peer_card(node, DiscoverySource::Manual, include_node_addr).await?; if let Some(path) = &out { std::fs::write(path, serde_json::to_string_pretty(&card)?)?; } Ok(ControlResponse::PeerCardExported { card, out, note: discovery_is_untrusted_note().to_owned(), }) } async fn local_peer_card( node: &LocalNode, source: DiscoverySource, include_node_addr: bool, ) -> Result { let endpoint_id = node .iroh_status .endpoint_id .clone() .ok_or(NodeError::IrohEndpointUnavailable)?; let endpoint = node .iroh_endpoint .lock() .map_err(|_| NodeError::RuntimeLockPoisoned)? .clone(); let node_addr = if include_node_addr { match endpoint { Some(endpoint) => endpoint.node_addr_snapshot().await.ok(), None => None, } } else { None }; let candidate = EndpointCandidate { endpoint_id, relay_url: node_addr.as_ref().and_then(|addr| addr.relay_url.clone()), direct_addresses: node_addr .map(|addr| addr.direct_addresses) .unwrap_or_default(), source, }; let key = AgentKey::load(&node.paths.agent_key())?; PeerCard::signed( NodeId::new(node.node_id.clone()), &key, vec![candidate], UnixMillis(geth_store::now_ms()), ) .map_err(NodeError::from) } async fn peer_ping(node: &LocalNode, peer_node: &str) -> Result { let store = Store::open(&node.paths.metadata_db())?; let stored = store .get_peer_card(peer_node)? .ok_or_else(|| NodeError::PeerNotFound(peer_node.to_owned()))?; let peer_card: PeerCard = serde_json::from_str(&stored.card_json)?; peer_card.validate_candidate()?; let candidate = peer_card .endpoints .first() .ok_or(geth_discovery::DiscoveryError::MissingEndpoint)?; ensure_peer_card_matches_endpoint(&peer_card, &candidate.endpoint_id)?; let node_addr = iroh_node_addr_from_candidate(candidate)?; let endpoint = node .iroh_endpoint .lock() .map_err(|_| NodeError::RuntimeLockPoisoned)? .clone() .ok_or(NodeError::IrohEndpointUnavailable)?; let self_card = local_peer_card(node, DiscoverySource::PeerExchange, true).await?; let nonce = geth_crypto::blake3_hex( format!("{}\0{}\0{}", node.node_id, peer_node, geth_store::now_ms()).as_bytes(), ); let request = PeerControlRequest::Ping { peer_card: self_card, nonce: nonce.clone(), }; let conn = endpoint .endpoint() .connect(node_addr, geth_iroh::ALPN_CONTROL) .await .map_err(|error| NodeError::IrohPeer(error.to_string()))?; let alpn = conn .alpn() .map(display_alpn) .unwrap_or_else(|| "unknown".to_owned()); let (mut send, mut recv) = conn .open_bi() .await .map_err(|error| NodeError::IrohPeer(error.to_string()))?; send.write_all(geth_control::encode_peer_request(&request)?.as_bytes()) .await .map_err(|error| NodeError::IrohPeer(error.to_string()))?; send.finish() .map_err(|error| NodeError::IrohPeer(error.to_string()))?; let bytes = recv .read_to_end(64 * 1024) .await .map_err(|error| NodeError::IrohPeer(error.to_string()))?; let text = std::str::from_utf8(&bytes).map_err(|error| NodeError::IrohPeer(error.to_string()))?; match geth_control::decode_peer_response(text)? { PeerControlResponse::Pong { node_id, agent_id, endpoint_id, alpn: remote_alpn, nonce: response_nonce, note, .. } if response_nonce == nonce => Ok(ControlResponse::PeerPinged { peer_node_id: node_id, peer_agent_id: agent_id, endpoint_id, alpn: if remote_alpn == "unknown" { alpn } else { remote_alpn }, note, }), PeerControlResponse::Pong { .. } => Err(NodeError::IrohPeer( "peer ping response nonce did not match request".to_owned(), )), PeerControlResponse::Error { message } => Err(NodeError::IrohPeer(message)), PeerControlResponse::AuthChecked { .. } => Err(NodeError::IrohPeer( "peer returned auth-check response to ping request".to_owned(), )), PeerControlResponse::SyncStatus { .. } | PeerControlResponse::CasFetched { .. } | PeerControlResponse::SshCertSynced { .. } | PeerControlResponse::SshRevocationSynced { .. } | PeerControlResponse::KvSynced { .. } | PeerControlResponse::PubsubPublished { .. } | PeerControlResponse::PubsubSubscribed { .. } | PeerControlResponse::PipeConnected { .. } | PeerControlResponse::SshProxyConnected { .. } | PeerControlResponse::DocumentSynced { .. } | PeerControlResponse::DbSynced { .. } => Err(NodeError::IrohPeer( "peer returned wrong response type to ping request".to_owned(), )), } } async fn peer_auth_check( node: &LocalNode, peer_node: &str, resource: String, capability: String, ) -> Result { let store = Store::open(&node.paths.metadata_db())?; let stored = store .get_peer_card(peer_node)? .ok_or_else(|| NodeError::PeerNotFound(peer_node.to_owned()))?; let peer_card: PeerCard = serde_json::from_str(&stored.card_json)?; peer_card.validate_candidate()?; let candidate = peer_card .endpoints .first() .ok_or(geth_discovery::DiscoveryError::MissingEndpoint)?; ensure_peer_card_matches_endpoint(&peer_card, &candidate.endpoint_id)?; let node_addr = iroh_node_addr_from_candidate(candidate)?; let endpoint = node .iroh_endpoint .lock() .map_err(|_| NodeError::RuntimeLockPoisoned)? .clone() .ok_or(NodeError::IrohEndpointUnavailable)?; let self_card = local_peer_card(node, DiscoverySource::PeerExchange, true).await?; let nonce = geth_crypto::blake3_hex( format!( "{}\0{}\0{}\0{}\0{}", node.node_id, peer_node, resource, capability, geth_store::now_ms() ) .as_bytes(), ); let request = PeerControlRequest::AuthCheck { peer_card: self_card, resource: resource.clone(), capability: capability.clone(), nonce: nonce.clone(), }; let conn = endpoint .endpoint() .connect(node_addr, geth_iroh::ALPN_CONTROL) .await .map_err(|error| NodeError::IrohPeer(error.to_string()))?; let (mut send, mut recv) = conn .open_bi() .await .map_err(|error| NodeError::IrohPeer(error.to_string()))?; send.write_all(geth_control::encode_peer_request(&request)?.as_bytes()) .await .map_err(|error| NodeError::IrohPeer(error.to_string()))?; send.finish() .map_err(|error| NodeError::IrohPeer(error.to_string()))?; let bytes = recv .read_to_end(64 * 1024) .await .map_err(|error| NodeError::IrohPeer(error.to_string()))?; let text = std::str::from_utf8(&bytes).map_err(|error| NodeError::IrohPeer(error.to_string()))?; match geth_control::decode_peer_response(text)? { PeerControlResponse::AuthChecked { node_id, agent_id, endpoint_id, resource: response_resource, capability: response_capability, allowed, reason, evaluated_ops, nonce: response_nonce, note, .. } if response_nonce == nonce && response_resource == resource && response_capability == capability => { Ok(ControlResponse::PeerAuthChecked { peer_node_id: node_id, peer_agent_id: agent_id, endpoint_id, resource: response_resource, capability: response_capability, allowed, reason, evaluated_ops, note, }) } PeerControlResponse::AuthChecked { .. } => Err(NodeError::IrohPeer( "peer auth-check response did not match request".to_owned(), )), PeerControlResponse::Error { message } => Err(NodeError::IrohPeer(message)), PeerControlResponse::Pong { .. } => Err(NodeError::IrohPeer( "peer returned pong to auth-check request".to_owned(), )), PeerControlResponse::SyncStatus { .. } | PeerControlResponse::CasFetched { .. } | PeerControlResponse::SshCertSynced { .. } | PeerControlResponse::SshRevocationSynced { .. } | PeerControlResponse::KvSynced { .. } | PeerControlResponse::PubsubPublished { .. } | PeerControlResponse::PubsubSubscribed { .. } | PeerControlResponse::PipeConnected { .. } | PeerControlResponse::SshProxyConnected { .. } | PeerControlResponse::DocumentSynced { .. } | PeerControlResponse::DbSynced { .. } => Err(NodeError::IrohPeer( "peer returned wrong response type to auth-check request".to_owned(), )), } } async fn cas_fetch_from_peer( node: &LocalNode, peer_node: &str, hash: BlobHash, ) -> Result { let store = Store::open(&node.paths.metadata_db())?; let stored = store .get_peer_card(peer_node)? .ok_or_else(|| NodeError::PeerNotFound(peer_node.to_owned()))?; let peer_card: PeerCard = serde_json::from_str(&stored.card_json)?; peer_card.validate_candidate()?; let candidate = peer_card .endpoints .first() .ok_or(geth_discovery::DiscoveryError::MissingEndpoint)?; ensure_peer_card_matches_endpoint(&peer_card, &candidate.endpoint_id)?; let node_addr = iroh_node_addr_from_candidate(candidate)?; let endpoint = node .iroh_endpoint .lock() .map_err(|_| NodeError::RuntimeLockPoisoned)? .clone() .ok_or(NodeError::IrohEndpointUnavailable)?; let self_card = local_peer_card(node, DiscoverySource::PeerExchange, true).await?; let nonce = geth_crypto::blake3_hex( format!( "{}\0{}\0{}\0{}", node.node_id, peer_node, hash.as_str(), geth_store::now_ms() ) .as_bytes(), ); let request = PeerControlRequest::CasFetch { peer_card: self_card, hash: hash.clone(), nonce: nonce.clone(), }; let conn = endpoint .endpoint() .connect(node_addr, geth_iroh::ALPN_CONTROL) .await .map_err(|error| NodeError::IrohPeer(error.to_string()))?; let (mut send, mut recv) = conn .open_bi() .await .map_err(|error| NodeError::IrohPeer(error.to_string()))?; send.write_all(geth_control::encode_peer_request(&request)?.as_bytes()) .await .map_err(|error| NodeError::IrohPeer(error.to_string()))?; send.finish() .map_err(|error| NodeError::IrohPeer(error.to_string()))?; let bytes = recv .read_to_end(64 * 1024 * 1024) .await .map_err(|error| NodeError::IrohPeer(error.to_string()))?; let text = std::str::from_utf8(&bytes).map_err(|error| NodeError::IrohPeer(error.to_string()))?; match geth_control::decode_peer_response(text)? { PeerControlResponse::CasFetched { node_id, agent_id, endpoint_id, hash: response_hash, size_bytes, content_base64, allowed, reason, nonce: response_nonce, note, .. } if response_nonce == nonce && response_hash == hash => { if !allowed { return Ok(ControlResponse::CasFetched { peer_node_id: node_id, peer_agent_id: agent_id, endpoint_id, hash: response_hash, size_bytes: 0, allowed, reason, note, }); } let content = content_base64 .ok_or_else(|| NodeError::IrohPeer("peer omitted CAS content".to_owned())) .and_then(|content| { base64::engine::general_purpose::STANDARD .decode(content) .map_err(|error| NodeError::IrohPeer(error.to_string())) })?; if content.len() as u64 != size_bytes { return Err(NodeError::IrohPeer(format!( "peer announced {size_bytes} CAS bytes but returned {} bytes", content.len() ))); } let cas = LocalCas::new(node.paths.cas_dir()); let info = cas.add_bytes(&content)?; if info.hash != response_hash { return Err(NodeError::IrohPeer(format!( "peer returned content hash {} for requested {}", info.hash, response_hash ))); } store.record_cas_object( info.hash.as_str(), info.size_bytes, &info.path.to_string_lossy(), )?; store.record_cas_provider(info.hash.as_str(), &node_id, &endpoint_id)?; Ok(ControlResponse::CasFetched { peer_node_id: node_id, peer_agent_id: agent_id, endpoint_id, hash: info.hash, size_bytes: info.size_bytes, allowed, reason, note, }) } PeerControlResponse::CasFetched { .. } => Err(NodeError::IrohPeer( "peer CAS fetch response did not match request".to_owned(), )), PeerControlResponse::Error { message } => Err(NodeError::IrohPeer(message)), PeerControlResponse::Pong { .. } | PeerControlResponse::AuthChecked { .. } | PeerControlResponse::SyncStatus { .. } | PeerControlResponse::SshCertSynced { .. } | PeerControlResponse::SshRevocationSynced { .. } | PeerControlResponse::KvSynced { .. } | PeerControlResponse::PubsubPublished { .. } | PeerControlResponse::PubsubSubscribed { .. } | PeerControlResponse::PipeConnected { .. } | PeerControlResponse::SshProxyConnected { .. } | PeerControlResponse::DocumentSynced { .. } | PeerControlResponse::DbSynced { .. } => Err(NodeError::IrohPeer( "peer returned wrong response type to CAS fetch".to_owned(), )), } } async fn ssh_cert_sync_from_peer( node: &LocalNode, peer_node: &str, ) -> Result { let since_ms = load_live_sync_cursor( &Store::open(&node.paths.metadata_db())?, peer_node, "ssh-certs", )?; let response = request_peer_control(node, peer_node, "ssh-cert-sync", |peer_card, nonce| { PeerControlRequest::SshCertSync { peer_card, since_ms, nonce, } }) .await?; match response { PeerControlResponse::SshCertSynced { node_id, agent_id, endpoint_id, requests, certificates, high_water_ms, allowed, reason, note, .. } => { if !allowed { return Ok(ControlResponse::SshCertSynced { peer_node_id: node_id, peer_agent_id: agent_id, endpoint_id, requests_imported: 0, certificates_imported: 0, allowed, reason, note, }); } let store = Store::open(&node.paths.metadata_db())?; let requests_imported = requests.len(); let certificates_imported = certificates.len(); for request in &requests { store.insert_ssh_cert_request(&stored_from_ssh_cert_request(request))?; } for certificate in &certificates { store.insert_ssh_certificate(&stored_from_ssh_certificate(certificate))?; } store_live_sync_cursor(&store, peer_node, "ssh-certs", high_water_ms)?; Ok(ControlResponse::SshCertSynced { peer_node_id: node_id, peer_agent_id: agent_id, endpoint_id, requests_imported, certificates_imported, allowed, reason, note, }) } PeerControlResponse::Error { message } => Err(NodeError::IrohPeer(message)), _ => Err(NodeError::IrohPeer( "peer returned wrong response type to SSH cert sync".to_owned(), )), } } async fn ssh_revocation_sync_from_peer( node: &LocalNode, peer_node: &str, ) -> Result { let since_ms = load_live_sync_cursor( &Store::open(&node.paths.metadata_db())?, peer_node, "ssh-revocations", )?; let response = request_peer_control( node, peer_node, "ssh-revocation-sync", |peer_card, nonce| PeerControlRequest::SshRevocationSync { peer_card, since_ms, nonce, }, ) .await?; match response { PeerControlResponse::SshRevocationSynced { node_id, agent_id, endpoint_id, revocations, high_water_ms, allowed, reason, note, .. } => { if !allowed { return Ok(ControlResponse::SshRevocationSynced { peer_node_id: node_id, peer_agent_id: agent_id, endpoint_id, revocations_imported: 0, allowed, reason, note, }); } let store = Store::open(&node.paths.metadata_db())?; let revocations_imported = revocations.len(); for revocation in &revocations { store.insert_ssh_revocation(&stored_from_ssh_revocation(revocation))?; } store_live_sync_cursor(&store, peer_node, "ssh-revocations", high_water_ms)?; Ok(ControlResponse::SshRevocationSynced { peer_node_id: node_id, peer_agent_id: agent_id, endpoint_id, revocations_imported, allowed, reason, note, }) } PeerControlResponse::Error { message } => Err(NodeError::IrohPeer(message)), _ => Err(NodeError::IrohPeer( "peer returned wrong response type to SSH revocation sync".to_owned(), )), } } async fn kv_sync_from_peer( node: &LocalNode, peer_node: &str, name: &str, ) -> Result { 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(), )), } } async fn pubsub_publish_to_peer( node: &LocalNode, peer_node: &str, topic: String, message: String, ) -> Result { geth_pubsub::validate_topic(&topic)?; geth_pubsub::validate_message(&message)?; let response = request_peer_control(node, peer_node, "pubsub-publish", |peer_card, nonce| { PeerControlRequest::PubsubPublish { peer_card, topic: topic.clone(), message: message.clone(), nonce, } }) .await?; match response { PeerControlResponse::PubsubPublished { node_id, agent_id, endpoint_id, message, allowed, reason, note, .. } => { let message = message.unwrap_or_else(|| PubsubMessage { topic: topic.into(), message: String::new(), published_at: UnixMillis(0), }); Ok(ControlResponse::PubsubRemotePublished { peer_node_id: node_id, peer_agent_id: agent_id, endpoint_id, message, allowed, reason, note, }) } PeerControlResponse::Error { message } => Err(NodeError::IrohPeer(message)), _ => Err(NodeError::IrohPeer( "peer returned wrong response type to pubsub publish".to_owned(), )), } } async fn pubsub_subscribe_from_peer( node: &LocalNode, peer_node: &str, topic: String, ) -> Result { geth_pubsub::validate_topic(&topic)?; let response = request_peer_control(node, peer_node, "pubsub-subscribe", |peer_card, nonce| { PeerControlRequest::PubsubSubscribe { peer_card, topic: topic.clone(), nonce, } }) .await?; match response { PeerControlResponse::PubsubSubscribed { node_id, agent_id, endpoint_id, topic: response_topic, messages, allowed, reason, note, .. } if response_topic == topic => Ok(ControlResponse::PubsubRemoteMessages { peer_node_id: node_id, peer_agent_id: agent_id, endpoint_id, topic: response_topic, messages, allowed, reason, note, }), PeerControlResponse::PubsubSubscribed { .. } => Err(NodeError::IrohPeer( "peer pubsub subscribe response did not match request".to_owned(), )), PeerControlResponse::Error { message } => Err(NodeError::IrohPeer(message)), _ => Err(NodeError::IrohPeer( "peer returned wrong response type to pubsub subscribe".to_owned(), )), } } async fn pipe_connect_to_peer( node: &LocalNode, peer_node: &str, target: String, ) -> Result { geth_pipe::validate_pipe_name(&target)?; let response = request_peer_control(node, peer_node, "pipe-connect", |peer_card, nonce| { PeerControlRequest::PipeConnect { peer_card, target: target.clone(), nonce, } }) .await?; match response { PeerControlResponse::PipeConnected { node_id, agent_id, endpoint_id, connection, allowed, reason, note, .. } => { let connection = connection.unwrap_or_else(|| PipeConnection { target, connected_at: UnixMillis(0), local_listener_found: false, note: geth_pipe::pipe_roadmap().to_owned(), }); Ok(ControlResponse::PipeRemoteConnected { peer_node_id: node_id, peer_agent_id: agent_id, endpoint_id, connection, allowed, reason, note, }) } PeerControlResponse::Error { message } => Err(NodeError::IrohPeer(message)), _ => Err(NodeError::IrohPeer( "peer returned wrong response type to pipe connect".to_owned(), )), } } async fn ssh_proxy_connect_to_peer( node: &LocalNode, peer_node: &str, ) -> Result { let response = request_peer_control(node, peer_node, "ssh-proxy-connect", |peer_card, nonce| { PeerControlRequest::SshProxyConnect { peer_card, nonce } }) .await?; match response { PeerControlResponse::SshProxyConnected { node_id, agent_id, endpoint_id, connection, allowed, reason, note, .. } => Ok(ControlResponse::SshProxyConnected { peer_node_id: node_id, peer_agent_id: agent_id, endpoint_id, connection, allowed, reason, note, }), PeerControlResponse::Error { message } => Err(NodeError::IrohPeer(message)), _ => Err(NodeError::IrohPeer( "peer returned wrong response type to SSH proxy connect".to_owned(), )), } } async fn document_sync_from_peer( node: &LocalNode, peer_node: &str, name: &str, ) -> Result { geth_document::validate_document_name(name) .map_err(|_| NodeError::InvalidDocumentName(name.to_owned()))?; let stream = format!("document:{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, "document-sync", |peer_card, nonce| { PeerControlRequest::DocumentSync { peer_card, name: name.to_owned(), since_ms, nonce, } }) .await?; match response { PeerControlResponse::DocumentSynced { node_id, agent_id, endpoint_id, name: response_name, state, high_water_ms, allowed, reason, note, .. } if response_name == name => { if !allowed { return Ok(ControlResponse::DocumentSynced { peer_node_id: node_id, peer_agent_id: agent_id, endpoint_id, name: response_name, updated: false, allowed, reason, note, }); } let store = Store::open(&node.paths.metadata_db())?; let mut updated = false; if let Some(state) = state { let local = ensure_local_document(&store, name)?; if state.updated_at.0 >= local.updated_at_ms { let normalized = geth_document::normalize_document_state(&state.state_json)?; store.insert_document_resource(&StoredDocumentResource { document_id: local.document_id, resource_id: local.resource_id, name: local.name, state_json: normalized, updated_at_ms: state.updated_at.0, })?; updated = true; } } store_live_sync_cursor(&store, peer_node, &stream, high_water_ms)?; Ok(ControlResponse::DocumentSynced { peer_node_id: node_id, peer_agent_id: agent_id, endpoint_id, name: response_name, updated, allowed, reason, note, }) } PeerControlResponse::DocumentSynced { .. } => Err(NodeError::IrohPeer( "peer document sync response did not match request".to_owned(), )), PeerControlResponse::Error { message } => Err(NodeError::IrohPeer(message)), _ => Err(NodeError::IrohPeer( "peer returned wrong response type to document sync".to_owned(), )), } } async fn sync_status_from_peer( node: &LocalNode, peer_node: &str, ) -> Result, 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, name: &str, limit: u32, ) -> Result { geth_db::validate_db_name(name).map_err(|_| NodeError::InvalidDbName(name.to_owned()))?; let store = Store::open(&node.paths.metadata_db())?; let local = store .get_db_resource_by_name(name)? .ok_or_else(|| NodeError::DbNotFound(name.to_owned()))?; let stream = format!("db:{name}"); let cursor = load_live_sync_cursor(&store, peer_node, &stream)?; let after_db_version = (cursor > 0).then_some(cursor); let response = request_peer_control(node, peer_node, "db-sync", |peer_card, nonce| { PeerControlRequest::DbSync { peer_card, name: name.to_owned(), after_db_version, limit, nonce, } }) .await?; match response { PeerControlResponse::DbSynced { node_id, agent_id, endpoint_id, name: response_name, batch, high_water_db_version, allowed, reason, note, .. } if response_name == name => { if !allowed { return Ok(ControlResponse::DbSynced { peer_node_id: node_id, peer_agent_id: agent_id, endpoint_id, name: response_name, changes_received: 0, max_db_version: None, schema_match: false, allowed, reason, note, }); } let local_schema = geth_db::schema_metadata(Path::new(&local.path))?; let (changes_received, max_db_version, schema_match) = if let Some(batch) = batch { let schema_match = batch.schema_metadata == local_schema; let changes_received = batch.changes.len(); let max_db_version = batch.max_db_version; if schema_match { if let Some(next_cursor) = high_water_db_version.or(max_db_version) { store_live_sync_cursor(&store, peer_node, &stream, next_cursor)?; } } (changes_received, max_db_version, schema_match) } else { (0, high_water_db_version, false) }; Ok(ControlResponse::DbSynced { peer_node_id: node_id, peer_agent_id: agent_id, endpoint_id, name: response_name, changes_received, max_db_version, schema_match, allowed, reason, note, }) } PeerControlResponse::DbSynced { .. } => Err(NodeError::IrohPeer( "peer DB sync response did not match request".to_owned(), )), PeerControlResponse::Error { message } => Err(NodeError::IrohPeer(message)), _ => Err(NodeError::IrohPeer( "peer returned wrong response type to DB sync".to_owned(), )), } } fn live_sync_cursor_key(peer_node: &str, stream: &str) -> String { format!("live-sync:{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 { return Ok(0); }; let cursor: LiveSyncCursor = serde_json::from_str(&state.state_json)?; Ok(cursor.cursor_ms) } fn store_live_sync_cursor( store: &Store, peer_node: &str, stream: &str, cursor_ms: i64, ) -> Result<(), NodeError> { store.put_module_state(&StoredModuleState { module: live_sync_cursor_key(peer_node, stream), state_json: serde_json::to_string(&LiveSyncCursor { cursor_ms })?, updated_at_ms: geth_store::now_ms(), })?; Ok(()) } fn should_live_sync_stream( store: &Store, peer_node: &str, stream: &str, remote_watermarks: Option<&BTreeMap>, ) -> Result { 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, 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 { 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, operation: &str, build_request: impl FnOnce(PeerCard, String) -> PeerControlRequest, ) -> Result { let store = Store::open(&node.paths.metadata_db())?; let stored = store .get_peer_card(peer_node)? .ok_or_else(|| NodeError::PeerNotFound(peer_node.to_owned()))?; let peer_card: PeerCard = serde_json::from_str(&stored.card_json)?; peer_card.validate_candidate()?; let candidate = peer_card .endpoints .first() .ok_or(geth_discovery::DiscoveryError::MissingEndpoint)?; ensure_peer_card_matches_endpoint(&peer_card, &candidate.endpoint_id)?; let node_addr = iroh_node_addr_from_candidate(candidate)?; let endpoint = node .iroh_endpoint .lock() .map_err(|_| NodeError::RuntimeLockPoisoned)? .clone() .ok_or(NodeError::IrohEndpointUnavailable)?; let self_card = local_peer_card(node, DiscoverySource::PeerExchange, true).await?; let nonce = geth_crypto::blake3_hex( format!( "{}\0{}\0{}\0{}", node.node_id, peer_node, operation, geth_store::now_ms() ) .as_bytes(), ); let request = build_request(self_card, nonce.clone()); let conn = endpoint .endpoint() .connect(node_addr, geth_iroh::ALPN_CONTROL) .await .map_err(|error| NodeError::IrohPeer(error.to_string()))?; let (mut send, mut recv) = conn .open_bi() .await .map_err(|error| NodeError::IrohPeer(error.to_string()))?; send.write_all(geth_control::encode_peer_request(&request)?.as_bytes()) .await .map_err(|error| NodeError::IrohPeer(error.to_string()))?; send.finish() .map_err(|error| NodeError::IrohPeer(error.to_string()))?; let bytes = recv .read_to_end(16 * 1024 * 1024) .await .map_err(|error| NodeError::IrohPeer(error.to_string()))?; let text = std::str::from_utf8(&bytes).map_err(|error| NodeError::IrohPeer(error.to_string()))?; let response = geth_control::decode_peer_response(text)?; match &response { PeerControlResponse::SshCertSynced { nonce: response_nonce, .. } | PeerControlResponse::SshRevocationSynced { nonce: response_nonce, .. } | PeerControlResponse::KvSynced { nonce: response_nonce, .. } | PeerControlResponse::PubsubPublished { nonce: response_nonce, .. } | PeerControlResponse::PubsubSubscribed { nonce: response_nonce, .. } | PeerControlResponse::PipeConnected { nonce: response_nonce, .. } | PeerControlResponse::SshProxyConnected { nonce: response_nonce, .. } | PeerControlResponse::DocumentSynced { nonce: response_nonce, .. } | PeerControlResponse::DbSynced { nonce: response_nonce, .. } | PeerControlResponse::SyncStatus { nonce: response_nonce, .. } if response_nonce == &nonce => Ok(response), PeerControlResponse::Error { .. } => Ok(response), _ => Err(NodeError::IrohPeer(format!( "peer {operation} response did not match request" ))), } } fn spawn_iroh_control_accept_loop(node: LocalNode, endpoint: GethIrohEndpoint) { let raw_endpoint = endpoint.endpoint(); tokio::spawn(async move { while let Some(incoming) = raw_endpoint.accept().await { let node = node.clone(); tokio::spawn(async move { if let Err(error) = handle_iroh_control_connection(node, incoming).await { tracing::warn!(%error, "iroh control request failed"); } }); } }); } fn spawn_background_live_sync(node: LocalNode) { tokio::spawn(async move { if let Err(error) = run_live_sync_once(&node).await { tracing::debug!(%error, "initial live sync tick failed"); } let mut interval = tokio::time::interval(LIVE_SYNC_INTERVAL); loop { interval.tick().await; if let Err(error) = run_live_sync_once(&node).await { tracing::debug!(%error, "live sync tick failed"); } } }); } 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()?; 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 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 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 { 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 { 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 { 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"); } } } } Ok(()) } async fn handle_iroh_control_connection( node: LocalNode, incoming: iroh::endpoint::Incoming, ) -> Result<(), NodeError> { let conn = incoming .await .map_err(|error| NodeError::IrohPeer(error.to_string()))?; let remote_endpoint_id = conn .remote_node_id() .map(|node_id| node_id.to_string()) .map_err(|error| NodeError::IrohPeer(error.to_string()))?; let alpn = conn .alpn() .map(display_alpn) .unwrap_or_else(|| "unknown".to_owned()); let (mut send, mut recv) = conn .accept_bi() .await .map_err(|error| NodeError::IrohPeer(error.to_string()))?; let bytes = recv .read_to_end(64 * 1024) .await .map_err(|error| NodeError::IrohPeer(error.to_string()))?; let text = std::str::from_utf8(&bytes).map_err(|error| NodeError::IrohPeer(error.to_string()))?; let response = match geth_control::decode_peer_request(text)? { PeerControlRequest::Ping { 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, })?; PeerControlResponse::Pong { 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, alpn, nonce, 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, capability, 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 explanation = geth_auth::explain_auth_ops( &load_auth_ops_for_resource(&store, &resource)?, PrincipalId::new(peer_card.node_id.to_string()), ResourceId::new(resource.clone()), Capability::new(capability.clone()), ); PeerControlResponse::AuthChecked { 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, resource, capability, allowed: explanation.allowed, reason: explanation.reason, evaluated_ops: explanation.evaluated_ops, nonce, note: "protected peer request authenticated endpoint/card binding before resource capability evaluation".to_owned(), } } PeerControlRequest::CasFetch { peer_card, hash, 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 resource = "resource:cas:local".to_owned(); let capability = "cas.fetch".to_owned(); let explanation = geth_auth::explain_auth_ops( &load_auth_ops_for_resource(&store, &resource)?, PrincipalId::new(peer_card.node_id.to_string()), ResourceId::new(resource), Capability::new(capability), ); if explanation.allowed { match LocalCas::new(node.paths.cas_dir()).read_bytes(&hash) { Ok(content) => { let size_bytes = content.len() as u64; PeerControlResponse::CasFetched { 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, hash, size_bytes, content_base64: Some( base64::engine::general_purpose::STANDARD.encode(content), ), allowed: true, reason: explanation.reason, evaluated_ops: explanation.evaluated_ops, nonce, note: "CAS fetch authenticated endpoint/card binding and required cas.fetch on resource:cas:local".to_owned(), } } Err(error) => PeerControlResponse::Error { message: format!("CAS blob {hash} is not available: {error}"), }, } } else { PeerControlResponse::CasFetched { 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, hash, size_bytes: 0, content_base64: None, allowed: false, reason: explanation.reason, evaluated_ops: explanation.evaluated_ops, nonce, note: "CAS fetch authenticated endpoint/card binding and required cas.fetch on resource:cas:local".to_owned(), } } } PeerControlRequest::SshCertSync { peer_card, since_ms, 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 resource = "resource:ssh:certs".to_owned(); let capability = "ssh_cert.sync".to_owned(); let high_water_ms = geth_store::now_ms(); let explanation = geth_auth::explain_auth_ops( &load_auth_ops_for_resource(&store, &resource)?, PrincipalId::new(peer_card.node_id.to_string()), ResourceId::new(resource), Capability::new(capability), ); let (requests, certificates) = if explanation.allowed { ( store .list_ssh_cert_requests_since(since_ms)? .into_iter() .map(ssh_cert_request_from_stored) .collect::, _>>()?, store .list_ssh_certificates_since(since_ms)? .into_iter() .map(ssh_certificate_from_stored) .collect(), ) } else { (Vec::new(), Vec::new()) }; PeerControlResponse::SshCertSynced { 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, requests, certificates, high_water_ms, allowed: explanation.allowed, reason: explanation.reason, evaluated_ops: explanation.evaluated_ops, nonce, note: "SSH certificate metadata sync authenticated endpoint/card binding and required ssh_cert.sync on resource:ssh:certs".to_owned(), } } PeerControlRequest::SshRevocationSync { peer_card, since_ms, 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 resource = "resource:ssh:revocations".to_owned(); let capability = "ssh_revocation.sync".to_owned(); let high_water_ms = geth_store::now_ms(); let explanation = geth_auth::explain_auth_ops( &load_auth_ops_for_resource(&store, &resource)?, PrincipalId::new(peer_card.node_id.to_string()), ResourceId::new(resource), Capability::new(capability), ); let revocations = if explanation.allowed { store .list_ssh_revocations_since(since_ms)? .into_iter() .map(ssh_revocation_from_stored) .collect::, _>>()? } else { Vec::new() }; PeerControlResponse::SshRevocationSynced { 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, revocations, high_water_ms, allowed: explanation.allowed, reason: explanation.reason, evaluated_ops: explanation.evaluated_ops, nonce, 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}"), } } } PeerControlRequest::PubsubPublish { peer_card, topic, message, nonce, } => { geth_pubsub::validate_topic(&topic)?; geth_pubsub::validate_message(&message)?; 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 resource = format!("resource:pubsub:{topic}"); let capability = "pubsub.publish".to_owned(); let explanation = geth_auth::explain_auth_ops( &load_auth_ops_for_resource(&store, &resource)?, PrincipalId::new(peer_card.node_id.to_string()), ResourceId::new(resource), Capability::new(capability), ); let published = if explanation.allowed { Some(record_pubsub_message(&node, topic, message)?) } else { None }; PeerControlResponse::PubsubPublished { 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, message: published, allowed: explanation.allowed, reason: explanation.reason, evaluated_ops: explanation.evaluated_ops, nonce, note: "pubsub publish authenticated endpoint/card binding and required pubsub.publish on the remote topic resource; pubsub is lossy".to_owned(), } } PeerControlRequest::PubsubSubscribe { peer_card, topic, nonce, } => { geth_pubsub::validate_topic(&topic)?; 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 resource = format!("resource:pubsub:{topic}"); let capability = "pubsub.subscribe".to_owned(); let explanation = geth_auth::explain_auth_ops( &load_auth_ops_for_resource(&store, &resource)?, PrincipalId::new(peer_card.node_id.to_string()), ResourceId::new(resource), Capability::new(capability), ); let messages = if explanation.allowed { pubsub_messages_for_topic(&node, &topic)? } else { Vec::new() }; PeerControlResponse::PubsubSubscribed { 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, topic, messages, allowed: explanation.allowed, reason: explanation.reason, evaluated_ops: explanation.evaluated_ops, nonce, note: "pubsub subscribe authenticated endpoint/card binding and required pubsub.subscribe on the remote topic resource; pubsub is lossy daemon-lifetime state".to_owned(), } } PeerControlRequest::PipeConnect { peer_card, target, nonce, } => { geth_pipe::validate_pipe_name(&target)?; 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 resource = format!("resource:pipe:{target}"); let capability = "pipe.connect".to_owned(); let explanation = geth_auth::explain_auth_ops( &load_auth_ops_for_resource(&store, &resource)?, PrincipalId::new(peer_card.node_id.to_string()), ResourceId::new(resource), Capability::new(capability), ); let connection = if explanation.allowed { Some(record_pipe_connection( &node, target, "remote pipe connect over protected Iroh control path; byte streams are not implemented yet".to_owned(), )?) } else { None }; PeerControlResponse::PipeConnected { 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, connection, allowed: explanation.allowed, reason: explanation.reason, evaluated_ops: explanation.evaluated_ops, nonce, note: "pipe connect authenticated endpoint/card binding and required pipe.connect on the remote pipe resource; byte streams are not implemented yet".to_owned(), } } PeerControlRequest::SshProxyConnect { 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 resource = "resource:ssh-proxy:local".to_owned(); let capability = "ssh_proxy.connect".to_owned(); let explanation = geth_auth::explain_auth_ops( &load_auth_ops_for_resource(&store, &resource)?, PrincipalId::new(peer_card.node_id.to_string()), ResourceId::new(resource), Capability::new(capability), ); let connection = if explanation.allowed { Some(SshProxyConnection { target_node: NodeId::new(node.node_id.clone()), connected_at: UnixMillis(geth_store::now_ms()), local_sshd_target: Some("127.0.0.1:22".to_owned()), admin_shell_available: false, note: "authorized SSH proxy control-plane check passed; byte proxying and sshd/admin-shell connection are not implemented yet".to_owned(), }) } else { None }; PeerControlResponse::SshProxyConnected { 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, connection, allowed: explanation.allowed, reason: explanation.reason, evaluated_ops: explanation.evaluated_ops, nonce, note: "SSH proxy authenticated endpoint/card binding and required ssh_proxy.connect on resource:ssh-proxy:local; SSH is not a geth transport and byte proxying is not implemented yet".to_owned(), } } PeerControlRequest::DocumentSync { peer_card, name, since_ms, nonce, } => { geth_document::validate_document_name(&name) .map_err(|_| NodeError::InvalidDocumentName(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(document) = store.get_document_resource_by_name(&name)? { let capability = "document.read".to_owned(); let high_water_ms = geth_store::now_ms(); let explanation = geth_auth::explain_auth_ops( &load_auth_ops_for_resource(&store, &document.resource_id)?, PrincipalId::new(peer_card.node_id.to_string()), ResourceId::new(document.resource_id.clone()), Capability::new(capability), ); let state = if explanation.allowed && document.updated_at_ms >= since_ms { Some(document_state_from_stored(&document)) } else { None }; PeerControlResponse::DocumentSynced { 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, state, high_water_ms, allowed: explanation.allowed, reason: explanation.reason, evaluated_ops: explanation.evaluated_ops, nonce, note: "document sync authenticated endpoint/card binding and required document.read on the remote document resource; JSON LWW sync is a bootstrap before Automerge".to_owned(), } } else { PeerControlResponse::Error { message: format!("document not found: {name}"), } } } PeerControlRequest::DbSync { peer_card, name, after_db_version, limit, nonce, } => { geth_db::validate_db_name(&name).map_err(|_| NodeError::InvalidDbName(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(db) = store.get_db_resource_by_name(&name)? { let capability = "db.sync".to_owned(); let explanation = geth_auth::explain_auth_ops( &load_auth_ops_for_resource(&store, &db.resource_id)?, PrincipalId::new(peer_card.node_id.to_string()), ResourceId::new(db.resource_id.clone()), Capability::new(capability), ); let sync_result: Result<_, String> = if explanation.allowed { let path = Path::new(&db.path); match geth_db::extract_crsqlite_changes(path, after_db_version, limit) { Ok(batch) => match geth_db::crsqlite_change_metadata(path) { Ok(metadata) => { let high_water_db_version = metadata.max_db_version.or(batch.max_db_version); Ok((Some(batch), high_water_db_version)) } Err(error) => Err(format!( "db sync cannot read crsql_changes metadata for {name}: {error}" )), }, Err(error) => Err(format!( "db sync cannot read crsql_changes for {name}: {error}" )), } } else { Ok((None, None)) }; match sync_result { Ok((batch, high_water_db_version)) => PeerControlResponse::DbSynced { 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, batch, high_water_db_version, allowed: explanation.allowed, reason: explanation.reason, evaluated_ops: explanation.evaluated_ops, nonce, note: "DB sync authenticated endpoint/card binding and required db.sync on the remote DB resource; bootstrap exchanges typed crsql_changes and cursors them, but applying remote changes is not implemented yet".to_owned(), }, Err(message) => PeerControlResponse::Error { message }, } } else { PeerControlResponse::Error { message: format!("db resource not found: {name}"), } } } }; send.write_all(geth_control::encode_peer_response(&response)?.as_bytes()) .await .map_err(|error| NodeError::IrohPeer(error.to_string()))?; send.finish() .map_err(|error| NodeError::IrohPeer(error.to_string()))?; Ok(()) } fn iroh_node_addr_from_candidate( candidate: &EndpointCandidate, ) -> Result { let node_id = candidate .endpoint_id .parse::() .map_err(|error| NodeError::IrohPeer(error.to_string()))?; let direct_addresses = candidate .direct_addresses .iter() .map(|addr| { addr.parse::() .map_err(|error| NodeError::IrohPeer(error.to_string())) }) .collect::, _>>()?; let mut node_addr = iroh::NodeAddr::new(node_id).with_direct_addresses(direct_addresses); if let Some(relay_url) = &candidate.relay_url { node_addr = node_addr.with_relay_url( relay_url .parse::() .map_err(|error| NodeError::IrohPeer(error.to_string()))?, ); } Ok(node_addr) } fn ensure_peer_card_matches_endpoint(card: &PeerCard, endpoint_id: &str) -> Result<(), NodeError> { if card .endpoints .iter() .any(|candidate| candidate.endpoint_id == endpoint_id) { Ok(()) } else { Err(NodeError::IrohPeer(format!( "signed peer card for {} does not bind Iroh endpoint {}", card.node_id, endpoint_id ))) } } fn display_alpn(alpn: Vec) -> String { String::from_utf8_lossy(&alpn).into_owned() } pub fn handle_request( node: &LocalNode, request: ControlRequest, ) -> Result { let store = Store::open(&node.paths.metadata_db())?; match request { ControlRequest::Status => Ok(ControlResponse::Status(StatusResponse { home: node.paths.home().to_path_buf(), socket: node.paths.socket_path(), agent_id: node.agent_id.clone(), node_id: node.node_id.clone(), iroh_enabled: node.iroh_status.enabled, endpoint_id: node.iroh_status.endpoint_id.clone(), iroh_relay_mode: node.iroh_status.relay_mode.clone(), iroh_local_discovery: node.iroh_status.local_discovery, iroh: node.iroh_status.note.clone(), })), ControlRequest::NodeId => Ok(ControlResponse::NodeId(NodeIdResponse { agent_id: node.agent_id.clone(), node_id: node.node_id.clone(), endpoint_id: node.iroh_status.endpoint_id.clone(), })), ControlRequest::PeerCardExport { out } => { let endpoint_id = node .iroh_status .endpoint_id .clone() .ok_or(NodeError::IrohEndpointUnavailable)?; let key = AgentKey::load(&node.paths.agent_key())?; let card = PeerCard::signed( NodeId::new(node.node_id.clone()), &key, vec![EndpointCandidate { endpoint_id, relay_url: None, direct_addresses: Vec::new(), source: DiscoverySource::Manual, }], UnixMillis(geth_store::now_ms()), )?; if let Some(path) = &out { std::fs::write(path, serde_json::to_string_pretty(&card)?)?; } Ok(ControlResponse::PeerCardExported { card, out, note: discovery_is_untrusted_note().to_owned(), }) } ControlRequest::PeerCardImport { path } => { let card_json = std::fs::read_to_string(&path)?; let card: PeerCard = serde_json::from_str(&card_json)?; card.validate_candidate()?; let discovered = DiscoveredPeer::candidate( card.clone(), UnixMillis(geth_store::now_ms()), DiscoverySource::Imported, )?; store.upsert_peer_card(&StoredPeerCard { peer_id: card.node_id.to_string(), card_json: serde_json::to_string(&card)?, updated_at_ms: discovered.discovered_at.0, })?; Ok(ControlResponse::PeerCardImported { peer: discovered, note: discovery_is_untrusted_note().to_owned(), }) } ControlRequest::PeerCardList => { let peers = store .list_peer_cards()? .into_iter() .map(discovered_peer_from_stored) .collect::, _>>()?; Ok(ControlResponse::PeerCardList { peers, note: discovery_is_untrusted_note().to_owned(), }) } ControlRequest::PeerPing { .. } => Err(NodeError::IrohEndpointUnavailable), ControlRequest::PeerAuthCheck { .. } => Err(NodeError::IrohEndpointUnavailable), ControlRequest::CasFetch { .. } => Err(NodeError::IrohEndpointUnavailable), ControlRequest::SshCertSync { .. } => Err(NodeError::IrohEndpointUnavailable), ControlRequest::SshRevocationSync { .. } => Err(NodeError::IrohEndpointUnavailable), ControlRequest::KvSync { .. } => Err(NodeError::IrohEndpointUnavailable), ControlRequest::DbSync { .. } => Err(NodeError::IrohEndpointUnavailable), ControlRequest::DocumentSync { .. } => Err(NodeError::IrohEndpointUnavailable), ControlRequest::ResourceList => Ok(ControlResponse::ResourceList { resources: store .list_resources()? .into_iter() .map(stored_resource_to_descriptor) .collect::, _>>()?, }), ControlRequest::ResourceCreate { kind, name } => { let kind = kind .parse::() .map_err(|_| NodeError::InvalidResourceKind(kind.clone()))?; let id = format!("resource:{}:{}", kind, name); let stored = StoredResource { resource_id: id, kind: kind.to_string(), name, status: "active".to_owned(), }; store.insert_resource(&stored)?; Ok(ControlResponse::ResourceCreated { resource: stored_resource_to_descriptor(stored)?, }) } ControlRequest::CasAdd { path } => { let cas = LocalCas::new(node.paths.cas_dir()); let info = cas.add_path(&path)?; store.record_cas_object( info.hash.as_str(), info.size_bytes, &info.path.to_string_lossy(), )?; Ok(ControlResponse::CasAdded { hash: info.hash, size_bytes: info.size_bytes, }) } ControlRequest::CasGet { hash, out } => { let cas = LocalCas::new(node.paths.cas_dir()); let size_bytes = cas.get_to_path(&hash, &out)?; Ok(ControlResponse::CasGot { hash, out, size_bytes, }) } ControlRequest::CasHash { path } => Ok(ControlResponse::CasHash { hash: hash_path(&path)?, }), ControlRequest::CasHas { hash } => { let cas = LocalCas::new(node.paths.cas_dir()); let present = cas.has(&hash)?; Ok(ControlResponse::CasHas { hash, present }) } ControlRequest::CasPin { hash } => { let cas = LocalCas::new(node.paths.cas_dir()); if !cas.has(&hash)? { return Err(NodeError::Cas(geth_cas::CasError::NotFound( hash.to_string(), ))); } store.pin_cas_object(hash.as_str())?; Ok(ControlResponse::CasPinned { hash, pinned: true }) } ControlRequest::CasUnpin { hash } => { geth_cas::validate_hash(&hash)?; store.unpin_cas_object(hash.as_str())?; Ok(ControlResponse::CasPinned { hash, pinned: false, }) } ControlRequest::CasCleanup { dry_run } => { let cas = LocalCas::new(node.paths.cas_dir()); let mut removed = Vec::new(); let mut retained_pinned = Vec::new(); for blob in cas.list()? { if store.is_cas_object_pinned(blob.hash.as_str())? { retained_pinned.push(blob.hash); continue; } if !dry_run && cas.remove(&blob.hash)? { store.delete_cas_object(blob.hash.as_str())?; } removed.push(blob.hash); } Ok(ControlResponse::CasCleanup { removed, retained_pinned, dry_run, }) } ControlRequest::CasProviders { hash } => { geth_cas::validate_hash(&hash)?; Ok(ControlResponse::CasProviders { providers: store .list_cas_providers(hash.as_str())? .into_iter() .map(|provider| CasProvider { peer_node_id: provider.peer_node_id, endpoint_id: provider.endpoint_id, last_seen_ms: provider.last_seen_ms, }) .collect(), hash, }) } ControlRequest::CasList => { let cas = LocalCas::new(node.paths.cas_dir()); let blobs = cas .list()? .into_iter() .map(|blob| { Ok(CasBlob { pinned: store.is_cas_object_pinned(blob.hash.as_str())?, hash: blob.hash, size_bytes: blob.size_bytes, }) }) .collect::, NodeError>>()?; Ok(ControlResponse::CasList { blobs }) } ControlRequest::CasRootAdd { name, path } => { geth_cas::validate_file_root_name(&name)?; if !path.is_dir() { return Err(NodeError::Cas(geth_cas::CasError::TreeRootNotDirectory( path.display().to_string(), ))); } let path = std::fs::canonicalize(path)?; let resource_id = format!("resource:cas-tree:{name}"); store.insert_resource(&StoredResource { resource_id: resource_id.clone(), kind: ResourceKind::Cas.to_string(), name: format!("file-root:{name}"), status: "active".to_owned(), })?; let stored = StoredFileRoot { root_id: format!("file-root:{name}"), resource_id, name, path: path.display().to_string(), latest_tree_hash: None, latest_tree_json: None, updated_at_ms: geth_store::now_ms(), }; store.upsert_file_root(&stored)?; Ok(ControlResponse::CasRootAdded { root: file_root_from_stored(&stored), }) } ControlRequest::CasRootList => Ok(ControlResponse::CasRootList { roots: store .list_file_roots()? .iter() .map(file_root_from_stored) .collect(), }), ControlRequest::CasRootScan { name } => { geth_cas::validate_file_root_name(&name)?; let mut stored = store .get_file_root_by_name(&name)? .ok_or_else(|| NodeError::ResourceNotFound(format!("file-root:{name}")))?; let previous_tree = stored .latest_tree_json .as_deref() .map(serde_json::from_str) .transpose()?; let cas = LocalCas::new(node.paths.cas_dir()); let scanned = cas.add_tree_path(Path::new(&stored.path))?; store.record_cas_object( scanned.object.hash.as_str(), scanned.object.size_bytes, &scanned.object.path.to_string_lossy(), )?; let changes = geth_cas::diff_tree_objects(previous_tree.as_ref(), &scanned.tree); stored.latest_tree_hash = Some(scanned.object.hash.to_string()); stored.latest_tree_json = Some(serde_json::to_string(&scanned.tree)?); stored.updated_at_ms = geth_store::now_ms(); store.upsert_file_root(&stored)?; Ok(ControlResponse::CasRootScanned { scan: FileRootScan { root: file_root_from_stored(&stored), tree: BlobInfoSummary { hash: scanned.object.hash, size_bytes: scanned.object.size_bytes, }, changes, note: geth_cas::file_root_scan_note().to_owned(), }, }) } ControlRequest::CasConflictRecord { root, path, kind, detail, base_tree, local_tree, remote_tree, } => { geth_cas::validate_file_root_name(&root)?; if let Some(hash) = &base_tree { geth_cas::validate_hash(hash)?; } if let Some(hash) = &local_tree { geth_cas::validate_hash(hash)?; } if let Some(hash) = &remote_tree { geth_cas::validate_hash(hash)?; } let root_record = store .get_file_root_by_name(&root)? .ok_or_else(|| NodeError::ResourceNotFound(format!("file-root:{root}")))?; let kind = FileConflictKind::parse(&kind)?; let created_at = geth_store::now_ms(); let conflict_id = generated_file_conflict_id(&root, &path, kind.as_str(), created_at); let conflict = StoredFileConflict { conflict_id, root_name: root, resource_id: root_record.resource_id, path, kind: kind.as_str().to_owned(), status: FileConflictStatus::Open.as_str().to_owned(), base_tree_hash: base_tree.map(|hash| hash.to_string()), local_tree_hash: local_tree.map(|hash| hash.to_string()), remote_tree_hash: remote_tree.map(|hash| hash.to_string()), detail, resolution: None, resolution_note: None, created_at_ms: created_at, resolved_at_ms: None, }; store.upsert_file_conflict(&conflict)?; Ok(ControlResponse::CasConflictRecorded { conflict: file_conflict_from_stored(conflict)?, }) } ControlRequest::CasConflictList { root } => { if let Some(root) = root.as_deref() { geth_cas::validate_file_root_name(root)?; } let conflicts = store .list_file_conflicts(root.as_deref())? .into_iter() .map(file_conflict_from_stored) .collect::, _>>()?; Ok(ControlResponse::CasConflictList { conflicts }) } ControlRequest::CasConflictResolve { conflict_id, resolution, note, } => { let resolution = FileConflictResolution::parse(&resolution)?; let mut conflict = store .get_file_conflict(&conflict_id)? .ok_or_else(|| NodeError::ResourceNotFound(conflict_id.clone()))?; conflict.status = FileConflictStatus::Resolved.as_str().to_owned(); conflict.resolution = Some(resolution.as_str().to_owned()); conflict.resolution_note = note; conflict.resolved_at_ms = Some(geth_store::now_ms()); store.upsert_file_conflict(&conflict)?; Ok(ControlResponse::CasConflictResolved { conflict: file_conflict_from_stored(conflict)?, }) } ControlRequest::KeychainInit { admin_key_path } => { let mut ops = Vec::new(); let created_at = UnixMillis(geth_store::now_ms()); let init = KeychainOp { id: generated_keychain_op_id("keychain-init", "local", created_at), created_at, kind: KeychainOpKind::KeychainInit, }; store_keychain_op(&store, &init)?; ops.push(init); if let Some(admin_key_path) = admin_key_path { let public_key = std::fs::read_to_string(admin_key_path)?; let created_at = UnixMillis(geth_store::now_ms()); let admin_key = KeyId::new(ssh_public_key_fingerprint(&public_key)); let op = KeychainOp { id: generated_keychain_op_id("admin-key-add", admin_key.as_str(), created_at), created_at, kind: KeychainOpKind::AdminKeyAdd { key: admin_key }, }; store_keychain_op(&store, &op)?; ops.push(op); } Ok(ControlResponse::KeychainInitialized { ops }) } ControlRequest::KeychainStatus => { let view = geth_keychain::reduce_keychain_ops(&load_keychain_ops(&store)?); Ok(ControlResponse::KeychainStatus(KeychainStatusResponse { initialized: view.initialized, admin_keys: view.admin_keys.len(), users: view.users.len(), devices: view.devices.len(), nodes: view.nodes.len(), })) } ControlRequest::SecretStatus => Ok(ControlResponse::SecretStatus { secrets: store .list_resource_secrets()? .into_iter() .map(resource_secret_from_stored) .collect(), }), ControlRequest::SecretCreate { resource } => { ensure_resource_exists(&store, &resource)?; let secret = create_resource_secret(&store, &resource, 1)?; Ok(ControlResponse::SecretCreated { secret }) } ControlRequest::SecretRotate { resource } => { ensure_resource_exists(&store, &resource)?; let next_epoch = store .latest_resource_secret(&resource)? .map(|secret| secret.epoch + 1) .unwrap_or(1); let secret = create_resource_secret(&store, &resource, next_epoch)?; Ok(ControlResponse::SecretCreated { secret }) } ControlRequest::SecretBearerCreate { resource, capabilities, expires_at_ms, } => { ensure_resource_exists(&store, &resource)?; let capabilities = capabilities .into_iter() .map(Capability::new) .collect::>(); geth_secrets::validate_bearer_capabilities(&capabilities)?; let created_at = UnixMillis(geth_store::now_ms()); let secret = geth_types::SecretId::new(format!( "bearer:{}", geth_crypto::blake3_hex( format!( "{resource}\0{}\0{}", capabilities .iter() .map(ToString::to_string) .collect::>() .join(","), created_at.0 ) .as_bytes() ) )); let access = BearerAccess::resource_scoped( secret.clone(), ResourceId::new(resource.clone()), capabilities.clone(), ); let op = AuthOp { id: generated_auth_op_id("bearer-create", &resource, secret.as_str(), created_at), resource: ResourceId::new(resource), created_at, kind: AuthOpKind::BearerAccessCreate { secret, capabilities, expires_at: expires_at_ms.map(UnixMillis), }, }; store_auth_op(&store, &op)?; Ok(ControlResponse::SecretBearerCreated { access: BearerAccess { expires_at: expires_at_ms.map(UnixMillis), ..access }, }) } ControlRequest::SecretBearerList => Ok(ControlResponse::SecretBearerList { access: load_bearer_access(&store)?, }), ControlRequest::SecretBearerRevoke { resource, secret } => { ensure_resource_exists(&store, &resource)?; let created_at = UnixMillis(geth_store::now_ms()); let op = AuthOp { id: generated_auth_op_id("bearer-revoke", &resource, &secret, created_at), resource: ResourceId::new(resource.clone()), created_at, kind: AuthOpKind::BearerAccessRevoke { secret: secret.clone().into(), }, }; store_auth_op(&store, &op)?; Ok(ControlResponse::SecretBearerRevoked { resource, secret }) } ControlRequest::AuthExplain { subject, resource, capability, } => { let ops = load_auth_ops_for_resource(&store, &resource)?; let discovered = store.get_peer_card(&subject)?.is_some(); if ops.is_empty() { if discovered { Ok(ControlResponse::AuthExplain( AuthExplanation::discovered_candidate(subject, resource, capability), )) } else { Ok(ControlResponse::AuthExplain(AuthExplanation::stub( subject, resource, capability, ))) } } else { let mut explanation = geth_auth::explain_auth_ops( &ops, PrincipalId::new(subject.clone()), ResourceId::new(resource.clone()), Capability::new(capability.clone()), ); if discovered && !explanation.allowed { explanation.reason = format!( "subject is a discovered peer candidate only; discovery does not grant trust or authorization; {}", explanation.reason ); } Ok(ControlResponse::AuthExplain(explanation)) } } ControlRequest::AuthGrant { subject, resource, capability, grant_id, } => { let created_at = UnixMillis(geth_store::now_ms()); let grant_id = grant_id.unwrap_or_else(|| generated_grant_id(&subject, &resource, &capability)); let op = AuthOp { id: generated_auth_op_id("grant-create", &resource, &grant_id, created_at), resource: ResourceId::new(resource), created_at, kind: AuthOpKind::GrantCreate { grant_id, principal: PrincipalId::new(subject), capabilities: vec![Capability::new(capability)], }, }; store_auth_op(&store, &op)?; Ok(ControlResponse::AuthOpRecorded { op }) } ControlRequest::AuthRevoke { resource, grant_id } => { let created_at = UnixMillis(geth_store::now_ms()); let op = AuthOp { id: generated_auth_op_id("grant-revoke", &resource, &grant_id, created_at), resource: ResourceId::new(resource), created_at, kind: AuthOpKind::GrantRevoke { grant_id }, }; store_auth_op(&store, &op)?; Ok(ControlResponse::AuthOpRecorded { op }) } ControlRequest::SshCertRequest { public_key_path, cert_kind, principals, requested_validity, renewal_of, reason, subject, } => { ensure_subject_authorized( &store, node, subject.as_deref(), "resource:ssh:certs", "ssh_cert.request", )?; if principals.is_empty() { return Err(NodeError::MissingSshCertPrincipal); } let cert_kind = cert_kind .parse::() .map_err(|_| NodeError::InvalidSshCertKind(cert_kind.clone()))?; let public_key = std::fs::read_to_string(&public_key_path)?; let created_at = UnixMillis(geth_store::now_ms()); let request = SshCertRequest { id: cert_request_id( &NodeId::new(node.node_id.clone()), &public_key, &principals, created_at, ), requester_node: NodeId::new(node.node_id.clone()), public_key_fingerprint: ssh_public_key_fingerprint(&public_key), public_key, cert_kind, principals, requested_validity, renewal_of: renewal_of.map(SshCertId::new), reason, status: SshCertRequestStatus::Pending, created_at, }; store.insert_ssh_cert_request(&stored_from_ssh_cert_request(&request))?; Ok(ControlResponse::SshCertRequested { request }) } ControlRequest::SshCertRequests { subject } => { ensure_subject_authorized( &store, node, subject.as_deref(), "resource:ssh:certs", "ssh_cert.read", )?; Ok(ControlResponse::SshCertRequests { requests: store .list_ssh_cert_requests()? .into_iter() .map(ssh_cert_request_from_stored) .collect::, _>>()?, }) } ControlRequest::SshCertApprove { request_id, ca_key_path, valid_for, serial, out, sign, subject, } => { ensure_subject_authorized( &store, node, subject.as_deref(), "resource:ssh:certs", "ssh_cert.approve", )?; let stored = store .get_ssh_cert_request(&request_id)? .ok_or_else(|| NodeError::SshCertRequestNotFound(request_id.clone()))?; let mut request = ssh_cert_request_from_stored(stored)?; request.status = SshCertRequestStatus::Approved; store.update_ssh_cert_request_status(request.id.as_str(), request.status.as_str())?; let public_key_path = out.clone().unwrap_or_else(|| { node.paths .home() .join("ssh-cert-requests") .join(format!("{}.pub", request.id)) }); if let Some(parent) = public_key_path.parent() { std::fs::create_dir_all(parent)?; } std::fs::write(&public_key_path, &request.public_key)?; let valid_for = valid_for .or_else(|| request.requested_validity.clone()) .unwrap_or_else(|| "+52w".to_owned()); let signing_command = build_ssh_cert_sign_command( &request, &ca_key_path, &public_key_path, &valid_for, serial, )?; let expected_certificate_path = expected_openssh_cert_path(&public_key_path); let mut signed = false; let mut certificate_id_value = None; let mut note = "request approved; run the signing command on the CA/YubiKey machine, then import the resulting -cert.pub file".to_owned(); if sign { geth_ssh_identity::ensure_ssh_keygen_available()?; let output = std::process::Command::new(&signing_command[0]) .args(&signing_command[1..]) .output()?; if !output.status.success() { return Err(geth_ssh_identity::SshIdentityError::SshKeygenFailed( String::from_utf8_lossy(&output.stderr).trim().to_owned(), ) .into()); } let certificate = std::fs::read_to_string(&expected_certificate_path)?; let record = SshCertificateRecord { id: certificate_id(&certificate), request_id: request.id.clone(), certificate_fingerprint: ssh_public_key_fingerprint(&certificate), certificate, imported_at: UnixMillis(geth_store::now_ms()), }; store.insert_ssh_certificate(&stored_from_ssh_certificate(&record))?; store.update_ssh_cert_request_status( request.id.as_str(), SshCertRequestStatus::Signed.as_str(), )?; signed = true; certificate_id_value = Some(record.id); note = "request approved, signed with ssh-keygen, and imported into local certificate metadata".to_owned(); } let approval = SshCertApproval { request_id: request.id, approved_by_node: NodeId::new(node.node_id.clone()), ca_key_path: ca_key_path.display().to_string(), key_id: request_id, valid_for, serial, output_path: Some(expected_certificate_path), signing_command, signed, certificate_id: certificate_id_value, note, }; Ok(ControlResponse::SshCertApproved { approval }) } ControlRequest::SshCertImport { request_id, cert_path, subject, } => { ensure_subject_authorized( &store, node, subject.as_deref(), "resource:ssh:certs", "ssh_cert.import", )?; let certificate = std::fs::read_to_string(&cert_path)?; let record = SshCertificateRecord { id: certificate_id(&certificate), request_id: SshCertRequestId::new(request_id.clone()), certificate_fingerprint: ssh_public_key_fingerprint(&certificate), certificate, imported_at: UnixMillis(geth_store::now_ms()), }; store.insert_ssh_certificate(&stored_from_ssh_certificate(&record))?; store.update_ssh_cert_request_status( &request_id, SshCertRequestStatus::Signed.as_str(), )?; Ok(ControlResponse::SshCertImported { certificate: record, }) } ControlRequest::SshCertList { subject } => { ensure_subject_authorized( &store, node, subject.as_deref(), "resource:ssh:certs", "ssh_cert.read", )?; Ok(ControlResponse::SshCertList { requests: store .list_ssh_cert_requests()? .into_iter() .map(ssh_cert_request_from_stored) .collect::, _>>()?, certificates: store .list_ssh_certificates()? .into_iter() .map(ssh_certificate_from_stored) .collect(), }) } ControlRequest::SshRevocationAdd { kind, target, reason, subject, } => { ensure_subject_authorized( &store, node, subject.as_deref(), "resource:ssh:revocations", "ssh_revocation.publish", )?; let kind = kind .parse::() .map_err(|_| NodeError::InvalidSshRevocationKind(kind.clone()))?; let created_at = UnixMillis(geth_store::now_ms()); let revocation = SshRevocationEntry { id: revocation_id(&kind, &target, created_at), kind, target, reason, created_at, published: true, }; store.insert_ssh_revocation(&stored_from_ssh_revocation(&revocation))?; Ok(ControlResponse::SshRevocationAdded { revocation }) } ControlRequest::SshRevocationList { subject } => { ensure_subject_authorized( &store, node, subject.as_deref(), "resource:ssh:revocations", "ssh_revocation.read", )?; Ok(ControlResponse::SshRevocationList { revocations: store .list_ssh_revocations()? .into_iter() .map(ssh_revocation_from_stored) .collect::, _>>()?, }) } ControlRequest::SshRevocationExport { out, format, ca_public, subject, } => { ensure_subject_authorized( &store, node, subject.as_deref(), "resource:ssh:revocations", "ssh_revocation.read", )?; let revocations = store .list_ssh_revocations()? .into_iter() .map(ssh_revocation_from_stored) .collect::, _>>()?; let format = format .parse::() .map_err(NodeError::SshIdentity)?; if let Some(parent) = out.parent() { std::fs::create_dir_all(parent)?; } let (body, note) = match format { SshRevocationExportFormat::Jsonl => { let mut body = String::new(); for revocation in &revocations { body.push_str(&serde_json::to_string(revocation)?); body.push('\n'); } ( body, "JSONL geth revocation metadata; not an OpenSSH KRL binary".to_owned(), ) } SshRevocationExportFormat::OpenSshKrlSpec => ( openssh_krl_spec(&revocations)?, "OpenSSH KRL specification; generate a binary KRL with ssh-keygen -k -f [-s ] ".to_owned(), ), SshRevocationExportFormat::OpenSshKrl => { write_openssh_krl(&revocations, &out, ca_public.as_deref())?; ( String::new(), "OpenSSH binary KRL generated with ssh-keygen; use ssh-keygen -Q -f to query it".to_owned(), ) } }; if format != SshRevocationExportFormat::OpenSshKrl { std::fs::write(&out, body)?; } Ok(ControlResponse::SshRevocationExported { out, format: format.to_string(), count: revocations.len(), note, }) } ControlRequest::SshRevocationImport { path, format, subject, } => { ensure_subject_authorized( &store, node, subject.as_deref(), "resource:ssh:revocations", "ssh_revocation.import", )?; let body = std::fs::read_to_string(&path)?; let created_at = UnixMillis(geth_store::now_ms()); let revocations = match format.as_str() { "jsonl" => body .lines() .filter(|line| !line.trim().is_empty()) .map(|line| serde_json::from_str(line).map_err(NodeError::from)) .collect::, NodeError>>()?, "openssh-krl-spec" | "krl-spec" => parse_openssh_krl_spec(&body)? .into_iter() .enumerate() .map(|(index, (kind, target))| { let created_at = UnixMillis(created_at.0 + index as i64); SshRevocationEntry { id: revocation_id(&kind, &target, created_at), kind, target, reason: Some(format!("imported from {}", path.display())), created_at, published: true, } }) .collect(), "openssh-krl" | "krl" => { return Err(NodeError::SshIdentity( geth_ssh_identity::SshIdentityError::BinaryKrlImportUnsupported, )); } _ => { return Err(NodeError::SshIdentity( geth_ssh_identity::SshIdentityError::InvalidRevocationExportFormat( format.clone(), ), )); } }; for revocation in &revocations { store.insert_ssh_revocation(&stored_from_ssh_revocation(revocation))?; } Ok(ControlResponse::SshRevocationImported { count: revocations.len(), revocations, format, note: "imported revocation metadata; binary OpenSSH KRL files cannot be enumerated, import JSONL or the KRL spec source instead".to_owned(), }) } ControlRequest::DbAdd { name, path } => { geth_db::validate_db_name(&name).map_err(|_| NodeError::InvalidDbName(name.clone()))?; if !path.is_file() { return Err(NodeError::InvalidDbPath(path.display().to_string())); } let path = std::fs::canonicalize(path)?; let schema_metadata = geth_db::schema_metadata(&path)?; let crsqlite_changes = geth_db::crsqlite_change_metadata(&path)?; let resource_id = format!("resource:db:{name}"); let db_id = format!("db:{name}"); let resource = StoredResource { resource_id: resource_id.clone(), kind: ResourceKind::Db.to_string(), name: name.clone(), status: "active".to_owned(), }; store.insert_resource(&resource)?; let stored = StoredDbResource { db_id, resource_id, name, path: path.display().to_string(), }; store.insert_db_resource(&stored)?; Ok(ControlResponse::DbAdded { db: db_resource_from_stored_with_schema(&stored, schema_metadata, crsqlite_changes), }) } ControlRequest::DbStatus { name } => { geth_db::validate_db_name(&name).map_err(|_| NodeError::InvalidDbName(name.clone()))?; let stored = store .get_db_resource_by_name(&name)? .ok_or_else(|| NodeError::DbNotFound(name.clone()))?; Ok(ControlResponse::DbStatus { db: db_resource_from_stored(&stored)?, }) } ControlRequest::DbChanges { name, after_db_version, limit, } => { geth_db::validate_db_name(&name).map_err(|_| NodeError::InvalidDbName(name.clone()))?; let stored = store .get_db_resource_by_name(&name)? .ok_or_else(|| NodeError::DbNotFound(name.clone()))?; let batch = geth_db::extract_crsqlite_changes( Path::new(&stored.path), after_db_version, limit, )?; Ok(ControlResponse::DbChanges { db: db_resource_from_stored(&stored)?, batch, }) } ControlRequest::KvCreate { name } => { geth_kv::validate_kv_name(&name).map_err(|_| NodeError::InvalidKvName(name.clone()))?; let resource_id = format!("resource:kv:{name}"); let kv_id = format!("kv:{name}"); let resource = StoredResource { resource_id: resource_id.clone(), kind: ResourceKind::Kv.to_string(), name: name.clone(), status: "active".to_owned(), }; store.insert_resource(&resource)?; let stored = StoredKvStore { kv_id, resource_id, name, }; store.insert_kv_store(&stored)?; Ok(ControlResponse::KvCreated { kv: kv_resource_from_stored(&stored), }) } ControlRequest::KvSet { name, key, value, subject, } => { geth_kv::validate_kv_name(&name).map_err(|_| NodeError::InvalidKvName(name.clone()))?; geth_kv::validate_kv_key(&key).map_err(|_| NodeError::InvalidKvKey(key.clone()))?; let kv = store .get_kv_store_by_name(&name)? .ok_or_else(|| NodeError::KvNotFound(name.clone()))?; ensure_kv_write_authorized(&store, node, subject, &kv.resource_id, &key)?; let stored = StoredKvEntry { kv_id: kv.kv_id, key, value, updated_at_ms: geth_store::now_ms(), }; store.set_kv_entry(&stored)?; Ok(ControlResponse::KvSet { entry: kv_entry_from_stored(stored), }) } ControlRequest::KvGet { name, key } => { geth_kv::validate_kv_name(&name).map_err(|_| NodeError::InvalidKvName(name.clone()))?; geth_kv::validate_kv_key(&key).map_err(|_| NodeError::InvalidKvKey(key.clone()))?; let kv = store .get_kv_store_by_name(&name)? .ok_or_else(|| NodeError::KvNotFound(name.clone()))?; let entry = store .get_kv_entry(&kv.kv_id, &key)? .map(kv_entry_from_stored); Ok(ControlResponse::KvGet { entry }) } ControlRequest::DocumentCreate { name } => { geth_document::validate_document_name(&name) .map_err(|_| NodeError::InvalidDocumentName(name.clone()))?; let resource_id = format!("resource:document:{name}"); let document_id = format!("document:{name}"); let resource = StoredResource { resource_id: resource_id.clone(), kind: ResourceKind::Document.to_string(), name: name.clone(), status: "active".to_owned(), }; store.insert_resource(&resource)?; let stored = StoredDocumentResource { document_id, resource_id, name, state_json: "{}".to_owned(), updated_at_ms: geth_store::now_ms(), }; store.insert_document_resource(&stored)?; Ok(ControlResponse::DocumentCreated { document: document_resource_from_stored(&stored), }) } ControlRequest::DocumentStatus { name } => { geth_document::validate_document_name(&name) .map_err(|_| NodeError::InvalidDocumentName(name.clone()))?; let stored = store .get_document_resource_by_name(&name)? .ok_or_else(|| NodeError::DocumentNotFound(name.clone()))?; Ok(ControlResponse::DocumentStatus { document: document_resource_from_stored(&stored), }) } ControlRequest::DocumentSet { name, state_json } => { geth_document::validate_document_name(&name) .map_err(|_| NodeError::InvalidDocumentName(name.clone()))?; let mut stored = store .get_document_resource_by_name(&name)? .ok_or_else(|| NodeError::DocumentNotFound(name.clone()))?; stored.state_json = geth_document::normalize_document_state(&state_json)?; stored.updated_at_ms = geth_store::now_ms(); store.insert_document_resource(&stored)?; Ok(ControlResponse::DocumentSet { state: document_state_from_stored(&stored), }) } ControlRequest::DocumentGet { name } => { geth_document::validate_document_name(&name) .map_err(|_| NodeError::InvalidDocumentName(name.clone()))?; let stored = store .get_document_resource_by_name(&name)? .ok_or_else(|| NodeError::DocumentNotFound(name.clone()))?; Ok(ControlResponse::DocumentGet { state: document_state_from_stored(&stored), }) } ControlRequest::PubsubPub { topic, message, node: None, } => { geth_pubsub::validate_topic(&topic)?; geth_pubsub::validate_message(&message)?; let message = record_pubsub_message(node, topic, message)?; Ok(ControlResponse::PubsubPublished { message }) } ControlRequest::PubsubPub { node: Some(_), .. } => Err(NodeError::IrohEndpointUnavailable), ControlRequest::PubsubSub { topic, node: None } => { geth_pubsub::validate_topic(&topic)?; let messages = pubsub_messages_for_topic(node, &topic)?; Ok(ControlResponse::PubsubMessages { topic, messages, note: geth_pubsub::pubsub_storage_warning().to_owned(), }) } ControlRequest::PubsubSub { node: Some(_), .. } => Err(NodeError::IrohEndpointUnavailable), ControlRequest::PipeListen { name } => { geth_pipe::validate_pipe_name(&name)?; let listener = PipeListener { id: format!("pipe:{name}").into(), name: name.clone(), listened_at: UnixMillis(geth_store::now_ms()), note: geth_pipe::local_pipe_runtime_note().to_owned(), }; let mut runtime = node .runtime .pipes .lock() .map_err(|_| NodeError::RuntimeLockPoisoned)?; runtime.listeners.insert(name, listener.clone()); Ok(ControlResponse::PipeListening { listener }) } ControlRequest::PipeConnect { target, node: None } => { geth_pipe::validate_pipe_name(&target)?; let connection = record_pipe_connection( node, target, geth_pipe::local_pipe_runtime_note().to_owned(), )?; Ok(ControlResponse::PipeConnected { connection }) } ControlRequest::PipeConnect { node: Some(_), .. } => { Err(NodeError::IrohEndpointUnavailable) } ControlRequest::SshProxyConnect { .. } => Err(NodeError::IrohEndpointUnavailable), ControlRequest::ModuleStub { module, command } => { Ok(ControlResponse::NotImplemented { module, command }) } } } fn stored_resource_to_descriptor(stored: StoredResource) -> Result { let kind = stored .kind .parse::() .map_err(|_| NodeError::InvalidResourceKind(stored.kind.clone()))?; Ok(ResourceDescriptor::local( ResourceId::new(stored.resource_id), kind, ResourceName::new(stored.name), )) } fn db_resource_from_stored(stored: &StoredDbResource) -> Result { let path = Path::new(&stored.path); let metadata = path.metadata().ok(); let (schema_metadata, crsqlite_changes) = if metadata .as_ref() .is_some_and(std::fs::Metadata::is_file) { ( geth_db::schema_metadata(path).unwrap_or_else(|error| format!("unavailable:{error}")), geth_db::crsqlite_change_metadata(path) .unwrap_or_else(|error| geth_db::CrSqliteChangeMetadata::error(error.to_string())), ) } else { ( "unavailable:path-missing".to_owned(), geth_db::CrSqliteChangeMetadata::error("path-missing"), ) }; Ok(db_resource_from_stored_with_schema( stored, schema_metadata, crsqlite_changes, )) } fn db_resource_from_stored_with_schema( stored: &StoredDbResource, schema_metadata: String, crsqlite_changes: geth_db::CrSqliteChangeMetadata, ) -> DbResource { let path = Path::new(&stored.path); let metadata = path.metadata().ok(); DbResource { id: stored.db_id.clone().into(), resource: stored.resource_id.clone().into(), name: stored.name.clone(), path: stored.path.clone(), path_exists: metadata.as_ref().is_some_and(std::fs::Metadata::is_file), size_bytes: metadata.map(|metadata| metadata.len()), schema_metadata, crsqlite_changes, sync_status: "local-only".to_owned(), } } fn kv_resource_from_stored(stored: &StoredKvStore) -> KvResource { KvResource { id: stored.kv_id.clone().into(), resource: stored.resource_id.clone().into(), name: stored.name.clone(), sync_status: "local-only".to_owned(), } } fn kv_entry_from_stored(stored: StoredKvEntry) -> KvEntry { KvEntry { store: stored.kv_id.into(), key: stored.key, value: stored.value, } } fn ensure_local_kv_store(store: &Store, name: &str) -> Result { 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 record_pubsub_message( node: &LocalNode, topic: String, message: String, ) -> Result { let message = PubsubMessage { topic: topic.into(), message, published_at: UnixMillis(geth_store::now_ms()), }; let mut runtime = node .runtime .pubsub .lock() .map_err(|_| NodeError::RuntimeLockPoisoned)?; runtime.messages.push_back(message.clone()); while runtime.messages.len() > PUBSUB_RING_LIMIT { runtime.messages.pop_front(); } Ok(message) } fn pubsub_messages_for_topic( node: &LocalNode, topic: &str, ) -> Result, NodeError> { let runtime = node .runtime .pubsub .lock() .map_err(|_| NodeError::RuntimeLockPoisoned)?; Ok(runtime .messages .iter() .filter(|message| message.topic.as_str() == topic) .cloned() .collect()) } fn record_pipe_connection( node: &LocalNode, target: String, note: String, ) -> Result { let mut runtime = node .runtime .pipes .lock() .map_err(|_| NodeError::RuntimeLockPoisoned)?; let local_listener_found = runtime.listeners.contains_key(&target); let connection = PipeConnection { target, connected_at: UnixMillis(geth_store::now_ms()), local_listener_found, note, }; runtime.connections.push_back(connection.clone()); while runtime.connections.len() > PIPE_CONNECTION_RING_LIMIT { runtime.connections.pop_front(); } Ok(connection) } fn ensure_local_document(store: &Store, name: &str) -> Result { if let Some(document) = store.get_document_resource_by_name(name)? { return Ok(document); } let stored = StoredDocumentResource { document_id: format!("document:{name}"), resource_id: format!("resource:document:{name}"), name: name.to_owned(), state_json: "{}".to_owned(), updated_at_ms: 0, }; store.insert_document_resource(&stored)?; Ok(stored) } fn document_resource_from_stored(stored: &StoredDocumentResource) -> DocumentResource { DocumentResource { id: stored.document_id.clone().into(), resource: stored.resource_id.clone().into(), name: stored.name.clone(), sync_status: "local-only".to_owned(), state_bytes: stored.state_json.len() as u64, } } fn document_state_from_stored(stored: &StoredDocumentResource) -> DocumentState { DocumentState { document: document_resource_from_stored(stored), state_json: stored.state_json.clone(), updated_at: UnixMillis(stored.updated_at_ms), } } fn file_root_from_stored(stored: &StoredFileRoot) -> FileRoot { FileRoot { id: stored.root_id.clone(), resource: stored.resource_id.clone(), name: stored.name.clone(), path: stored.path.clone(), latest_tree: stored.latest_tree_hash.clone().map(Into::into), updated_at_ms: stored.updated_at_ms, } } fn discovered_peer_from_stored(stored: StoredPeerCard) -> Result { let card: PeerCard = serde_json::from_str(&stored.card_json)?; Ok(DiscoveredPeer::candidate( card, UnixMillis(stored.updated_at_ms), DiscoverySource::Imported, )?) } fn file_conflict_from_stored(stored: StoredFileConflict) -> Result { Ok(FileConflict { id: stored.conflict_id, root: stored.root_name, resource: stored.resource_id, path: stored.path, kind: FileConflictKind::parse(&stored.kind)?, status: FileConflictStatus::parse(&stored.status)?, base_tree: stored.base_tree_hash.map(Into::into), local_tree: stored.local_tree_hash.map(Into::into), remote_tree: stored.remote_tree_hash.map(Into::into), detail: stored.detail, resolution: stored .resolution .as_deref() .map(FileConflictResolution::parse) .transpose()?, resolution_note: stored.resolution_note, created_at_ms: stored.created_at_ms, resolved_at_ms: stored.resolved_at_ms, }) } fn ensure_resource_exists(store: &Store, resource_id: &str) -> Result<(), NodeError> { if store .list_resources()? .into_iter() .any(|resource| resource.resource_id == resource_id) { Ok(()) } else { Err(NodeError::ResourceNotFound(resource_id.to_owned())) } } fn ensure_kv_write_authorized( store: &Store, node: &LocalNode, subject: Option, resource_id: &str, key: &str, ) -> Result<(), NodeError> { let Some(subject) = subject else { return Ok(()); }; let capability = format!("kv.write_key:{key}"); ensure_subject_authorized(store, node, Some(&subject), resource_id, &capability) } fn ensure_subject_authorized( store: &Store, node: &LocalNode, subject: Option<&str>, resource_id: &str, capability: &str, ) -> Result<(), NodeError> { let Some(subject) = subject else { return Ok(()); }; if subject == node.node_id || subject == node.agent_id { return Ok(()); } let ops = load_auth_ops_for_resource(store, resource_id)?; let explanation = geth_auth::explain_auth_ops( &ops, PrincipalId::new(subject.to_owned()), ResourceId::new(resource_id.to_owned()), Capability::new(capability.to_owned()), ); if explanation.allowed { Ok(()) } else { Err(NodeError::Unauthorized(format!( "{subject} lacks {capability} on {resource_id}: {}", explanation.reason ))) } } fn create_resource_secret( store: &Store, resource_id: &str, epoch: u64, ) -> Result { let created_at = UnixMillis(geth_store::now_ms()); let secret_id = format!( "secret:{}", geth_crypto::blake3_hex(format!("{resource_id}\0{epoch}\0{}", created_at.0).as_bytes()) ); let stored = StoredResourceSecret { secret_id, resource_id: resource_id.to_owned(), epoch, status: "active".to_owned(), created_at_ms: created_at.0, }; store.insert_resource_secret(&stored)?; Ok(resource_secret_from_stored(stored)) } fn resource_secret_from_stored(stored: StoredResourceSecret) -> ResourceMasterSecret { ResourceMasterSecret { id: stored.secret_id.into(), resource: stored.resource_id.into(), epoch: stored.epoch, created_at: UnixMillis(stored.created_at_ms), } } fn load_bearer_access(store: &Store) -> Result, NodeError> { let ops = store .list_auth_ops()? .into_iter() .map(|stored| serde_json::from_str(&stored.op_json).map_err(NodeError::from)) .collect::, NodeError>>()?; let view = geth_auth::reduce_auth_ops(&ops); Ok(view .bearer_access .into_values() .map(|record| BearerAccess { secret: record.secret, resource: record.resource, capabilities: record.capabilities, expires_at: record.expires_at, may_delegate: false, }) .collect()) } fn store_auth_op(store: &Store, op: &AuthOp) -> Result<(), NodeError> { store.insert_auth_op(&StoredAuthOp { op_id: op.id.to_string(), resource_id: op.resource.to_string(), op_json: serde_json::to_string(op)?, created_at_ms: op.created_at.0, })?; Ok(()) } fn load_auth_ops_for_resource(store: &Store, resource: &str) -> Result, NodeError> { store .list_auth_ops_for_resource(resource)? .into_iter() .map(|stored| serde_json::from_str(&stored.op_json).map_err(NodeError::from)) .collect() } fn store_keychain_op(store: &Store, op: &KeychainOp) -> Result<(), NodeError> { store.insert_keychain_op(&StoredKeychainOp { op_id: op.id.to_string(), op_json: serde_json::to_string(op)?, created_at_ms: op.created_at.0, })?; Ok(()) } fn load_keychain_ops(store: &Store) -> Result, NodeError> { store .list_keychain_ops()? .into_iter() .map(|stored| serde_json::from_str(&stored.op_json).map_err(NodeError::from)) .collect() } fn generated_grant_id(subject: &str, resource: &str, capability: &str) -> String { format!( "grant:{}", geth_crypto::blake3_hex(format!("{subject}\0{resource}\0{capability}").as_bytes()) ) } fn generated_file_conflict_id(root: &str, path: &str, kind: &str, created_at_ms: i64) -> String { format!( "file-conflict:{}", geth_crypto::blake3_hex(format!("{created_at_ms}\0{root}\0{path}\0{kind}").as_bytes()) ) } fn generated_auth_op_id( kind: &str, resource: &str, stable_id: &str, created_at: UnixMillis, ) -> AuthOpId { AuthOpId::new(format!( "auth-op:{}", geth_crypto::blake3_hex( format!("{}\0{kind}\0{resource}\0{stable_id}", created_at.0).as_bytes() ) )) } fn generated_keychain_op_id(kind: &str, stable_id: &str, created_at: UnixMillis) -> AuthOpId { AuthOpId::new(format!( "keychain-op:{}", geth_crypto::blake3_hex(format!("{}\0{kind}\0{stable_id}", created_at.0).as_bytes()) )) } fn stable_node_id(agent_id: &str) -> String { format!("node:{agent_id}") } async fn start_daemon_iroh_endpoint( node: &mut LocalNode, ) -> Result, NodeError> { let node_config = GethConfig::load(&node.paths.config_file())?; let relay_mode = node_config.iroh.relay_mode.clone(); let relay_mode_label = relay_mode.label(); let local_discovery = node_config.iroh.local_discovery; let iroh_relay_mode = config_relay_mode_to_iroh(&relay_mode, &node_config.iroh.relay_maps); let mut config = GethIrohConfig::local_with_relay(node.paths.iroh_key(), iroh_relay_mode); config.local_discovery = local_discovery; match geth_iroh::start_endpoint(&config).await { Ok(endpoint) => { let status = endpoint.status(); if let Some(endpoint_id) = &status.endpoint_id { let store = Store::open(&node.paths.metadata_db())?; store.upsert_node_endpoint(endpoint_id, &node.node_id, &node.agent_id, "iroh")?; } node.iroh_status = status; *node .iroh_endpoint .lock() .map_err(|_| NodeError::RuntimeLockPoisoned)? = Some(endpoint.clone()); Ok(Some(endpoint)) } Err(error) => { node.iroh_status = EndpointStatus { enabled: false, endpoint_id: None, relay_mode: relay_mode_label, local_discovery, note: format!("Iroh endpoint failed to start: {error}"), }; Ok(None) } } } fn config_relay_mode_to_iroh( mode: &RelayMode, relay_maps: &std::collections::BTreeMap, ) -> GethRelayMode { match mode { RelayMode::Disabled => GethRelayMode::Disabled, RelayMode::Default => GethRelayMode::Default, RelayMode::Staging => GethRelayMode::Staging, RelayMode::Custom { map } => { let relay_map = relay_maps .get(map) .expect("custom relay map was validated during config load"); GethRelayMode::Custom { name: map.clone(), relay_urls: relay_map.urls.clone(), } } } } fn expected_openssh_cert_path(public_key_path: &Path) -> String { let text = public_key_path.display().to_string(); if let Some(prefix) = text.strip_suffix(".pub") { format!("{prefix}-cert.pub") } else { format!("{text}-cert.pub") } } fn stored_from_ssh_cert_request(request: &SshCertRequest) -> StoredSshCertRequest { StoredSshCertRequest { request_id: request.id.to_string(), requester_node: request.requester_node.to_string(), public_key: request.public_key.clone(), public_key_fingerprint: request.public_key_fingerprint.clone(), cert_kind: request.cert_kind.to_string(), principals: request.principals.clone(), requested_validity: request.requested_validity.clone(), renewal_of: request.renewal_of.as_ref().map(ToString::to_string), reason: request.reason.clone(), status: request.status.to_string(), created_at_ms: request.created_at.0, } } fn ssh_cert_request_from_stored(stored: StoredSshCertRequest) -> Result { let cert_kind = stored .cert_kind .parse::() .map_err(|_| NodeError::InvalidSshCertKind(stored.cert_kind.clone()))?; let status = stored .status .parse::() .map_err(|_| NodeError::InvalidSshCertStatus(stored.status.clone()))?; Ok(SshCertRequest { id: SshCertRequestId::new(stored.request_id), requester_node: NodeId::new(stored.requester_node), public_key: stored.public_key, public_key_fingerprint: stored.public_key_fingerprint, cert_kind, principals: stored.principals, requested_validity: stored.requested_validity, renewal_of: stored.renewal_of.map(SshCertId::new), reason: stored.reason, status, created_at: UnixMillis(stored.created_at_ms), }) } fn stored_from_ssh_certificate(certificate: &SshCertificateRecord) -> StoredSshCertificate { StoredSshCertificate { cert_id: certificate.id.to_string(), request_id: certificate.request_id.to_string(), certificate: certificate.certificate.clone(), certificate_fingerprint: certificate.certificate_fingerprint.clone(), imported_at_ms: certificate.imported_at.0, } } fn ssh_certificate_from_stored(stored: StoredSshCertificate) -> SshCertificateRecord { SshCertificateRecord { id: SshCertId::new(stored.cert_id), request_id: SshCertRequestId::new(stored.request_id), certificate: stored.certificate, certificate_fingerprint: stored.certificate_fingerprint, imported_at: UnixMillis(stored.imported_at_ms), } } fn stored_from_ssh_revocation(revocation: &SshRevocationEntry) -> StoredSshRevocation { StoredSshRevocation { revocation_id: revocation.id.to_string(), kind: revocation.kind.to_string(), target: revocation.target.clone(), reason: revocation.reason.clone(), created_at_ms: revocation.created_at.0, published: revocation.published, } } fn ssh_revocation_from_stored( stored: StoredSshRevocation, ) -> Result { let kind = stored .kind .parse::() .map_err(|_| NodeError::InvalidSshRevocationKind(stored.kind.clone()))?; Ok(SshRevocationEntry { id: geth_types::SshRevocationId::new(stored.revocation_id), kind, target: stored.target, reason: stored.reason, created_at: UnixMillis(stored.created_at_ms), published: stored.published, }) } #[cfg(test)] mod tests { use super::*; fn write_offline_iroh_config(paths: &GethPaths) { std::fs::write( paths.config_file(), "[iroh]\nrelay_mode = \"disabled\"\nlocal_discovery = false\n", ) .expect("write config"); } fn create_mock_crsqlite_db(path: &Path, change_db_version: Option) { let conn = rusqlite::Connection::open(path).expect("open sqlite"); conn.execute("CREATE TABLE notes(id INTEGER PRIMARY KEY, body TEXT)", []) .expect("create notes"); conn.execute( r#"CREATE TABLE crsql_changes( table_name TEXT NOT NULL, pk BLOB NOT NULL, cid TEXT NOT NULL, val BLOB, col_version INTEGER NOT NULL, db_version INTEGER NOT NULL, site_id BLOB, cl INTEGER, seq INTEGER )"#, [], ) .expect("create crsql_changes"); if let Some(db_version) = change_db_version { conn.execute( "INSERT INTO crsql_changes(table_name, pk, cid, val, col_version, db_version, site_id, cl, seq) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)", ( "notes", vec![db_version as u8], "body", Vec::from("hello".as_bytes()), 1_i64, db_version, vec![1_u8], 1_i64, db_version, ), ) .expect("insert initial mock crsqlite change"); } } fn insert_mock_crsqlite_change(path: &Path, db_version: i64, body: &str) { let conn = rusqlite::Connection::open(path).expect("open sqlite"); conn.execute( "INSERT INTO crsql_changes(table_name, pk, cid, val, col_version, db_version, site_id, cl, seq) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)", ( "notes", vec![db_version as u8], "body", body.as_bytes().to_vec(), 1_i64, db_version, vec![1_u8], 1_i64, db_version, ), ) .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(); let card = PeerCard::signed( "node:test".into(), &key, vec![EndpointCandidate { endpoint_id: "endpoint:test".to_owned(), relay_url: None, direct_addresses: vec![ "127.0.0.1:1111".to_owned(), "127.0.0.2:1111".to_owned(), "127.0.0.3:2222".to_owned(), ], source: DiscoverySource::Mdns, }], UnixMillis(1), ) .expect("peer card"); let (port, addrs) = lan_discovery_addrs(&card).expect("lan addresses"); assert_eq!(port, 1111); assert_eq!(addrs.len(), 2); } #[tokio::test] async fn peer_ping_uses_signed_peer_card_over_iroh() { let left_home = tempfile::tempdir().expect("left home"); let right_home = tempfile::tempdir().expect("right home"); let left_paths = GethPaths::from_home(left_home.path()); let right_paths = GethPaths::from_home(right_home.path()); let mut left = init_node(&left_paths).expect("init left"); let mut right = init_node(&right_paths).expect("init right"); write_offline_iroh_config(&left_paths); write_offline_iroh_config(&right_paths); let Some(left_endpoint) = start_daemon_iroh_endpoint(&mut left) .await .expect("left iroh") else { eprintln!("skipping peer ping assertion; left Iroh endpoint unavailable"); return; }; let Some(right_endpoint) = start_daemon_iroh_endpoint(&mut right) .await .expect("right iroh") else { eprintln!("skipping peer ping assertion; right Iroh endpoint unavailable"); left_endpoint.shutdown().await; return; }; spawn_iroh_control_accept_loop(right.clone(), right_endpoint.clone()); let exported = handle_request_async(&right, ControlRequest::PeerCardExport { out: None }) .await .expect("export right peer card"); let right_card = match exported { ControlResponse::PeerCardExported { card, .. } => card, other => panic!("unexpected export response: {other:?}"), }; assert!(!right_card.endpoints[0].direct_addresses.is_empty()); Store::open(&left_paths.metadata_db()) .expect("open left store") .upsert_peer_card(&StoredPeerCard { peer_id: right_card.node_id.to_string(), card_json: serde_json::to_string(&right_card).expect("card json"), updated_at_ms: geth_store::now_ms(), }) .expect("insert right peer"); let right_blob = LocalCas::new(right.paths.cas_dir()) .add_bytes(b"remote cas bytes") .expect("right cas add"); let right_pubkey = right_home.path().join("request.pub"); std::fs::write( &right_pubkey, "ssh-ed25519 AAAAC3NzaC1lZDI1NTE5AAAAIGV0aA== node\n", ) .expect("write right public key"); let requested = handle_request( &right, ControlRequest::SshCertRequest { public_key_path: right_pubkey, cert_kind: "user".to_owned(), principals: vec!["eric".to_owned()], requested_validity: Some("+52w".to_owned()), renewal_of: None, reason: Some("test sync".to_owned()), subject: None, }, ) .expect("right ssh cert request"); let right_request_id = match requested { ControlResponse::SshCertRequested { request } => request.id, other => panic!("unexpected SSH cert request response: {other:?}"), }; let added_revocation = handle_request( &right, ControlRequest::SshRevocationAdd { kind: "key-id".to_owned(), target: "old-node-key".to_owned(), reason: Some("test sync".to_owned()), subject: None, }, ) .expect("right ssh revocation"); let right_revocation_id = match added_revocation { 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"); handle_request( &right, ControlRequest::PipeListen { name: "inbox".to_owned(), }, ) .expect("right pipe listen"); handle_request( &right, ControlRequest::DocumentCreate { name: "notes".to_owned(), }, ) .expect("right document create"); handle_request( &right, ControlRequest::DocumentSet { name: "notes".to_owned(), state_json: r#"{"title":"remote"}"#.to_owned(), }, ) .expect("right document set"); let left_db_path = left_home.path().join("notes.sqlite"); let right_db_path = right_home.path().join("notes.sqlite"); create_mock_crsqlite_db(&left_db_path, None); create_mock_crsqlite_db(&right_db_path, Some(7)); handle_request( &left, ControlRequest::DbAdd { name: "notes".to_owned(), path: left_db_path.clone(), }, ) .expect("left db add"); handle_request( &right, ControlRequest::DbAdd { name: "notes".to_owned(), path: right_db_path.clone(), }, ) .expect("right db add"); let ping = handle_request_async( &left, ControlRequest::PeerPing { node: right_card.node_id.to_string(), }, ) .await .expect("peer ping"); match ping { ControlResponse::PeerPinged { peer_node_id, peer_agent_id, alpn, note, .. } => { assert_eq!(peer_node_id, right.node_id); assert_eq!(peer_agent_id, right.agent_id); assert_eq!(alpn, "/geth/control/1"); assert!(note.contains("does not grant resource capabilities")); } other => panic!("unexpected ping response: {other:?}"), } let denied = handle_request_async( &left, ControlRequest::PeerAuthCheck { node: right_card.node_id.to_string(), resource: "resource:cas:local".to_owned(), capability: "cas.fetch".to_owned(), }, ) .await .expect("denied peer auth check"); match denied { ControlResponse::PeerAuthChecked { allowed, reason, note, .. } => { assert!(!allowed); assert!(reason.contains("no active direct or group grant")); assert!(note.contains("authenticated endpoint/card binding")); } other => panic!("unexpected denied auth response: {other:?}"), } let denied_fetch = handle_request_async( &left, ControlRequest::CasFetch { node: right_card.node_id.to_string(), hash: right_blob.hash.clone(), }, ) .await .expect("denied cas fetch"); match denied_fetch { ControlResponse::CasFetched { allowed, reason, size_bytes, .. } => { assert!(!allowed); assert_eq!(size_bytes, 0); assert!(reason.contains("no active direct or group grant")); assert!( !LocalCas::new(left.paths.cas_dir()) .has(&right_blob.hash) .expect("left cas has after denied fetch") ); } 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:?}"), } let denied_pubsub = handle_request_async( &left, ControlRequest::PubsubPub { topic: "presence/test".to_owned(), message: "hello".to_owned(), node: Some(right_card.node_id.to_string()), }, ) .await .expect("denied remote pubsub publish"); match denied_pubsub { ControlResponse::PubsubRemotePublished { allowed, reason, message, .. } => { assert!(!allowed); assert!(message.message.is_empty()); assert!(reason.contains("no active direct or group grant")); } other => panic!("unexpected denied pubsub publish response: {other:?}"), } let denied_pubsub_subscribe = handle_request_async( &left, ControlRequest::PubsubSub { topic: "presence/test".to_owned(), node: Some(right_card.node_id.to_string()), }, ) .await .expect("denied remote pubsub subscribe"); match denied_pubsub_subscribe { ControlResponse::PubsubRemoteMessages { allowed, messages, reason, .. } => { assert!(!allowed); assert!(messages.is_empty()); assert!(reason.contains("no active direct or group grant")); } other => panic!("unexpected denied pubsub subscribe response: {other:?}"), } let denied_pipe = handle_request_async( &left, ControlRequest::PipeConnect { target: "inbox".to_owned(), node: Some(right_card.node_id.to_string()), }, ) .await .expect("denied remote pipe connect"); match denied_pipe { ControlResponse::PipeRemoteConnected { allowed, reason, connection, .. } => { assert!(!allowed); assert!(!connection.local_listener_found); assert!(reason.contains("no active direct or group grant")); } other => panic!("unexpected denied pipe connect response: {other:?}"), } let denied_ssh_proxy = handle_request_async( &left, ControlRequest::SshProxyConnect { node: right_card.node_id.to_string(), }, ) .await .expect("denied SSH proxy connect"); match denied_ssh_proxy { ControlResponse::SshProxyConnected { allowed, reason, connection, .. } => { assert!(!allowed); assert!(connection.is_none()); assert!(reason.contains("no active direct or group grant")); } other => panic!("unexpected denied SSH proxy response: {other:?}"), } let denied_document = handle_request_async( &left, ControlRequest::DocumentSync { node: right_card.node_id.to_string(), name: "notes".to_owned(), }, ) .await .expect("denied document sync"); match denied_document { ControlResponse::DocumentSynced { allowed, updated, reason, .. } => { assert!(!allowed); assert!(!updated); assert!(reason.contains("no active direct or group grant")); } other => panic!("unexpected denied document sync response: {other:?}"), } let denied_db_sync = handle_request_async( &left, ControlRequest::DbSync { node: right_card.node_id.to_string(), name: "notes".to_owned(), limit: 10, }, ) .await .expect("denied db sync"); match denied_db_sync { ControlResponse::DbSynced { allowed, changes_received, schema_match, reason, .. } => { assert!(!allowed); assert_eq!(changes_received, 0); assert!(!schema_match); assert!(reason.contains("no active direct or group grant")); } other => panic!("unexpected denied DB sync response: {other:?}"), } handle_request( &right, ControlRequest::AuthGrant { subject: left.node_id.clone(), resource: "resource:cas:local".to_owned(), capability: "cas.fetch".to_owned(), grant_id: Some("grant:left-cas-fetch".to_owned()), }, ) .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"); handle_request( &right, ControlRequest::AuthGrant { subject: left.node_id.clone(), resource: "resource:pubsub:presence/test".to_owned(), capability: "pubsub.publish".to_owned(), grant_id: Some("grant:left-pubsub-publish".to_owned()), }, ) .expect("grant left pubsub publish"); handle_request( &right, ControlRequest::AuthGrant { subject: left.node_id.clone(), resource: "resource:pubsub:presence/test".to_owned(), capability: "pubsub.subscribe".to_owned(), grant_id: Some("grant:left-pubsub-subscribe".to_owned()), }, ) .expect("grant left pubsub subscribe"); handle_request( &right, ControlRequest::AuthGrant { subject: left.node_id.clone(), resource: "resource:pipe:inbox".to_owned(), capability: "pipe.connect".to_owned(), grant_id: Some("grant:left-pipe-connect".to_owned()), }, ) .expect("grant left pipe connect"); handle_request( &right, ControlRequest::AuthGrant { subject: left.node_id.clone(), resource: "resource:ssh-proxy:local".to_owned(), capability: "ssh_proxy.connect".to_owned(), grant_id: Some("grant:left-ssh-proxy-connect".to_owned()), }, ) .expect("grant left ssh proxy connect"); handle_request( &right, ControlRequest::AuthGrant { subject: left.node_id.clone(), resource: "resource:document:notes".to_owned(), capability: "document.read".to_owned(), grant_id: Some("grant:left-document-read".to_owned()), }, ) .expect("grant left document read"); handle_request( &right, ControlRequest::AuthGrant { subject: left.node_id.clone(), resource: "resource:db:notes".to_owned(), capability: "db.sync".to_owned(), grant_id: Some("grant:left-db-sync".to_owned()), }, ) .expect("grant left db sync"); let allowed = handle_request_async( &left, ControlRequest::PeerAuthCheck { node: right_card.node_id.to_string(), resource: "resource:cas:local".to_owned(), capability: "cas.fetch".to_owned(), }, ) .await .expect("allowed peer auth check"); match allowed { ControlResponse::PeerAuthChecked { allowed, reason, evaluated_ops, .. } => { assert!(allowed); assert!(reason.contains("direct grant")); assert_eq!(evaluated_ops, 1); } other => panic!("unexpected allowed auth response: {other:?}"), } let fetched = handle_request_async( &left, ControlRequest::CasFetch { node: right_card.node_id.to_string(), hash: right_blob.hash.clone(), }, ) .await .expect("allowed cas fetch"); match fetched { ControlResponse::CasFetched { allowed, hash, size_bytes, reason, note, .. } => { assert!(allowed); assert_eq!(hash, right_blob.hash); assert_eq!(size_bytes, right_blob.size_bytes); assert!(reason.contains("direct grant")); assert!(note.contains("required cas.fetch")); assert_eq!( LocalCas::new(left.paths.cas_dir()) .read_bytes(&right_blob.hash) .expect("left cas read fetched blob"), b"remote cas bytes" ); } other => panic!("unexpected allowed CAS fetch response: {other:?}"), } let providers = Store::open(&left_paths.metadata_db()) .expect("open left after cas fetch") .list_cas_providers(right_blob.hash.as_str()) .expect("list cas providers"); assert_eq!(providers.len(), 1); assert_eq!(providers[0].peer_node_id, right.node_id); 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 published = handle_request_async( &left, ControlRequest::PubsubPub { topic: "presence/test".to_owned(), message: "hello".to_owned(), node: Some(right_card.node_id.to_string()), }, ) .await .expect("allowed remote pubsub publish"); match published { ControlResponse::PubsubRemotePublished { allowed, message, reason, note, .. } => { assert!(allowed); assert_eq!(message.topic.as_str(), "presence/test"); assert_eq!(message.message, "hello"); assert!(reason.contains("direct grant")); assert!(note.contains("lossy")); } other => panic!("unexpected allowed pubsub publish response: {other:?}"), } let remote_messages = handle_request( &right, ControlRequest::PubsubSub { topic: "presence/test".to_owned(), node: None, }, ) .expect("right pubsub sub after remote publish"); match remote_messages { ControlResponse::PubsubMessages { messages, .. } => { assert_eq!(messages.len(), 1); assert_eq!(messages[0].message, "hello"); } other => panic!("unexpected remote pubsub messages response: {other:?}"), } let subscribed = handle_request_async( &left, ControlRequest::PubsubSub { topic: "presence/test".to_owned(), node: Some(right_card.node_id.to_string()), }, ) .await .expect("allowed remote pubsub subscribe"); match subscribed { ControlResponse::PubsubRemoteMessages { allowed, messages, reason, note, .. } => { assert!(allowed); assert_eq!(messages.len(), 1); assert_eq!(messages[0].message, "hello"); assert!(reason.contains("direct grant")); assert!(note.contains("pubsub.subscribe")); } other => panic!("unexpected allowed pubsub subscribe response: {other:?}"), } let remote_pipe = handle_request_async( &left, ControlRequest::PipeConnect { target: "inbox".to_owned(), node: Some(right_card.node_id.to_string()), }, ) .await .expect("allowed remote pipe connect"); match remote_pipe { ControlResponse::PipeRemoteConnected { allowed, connection, reason, note, .. } => { assert!(allowed); assert!(connection.local_listener_found); assert!(reason.contains("direct grant")); assert!(note.contains("byte streams are not implemented yet")); } other => panic!("unexpected allowed pipe connect response: {other:?}"), } let ssh_proxy = handle_request_async( &left, ControlRequest::SshProxyConnect { node: right_card.node_id.to_string(), }, ) .await .expect("allowed SSH proxy connect"); match ssh_proxy { ControlResponse::SshProxyConnected { allowed, connection, reason, note, .. } => { assert!(allowed); let connection = connection.expect("proxy connection metadata"); assert_eq!(connection.target_node.to_string(), right.node_id); assert_eq!( connection.local_sshd_target.as_deref(), Some("127.0.0.1:22") ); assert!(reason.contains("direct grant")); assert!(note.contains("SSH is not a geth transport")); } other => panic!("unexpected allowed SSH proxy response: {other:?}"), } let document_sync = handle_request_async( &left, ControlRequest::DocumentSync { node: right_card.node_id.to_string(), name: "notes".to_owned(), }, ) .await .expect("allowed document sync"); match document_sync { ControlResponse::DocumentSynced { allowed, updated, reason, note, .. } => { assert!(allowed); assert!(updated); assert!(reason.contains("direct grant")); assert!(note.contains("JSON LWW")); } other => panic!("unexpected allowed document sync response: {other:?}"), } let synced_document = handle_request( &left, ControlRequest::DocumentGet { name: "notes".to_owned(), }, ) .expect("left document get synced"); match synced_document { ControlResponse::DocumentGet { state } => { assert_eq!(state.state_json, r#"{"title":"remote"}"#); } other => panic!("unexpected synced document get response: {other:?}"), } let db_sync = handle_request_async( &left, ControlRequest::DbSync { node: right_card.node_id.to_string(), name: "notes".to_owned(), limit: 10, }, ) .await .expect("allowed db sync"); match db_sync { ControlResponse::DbSynced { allowed, changes_received, max_db_version, schema_match, reason, note, .. } => { assert!(allowed); assert_eq!(changes_received, 1); assert_eq!(max_db_version, Some(7)); assert!(schema_match); assert!(reason.contains("direct grant")); assert!(note.contains("db.sync")); assert!(note.contains("applying remote changes is not implemented yet")); } other => panic!("unexpected allowed DB sync response: {other:?}"), } let denied_cert_sync = handle_request_async( &left, ControlRequest::SshCertSync { node: right_card.node_id.to_string(), }, ) .await .expect("denied ssh cert sync"); match denied_cert_sync { ControlResponse::SshCertSynced { allowed, requests_imported, certificates_imported, reason, .. } => { assert!(!allowed); assert_eq!(requests_imported, 0); assert_eq!(certificates_imported, 0); assert!(reason.contains("no active direct or group grant")); } other => panic!("unexpected denied SSH cert sync response: {other:?}"), } handle_request( &right, ControlRequest::AuthGrant { subject: left.node_id.clone(), resource: "resource:ssh:certs".to_owned(), capability: "ssh_cert.sync".to_owned(), grant_id: Some("grant:left-ssh-cert-sync".to_owned()), }, ) .expect("grant left ssh cert sync"); handle_request( &right, ControlRequest::AuthGrant { subject: left.node_id.clone(), resource: "resource:ssh:revocations".to_owned(), capability: "ssh_revocation.sync".to_owned(), grant_id: Some("grant:left-ssh-revocation-sync".to_owned()), }, ) .expect("grant left ssh revocation sync"); let cert_sync = handle_request_async( &left, ControlRequest::SshCertSync { node: right_card.node_id.to_string(), }, ) .await .expect("allowed ssh cert sync"); match cert_sync { ControlResponse::SshCertSynced { allowed, requests_imported, certificates_imported, reason, note, .. } => { assert!(allowed); assert_eq!(requests_imported, 1); assert_eq!(certificates_imported, 0); assert!(reason.contains("direct grant")); assert!(note.contains("ssh_cert.sync")); } other => panic!("unexpected allowed SSH cert sync response: {other:?}"), } let left_requests = Store::open(&left_paths.metadata_db()) .expect("open left after cert sync") .list_ssh_cert_requests() .expect("list synced requests"); assert!( left_requests .iter() .any(|request| request.request_id == right_request_id.as_str()) ); let revocation_sync = handle_request_async( &left, ControlRequest::SshRevocationSync { node: right_card.node_id.to_string(), }, ) .await .expect("allowed ssh revocation sync"); match revocation_sync { ControlResponse::SshRevocationSynced { allowed, revocations_imported, reason, note, .. } => { assert!(allowed); assert_eq!(revocations_imported, 1); assert!(reason.contains("direct grant")); assert!(note.contains("ssh_revocation.sync")); } other => panic!("unexpected allowed SSH revocation sync response: {other:?}"), } let left_revocations = Store::open(&left_paths.metadata_db()) .expect("open left after revocation sync") .list_ssh_revocations() .expect("list synced revocations"); assert!( left_revocations .iter() .any(|revocation| revocation.revocation_id == right_revocation_id.as_str()) ); tokio::time::sleep(Duration::from_millis(2)).await; let second_pubkey = right_home.path().join("request-2.pub"); std::fs::write( &second_pubkey, "ssh-ed25519 AAAAC3NzaC1lZDI1NTE5AAAAIGV0aDI= node2\n", ) .expect("write second public key"); let second_requested = handle_request( &right, ControlRequest::SshCertRequest { public_key_path: second_pubkey, cert_kind: "user".to_owned(), principals: vec!["admin".to_owned()], requested_validity: Some("+4w".to_owned()), renewal_of: None, reason: Some("background sync test".to_owned()), subject: None, }, ) .expect("right second ssh cert request"); let second_request_id = match second_requested { ControlResponse::SshCertRequested { request } => request.id, other => panic!("unexpected second SSH cert request response: {other:?}"), }; let second_revocation = handle_request( &right, ControlRequest::SshRevocationAdd { kind: "key-id".to_owned(), target: "newly-revoked-key".to_owned(), reason: Some("background sync test".to_owned()), subject: None, }, ) .expect("right second ssh revocation"); let second_revocation_id = match second_revocation { 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"); handle_request( &right, ControlRequest::DocumentSet { name: "notes".to_owned(), state_json: r#"{"title":"live"}"#.to_owned(), }, ) .expect("right live document set"); insert_mock_crsqlite_change(&right_db_path, 8, "live"); run_live_sync_once(&left) .await .expect("background live sync tick"); let left_store = Store::open(&left_paths.metadata_db()).expect("open left after live sync"); assert!( left_store .list_ssh_cert_requests() .expect("list live-synced requests") .iter() .any(|request| request.request_id == second_request_id.as_str()) ); assert!( left_store .list_ssh_revocations() .expect("list live-synced revocations") .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()) ); assert_eq!( left_store .get_document_resource_by_name("notes") .expect("get live-synced document") .map(|document| document.state_json), Some(r#"{"title":"live"}"#.to_owned()) ); let db_cursor = left_store .get_module_state(&live_sync_cursor_key( right_card.node_id.as_str(), "db:notes", )) .expect("get live-synced db cursor") .expect("db cursor exists"); assert_eq!( serde_json::from_str::(&db_cursor.state_json) .expect("parse db cursor") .cursor_ms, 8 ); left_endpoint.shutdown().await; right_endpoint.shutdown().await; } }