refactor: centralize peer control caller auth

This commit is contained in:
Eric Wendland 2026-07-05 23:38:19 +02:00
commit 1a6d9db913
3 changed files with 153 additions and 384 deletions

View file

@ -222,6 +222,12 @@ struct PipeUnixConnectWire {
bearer_proof: Option<BearerProof>,
}
struct PeerControlCaller {
peer_card: PeerCard,
remote_endpoint_id: String,
store: Store,
}
#[derive(Clone, Debug)]
pub struct InitOwnerOptions {
pub admin_key_path: Option<PathBuf>,
@ -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::<Result<Vec<_>, _>>()?,
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<PeerControlCaller, NodeError> {
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,