diff --git a/AGENTS.md b/AGENTS.md index 3d4bf0a..1cdd8e2 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -179,8 +179,12 @@ Roadmap items should be actionable and checkable: resource before creating a daemon-lifetime listener on the peer. `geth pipe send [message|--in |--in -] --node ` carries a byte message over the dedicated `/geth/pipe/1` ALPN when `pipe.connect` is authorized, and - `geth pipe recv ` drains local daemon-lifetime messages. Long-lived - bidirectional streams and TCP/Unix forwarding are still roadmap work. + `geth pipe recv ` drains local daemon-lifetime messages. `geth pipe + forward-tcp --listen 127.0.0.1: --node --target + 127.0.0.1:` opens authorized bidirectional byte streams over + `/geth/pipe/1`; the remote daemon requires `pipe.forward` on + `resource:pipe-tcp:` before connecting to the loopback target. Unix + socket forwarding is still roadmap work. - `geth ssh proxy ` is a streaming OpenSSH ProxyCommand-style path. The CLI streams through the local daemon, the daemon uses `/geth/ssh-proxy/1` over Iroh, the remote daemon requires `ssh_proxy.connect` on diff --git a/README.md b/README.md index 9d5af52..f425bc6 100644 --- a/README.md +++ b/README.md @@ -146,7 +146,8 @@ The bootstrap implementation provides: `geth pipe listen [--node ] [--bearer-secret ]`, `geth pipe connect [--node ] [--bearer-secret ]`, `geth pipe send [message|--in |--in -] [--node ] [--bearer-secret ]`, - and `geth pipe recv [--peek]` + `geth pipe recv [--peek]`, and + `geth pipe forward-tcp --listen 127.0.0.1: --node --target 127.0.0.1:` `geth peer export/import/list` is for untrusted peer-card exchange. Peer cards include the Iroh EndpointID plus currently known relay/direct addresses. @@ -207,8 +208,14 @@ recv ` drains local daemon-lifetime messages. Remote pipe listen uses the same protected path: `geth pipe listen --node ` requires `pipe.listen` on `resource:pipe:` before registering a daemon-lifetime listener on the -peer. Long-lived stdin/stdout streaming and socket forwarding are still future -work. +peer. +`geth pipe forward-tcp --listen 127.0.0.1: --node --target +127.0.0.1:` starts a local loopback TCP listener. Each accepted +connection asks the local daemon to open an authorized `/geth/pipe/1` byte +stream to the peer. The remote daemon validates the signed endpoint/card binding +and requires `pipe.forward` on `resource:pipe-tcp:` before connecting to +the remote loopback TCP target. This is loopback-only in the prototype to avoid +turning geth into an accidental open proxy. `geth ssh proxy ` is usable as an OpenSSH `ProxyCommand`: the CLI opens a local daemon stream, the daemon opens the dedicated `/geth/ssh-proxy/1` Iroh ALPN, the remote daemon validates the caller's endpoint/card binding and diff --git a/crates/geth-cli/src/lib.rs b/crates/geth-cli/src/lib.rs index ca9e167..c778250 100644 --- a/crates/geth-cli/src/lib.rs +++ b/crates/geth-cli/src/lib.rs @@ -397,6 +397,16 @@ pub enum PipeCommand { #[arg(long)] bearer_secret: Option, }, + ForwardTcp { + #[arg(long)] + listen: String, + #[arg(long)] + node: String, + #[arg(long)] + target: String, + #[arg(long)] + bearer_secret: Option, + }, Send { target: String, message: Option, @@ -611,6 +621,20 @@ pub async fn run() -> Result<()> { .await .context("stream SSH proxy through geth daemon")?; } + Command::Pipe { + command: + PipeCommand::ForwardTcp { + listen, + node, + target, + bearer_secret, + }, + } if !cli.json && !cli.jsonl => { + println!("forwarding tcp {listen} -> {node}:{target}"); + geth_node::run_tcp_forward(&paths, listen, node, target, bearer_secret) + .await + .context("run TCP forward through geth daemon")?; + } command => { let request = request_for_command(command)?; let response = geth_node::send_control(&paths, request) @@ -891,6 +915,17 @@ fn request_for_command(command: Command) -> Result { node, bearer_secret, }, + PipeCommand::ForwardTcp { + listen, + node, + target, + bearer_secret, + } => ControlRequest::PipeTcpForward { + listen_addr: listen, + node, + target_addr: target, + bearer_secret, + }, PipeCommand::Send { target, message, diff --git a/crates/geth-control/src/lib.rs b/crates/geth-control/src/lib.rs index 1cc3ea0..5d59fa6 100644 --- a/crates/geth-control/src/lib.rs +++ b/crates/geth-control/src/lib.rs @@ -302,6 +302,17 @@ pub enum ControlRequest { node: Option, bearer_secret: Option, }, + PipeTcpForward { + node: String, + listen_addr: String, + target_addr: String, + bearer_secret: Option, + }, + PipeTcpStream { + node: String, + target_addr: String, + bearer_secret: Option, + }, PipeSend { target: String, data_base64: String, @@ -1033,6 +1044,12 @@ pub enum PipeWireRequest { nonce: String, bearer_proof: Option, }, + TcpConnect { + peer_card: PeerCard, + target_addr: String, + nonce: String, + bearer_proof: Option, + }, } #[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] @@ -1051,6 +1068,18 @@ pub enum PipeWireResponse { nonce: String, note: String, }, + Connected { + 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, }, @@ -2037,6 +2066,58 @@ mod tests { response ); + let request = PipeWireRequest::TcpConnect { + peer_card: PeerCard { + node_id: "node:caller".into(), + agent_id: "agent:caller".into(), + endpoints: Vec::new(), + issued_at: geth_types::UnixMillis(1), + signature: geth_discovery::SignatureMetadata { + namespace: "geth.peer-card.v1@geth.local".to_owned(), + signer: "agent:caller".to_owned(), + public_key: "key".to_owned(), + signature: "sig".to_owned(), + }, + }, + target_addr: "127.0.0.1:22".to_owned(), + nonce: "nonce".to_owned(), + bearer_proof: Some(BearerProof { + secret: "bearer:test".into(), + resource: "resource:pipe-tcp:127.0.0.1:22".into(), + capabilities: vec!["pipe.forward".into()], + nonce: "nonce".to_owned(), + response: "response".to_owned(), + }), + }; + assert_eq!( + decode_pipe_wire_request(&encode_pipe_wire_request(&request).expect("encode")) + .expect("decode"), + request + ); + + let response = PipeWireResponse::Connected { + 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: "127.0.0.1:22".to_owned(), + connected_at: geth_types::UnixMillis(1), + local_listener_found: true, + note: "tcp".to_owned(), + }), + allowed: true, + reason: "direct grant".to_owned(), + evaluated_ops: 1, + nonce: "nonce".to_owned(), + note: "pipe tcp".to_owned(), + }; + assert_eq!( + decode_pipe_wire_response(&encode_pipe_wire_response(&response).expect("encode")) + .expect("decode"), + response + ); + let response = PeerControlResponse::SshProxyConnected { node_id: "node:peer".to_owned(), agent_id: "agent:peer".to_owned(), diff --git a/crates/geth-node/src/lib.rs b/crates/geth-node/src/lib.rs index 4fd5402..6214b8f 100644 --- a/crates/geth-node/src/lib.rs +++ b/crates/geth-node/src/lib.rs @@ -48,7 +48,7 @@ use std::path::{Component, Path, PathBuf}; use std::sync::{Arc, Mutex}; use std::time::Duration; use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; -use tokio::net::{TcpStream, UnixListener, UnixStream}; +use tokio::net::{TcpListener, TcpStream, UnixListener, UnixStream}; #[derive(Debug, thiserror::Error)] pub enum NodeError { @@ -163,6 +163,13 @@ struct LiveSyncCursor { cursor_ms: i64, } +struct PipeTcpConnectWire { + peer_card: PeerCard, + target_addr: String, + nonce: String, + bearer_proof: Option, +} + pub fn init_node(paths: &GethPaths) -> Result { paths.ensure_base_dirs()?; if !paths.config_file().exists() { @@ -295,6 +302,87 @@ pub async fn stream_ssh_proxy( } } +pub async fn run_tcp_forward( + paths: &GethPaths, + listen_addr: String, + peer_node: String, + target_addr: String, + bearer_secret: Option, +) -> Result<(), NodeError> { + let listen_addr = geth_pipe::validate_tcp_forward_listen_addr(&listen_addr)?; + geth_pipe::validate_tcp_forward_target_addr(&target_addr)?; + let listener = TcpListener::bind(listen_addr).await?; + tracing::info!( + listen = %listen_addr, + peer = %peer_node, + target = %target_addr, + "tcp forward listening" + ); + loop { + let (client, client_addr) = listener.accept().await?; + let paths = paths.clone(); + let peer_node = peer_node.clone(); + let target_addr = target_addr.clone(); + let bearer_secret = bearer_secret.clone(); + tokio::spawn(async move { + if let Err(error) = + stream_pipe_tcp(&paths, peer_node, target_addr, bearer_secret, client).await + { + tracing::warn!(%error, %client_addr, "tcp forward connection failed"); + } + }); + } +} + +async fn stream_pipe_tcp( + paths: &GethPaths, + peer_node: String, + target_addr: String, + bearer_secret: Option, + client: TcpStream, +) -> Result<(), NodeError> { + let mut stream = UnixStream::connect(paths.socket_path()).await?; + stream + .write_all( + geth_control::encode_request(&ControlRequest::PipeTcpStream { + node: peer_node, + target_addr, + bearer_secret, + })? + .as_bytes(), + ) + .await?; + let mut reader = BufReader::new(stream); + let mut line = String::new(); + reader.read_line(&mut line).await?; + match geth_control::decode_response(&line)? { + ControlResponse::PipeRemoteConnected { allowed: true, .. } => { + let daemon_stream = reader.into_inner(); + let (mut daemon_read, mut daemon_write) = daemon_stream.into_split(); + let (mut client_read, mut client_write) = client.into_split(); + let upload = async { + tokio::io::copy(&mut client_read, &mut daemon_write).await?; + daemon_write.shutdown().await + }; + let download = async { + tokio::io::copy(&mut daemon_read, &mut client_write).await?; + client_write.shutdown().await + }; + tokio::try_join!(upload, download)?; + Ok(()) + } + ControlResponse::PipeRemoteConnected { + allowed: false, + reason, + .. + } => Err(NodeError::Unauthorized(reason)), + ControlResponse::Error { message } => Err(NodeError::IrohPeer(message)), + other => Err(NodeError::IrohPeer(format!( + "daemon returned unexpected pipe TCP stream response: {other:?}" + ))), + } +} + pub async fn handle_request_async( node: &LocalNode, request: ControlRequest, @@ -363,6 +451,12 @@ pub async fn handle_request_async( node: Some(peer_node), bearer_secret, } => pipe_send_to_peer(node, &peer_node, target, data_base64, bearer_secret).await, + ControlRequest::PipeTcpForward { .. } | ControlRequest::PipeTcpStream { .. } => { + Err(NodeError::IrohPeer( + "TCP forwarding is a streaming control command; run without --json/--jsonl" + .to_owned(), + )) + } ControlRequest::SshProxyConnect { node: peer_node, bearer_secret, @@ -389,6 +483,16 @@ async fn handle_stream(node: LocalNode, stream: UnixStream) -> Result<(), NodeEr let stream = reader.into_inner(); return handle_local_ssh_proxy_stream(node, &peer_node, bearer_secret, stream).await; } + if let ControlRequest::PipeTcpStream { + node: peer_node, + target_addr, + bearer_secret, + } = request + { + let stream = reader.into_inner(); + return handle_local_pipe_tcp_stream(node, &peer_node, target_addr, bearer_secret, stream) + .await; + } let response = match handle_request_async(&node, request).await { Ok(response) => response, Err(error) => ControlResponse::Error { @@ -1584,9 +1688,134 @@ async fn pipe_send_to_peer( note, }), PipeWireResponse::Error { message } => Err(NodeError::IrohPeer(message)), + PipeWireResponse::Connected { .. } => Err(NodeError::IrohPeer( + "peer returned TCP stream response to pipe send".to_owned(), + )), } } +async fn handle_local_pipe_tcp_stream( + node: LocalNode, + peer_node: &str, + target_addr: String, + bearer_secret: Option, + local_stream: UnixStream, +) -> Result<(), NodeError> { + geth_pipe::validate_tcp_forward_target_addr(&target_addr)?; + 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{}\0pipe-tcp-stream\0{}\0{}", + node.node_id, + peer_node, + target_addr, + geth_store::now_ms() + ) + .as_bytes(), + ); + let resource = format!("resource:pipe-tcp:{target_addr}"); + let request = PipeWireRequest::TcpConnect { + peer_card: self_card, + target_addr: target_addr.clone(), + nonce: nonce.clone(), + bearer_proof: bearer_proof(bearer_secret, &resource, "pipe.forward", &nonce), + }; + let conn = endpoint + .endpoint() + .connect(node_addr, geth_iroh::ALPN_PIPE) + .await + .map_err(|error| NodeError::IrohPeer(error.to_string()))?; + let (mut remote_send, remote_recv) = conn + .open_bi() + .await + .map_err(|error| NodeError::IrohPeer(error.to_string()))?; + remote_send + .write_all(geth_control::encode_pipe_wire_request(&request)?.as_bytes()) + .await + .map_err(|error| NodeError::IrohPeer(error.to_string()))?; + let mut remote_reader = BufReader::new(remote_recv); + let mut line = String::new(); + remote_reader + .read_line(&mut line) + .await + .map_err(|error| NodeError::IrohPeer(error.to_string()))?; + let response = geth_control::decode_pipe_wire_response(&line)?; + let local_response = match response { + PipeWireResponse::Connected { + node_id, + agent_id, + endpoint_id, + connection, + allowed, + reason, + note, + nonce: response_nonce, + .. + } if response_nonce == nonce => ControlResponse::PipeRemoteConnected { + peer_node_id: node_id, + peer_agent_id: agent_id, + endpoint_id, + connection: connection.unwrap_or_else(|| PipeConnection { + target: target_addr.clone(), + connected_at: UnixMillis(0), + local_listener_found: false, + note: "pipe TCP forward was denied before connecting".to_owned(), + }), + allowed, + reason, + note, + }, + PipeWireResponse::Error { message } => ControlResponse::Error { message }, + _ => ControlResponse::Error { + message: "peer returned wrong response type to pipe TCP stream".to_owned(), + }, + }; + let mut local_stream = local_stream; + local_stream + .write_all(geth_control::encode_response(&local_response)?.as_bytes()) + .await?; + let ControlResponse::PipeRemoteConnected { allowed: true, .. } = local_response else { + return Ok(()); + }; + let remote_recv = remote_reader.into_inner(); + let (mut local_read, mut local_write) = local_stream.into_split(); + let upload = async { + tokio::io::copy(&mut local_read, &mut remote_send) + .await + .map_err(|error| NodeError::IrohPeer(error.to_string()))?; + remote_send + .finish() + .map_err(|error| NodeError::IrohPeer(error.to_string())) + }; + let mut remote_recv = remote_recv; + let download = async { + tokio::io::copy(&mut remote_recv, &mut local_write) + .await + .map_err(|error| NodeError::IrohPeer(error.to_string()))?; + local_write.shutdown().await.map_err(NodeError::from) + }; + tokio::try_join!(upload, download)?; + Ok(()) +} + async fn ssh_proxy_connect_to_peer( node: &LocalNode, peer_node: &str, @@ -2577,21 +2806,15 @@ async fn handle_iroh_control_connection( if alpn == display_alpn(geth_iroh::ALPN_SSH_PROXY.to_vec()) { return handle_ssh_proxy_wire_connection(node, remote_endpoint_id, send, recv).await; } + if alpn == display_alpn(geth_iroh::ALPN_PIPE.to_vec()) { + return handle_pipe_wire_connection(node, remote_endpoint_id, send, recv).await; + } 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()))?; - if alpn == display_alpn(geth_iroh::ALPN_PIPE.to_vec()) { - let response = handle_pipe_wire_request(&node, &remote_endpoint_id, text)?; - send.write_all(geth_control::encode_pipe_wire_response(&response)?.as_bytes()) - .await - .map_err(|error| NodeError::IrohPeer(error.to_string()))?; - send.finish() - .map_err(|error| NodeError::IrohPeer(error.to_string()))?; - return Ok(()); - } let response = match geth_control::decode_peer_request(text)? { PeerControlRequest::Ping { peer_card, nonce } => { peer_card.validate_candidate()?; @@ -3526,12 +3749,19 @@ async fn handle_ssh_proxy_wire_connection( Ok(()) } -fn handle_pipe_wire_request( - node: &LocalNode, - remote_endpoint_id: &str, - text: &str, -) -> Result { - match geth_control::decode_pipe_wire_request(text)? { +async fn handle_pipe_wire_connection( + node: LocalNode, + remote_endpoint_id: String, + mut send: iroh::endpoint::SendStream, + recv: iroh::endpoint::RecvStream, +) -> Result<(), NodeError> { + let mut reader = BufReader::new(recv); + let mut line = String::new(); + reader + .read_line(&mut line) + .await + .map_err(|error| NodeError::IrohPeer(error.to_string()))?; + match geth_control::decode_pipe_wire_request(&line)? { PipeWireRequest::Send { peer_card, target, @@ -3539,61 +3769,219 @@ fn handle_pipe_wire_request( nonce, bearer_proof, } => { - geth_pipe::validate_pipe_name(&target)?; - 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 resource = format!("resource:pipe:{target}"); - let capability = "pipe.connect".to_owned(); - let explanation = explain_peer_or_bearer( - &store, - peer_card.node_id.as_str(), - &resource, - &capability, - &nonce, - bearer_proof.as_ref(), - )?; - let (message, listener_found) = if explanation.allowed { - record_pipe_message( - node, - target, - data_base64, - Some(peer_card.node_id.to_string()), - "remote pipe byte message over dedicated Iroh pipe ALPN".to_owned(), - )? - } else { - (None, false) - }; - Ok(PipeWireResponse::Sent { - 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(), - message: Box::new(message), - listener_found, - allowed: explanation.allowed, - reason: explanation.reason, - evaluated_ops: explanation.evaluated_ops, + let response = handle_pipe_send_wire_request( + &node, + &remote_endpoint_id, + peer_card, + target, + data_base64, nonce, - note: "pipe send authenticated endpoint/card binding and required pipe.connect on the remote pipe resource; payload used the dedicated Iroh pipe ALPN".to_owned(), - }) + bearer_proof, + )?; + send.write_all(geth_control::encode_pipe_wire_response(&response)?.as_bytes()) + .await + .map_err(|error| NodeError::IrohPeer(error.to_string()))?; + send.finish() + .map_err(|error| NodeError::IrohPeer(error.to_string()))?; + Ok(()) + } + PipeWireRequest::TcpConnect { + peer_card, + target_addr, + nonce, + bearer_proof, + } => { + handle_pipe_tcp_wire_connection( + node, + &remote_endpoint_id, + PipeTcpConnectWire { + peer_card, + target_addr, + nonce, + bearer_proof, + }, + send, + reader.into_inner(), + ) + .await } } } +fn handle_pipe_send_wire_request( + node: &LocalNode, + remote_endpoint_id: &str, + peer_card: PeerCard, + target: String, + data_base64: String, + nonce: String, + bearer_proof: Option, +) -> Result { + geth_pipe::validate_pipe_name(&target)?; + 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 resource = format!("resource:pipe:{target}"); + let capability = "pipe.connect".to_owned(); + let explanation = explain_peer_or_bearer( + &store, + peer_card.node_id.as_str(), + &resource, + &capability, + &nonce, + bearer_proof.as_ref(), + )?; + let (message, listener_found) = if explanation.allowed { + record_pipe_message( + node, + target, + data_base64, + Some(peer_card.node_id.to_string()), + "remote pipe byte message over dedicated Iroh pipe ALPN".to_owned(), + )? + } else { + (None, false) + }; + Ok(PipeWireResponse::Sent { + 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(), + message: Box::new(message), + listener_found, + allowed: explanation.allowed, + reason: explanation.reason, + evaluated_ops: explanation.evaluated_ops, + nonce, + note: "pipe send authenticated endpoint/card binding and required pipe.connect on the remote pipe resource; payload used the dedicated Iroh pipe ALPN".to_owned(), + }) +} + +async fn handle_pipe_tcp_wire_connection( + node: LocalNode, + remote_endpoint_id: &str, + request: PipeTcpConnectWire, + mut send: iroh::endpoint::SendStream, + recv: iroh::endpoint::RecvStream, +) -> Result<(), NodeError> { + let PipeTcpConnectWire { + peer_card, + target_addr, + nonce, + 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 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(), + &resource, + &capability, + &nonce, + bearer_proof.as_ref(), + )?; + if !explanation.allowed { + let response = PipeWireResponse::Connected { + 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(), + connection: None, + allowed: false, + reason: explanation.reason, + evaluated_ops: explanation.evaluated_ops, + nonce, + note: "pipe TCP forward denied before opening the remote TCP target".to_owned(), + }; + send.write_all(geth_control::encode_pipe_wire_response(&response)?.as_bytes()) + .await + .map_err(|error| NodeError::IrohPeer(error.to_string()))?; + return send + .finish() + .map_err(|error| NodeError::IrohPeer(error.to_string())); + } + let tcp = match TcpStream::connect(target).await { + Ok(tcp) => tcp, + Err(error) => { + let response = PipeWireResponse::Error { + message: format!( + "authorized pipe TCP forward could not connect to {target_addr}: {error}" + ), + }; + send.write_all(geth_control::encode_pipe_wire_response(&response)?.as_bytes()) + .await + .map_err(|error| NodeError::IrohPeer(error.to_string()))?; + return send + .finish() + .map_err(|error| NodeError::IrohPeer(error.to_string())); + } + }; + let response = PipeWireResponse::Connected { + 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(), + connection: Some(PipeConnection { + target: target_addr.clone(), + connected_at: UnixMillis(geth_store::now_ms()), + local_listener_found: true, + note: "authorized TCP byte stream over the dedicated Iroh pipe ALPN".to_owned(), + }), + allowed: true, + reason: explanation.reason, + evaluated_ops: explanation.evaluated_ops, + nonce, + note: "pipe TCP forward authenticated endpoint/card binding and required pipe.forward before connecting to the loopback TCP target".to_owned(), + }; + send.write_all(geth_control::encode_pipe_wire_response(&response)?.as_bytes()) + .await + .map_err(|error| NodeError::IrohPeer(error.to_string()))?; + let (mut tcp_read, mut tcp_write) = tcp.into_split(); + let mut recv = recv; + let inbound = async { + tokio::io::copy(&mut recv, &mut tcp_write) + .await + .map_err(|error| NodeError::IrohPeer(error.to_string()))?; + tcp_write.shutdown().await.map_err(NodeError::from) + }; + let outbound = async { + tokio::io::copy(&mut tcp_read, &mut send) + .await + .map_err(|error| NodeError::IrohPeer(error.to_string()))?; + send.finish() + .map_err(|error| NodeError::IrohPeer(error.to_string())) + }; + tokio::try_join!(inbound, outbound)?; + Ok(()) +} + fn iroh_node_addr_from_candidate( candidate: &EndpointCandidate, ) -> Result { @@ -4929,6 +5317,9 @@ pub fn handle_request( }) } ControlRequest::PipeSend { node: Some(_), .. } => Err(NodeError::IrohEndpointUnavailable), + ControlRequest::PipeTcpForward { .. } | ControlRequest::PipeTcpStream { .. } => { + Err(NodeError::IrohEndpointUnavailable) + } ControlRequest::PipeRecv { name, peek } => { geth_pipe::validate_pipe_name(&name)?; let messages = pipe_messages(node, &name, !peek)?; diff --git a/crates/geth-pipe/src/lib.rs b/crates/geth-pipe/src/lib.rs index 36e86d4..5e0a01c 100644 --- a/crates/geth-pipe/src/lib.rs +++ b/crates/geth-pipe/src/lib.rs @@ -1,5 +1,6 @@ use geth_types::{PipeId, ResourceId, UnixMillis}; use serde::{Deserialize, Serialize}; +use std::net::SocketAddr; #[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] pub struct PipeResource { @@ -37,6 +38,10 @@ pub struct PipeMessage { pub enum PipeError { #[error("invalid pipe name or target: {0}")] InvalidName(String), + #[error("invalid TCP address: {0}")] + InvalidTcpAddress(String), + #[error("TCP forwarding addresses must be loopback addresses: {0}")] + NonLoopbackTcpAddress(String), } pub fn validate_pipe_name(name: &str) -> Result<(), PipeError> { @@ -50,6 +55,26 @@ pub fn validate_pipe_name(name: &str) -> Result<(), PipeError> { Ok(()) } +pub fn validate_tcp_forward_listen_addr(addr: &str) -> Result { + let addr = addr + .parse::() + .map_err(|_| PipeError::InvalidTcpAddress(addr.to_owned()))?; + if !addr.ip().is_loopback() { + return Err(PipeError::NonLoopbackTcpAddress(addr.to_string())); + } + Ok(addr) +} + +pub fn validate_tcp_forward_target_addr(addr: &str) -> Result { + let addr = addr + .parse::() + .map_err(|_| PipeError::InvalidTcpAddress(addr.to_owned()))?; + if !addr.ip().is_loopback() { + return Err(PipeError::NonLoopbackTcpAddress(addr.to_string())); + } + Ok(addr) +} + #[must_use] pub fn pipe_roadmap() -> &'static str { "future pipes are authorized Iroh bidirectional streams for stdin/stdout and forwarding" @@ -73,4 +98,13 @@ mod tests { assert!(validate_pipe_name("inbox/main").is_err()); assert!(validate_pipe_name("inbox main").is_err()); } + + #[test] + fn tcp_forward_addresses_must_be_explicit_loopback_socket_addrs() { + assert!(validate_tcp_forward_listen_addr("127.0.0.1:9000").is_ok()); + assert!(validate_tcp_forward_target_addr("[::1]:22").is_ok()); + assert!(validate_tcp_forward_listen_addr("localhost:9000").is_err()); + assert!(validate_tcp_forward_target_addr("0.0.0.0:22").is_err()); + assert!(validate_tcp_forward_target_addr("192.0.2.10:22").is_err()); + } } diff --git a/docs/architecture.md b/docs/architecture.md index f16f65c..c68a29e 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -213,8 +213,13 @@ uses the dedicated `/geth/pipe/1` ALPN to write a byte message to a peer listener after the same endpoint/card and capability checks. `geth pipe recv ` drains local daemon-lifetime messages. `geth pipe listen --node ` can also ask a peer to register a daemon-lifetime listener after -checking `pipe.listen` on the same resource. Long-lived stdin/stdout streams and -TCP/Unix forwarding are still future work. +checking `pipe.listen` on the same resource. `geth pipe forward-tcp --listen +127.0.0.1: --node --target 127.0.0.1:` runs a +local loopback listener and opens one authorized `/geth/pipe/1` byte stream per +accepted connection. The remote daemon validates endpoint/card binding and +requires `pipe.forward` on `resource:pipe-tcp:` before connecting to the +remote loopback TCP target. TCP forwarding is loopback-only in the prototype; +Unix socket forwarding is still future work. `geth-ssh-proxy` defines proxy target and connection metadata. `geth ssh proxy ` is a streaming command intended for OpenSSH `ProxyCommand`: the CLI diff --git a/docs/roadmap.md b/docs/roadmap.md index 9d5f495..a1d572d 100644 --- a/docs/roadmap.md +++ b/docs/roadmap.md @@ -369,14 +369,21 @@ Goal: add authorized stream-oriented management workflows over Iroh. - `[x]` `geth pipe recv ` drains daemon-lifetime pipe messages. - `[x]` Remote pipe send requires `pipe.connect` on `resource:pipe:`. - - `[ ]` Pipe connect carries bidirectional byte streams over Iroh. - - `[ ]` Streams close cleanly and propagate errors. + - `[x]` TCP forwarding carries bidirectional byte streams over Iroh. + - `[x]` TCP streams close cleanly and propagate errors through the local + forwarder logs. -- `[ ]` TCP forwarding. +- `[~]` TCP forwarding. Acceptance criteria: - - A local TCP listener can forward over an authorized Iroh pipe. - - Tests cover basic request/response forwarding. - - Forwarding is resource-scoped and can be disabled by auth. + - `[x]` A local loopback TCP listener can forward over an authorized Iroh + pipe with `geth pipe forward-tcp`. + - `[x]` The remote side connects only to explicit loopback socket addresses + in the prototype. + - `[x]` Forwarding is resource-scoped with `pipe.forward` on + `resource:pipe-tcp:`. + - `[x]` Tests cover address validation and TCP pipe wire + request/response serialization. + - `[ ]` Tests cover a full two-node request/response forwarding exchange. - `[ ]` Unix socket forwarding where supported. Acceptance criteria: