diff --git a/AGENTS.md b/AGENTS.md index b0287da..ffbb06b 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -150,9 +150,10 @@ Roadmap items should be actionable and checkable: bounded in-memory ring buffer. Remote publish over the protected Iroh control ALPN requires `pubsub.publish` on `resource:pubsub:`. Iroh-gossip replication, private topics, and remote subscribe are still roadmap work. -- Pipe listen/connect supports a local daemon-lifetime registry only. Iroh byte - streams, TCP/Unix forwarding, and pipe capability enforcement are still - roadmap work. +- Pipe listen/connect supports a daemon-lifetime registry. Remote + `geth pipe connect --node ` uses the protected Iroh control + ALPN and requires `pipe.connect` on `resource:pipe:`. Iroh byte streams, + TCP/Unix forwarding, and remote listener creation are still roadmap work. - Resource secret epoch metadata can be created, rotated, and listed locally. Bearer access metadata can be created/listed/revoked as resource-scoped auth ops and must not allow trust graph mutation capabilities. Payload encryption, diff --git a/README.md b/README.md index 7fca023..3fdffbd 100644 --- a/README.md +++ b/README.md @@ -120,7 +120,8 @@ The bootstrap implementation provides: - `geth ssh revocation export --out [--format jsonl|openssh-krl-spec|openssh-krl]` - `geth ssh revocation import [--format jsonl|openssh-krl-spec]` - `geth ssh revocation sync ` -- local pipe registry commands: `geth pipe listen/connect` +- pipe registry/connect commands: `geth pipe listen ` and + `geth pipe connect [--node ]` `geth peer export/import/list` is for untrusted peer-card exchange. Peer cards include the Iroh EndpointID plus currently known relay/direct addresses. @@ -151,6 +152,10 @@ Remote pubsub publish uses the protected Iroh control path too. The remote peer requires `pubsub.publish` on `resource:pubsub:` before recording the message in its local daemon-lifetime ring buffer. Pubsub remains lossy and is not durable storage. +Remote pipe connect uses the same protected Iroh control path and requires +`pipe.connect` on `resource:pipe:`. The current prototype records a remote +connection attempt and whether a listener exists; byte streaming and forwarding +are still future work. Importing or pinging a peer card never grants capabilities by itself. When `[iroh].local_discovery = true`, the daemon also advertises and discovers signed peer cards on LAN using a geth-specific mDNS TXT payload. That payload is diff --git a/crates/geth-cli/src/lib.rs b/crates/geth-cli/src/lib.rs index 398882c..f5695d4 100644 --- a/crates/geth-cli/src/lib.rs +++ b/crates/geth-cli/src/lib.rs @@ -325,8 +325,14 @@ pub enum PubsubCommand { #[derive(Debug, Subcommand)] pub enum PipeCommand { - Listen { name: String }, - Connect { target: String }, + Listen { + name: String, + }, + Connect { + target: String, + #[arg(long)] + node: Option, + }, } #[derive(Debug, Subcommand)] @@ -644,7 +650,7 @@ fn request_for_command(command: Command) -> Result { }, Command::Pipe { command } => match command { PipeCommand::Listen { name } => ControlRequest::PipeListen { name }, - PipeCommand::Connect { target } => ControlRequest::PipeConnect { target }, + PipeCommand::Connect { target, node } => ControlRequest::PipeConnect { target, node }, }, Command::Db { command } => match command { DbCommand::Add { name, path } => ControlRequest::DbAdd { name, path }, @@ -1422,6 +1428,29 @@ fn print_response(response: ControlResponse, json: bool) -> Result<()> { println!("connected_at_ms: {}", connection.connected_at.0); println!("note: {}", connection.note); } + ControlResponse::PipeRemoteConnected { + peer_node_id, + peer_agent_id, + endpoint_id, + connection, + allowed, + reason, + note, + } => { + if allowed { + println!("pipe target: {}", connection.target); + println!("peer: {peer_node_id}"); + println!("remote_listener_found: {}", connection.local_listener_found); + println!("connected_at_ms: {}", connection.connected_at.0); + } else { + println!("pipe connect denied by {peer_node_id}"); + } + println!("agent: {peer_agent_id}"); + println!("endpoint: {endpoint_id}"); + println!("allowed: {allowed}"); + println!("reason: {reason}"); + println!("note: {note}"); + } ControlResponse::NotImplemented { module, command } => { println!("{module} {command}: not implemented yet"); } diff --git a/crates/geth-control/src/lib.rs b/crates/geth-control/src/lib.rs index 75b87ce..5c1f7a7 100644 --- a/crates/geth-control/src/lib.rs +++ b/crates/geth-control/src/lib.rs @@ -226,6 +226,7 @@ pub enum ControlRequest { }, PipeConnect { target: String, + node: Option, }, ModuleStub { module: String, @@ -472,6 +473,15 @@ pub enum ControlResponse { PipeConnected { connection: PipeConnection, }, + PipeRemoteConnected { + peer_node_id: String, + peer_agent_id: String, + endpoint_id: String, + connection: PipeConnection, + allowed: bool, + reason: String, + note: String, + }, NotImplemented { module: String, command: String, @@ -557,6 +567,11 @@ pub enum PeerControlRequest { message: String, nonce: String, }, + PipeConnect { + peer_card: PeerCard, + target: String, + nonce: String, + }, } #[derive(Clone, Debug, PartialEq, Serialize, Deserialize)] @@ -651,6 +666,18 @@ pub enum PeerControlResponse { nonce: String, note: String, }, + PipeConnected { + node_id: String, + agent_id: String, + endpoint_id: String, + remote_endpoint_id: String, + connection: Option, + allowed: bool, + reason: String, + evaluated_ops: usize, + nonce: String, + note: String, + }, Error { message: String, }, @@ -850,6 +877,15 @@ mod tests { request ); + let request = ControlRequest::PipeConnect { + target: "inbox".to_owned(), + node: Some("node:peer".to_owned()), + }; + assert_eq!( + decode_request(&encode_request(&request).expect("encode")).expect("decode"), + request + ); + let request = ControlRequest::PubsubPub { topic: "presence/test".to_owned(), message: "online".to_owned(), @@ -878,6 +914,25 @@ mod tests { response ); + let response = ControlResponse::PipeRemoteConnected { + peer_node_id: "node:peer".to_owned(), + peer_agent_id: "agent:peer".to_owned(), + endpoint_id: "endpoint:peer".to_owned(), + connection: PipeConnection { + target: "inbox".to_owned(), + connected_at: geth_types::UnixMillis(1), + local_listener_found: true, + note: "remote".to_owned(), + }, + allowed: true, + reason: "direct grant".to_owned(), + note: "pipe connect".to_owned(), + }; + assert_eq!( + decode_response(&encode_response(&response).expect("encode")).expect("decode"), + response + ); + let request = ControlRequest::DbChanges { name: "notes".to_owned(), after_db_version: Some(7), @@ -1068,5 +1123,28 @@ mod tests { .expect("decode"), response ); + + let response = PeerControlResponse::PipeConnected { + node_id: "node:peer".to_owned(), + agent_id: "agent:peer".to_owned(), + endpoint_id: "endpoint:peer".to_owned(), + remote_endpoint_id: "endpoint:caller".to_owned(), + connection: Some(PipeConnection { + target: "inbox".to_owned(), + connected_at: geth_types::UnixMillis(1), + local_listener_found: true, + note: "remote".to_owned(), + }), + allowed: true, + reason: "direct grant".to_owned(), + evaluated_ops: 1, + nonce: "nonce".to_owned(), + note: "pipe connect".to_owned(), + }; + assert_eq!( + decode_peer_response(&encode_peer_response(&response).expect("encode")) + .expect("decode"), + response + ); } } diff --git a/crates/geth-node/src/lib.rs b/crates/geth-node/src/lib.rs index e2bfea7..9edeba4 100644 --- a/crates/geth-node/src/lib.rs +++ b/crates/geth-node/src/lib.rs @@ -268,6 +268,10 @@ pub async fn handle_request_async( message, node: Some(peer_node), } => pubsub_publish_to_peer(node, &peer_node, topic, message).await, + ControlRequest::PipeConnect { + target, + node: Some(peer_node), + } => pipe_connect_to_peer(node, &peer_node, target).await, other => handle_request(node, other), } } @@ -545,7 +549,8 @@ async fn peer_ping(node: &LocalNode, peer_node: &str) -> Result Err(NodeError::IrohPeer( + | PeerControlResponse::PubsubPublished { .. } + | PeerControlResponse::PipeConnected { .. } => Err(NodeError::IrohPeer( "peer returned wrong response type to ping request".to_owned(), )), } @@ -654,7 +659,8 @@ async fn peer_auth_check( | PeerControlResponse::SshCertSynced { .. } | PeerControlResponse::SshRevocationSynced { .. } | PeerControlResponse::KvSynced { .. } - | PeerControlResponse::PubsubPublished { .. } => Err(NodeError::IrohPeer( + | PeerControlResponse::PubsubPublished { .. } + | PeerControlResponse::PipeConnected { .. } => Err(NodeError::IrohPeer( "peer returned wrong response type to auth-check request".to_owned(), )), } @@ -792,7 +798,8 @@ async fn cas_fetch_from_peer( | PeerControlResponse::SshCertSynced { .. } | PeerControlResponse::SshRevocationSynced { .. } | PeerControlResponse::KvSynced { .. } - | PeerControlResponse::PubsubPublished { .. } => Err(NodeError::IrohPeer( + | PeerControlResponse::PubsubPublished { .. } + | PeerControlResponse::PipeConnected { .. } => Err(NodeError::IrohPeer( "peer returned wrong response type to CAS fetch".to_owned(), )), } @@ -1068,6 +1075,54 @@ async fn pubsub_publish_to_peer( } } +async fn pipe_connect_to_peer( + node: &LocalNode, + peer_node: &str, + target: String, +) -> Result { + geth_pipe::validate_pipe_name(&target)?; + let response = request_peer_control(node, peer_node, "pipe-connect", |peer_card, nonce| { + PeerControlRequest::PipeConnect { + peer_card, + target: target.clone(), + nonce, + } + }) + .await?; + match response { + PeerControlResponse::PipeConnected { + node_id, + agent_id, + endpoint_id, + connection, + allowed, + reason, + note, + .. + } => { + let connection = connection.unwrap_or_else(|| PipeConnection { + target, + connected_at: UnixMillis(0), + local_listener_found: false, + note: geth_pipe::pipe_roadmap().to_owned(), + }); + Ok(ControlResponse::PipeRemoteConnected { + peer_node_id: node_id, + peer_agent_id: agent_id, + endpoint_id, + connection, + allowed, + reason, + note, + }) + } + PeerControlResponse::Error { message } => Err(NodeError::IrohPeer(message)), + _ => Err(NodeError::IrohPeer( + "peer returned wrong response type to pipe connect".to_owned(), + )), + } +} + fn live_sync_cursor_key(peer_node: &str, stream: &str) -> String { format!("live-sync:{peer_node}:{stream}") } @@ -1168,6 +1223,10 @@ async fn request_peer_control( | PeerControlResponse::PubsubPublished { nonce: response_nonce, .. + } + | PeerControlResponse::PipeConnected { + nonce: response_nonce, + .. } if response_nonce == &nonce => Ok(response), PeerControlResponse::Error { .. } => Ok(response), _ => Err(NodeError::IrohPeer(format!( @@ -1606,6 +1665,55 @@ async fn handle_iroh_control_connection( note: "pubsub publish authenticated endpoint/card binding and required pubsub.publish on the remote topic resource; pubsub is lossy".to_owned(), } } + PeerControlRequest::PipeConnect { + peer_card, + target, + nonce, + } => { + 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 resource = format!("resource:pipe:{target}"); + let capability = "pipe.connect".to_owned(); + let explanation = geth_auth::explain_auth_ops( + &load_auth_ops_for_resource(&store, &resource)?, + PrincipalId::new(peer_card.node_id.to_string()), + ResourceId::new(resource), + Capability::new(capability), + ); + let connection = if explanation.allowed { + Some(record_pipe_connection( + &node, + target, + "remote pipe connect over protected Iroh control path; byte streams are not implemented yet".to_owned(), + )?) + } else { + None + }; + PeerControlResponse::PipeConnected { + 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, + connection, + allowed: explanation.allowed, + reason: explanation.reason, + evaluated_ops: explanation.evaluated_ops, + nonce, + note: "pipe connect authenticated endpoint/card binding and required pipe.connect on the remote pipe resource; byte streams are not implemented yet".to_owned(), + } + } }; send.write_all(geth_control::encode_peer_response(&response)?.as_bytes()) .await @@ -2647,26 +2755,18 @@ pub fn handle_request( runtime.listeners.insert(name, listener.clone()); Ok(ControlResponse::PipeListening { listener }) } - ControlRequest::PipeConnect { target } => { + ControlRequest::PipeConnect { target, node: None } => { geth_pipe::validate_pipe_name(&target)?; - let mut runtime = node - .runtime - .pipes - .lock() - .map_err(|_| NodeError::RuntimeLockPoisoned)?; - let local_listener_found = runtime.listeners.contains_key(&target); - let connection = PipeConnection { + let connection = record_pipe_connection( + node, target, - connected_at: UnixMillis(geth_store::now_ms()), - local_listener_found, - note: geth_pipe::local_pipe_runtime_note().to_owned(), - }; - runtime.connections.push_back(connection.clone()); - while runtime.connections.len() > PIPE_CONNECTION_RING_LIMIT { - runtime.connections.pop_front(); - } + geth_pipe::local_pipe_runtime_note().to_owned(), + )?; Ok(ControlResponse::PipeConnected { connection }) } + ControlRequest::PipeConnect { node: Some(_), .. } => { + Err(NodeError::IrohEndpointUnavailable) + } ControlRequest::ModuleStub { module, command } => { Ok(ControlResponse::NotImplemented { module, command }) } @@ -2782,6 +2882,30 @@ fn record_pubsub_message( Ok(message) } +fn record_pipe_connection( + node: &LocalNode, + target: String, + note: String, +) -> Result { + let mut runtime = node + .runtime + .pipes + .lock() + .map_err(|_| NodeError::RuntimeLockPoisoned)?; + let local_listener_found = runtime.listeners.contains_key(&target); + let connection = PipeConnection { + target, + connected_at: UnixMillis(geth_store::now_ms()), + local_listener_found, + note, + }; + runtime.connections.push_back(connection.clone()); + while runtime.connections.len() > PIPE_CONNECTION_RING_LIMIT { + runtime.connections.pop_front(); + } + Ok(connection) +} + fn document_resource_from_stored(stored: &StoredDocumentResource) -> DocumentResource { DocumentResource { id: stored.document_id.clone().into(), @@ -3302,6 +3426,13 @@ mod tests { }, ) .expect("right kv set"); + handle_request( + &right, + ControlRequest::PipeListen { + name: "inbox".to_owned(), + }, + ) + .expect("right pipe listen"); let ping = handle_request_async( &left, @@ -3427,6 +3558,29 @@ mod tests { other => panic!("unexpected denied pubsub publish response: {other:?}"), } + let denied_pipe = handle_request_async( + &left, + ControlRequest::PipeConnect { + target: "inbox".to_owned(), + node: Some(right_card.node_id.to_string()), + }, + ) + .await + .expect("denied remote pipe connect"); + match denied_pipe { + ControlResponse::PipeRemoteConnected { + allowed, + reason, + connection, + .. + } => { + assert!(!allowed); + assert!(!connection.local_listener_found); + assert!(reason.contains("no active direct or group grant")); + } + other => panic!("unexpected denied pipe connect response: {other:?}"), + } + handle_request( &right, ControlRequest::AuthGrant { @@ -3457,6 +3611,16 @@ mod tests { }, ) .expect("grant left pubsub publish"); + handle_request( + &right, + ControlRequest::AuthGrant { + subject: left.node_id.clone(), + resource: "resource:pipe:inbox".to_owned(), + capability: "pipe.connect".to_owned(), + grant_id: Some("grant:left-pipe-connect".to_owned()), + }, + ) + .expect("grant left pipe connect"); let allowed = handle_request_async( &left, @@ -3595,6 +3759,31 @@ mod tests { other => panic!("unexpected remote pubsub messages response: {other:?}"), } + let remote_pipe = handle_request_async( + &left, + ControlRequest::PipeConnect { + target: "inbox".to_owned(), + node: Some(right_card.node_id.to_string()), + }, + ) + .await + .expect("allowed remote pipe connect"); + match remote_pipe { + ControlResponse::PipeRemoteConnected { + allowed, + connection, + reason, + note, + .. + } => { + assert!(allowed); + assert!(connection.local_listener_found); + assert!(reason.contains("direct grant")); + assert!(note.contains("byte streams are not implemented yet")); + } + other => panic!("unexpected allowed pipe connect response: {other:?}"), + } + let denied_cert_sync = handle_request_async( &left, ControlRequest::SshCertSync { diff --git a/crates/geth/tests/bootstrap.rs b/crates/geth/tests/bootstrap.rs index 65896bd..f093707 100644 --- a/crates/geth/tests/bootstrap.rs +++ b/crates/geth/tests/bootstrap.rs @@ -1312,6 +1312,7 @@ fn pipe_listen_connect_uses_local_runtime_registry() { &node, geth_control::ControlRequest::PipeConnect { target: "inbox".to_owned(), + node: None, }, ) .expect("connect pipe"); @@ -1333,6 +1334,7 @@ fn pipe_listen_connect_uses_local_runtime_registry() { &reopened, geth_control::ControlRequest::PipeConnect { target: "inbox".to_owned(), + node: None, }, ) .expect("connect after reopen"); diff --git a/docs/architecture.md b/docs/architecture.md index 10008cd..0b9afd9 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -165,10 +165,13 @@ control ALPN. The remote daemon validates endpoint/card binding and requires its local ring buffer. Iroh-gossip replication and private topics are future work. -`geth-pipe` currently supports `pipe listen/connect` against a local -daemon-lifetime registry. This is a control-plane scaffold for names and -connection attempts only; it does not carry bytes, forward sockets, or use Iroh -streams yet. +`geth-pipe` currently supports `pipe listen/connect` against a daemon-lifetime +registry. `geth pipe connect --node ` sends an authorized remote +connect request over the protected Iroh control ALPN. The remote daemon validates +endpoint/card binding and requires `pipe.connect` on `resource:pipe:` +before recording the connection attempt and reporting whether a listener exists. +This is still a control-plane scaffold for names and connection attempts only; +it does not carry bytes or forward sockets yet. `geth-ssh-proxy` currently defines types, command shape, and roadmap stubs. diff --git a/docs/roadmap.md b/docs/roadmap.md index 0d8196f..f406a99 100644 --- a/docs/roadmap.md +++ b/docs/roadmap.md @@ -296,8 +296,13 @@ Goal: add authorized stream-oriented management workflows over Iroh. registry. - `[x]` Tests cover local listener registration, local connect matching, and daemon restart behavior. - - `[ ]` `geth pipe listen/connect` can connect two local test nodes over Iroh. - - `[ ]` Pipe access requires `pipe.listen` or `pipe.connect`. + - `[x]` `geth pipe connect --node ` can connect to an + imported peer over Iroh at the control-plane level. + - `[x]` Remote pipe connect requires `pipe.connect` on + `resource:pipe:`. + - `[ ]` Remote pipe listen requires `pipe.listen` when remote listener + creation is added. + - `[ ]` Pipe connect carries bidirectional byte streams over Iroh. - `[ ]` Streams close cleanly and propagate errors. - `[ ]` TCP forwarding.