Add Iroh peer ping

This commit is contained in:
Eric Wendland 2026-05-18 12:09:50 +02:00
commit 679eeb48a3
15 changed files with 619 additions and 18 deletions

View file

@ -28,3 +28,7 @@ geth-secrets = { path = "../geth-secrets" }
geth-ssh-identity = { path = "../geth-ssh-identity" }
geth-store = { path = "../geth-store" }
geth-types = { path = "../geth-types" }
iroh.workspace = true
[dev-dependencies]
tempfile.workspace = true

View file

@ -8,7 +8,7 @@ use geth_cas::{
use geth_config::{GethConfig, GethPaths, RelayMode};
use geth_control::{
CasBlob, ControlRequest, ControlResponse, KeychainStatusResponse, NodeIdResponse,
StatusResponse,
PeerControlRequest, PeerControlResponse, StatusResponse,
};
use geth_crypto::AgentKey;
use geth_db::DbResource;
@ -112,6 +112,10 @@ pub enum NodeError {
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)]
@ -120,6 +124,7 @@ pub struct LocalNode {
pub agent_id: String,
pub node_id: String,
pub iroh_status: EndpointStatus,
iroh_endpoint: Arc<Mutex<Option<GethIrohEndpoint>>>,
runtime: Arc<NodeRuntime>,
}
@ -165,6 +170,7 @@ pub fn init_node(paths: &GethPaths) -> Result<LocalNode, NodeError> {
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()),
@ -179,6 +185,14 @@ pub fn open_node(paths: &GethPaths) -> Result<LocalNode, NodeError> {
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);
}
if Path::new(&paths.socket_path()).exists() {
std::fs::remove_file(paths.socket_path())?;
}
@ -211,12 +225,23 @@ pub async fn send_control(
Ok(geth_control::decode_response(&line)?)
}
pub async fn handle_request_async(
node: &LocalNode,
request: ControlRequest,
) -> Result<ControlResponse, NodeError> {
match request {
ControlRequest::PeerCardExport { out } => export_peer_card(node, out, true).await,
ControlRequest::PeerPing { node: peer_node } => peer_ping(node, &peer_node).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(&node, request) {
let response = match handle_request_async(&node, request).await {
Ok(response) => response,
Err(error) => ControlResponse::Error {
message: error.to_string(),
@ -229,6 +254,243 @@ async fn handle_stream(node: LocalNode, stream: UnixStream) -> Result<(), NodeEr
Ok(())
}
async fn export_peer_card(
node: &LocalNode,
out: Option<std::path::PathBuf>,
include_node_addr: bool,
) -> Result<ControlResponse, NodeError> {
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<PeerCard, NodeError> {
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<ControlResponse, 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)?;
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)),
}
}
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");
}
});
}
});
}
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()?;
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(),
}
}
};
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<iroh::NodeAddr, NodeError> {
let node_id = candidate
.endpoint_id
.parse::<iroh::NodeId>()
.map_err(|error| NodeError::IrohPeer(error.to_string()))?;
let direct_addresses = candidate
.direct_addresses
.iter()
.map(|addr| {
addr.parse::<std::net::SocketAddr>()
.map_err(|error| NodeError::IrohPeer(error.to_string()))
})
.collect::<Result<Vec<_>, _>>()?;
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::<iroh::RelayUrl>()
.map_err(|error| NodeError::IrohPeer(error.to_string()))?,
);
}
Ok(node_addr)
}
fn display_alpn(alpn: Vec<u8>) -> String {
String::from_utf8_lossy(&alpn).into_owned()
}
pub fn handle_request(
node: &LocalNode,
request: ControlRequest,
@ -264,6 +526,7 @@ pub fn handle_request(
vec![EndpointCandidate {
endpoint_id,
relay_url: None,
direct_addresses: Vec::new(),
source: DiscoverySource::Manual,
}],
UnixMillis(geth_store::now_ms()),
@ -307,6 +570,7 @@ pub fn handle_request(
note: discovery_is_untrusted_note().to_owned(),
})
}
ControlRequest::PeerPing { .. } => Err(NodeError::IrohEndpointUnavailable),
ControlRequest::ResourceList => Ok(ControlResponse::ResourceList {
resources: store
.list_resources()?
@ -1564,6 +1828,10 @@ async fn start_daemon_iroh_endpoint(
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) => {
@ -1695,3 +1963,90 @@ fn ssh_revocation_from_stored(
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");
}
#[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 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:?}"),
}
left_endpoint.shutdown().await;
right_endpoint.shutdown().await;
}
}