refactor: extract outbound peer client

This commit is contained in:
Eric Wendland 2026-07-05 17:43:24 +02:00
commit 5ccebe1978
3 changed files with 264 additions and 263 deletions

View file

@ -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<PeerControlResponse, NodeError> {
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<PipeWireResponse, NodeError> {
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<OverlayWireResponse, NodeError> {
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<SyncPeerRun, NodeError> {
let kv_stores = Store::open(&node.paths.metadata_db())?.list_kv_stores()?;
let documents = Store::open(&node.paths.metadata_db())?.list_document_resources()?;

View file

@ -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<PeerControlResponse, NodeError> {
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<PipeWireResponse, NodeError> {
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<OverlayWireResponse, NodeError> {
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(),
)
}

View file

@ -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