diff --git a/crates/geth-node/src/lib.rs b/crates/geth-node/src/lib.rs index f784871..10c6e6b 100644 --- a/crates/geth-node/src/lib.rs +++ b/crates/geth-node/src/lib.rs @@ -1,4 +1,5 @@ mod daemon; +mod peer_client; mod runtime; pub mod service; mod sync; @@ -59,6 +60,7 @@ use geth_types::{ }; use iroh::protocol::ProtocolHandler; use iroh_docs::api::protocol::{AddrInfoOptions, ShareMode}; +use peer_client::{request_overlay_wire, request_peer_control, request_pipe_wire}; use runtime::{ NodeRuntime, OverlayTunCounters, OverlayTunRuntime, PipeRuntime, PubsubGossipTopicRuntime, PubsubRuntime, @@ -4073,269 +4075,6 @@ fn explain_peer_or_bearer( }) } -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 peer_node = resolve_peer_node_for_control(&store, peer_node); - let stored = store - .get_peer_card(&peer_node)? - .ok_or_else(|| NodeError::PeerNotFound(peer_node.clone()))?; - 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, 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()))?; - send.stopped() - .await - .map_err(|error| NodeError::IrohPeer(error.to_string()))?; - drop(send); - let mut response_line = String::new(); - let mut reader = BufReader::new(recv); - reader - .read_line(&mut response_line) - .await - .map_err(|error| NodeError::IrohPeer(error.to_string()))?; - let response = geth_control::decode_peer_response(&response_line)?; - match &response { - PeerControlResponse::SshCertSynced { - nonce: response_nonce, - .. - } - | PeerControlResponse::KeychainSynced { - nonce: response_nonce, - .. - } - | PeerControlResponse::AuthSynced { - nonce: response_nonce, - .. - } - | PeerControlResponse::NodeEnrollmentSubmitted { - 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::PipeListening { - nonce: response_nonce, - .. - } - | PeerControlResponse::CasRootSynced { - nonce: response_nonce, - .. - } - | PeerControlResponse::SshProxyConnected { - nonce: response_nonce, - .. - } - | PeerControlResponse::SshAdminShellOutput { - 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" - ))), - } -} - -async fn request_pipe_wire( - node: &LocalNode, - peer_node: &str, - operation: &str, - build_request: impl FnOnce(PeerCard, String) -> PipeWireRequest, -) -> 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_PIPE) - .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_pipe_wire_request(&request)?.as_bytes()) - .await - .map_err(|error| NodeError::IrohPeer(error.to_string()))?; - finish_iroh_send(&mut send).await?; - drop(send); - let text = read_iroh_line(&mut recv, 16 * 1024 * 1024).await?; - let response = geth_control::decode_pipe_wire_response(&text)?; - match &response { - PipeWireResponse::Sent { - nonce: response_nonce, - .. - } if response_nonce == &nonce => Ok(response), - PipeWireResponse::Error { .. } => Ok(response), - _ => Err(NodeError::IrohPeer(format!( - "peer {operation} response did not match request" - ))), - } -} - -async fn request_overlay_wire( - node: &LocalNode, - peer_node: &str, - operation: &str, - build_request: impl FnOnce(PeerCard, String) -> OverlayWireRequest, -) -> Result { - let store = Store::open(&node.paths.metadata_db())?; - let peer_node = resolve_peer_node_for_control(&store, peer_node); - let stored = store - .get_peer_card(&peer_node)? - .ok_or_else(|| NodeError::PeerNotFound(peer_node.clone()))?; - 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_OVERLAY) - .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_overlay_wire_request(&request)?.as_bytes()) - .await - .map_err(|error| NodeError::IrohPeer(error.to_string()))?; - finish_iroh_send(&mut send).await?; - drop(send); - let text = read_iroh_line(&mut recv, 16 * 1024 * 1024).await?; - let response = geth_control::decode_overlay_wire_response(&text)?; - match &response { - OverlayWireResponse::PacketAccepted { - nonce: response_nonce, - .. - } if response_nonce == &nonce => Ok(response), - OverlayWireResponse::Error { .. } => Ok(response), - _ => Err(NodeError::IrohPeer(format!( - "peer {operation} response did not match request" - ))), - } -} - async fn run_sync_for_peer(node: &LocalNode, peer_node: &str) -> Result { let kv_stores = Store::open(&node.paths.metadata_db())?.list_kv_stores()?; let documents = Store::open(&node.paths.metadata_db())?.list_document_resources()?; diff --git a/crates/geth-node/src/peer_client.rs b/crates/geth-node/src/peer_client.rs new file mode 100644 index 0000000..1ecde06 --- /dev/null +++ b/crates/geth-node/src/peer_client.rs @@ -0,0 +1,260 @@ +//! Outbound Iroh peer request helpers. + +use crate::wire::{finish_iroh_send, read_iroh_line}; +use crate::{LocalNode, NodeError}; +use geth_control::{ + OverlayWireRequest, OverlayWireResponse, PeerControlRequest, PeerControlResponse, + PipeWireRequest, PipeWireResponse, +}; +use geth_discovery::{DiscoverySource, PeerCard}; +use geth_store::Store; +use tokio::io::{AsyncBufReadExt, BufReader}; + +pub(crate) 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 peer_node = super::resolve_peer_node_for_control(&store, peer_node); + let stored = store + .get_peer_card(&peer_node)? + .ok_or_else(|| NodeError::PeerNotFound(peer_node.clone()))?; + 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)?; + super::ensure_peer_card_matches_endpoint(&peer_card, &candidate.endpoint_id)?; + let node_addr = super::iroh_node_addr_from_candidate(candidate)?; + let endpoint = node + .iroh_endpoint + .lock() + .map_err(|_| NodeError::RuntimeLockPoisoned)? + .clone() + .ok_or(NodeError::IrohEndpointUnavailable)?; + let self_card = super::local_peer_card(node, DiscoverySource::PeerExchange, true).await?; + let nonce = request_nonce(&node.node_id, &peer_node, operation); + 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, 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()))?; + send.stopped() + .await + .map_err(|error| NodeError::IrohPeer(error.to_string()))?; + drop(send); + let mut response_line = String::new(); + let mut reader = BufReader::new(recv); + reader + .read_line(&mut response_line) + .await + .map_err(|error| NodeError::IrohPeer(error.to_string()))?; + let response = geth_control::decode_peer_response(&response_line)?; + match &response { + PeerControlResponse::SshCertSynced { + nonce: response_nonce, + .. + } + | PeerControlResponse::KeychainSynced { + nonce: response_nonce, + .. + } + | PeerControlResponse::AuthSynced { + nonce: response_nonce, + .. + } + | PeerControlResponse::NodeEnrollmentSubmitted { + 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::PipeListening { + nonce: response_nonce, + .. + } + | PeerControlResponse::CasRootSynced { + nonce: response_nonce, + .. + } + | PeerControlResponse::SshProxyConnected { + nonce: response_nonce, + .. + } + | PeerControlResponse::SshAdminShellOutput { + 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" + ))), + } +} + +pub(crate) async fn request_pipe_wire( + node: &LocalNode, + peer_node: &str, + operation: &str, + build_request: impl FnOnce(PeerCard, String) -> PipeWireRequest, +) -> 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)?; + super::ensure_peer_card_matches_endpoint(&peer_card, &candidate.endpoint_id)?; + let node_addr = super::iroh_node_addr_from_candidate(candidate)?; + let endpoint = node + .iroh_endpoint + .lock() + .map_err(|_| NodeError::RuntimeLockPoisoned)? + .clone() + .ok_or(NodeError::IrohEndpointUnavailable)?; + let self_card = super::local_peer_card(node, DiscoverySource::PeerExchange, true).await?; + let nonce = request_nonce(&node.node_id, peer_node, operation); + let request = build_request(self_card, nonce.clone()); + let conn = endpoint + .endpoint() + .connect(node_addr, geth_iroh::ALPN_PIPE) + .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_pipe_wire_request(&request)?.as_bytes()) + .await + .map_err(|error| NodeError::IrohPeer(error.to_string()))?; + finish_iroh_send(&mut send).await?; + drop(send); + let text = read_iroh_line(&mut recv, 16 * 1024 * 1024).await?; + let response = geth_control::decode_pipe_wire_response(&text)?; + match &response { + PipeWireResponse::Sent { + nonce: response_nonce, + .. + } if response_nonce == &nonce => Ok(response), + PipeWireResponse::Error { .. } => Ok(response), + _ => Err(NodeError::IrohPeer(format!( + "peer {operation} response did not match request" + ))), + } +} + +pub(crate) async fn request_overlay_wire( + node: &LocalNode, + peer_node: &str, + operation: &str, + build_request: impl FnOnce(PeerCard, String) -> OverlayWireRequest, +) -> Result { + let store = Store::open(&node.paths.metadata_db())?; + let peer_node = super::resolve_peer_node_for_control(&store, peer_node); + let stored = store + .get_peer_card(&peer_node)? + .ok_or_else(|| NodeError::PeerNotFound(peer_node.clone()))?; + 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)?; + super::ensure_peer_card_matches_endpoint(&peer_card, &candidate.endpoint_id)?; + let node_addr = super::iroh_node_addr_from_candidate(candidate)?; + let endpoint = node + .iroh_endpoint + .lock() + .map_err(|_| NodeError::RuntimeLockPoisoned)? + .clone() + .ok_or(NodeError::IrohEndpointUnavailable)?; + let self_card = super::local_peer_card(node, DiscoverySource::PeerExchange, true).await?; + let nonce = request_nonce(&node.node_id, &peer_node, operation); + let request = build_request(self_card, nonce.clone()); + let conn = endpoint + .endpoint() + .connect(node_addr, geth_iroh::ALPN_OVERLAY) + .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_overlay_wire_request(&request)?.as_bytes()) + .await + .map_err(|error| NodeError::IrohPeer(error.to_string()))?; + finish_iroh_send(&mut send).await?; + drop(send); + let text = read_iroh_line(&mut recv, 16 * 1024 * 1024).await?; + let response = geth_control::decode_overlay_wire_response(&text)?; + match &response { + OverlayWireResponse::PacketAccepted { + nonce: response_nonce, + .. + } if response_nonce == &nonce => Ok(response), + OverlayWireResponse::Error { .. } => Ok(response), + _ => Err(NodeError::IrohPeer(format!( + "peer {operation} response did not match request" + ))), + } +} + +fn request_nonce(local_node_id: &str, peer_node: &str, operation: &str) -> String { + geth_crypto::blake3_hex( + format!( + "{}\0{}\0{}\0{}", + local_node_id, + peer_node, + operation, + geth_store::now_ms() + ) + .as_bytes(), + ) +} diff --git a/docs/production-readiness-roadmap.md b/docs/production-readiness-roadmap.md index d8a7390..b42fd1b 100644 --- a/docs/production-readiness-roadmap.md +++ b/docs/production-readiness-roadmap.md @@ -56,6 +56,8 @@ behavior. 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 endpoint-binding validation are centralized. - `[ ]` Feature handlers receive authenticated caller context rather than