diff --git a/crates/geth-node/src/lib.rs b/crates/geth-node/src/lib.rs index 614d550..f615a20 100644 --- a/crates/geth-node/src/lib.rs +++ b/crates/geth-node/src/lib.rs @@ -222,6 +222,12 @@ struct PipeUnixConnectWire { bearer_proof: Option, } +struct PeerControlCaller { + peer_card: PeerCard, + remote_endpoint_id: String, + store: Store, +} + #[derive(Clone, Debug)] pub struct InitOwnerOptions { pub admin_key_path: Option, @@ -4114,49 +4120,26 @@ async fn handle_iroh_control_connection( tracing::debug!(remote_endpoint_id = %log_remote_endpoint_id, alpn = %log_alpn, request_bytes, "read iroh peer-control request line"); let response = match geth_control::decode_peer_request(&request_line)? { 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, - })?; + let caller = authenticate_peer_control_caller(&node, peer_card, &remote_endpoint_id)?; 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, + remote_endpoint_id: caller.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())?; + let caller = authenticate_peer_control_caller(&node, peer_card, &remote_endpoint_id)?; + let watermarks = + sync_watermarks_for_peer(&caller.store, caller.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, + remote_endpoint_id: caller.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(), @@ -4167,21 +4150,9 @@ async fn handle_iroh_control_connection( 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 ops = load_keychain_ops(&store)?; - let signatures = load_keychain_signatures(&store)?; + let caller = authenticate_peer_control_caller(&node, peer_card, &remote_endpoint_id)?; + let ops = load_keychain_ops(&caller.store)?; + let signatures = load_keychain_signatures(&caller.store)?; let returned_op_ids = ops .iter() .filter(|op| op.created_at.0 >= since_ms) @@ -4197,7 +4168,7 @@ async fn handle_iroh_control_connection( 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, + remote_endpoint_id: caller.remote_endpoint_id, ops: ops .into_iter() .filter(|op| op.created_at.0 >= since_ms) @@ -4219,21 +4190,9 @@ async fn handle_iroh_control_connection( 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 ops = load_auth_ops(&store)?; - let signatures = load_auth_signatures(&store)?; + let caller = authenticate_peer_control_caller(&node, peer_card, &remote_endpoint_id)?; + let ops = load_auth_ops(&caller.store)?; + let signatures = load_auth_signatures(&caller.store)?; let returned_op_ids = ops .iter() .filter(|op| op.created_at.0 >= since_ms) @@ -4249,7 +4208,7 @@ async fn handle_iroh_control_connection( 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, + remote_endpoint_id: caller.remote_endpoint_id, ops: ops .into_iter() .filter(|op| op.created_at.0 >= since_ms) @@ -4271,25 +4230,14 @@ async fn handle_iroh_control_connection( request, 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 accepted = insert_node_enrollment_request_if_not_conflicting(&store, &request)?; + let caller = authenticate_peer_control_caller(&node, peer_card, &remote_endpoint_id)?; + let accepted = + insert_node_enrollment_request_if_not_conflicting(&caller.store, &request)?; PeerControlResponse::NodeEnrollmentSubmitted { 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, + remote_endpoint_id: caller.remote_endpoint_id, request_id: request.id.to_string(), accepted, nonce, @@ -4302,22 +4250,10 @@ async fn handle_iroh_control_connection( 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 caller = authenticate_peer_control_caller(&node, peer_card, &remote_endpoint_id)?; let explanation = geth_auth::explain_auth_ops( - &load_auth_ops_for_resource(&store, &resource)?, - PrincipalId::new(peer_card.node_id.to_string()), + &load_auth_ops_for_resource(&caller.store, &resource)?, + PrincipalId::new(caller.peer_card.node_id.to_string()), ResourceId::new(resource.clone()), Capability::new(capability.clone()), ); @@ -4325,7 +4261,7 @@ async fn handle_iroh_control_connection( 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, + remote_endpoint_id: caller.remote_endpoint_id, resource, capability, allowed: explanation.allowed, @@ -4341,24 +4277,12 @@ async fn handle_iroh_control_connection( nonce, bearer_proof, } => { - 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 caller = authenticate_peer_control_caller(&node, peer_card, &remote_endpoint_id)?; let resource = "resource:cas:local".to_owned(); let capability = "cas.fetch".to_owned(); let explanation = explain_peer_or_bearer( - &store, - peer_card.node_id.as_str(), + &caller.store, + caller.peer_card.node_id.as_str(), &resource, &capability, &nonce, @@ -4372,7 +4296,7 @@ async fn handle_iroh_control_connection( 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, + remote_endpoint_id: caller.remote_endpoint_id, hash, size_bytes, content_base64: Some( @@ -4394,7 +4318,7 @@ async fn handle_iroh_control_connection( 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, + remote_endpoint_id: caller.remote_endpoint_id, hash, size_bytes: 0, content_base64: None, @@ -4413,31 +4337,19 @@ async fn handle_iroh_control_connection( bearer_proof, } => { geth_cas::validate_file_root_name(&name)?; - 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 caller = authenticate_peer_control_caller(&node, peer_card, &remote_endpoint_id)?; let resource = format!("resource:cas-tree:{name}"); let capability = "cas.fetch".to_owned(); let explanation = explain_peer_or_bearer( - &store, - peer_card.node_id.as_str(), + &caller.store, + caller.peer_card.node_id.as_str(), &resource, &capability, &nonce, bearer_proof.as_ref(), )?; let (root, tree_content_base64) = if explanation.allowed { - if let Some(root) = store.get_file_root_by_name(&name)? { + if let Some(root) = caller.store.get_file_root_by_name(&name)? { let root_view = file_root_from_stored(&root); let tree_content_base64 = root .latest_tree_hash @@ -4462,7 +4374,7 @@ async fn handle_iroh_control_connection( 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, + remote_endpoint_id: caller.remote_endpoint_id, name, root, tree_content_base64, @@ -4479,25 +4391,13 @@ async fn handle_iroh_control_connection( nonce, bearer_proof, } => { - 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 caller = authenticate_peer_control_caller(&node, peer_card, &remote_endpoint_id)?; let resource = "resource:ssh:certs".to_owned(); let capability = "ssh_cert.sync".to_owned(); let high_water_ms = geth_store::now_ms(); let explanation = explain_peer_or_bearer( - &store, - peer_card.node_id.as_str(), + &caller.store, + caller.peer_card.node_id.as_str(), &resource, &capability, &nonce, @@ -4505,12 +4405,14 @@ async fn handle_iroh_control_connection( )?; let (requests, certificates, log_entries) = if explanation.allowed { ( - store + caller + .store .list_ssh_cert_requests_since(since_ms)? .into_iter() .map(ssh_cert_request_from_stored) .collect::, _>>()?, - store + caller + .store .list_ssh_certificates_since(since_ms)? .into_iter() .map(ssh_certificate_from_stored) @@ -4529,7 +4431,7 @@ async fn handle_iroh_control_connection( 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, + remote_endpoint_id: caller.remote_endpoint_id, requests, certificates, log_entries, @@ -4547,32 +4449,21 @@ async fn handle_iroh_control_connection( nonce, bearer_proof, } => { - 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 caller = authenticate_peer_control_caller(&node, peer_card, &remote_endpoint_id)?; let resource = "resource:ssh:revocations".to_owned(); let capability = "ssh_revocation.sync".to_owned(); let high_water_ms = geth_store::now_ms(); let explanation = explain_peer_or_bearer( - &store, - peer_card.node_id.as_str(), + &caller.store, + caller.peer_card.node_id.as_str(), &resource, &capability, &nonce, bearer_proof.as_ref(), )?; let (revocations, log_entries) = if explanation.allowed { - store + caller + .store .list_ssh_revocations_since(since_ms)? .into_iter() .map(ssh_revocation_from_stored) @@ -4589,7 +4480,7 @@ async fn handle_iroh_control_connection( 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, + remote_endpoint_id: caller.remote_endpoint_id, revocations, log_entries, high_water_ms, @@ -4608,19 +4499,11 @@ async fn handle_iroh_control_connection( bearer_proof, } => { 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, - })?; + let PeerControlCaller { + peer_card, + remote_endpoint_id, + store, + } = authenticate_peer_control_caller(&node, peer_card, &remote_endpoint_id)?; if let Some(kv) = store.get_kv_store_by_name(&name)? { let capability = "kv.read".to_owned(); let explanation = explain_peer_or_bearer( @@ -4686,19 +4569,11 @@ async fn handle_iroh_control_connection( } => { 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 PeerControlCaller { + peer_card, + remote_endpoint_id, + store, + } = authenticate_peer_control_caller(&node, peer_card, &remote_endpoint_id)?; let resource = format!("resource:pubsub:{topic}"); let capability = "pubsub.publish".to_owned(); let explanation = explain_peer_or_bearer( @@ -4746,19 +4621,11 @@ async fn handle_iroh_control_connection( bearer_proof, } => { 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 PeerControlCaller { + peer_card, + remote_endpoint_id, + store, + } = authenticate_peer_control_caller(&node, peer_card, &remote_endpoint_id)?; let resource = format!("resource:pubsub:{topic}"); let capability = "pubsub.subscribe".to_owned(); let explanation = explain_peer_or_bearer( @@ -4800,24 +4667,12 @@ async fn handle_iroh_control_connection( bearer_proof, } => { 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 caller = authenticate_peer_control_caller(&node, peer_card, &remote_endpoint_id)?; let resource = format!("resource:pipe:{target}"); let capability = "pipe.connect".to_owned(); let explanation = explain_peer_or_bearer( - &store, - peer_card.node_id.as_str(), + &caller.store, + caller.peer_card.node_id.as_str(), &resource, &capability, &nonce, @@ -4836,7 +4691,7 @@ async fn handle_iroh_control_connection( 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, + remote_endpoint_id: caller.remote_endpoint_id, connection, allowed: explanation.allowed, reason: explanation.reason, @@ -4852,24 +4707,12 @@ async fn handle_iroh_control_connection( bearer_proof, } => { geth_pipe::validate_pipe_name(&name)?; - 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 caller = authenticate_peer_control_caller(&node, peer_card, &remote_endpoint_id)?; let resource = format!("resource:pipe:{name}"); let capability = "pipe.listen".to_owned(); let explanation = explain_peer_or_bearer( - &store, - peer_card.node_id.as_str(), + &caller.store, + caller.peer_card.node_id.as_str(), &resource, &capability, &nonce, @@ -4884,7 +4727,7 @@ async fn handle_iroh_control_connection( 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, + remote_endpoint_id: caller.remote_endpoint_id, listener, allowed: explanation.allowed, reason: explanation.reason, @@ -4898,24 +4741,12 @@ async fn handle_iroh_control_connection( nonce, bearer_proof, } => { - 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 caller = authenticate_peer_control_caller(&node, peer_card, &remote_endpoint_id)?; let resource = "resource:ssh-proxy:local".to_owned(); let capability = "ssh_proxy.connect".to_owned(); let explanation = explain_peer_or_bearer( - &store, - peer_card.node_id.as_str(), + &caller.store, + caller.peer_card.node_id.as_str(), &resource, &capability, &nonce, @@ -4936,7 +4767,7 @@ async fn handle_iroh_control_connection( 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, + remote_endpoint_id: caller.remote_endpoint_id, connection, allowed: explanation.allowed, reason: explanation.reason, @@ -4951,24 +4782,12 @@ async fn handle_iroh_control_connection( nonce, bearer_proof, } => { - 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 caller = authenticate_peer_control_caller(&node, peer_card, &remote_endpoint_id)?; let resource = "resource:ssh-proxy:local".to_owned(); let capability = "ssh_proxy.admin_shell".to_owned(); let explanation = explain_peer_or_bearer( - &store, - peer_card.node_id.as_str(), + &caller.store, + caller.peer_card.node_id.as_str(), &resource, &capability, &nonce, @@ -4983,7 +4802,7 @@ async fn handle_iroh_control_connection( 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, + remote_endpoint_id: caller.remote_endpoint_id, command, output, allowed: explanation.allowed, @@ -5002,25 +4821,13 @@ async fn handle_iroh_control_connection( } => { 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 caller = authenticate_peer_control_caller(&node, peer_card, &remote_endpoint_id)?; + if let Some(document) = caller.store.get_document_resource_by_name(&name)? { let capability = "document.read".to_owned(); let high_water_ms = document.updated_at_ms; let explanation = explain_peer_or_bearer( - &store, - peer_card.node_id.as_str(), + &caller.store, + caller.peer_card.node_id.as_str(), &document.resource_id, &capability, &nonce, @@ -5035,7 +4842,7 @@ async fn handle_iroh_control_connection( 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, + remote_endpoint_id: caller.remote_endpoint_id, name, state, high_water_ms, @@ -5060,24 +4867,12 @@ async fn handle_iroh_control_connection( bearer_proof, } => { 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 caller = authenticate_peer_control_caller(&node, peer_card, &remote_endpoint_id)?; + if let Some(db) = caller.store.get_db_resource_by_name(&name)? { let capability = "db.sync".to_owned(); let explanation = explain_peer_or_bearer( - &store, - peer_card.node_id.as_str(), + &caller.store, + caller.peer_card.node_id.as_str(), &db.resource_id, &capability, &nonce, @@ -5108,7 +4903,7 @@ async fn handle_iroh_control_connection( 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, + remote_endpoint_id: caller.remote_endpoint_id, name, batch, high_water_db_version, @@ -5164,19 +4959,11 @@ async fn handle_ssh_proxy_wire_connection( return Ok(()); }; - 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 PeerControlCaller { + peer_card, + remote_endpoint_id, + store, + } = authenticate_peer_control_caller(&node, peer_card, &remote_endpoint_id)?; let resource = "resource:ssh-proxy:local".to_owned(); let capability = "ssh_proxy.connect".to_owned(); let explanation = explain_peer_or_bearer( @@ -5389,24 +5176,12 @@ async fn handle_overlay_packet_wire_request( .decode(&packet_base64) .map_err(|error| NodeError::IrohPeer(format!("invalid overlay packet base64: {error}")))?; geth_overlay::validate_ipv4_packet(&packet)?; - 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 caller = authenticate_peer_control_caller(node, peer_card, remote_endpoint_id)?; let resource = geth_overlay::overlay_resource_id(&network).to_string(); let capability = "overlay.route".to_owned(); let explanation = explain_peer_or_bearer( - &store, - peer_card.node_id.as_str(), + &caller.store, + caller.peer_card.node_id.as_str(), &resource, &capability, &nonce, @@ -5415,9 +5190,9 @@ async fn handle_overlay_packet_wire_request( let packet = if explanation.allowed { let injected = overlay_runtime_inject(node, &network, packet.clone())?; let mut packet = record_overlay_packet( - &store, + &caller.store, &network, - peer_card.node_id.as_str(), + caller.peer_card.node_id.as_str(), &node.node_id, packet_base64, packet.len(), @@ -5434,7 +5209,7 @@ async fn handle_overlay_packet_wire_request( 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: remote_endpoint_id.to_owned(), + remote_endpoint_id: caller.remote_endpoint_id, packet, allowed: explanation.allowed, reason: explanation.reason, @@ -5457,24 +5232,12 @@ fn handle_pipe_send_wire_request( base64::engine::general_purpose::STANDARD .decode(&data_base64) .map_err(|error| NodeError::IrohPeer(format!("invalid pipe payload: {error}")))?; - 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 caller = authenticate_peer_control_caller(node, peer_card, remote_endpoint_id)?; let resource = format!("resource:pipe:{target}"); let capability = "pipe.connect".to_owned(); let explanation = explain_peer_or_bearer( - &store, - peer_card.node_id.as_str(), + &caller.store, + caller.peer_card.node_id.as_str(), &resource, &capability, &nonce, @@ -5485,7 +5248,7 @@ fn handle_pipe_send_wire_request( node, target, data_base64, - Some(peer_card.node_id.to_string()), + Some(caller.peer_card.node_id.to_string()), "remote pipe byte message over dedicated Iroh pipe ALPN".to_owned(), )? } else { @@ -5495,7 +5258,7 @@ fn handle_pipe_send_wire_request( 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: remote_endpoint_id.to_owned(), + remote_endpoint_id: caller.remote_endpoint_id, message: Box::new(message), listener_found, allowed: explanation.allowed, @@ -5520,24 +5283,12 @@ async fn handle_pipe_tcp_wire_connection( bearer_proof, } = request; let target = geth_pipe::validate_tcp_forward_target_addr(&target_addr)?; - 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 caller = authenticate_peer_control_caller(&node, peer_card, remote_endpoint_id)?; let resource = format!("resource:pipe-tcp:{target_addr}"); let capability = "pipe.forward".to_owned(); let explanation = explain_peer_or_bearer( - &store, - peer_card.node_id.as_str(), + &caller.store, + caller.peer_card.node_id.as_str(), &resource, &capability, &nonce, @@ -5548,7 +5299,7 @@ async fn handle_pipe_tcp_wire_connection( 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: remote_endpoint_id.to_owned(), + remote_endpoint_id: caller.remote_endpoint_id, connection: None, allowed: false, reason: explanation.reason, @@ -5583,7 +5334,7 @@ async fn handle_pipe_tcp_wire_connection( 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: remote_endpoint_id.to_owned(), + remote_endpoint_id: caller.remote_endpoint_id, connection: Some(PipeConnection { target: target_addr.clone(), connected_at: UnixMillis(geth_store::now_ms()), @@ -5632,24 +5383,12 @@ async fn handle_pipe_unix_wire_connection( } = request; let target_path = geth_pipe::validate_unix_forward_path(&target_path)?; let target_display = target_path.display().to_string(); - 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 caller = authenticate_peer_control_caller(&node, peer_card, remote_endpoint_id)?; let resource = format!("resource:pipe-unix:{target_display}"); let capability = "pipe.forward".to_owned(); let explanation = explain_peer_or_bearer( - &store, - peer_card.node_id.as_str(), + &caller.store, + caller.peer_card.node_id.as_str(), &resource, &capability, &nonce, @@ -5660,7 +5399,7 @@ async fn handle_pipe_unix_wire_connection( 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: remote_endpoint_id.to_owned(), + remote_endpoint_id: caller.remote_endpoint_id, connection: None, allowed: false, reason: explanation.reason, @@ -5696,7 +5435,7 @@ async fn handle_pipe_unix_wire_connection( 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: remote_endpoint_id.to_owned(), + remote_endpoint_id: caller.remote_endpoint_id, connection: Some(PipeConnection { target: target_display.clone(), connected_at: UnixMillis(geth_store::now_ms()), @@ -5774,6 +5513,31 @@ fn ensure_peer_card_matches_endpoint(card: &PeerCard, endpoint_id: &str) -> Resu } } +fn authenticate_peer_control_caller( + node: &LocalNode, + peer_card: PeerCard, + remote_endpoint_id: &str, +) -> Result { + 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, + })?; + Ok(PeerControlCaller { + peer_card, + remote_endpoint_id: remote_endpoint_id.to_owned(), + store, + }) +} + fn peer_gossip_endpoint_id( node: &LocalNode, peer_node: &str, diff --git a/docs/architecture.md b/docs/architecture.md index 99a18df..703e435 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -101,6 +101,11 @@ authenticates the Iroh endpoint and peer-card signature, but it does not authorize any resource module. Protected peer control requests must also prove that the signed peer card binds the observed Iroh EndpointID, then reduce resource auth ops; an EndpointID alone is not accepted as a resource principal. +Inbound protected control, pipe-wire, SSH-proxy, and overlay-wire handlers use a +shared authenticated caller context that validates the peer card, checks the +observed endpoint binding, records the peer card as untrusted candidate +metadata, and then passes the verified caller identity plus store handle to the +resource-specific authorization step. Remote geth JSONL messages over Iroh are bounded before decoding: peer-control and module wire requests/responses are limited to 16 MiB, streaming handshakes for SSH proxy and TCP/Unix pipe forwarding are limited to diff --git a/docs/production-readiness-roadmap.md b/docs/production-readiness-roadmap.md index 0031608..d40aa1e 100644 --- a/docs/production-readiness-roadmap.md +++ b/docs/production-readiness-roadmap.md @@ -52,17 +52,17 @@ behavior. - `[ ]` Each command family has a small handler module or function group. - `[x]` Local-only behavior remains covered by existing integration tests. -- `[ ]` Extract protected peer-control routing. +- `[~]` Extract protected peer-control routing. Acceptance criteria: - `[x]` Shared bounded Iroh line-read and send-finish helpers live outside the main feature handler module. - `[x]` Outbound peer-control, pipe-wire, and overlay-wire request helpers live outside the main feature handler module. - - `[ ]` Iroh control ALPN handling, nonce checks, peer-card validation, and + - `[~]` Iroh control ALPN handling, nonce checks, peer-card validation, and endpoint-binding validation are centralized. - - `[ ]` Feature handlers receive authenticated caller context rather than + - `[x]` Feature handlers receive authenticated caller context rather than repeating peer-card boilerplate. - - `[ ]` Remote request tests still prove discovery alone grants no access. + - `[x]` Remote request tests still prove discovery alone grants no access. - `[ ]` Extract resource module handlers. Acceptance criteria: