diff --git a/crates/geth-node/src/lib.rs b/crates/geth-node/src/lib.rs index 23730e0..614d550 100644 --- a/crates/geth-node/src/lib.rs +++ b/crates/geth-node/src/lib.rs @@ -1,6 +1,7 @@ pub mod backup; mod daemon; pub mod doctor; +mod local_control; mod peer_client; mod runtime; pub mod service; @@ -64,6 +65,7 @@ use geth_types::{ }; use iroh::protocol::ProtocolHandler; use iroh_docs::api::protocol::{AddrInfoOptions, ShareMode}; +pub use local_control::handle_request_async; use peer_client::{request_overlay_wire, request_peer_control, request_pipe_wire}; use runtime::{ NodeRuntime, OverlayTunCounters, OverlayTunRuntime, PipeRuntime, PubsubGossipTopicRuntime, @@ -644,488 +646,6 @@ async fn stream_pipe_unix( } } -pub async fn handle_request_async( - node: &LocalNode, - request: ControlRequest, -) -> Result { - let trace = control_request_trace_fields(&request); - tracing::info!( - command = trace.command, - peer_node = trace.peer_node.as_deref().unwrap_or(""), - resource = trace.resource.as_deref().unwrap_or(""), - capability = trace.capability.as_deref().unwrap_or(""), - stream = trace.stream.as_deref().unwrap_or(""), - "control request started" - ); - let result = match request { - ControlRequest::PeerCardExport { out } => export_peer_card(node, out, true).await, - ControlRequest::PeerPing { node: peer_node } => peer_ping(node, &peer_node).await, - ControlRequest::PeerAuthCheck { - node: peer_node, - resource, - capability, - } => peer_auth_check(node, &peer_node, resource, capability).await, - ControlRequest::KeychainSync { node: peer_node } => { - keychain_sync_from_peer(node, &peer_node).await - } - ControlRequest::AuthSync { node: peer_node } => auth_sync_from_peer(node, &peer_node).await, - ControlRequest::SyncNow { node: peer_node } => sync_now(node, peer_node).await, - ControlRequest::NodeEnrollSubmit { - owner_node, - request_id, - path, - } => node_enrollment_submit_to_peer(node, &owner_node, request_id, path).await, - ControlRequest::NodeEnrollSync { owner_node } => { - node_enrollment_sync_from_owner(node, &owner_node).await - } - ControlRequest::CasFetch { - node: peer_node, - hash, - bearer_secret, - } => cas_fetch_from_peer(node, &peer_node, hash, bearer_secret).await, - ControlRequest::CasRootSync { - node: peer_node, - name, - bearer_secret, - } => cas_root_sync_from_peer(node, &peer_node, &name, bearer_secret).await, - ControlRequest::SshCertSync { - node: peer_node, - bearer_secret, - } => ssh_cert_sync_from_peer(node, &peer_node, bearer_secret).await, - ControlRequest::SshRevocationSync { - node: peer_node, - bearer_secret, - } => ssh_revocation_sync_from_peer(node, &peer_node, bearer_secret).await, - ControlRequest::KvSync { - node: peer_node, - name, - bearer_secret, - } => kv_sync_from_peer(node, &peer_node, &name, bearer_secret).await, - ControlRequest::DbSync { - node: peer_node, - name, - limit, - bearer_secret, - } => db_sync_from_peer(node, &peer_node, &name, limit, bearer_secret).await, - ControlRequest::PubsubPub { - topic, - message, - node: Some(peer_node), - bearer_secret, - } => pubsub_publish_to_peer(node, &peer_node, topic, message, bearer_secret).await, - ControlRequest::PubsubSub { - topic, - node: Some(peer_node), - bearer_secret, - } => pubsub_subscribe_from_peer(node, &peer_node, topic, bearer_secret).await, - ControlRequest::PipeListen { - name, - node: Some(peer_node), - bearer_secret, - } => pipe_listen_on_peer(node, &peer_node, name, bearer_secret).await, - ControlRequest::PipeConnect { - target, - node: Some(peer_node), - bearer_secret, - } => pipe_connect_to_peer(node, &peer_node, target, bearer_secret).await, - ControlRequest::PipeSend { - target, - data_base64, - node: Some(peer_node), - bearer_secret, - } => pipe_send_to_peer(node, &peer_node, target, data_base64, bearer_secret).await, - ControlRequest::PipeTcpForward { .. } - | ControlRequest::PipeTcpStream { .. } - | ControlRequest::PipeUnixForward { .. } - | ControlRequest::PipeUnixStream { .. } => Err(NodeError::IrohPeer( - "pipe forwarding is a streaming control command; run without --json/--jsonl".to_owned(), - )), - ControlRequest::SshProxyConnect { - node: peer_node, - bearer_secret, - } => ssh_proxy_connect_to_peer(node, &peer_node, bearer_secret).await, - ControlRequest::SshAdminShell { - node: peer_node, - command, - bearer_secret, - } => ssh_admin_shell_to_peer(node, &peer_node, command, bearer_secret).await, - ControlRequest::DocumentSync { - node: peer_node, - name, - bearer_secret, - } => document_sync_from_peer(node, &peer_node, &name, bearer_secret).await, - ControlRequest::OverlaySend { - name, - node: peer_node, - packet_base64, - bearer_secret, - } => overlay_send_to_peer(node, &peer_node, name, packet_base64, bearer_secret).await, - ControlRequest::OverlayUp { - name, - bearer_secret, - mtu, - } => overlay_runtime_up(node, name, bearer_secret, mtu).await, - ControlRequest::OverlayDown { name } => overlay_runtime_down(node, &name), - other => { - let response = handle_request(node, other)?; - match &response { - ControlResponse::CasAdded { hash, .. } => { - mirror_local_blob_to_iroh_blobs(node, hash).await?; - } - ControlResponse::CasPrivateAdded { encrypted_hash, .. } => { - mirror_local_blob_to_iroh_blobs(node, encrypted_hash).await?; - } - ControlResponse::KvCreated { kv } => { - mirror_kv_store_to_iroh_docs(node, kv.name.as_str()).await?; - } - ControlResponse::KvSet { entry } => { - let name = entry - .store - .as_str() - .strip_prefix("kv:") - .unwrap_or(entry.store.as_str()); - mirror_kv_store_to_iroh_docs(node, name).await?; - } - ControlResponse::PubsubPublished { message } => { - broadcast_pubsub_gossip( - node, - message.topic.as_str(), - message.message.as_str(), - Vec::new(), - ) - .await?; - } - _ => {} - } - Ok(response) - } - }; - match &result { - Ok(_) => tracing::info!( - command = trace.command, - peer_node = trace.peer_node.as_deref().unwrap_or(""), - resource = trace.resource.as_deref().unwrap_or(""), - capability = trace.capability.as_deref().unwrap_or(""), - stream = trace.stream.as_deref().unwrap_or(""), - "control request completed" - ), - Err(error) => tracing::warn!( - command = trace.command, - peer_node = trace.peer_node.as_deref().unwrap_or(""), - resource = trace.resource.as_deref().unwrap_or(""), - capability = trace.capability.as_deref().unwrap_or(""), - stream = trace.stream.as_deref().unwrap_or(""), - error_code = node_error_code(error), - %error, - "control request failed" - ), - } - result -} - -#[derive(Debug, PartialEq, Eq)] -struct ControlTraceFields { - command: &'static str, - peer_node: Option, - resource: Option, - capability: Option, - stream: Option, -} - -fn control_request_trace_fields(request: &ControlRequest) -> ControlTraceFields { - let mut trace = ControlTraceFields { - command: control_request_command(request), - peer_node: None, - resource: None, - capability: None, - stream: None, - }; - match request { - ControlRequest::PeerPing { node } - | ControlRequest::KeychainSync { node } - | ControlRequest::AuthSync { node } => { - trace.peer_node = Some(node.clone()); - } - ControlRequest::PeerAuthCheck { - node, - resource, - capability, - } => { - trace.peer_node = Some(node.clone()); - trace.resource = Some(resource.clone()); - trace.capability = Some(capability.clone()); - } - ControlRequest::CasRootSync { node, name, .. } => { - trace.peer_node = Some(node.clone()); - trace.resource = Some(format!("resource:cas-tree:{name}")); - trace.capability = Some("cas.fetch".to_owned()); - trace.stream = Some(format!("cas-tree:{name}")); - } - ControlRequest::SyncNow { node } => { - trace.peer_node = node.clone(); - trace.stream = Some("all".to_owned()); - } - ControlRequest::NodeGrant { - node, - resource, - capability, - .. - } => { - trace.peer_node = Some(node.clone()); - trace.resource = Some(resource.clone()); - trace.capability = Some(capability.clone()); - } - ControlRequest::NodeRevokeGrant { resource, .. } => { - trace.resource = Some(resource.clone()); - } - ControlRequest::AuthExplain { - resource, - capability, - .. - } - | ControlRequest::AuthGrant { - resource, - capability, - .. - } => { - trace.resource = Some(resource.clone()); - trace.capability = Some(capability.clone()); - } - ControlRequest::AuthRevoke { resource, .. } - | ControlRequest::SecretCreate { resource } - | ControlRequest::SecretRotate { resource } - | ControlRequest::SecretBearerCreate { resource, .. } - | ControlRequest::SecretBearerChallenge { resource, .. } - | ControlRequest::SecretBearerProve { resource, .. } - | ControlRequest::SecretBearerVerify { resource, .. } - | ControlRequest::SecretBearerRevoke { resource, .. } - | ControlRequest::CasAddPrivate { resource, .. } - | ControlRequest::CasGetPrivate { resource, .. } => { - trace.resource = Some(resource.clone()); - } - ControlRequest::CasFetch { node, .. } => { - trace.peer_node = Some(node.clone()); - trace.resource = Some("resource:cas:local".to_owned()); - trace.capability = Some("cas.fetch".to_owned()); - } - ControlRequest::SshCertSync { node, .. } => { - trace.peer_node = Some(node.clone()); - trace.resource = Some("resource:ssh:certs".to_owned()); - trace.capability = Some("ssh_cert.sync".to_owned()); - trace.stream = Some("ssh-certs".to_owned()); - } - ControlRequest::SshRevocationSync { node, .. } => { - trace.peer_node = Some(node.clone()); - trace.resource = Some("resource:ssh:revocations".to_owned()); - trace.capability = Some("ssh_revocation.sync".to_owned()); - trace.stream = Some("ssh-revocations".to_owned()); - } - ControlRequest::KvSet { name, key, .. } => { - trace.resource = Some(format!("resource:kv:{name}")); - trace.capability = Some(format!("kv.write_key:{key}")); - trace.stream = Some(format!("kv:{name}")); - } - ControlRequest::KvSync { node, name, .. } => { - trace.peer_node = Some(node.clone()); - trace.resource = Some(format!("resource:kv:{name}")); - trace.capability = Some("kv.read".to_owned()); - trace.stream = Some(format!("kv:{name}")); - } - ControlRequest::KvCreate { name } | ControlRequest::KvGet { name, .. } => { - trace.resource = Some(format!("resource:kv:{name}")); - trace.stream = Some(format!("kv:{name}")); - } - ControlRequest::DbSync { node, name, .. } => { - trace.peer_node = Some(node.clone()); - trace.resource = Some(format!("resource:db:{name}")); - trace.capability = Some("db.sync".to_owned()); - trace.stream = Some(format!("db:{name}")); - } - ControlRequest::DbAdd { name, .. } - | ControlRequest::DbStatus { name } - | ControlRequest::DbChanges { name, .. } => { - trace.resource = Some(format!("resource:db:{name}")); - trace.stream = Some(format!("db:{name}")); - } - ControlRequest::DocumentSync { node, name, .. } => { - trace.peer_node = Some(node.clone()); - trace.resource = Some(format!("resource:document:{name}")); - trace.capability = Some("document.read".to_owned()); - trace.stream = Some(format!("document:{name}")); - } - ControlRequest::DocumentCreate { name } - | ControlRequest::DocumentStatus { name } - | ControlRequest::DocumentSet { name, .. } - | ControlRequest::DocumentGet { name } => { - trace.resource = Some(format!("resource:document:{name}")); - trace.stream = Some(format!("document:{name}")); - } - ControlRequest::PubsubPub { topic, node, .. } => { - trace.peer_node = node.clone(); - trace.resource = Some(format!("resource:pubsub:{topic}")); - trace.capability = Some("pubsub.publish".to_owned()); - trace.stream = Some(format!("pubsub:{topic}")); - } - ControlRequest::PubsubSub { topic, node, .. } => { - trace.peer_node = node.clone(); - trace.resource = Some(format!("resource:pubsub:{topic}")); - trace.capability = Some("pubsub.subscribe".to_owned()); - trace.stream = Some(format!("pubsub:{topic}")); - } - ControlRequest::PipeListen { name, node, .. } => { - trace.peer_node = node.clone(); - trace.resource = Some(format!("resource:pipe:{name}")); - trace.capability = Some("pipe.listen".to_owned()); - trace.stream = Some(format!("pipe:{name}")); - } - ControlRequest::PipeConnect { target, node, .. } - | ControlRequest::PipeSend { target, node, .. } => { - trace.peer_node = node.clone(); - trace.resource = Some(format!("resource:pipe:{target}")); - trace.capability = Some("pipe.connect".to_owned()); - trace.stream = Some(format!("pipe:{target}")); - } - ControlRequest::PipeTcpForward { - node, target_addr, .. - } - | ControlRequest::PipeTcpStream { - node, target_addr, .. - } => { - trace.peer_node = Some(node.clone()); - trace.resource = Some(format!("resource:pipe-tcp:{target_addr}")); - trace.capability = Some("pipe.forward".to_owned()); - trace.stream = Some("pipe-tcp".to_owned()); - } - ControlRequest::PipeUnixForward { - node, target_path, .. - } - | ControlRequest::PipeUnixStream { - node, target_path, .. - } => { - trace.peer_node = Some(node.clone()); - trace.resource = Some(format!("resource:pipe-unix:{}", target_path.display())); - trace.capability = Some("pipe.forward".to_owned()); - trace.stream = Some("pipe-unix".to_owned()); - } - ControlRequest::SshProxyConnect { node, .. } - | ControlRequest::SshProxyStream { node, .. } => { - trace.peer_node = Some(node.clone()); - trace.resource = Some("resource:ssh-proxy:local".to_owned()); - trace.capability = Some("ssh_proxy.connect".to_owned()); - trace.stream = Some("ssh-proxy".to_owned()); - } - ControlRequest::SshAdminShell { node, .. } => { - trace.peer_node = Some(node.clone()); - trace.resource = Some("resource:ssh-proxy:local".to_owned()); - trace.capability = Some("ssh_proxy.admin_shell".to_owned()); - trace.stream = Some("ssh-admin-shell".to_owned()); - } - ControlRequest::OverlaySend { name, node, .. } => { - trace.peer_node = Some(node.clone()); - trace.resource = Some(format!("resource:overlay:{name}")); - trace.capability = Some("overlay.route".to_owned()); - trace.stream = Some(format!("overlay:{name}")); - } - ControlRequest::OverlayUp { name, .. } => { - trace.resource = Some(format!("resource:overlay:{name}")); - trace.capability = Some("overlay.route".to_owned()); - trace.stream = Some(format!("overlay:{name}")); - } - _ => {} - } - trace -} - -fn control_request_command(request: &ControlRequest) -> &'static str { - match request { - ControlRequest::Status => "status", - ControlRequest::NodeId => "node-id", - ControlRequest::PeerCardExport { .. } => "peer.export", - ControlRequest::PeerCardImport { .. } => "peer.import", - ControlRequest::PeerCardList => "peer.list", - ControlRequest::PeerPing { .. } => "peer.ping", - ControlRequest::PeerAuthCheck { .. } => "peer.auth-check", - ControlRequest::ResourceList => "resource.list", - ControlRequest::ResourceCreate { .. } => "resource.create", - ControlRequest::KeychainSync { .. } => "keychain.sync", - ControlRequest::AuthSync { .. } => "auth.sync", - ControlRequest::SyncStatus => "sync.status", - ControlRequest::SyncNow { .. } => "sync.now", - ControlRequest::AuthExplain { .. } => "auth.explain", - ControlRequest::AuthGrant { .. } => "auth.grant", - ControlRequest::AuthRevoke { .. } => "auth.revoke", - ControlRequest::CasFetch { .. } => "cas.fetch", - ControlRequest::CasRootSync { .. } => "cas.root.sync", - ControlRequest::CasAddPrivate { .. } => "cas.add-private", - ControlRequest::CasGetPrivate { .. } => "cas.get-private", - ControlRequest::SshCertSync { .. } => "ssh.cert.sync", - ControlRequest::SshRevocationSync { .. } => "ssh.revocation.sync", - ControlRequest::KvCreate { .. } => "kv.create", - ControlRequest::KvSet { .. } => "kv.set", - ControlRequest::KvGet { .. } => "kv.get", - ControlRequest::KvSync { .. } => "kv.sync", - ControlRequest::DbAdd { .. } => "db.add", - ControlRequest::DbStatus { .. } => "db.status", - ControlRequest::DbChanges { .. } => "db.changes", - ControlRequest::DbSync { .. } => "db.sync", - ControlRequest::DocumentCreate { .. } => "document.create", - ControlRequest::DocumentStatus { .. } => "document.status", - ControlRequest::DocumentSet { .. } => "document.set", - ControlRequest::DocumentGet { .. } => "document.get", - ControlRequest::DocumentSync { .. } => "document.sync", - ControlRequest::PubsubPub { .. } => "pubsub.pub", - ControlRequest::PubsubSub { .. } => "pubsub.sub", - ControlRequest::PipeListen { .. } => "pipe.listen", - ControlRequest::PipeConnect { .. } => "pipe.connect", - ControlRequest::PipeSend { .. } => "pipe.send", - ControlRequest::PipeRecv { .. } => "pipe.recv", - ControlRequest::PipeTcpForward { .. } => "pipe.forward-tcp", - ControlRequest::PipeTcpStream { .. } => "pipe.tcp-stream", - ControlRequest::PipeUnixForward { .. } => "pipe.forward-unix", - ControlRequest::PipeUnixStream { .. } => "pipe.unix-stream", - ControlRequest::SshProxyConnect { .. } => "ssh.proxy", - ControlRequest::SshProxyStream { .. } => "ssh.proxy-stream", - ControlRequest::SshAdminShell { .. } => "ssh.admin-shell", - ControlRequest::OverlayStatus => "overlay.status", - ControlRequest::OverlayPlan { .. } => "overlay.plan", - ControlRequest::OverlayJoin { .. } => "overlay.join", - ControlRequest::OverlayLeave { .. } => "overlay.leave", - ControlRequest::OverlayInterfacePlan { .. } => "overlay.interface-plan", - ControlRequest::OverlayUp { .. } => "overlay.up", - ControlRequest::OverlayDown { .. } => "overlay.down", - ControlRequest::OverlayPeers { .. } => "overlay.peers", - ControlRequest::OverlaySend { .. } => "overlay.send", - ControlRequest::OverlayRecv { .. } => "overlay.recv", - ControlRequest::SecretStatus => "secret.status", - ControlRequest::SecretCreate { .. } => "secret.create", - ControlRequest::SecretRotate { .. } => "secret.rotate", - ControlRequest::SecretBearerCreate { .. } => "secret.bearer.create", - ControlRequest::SecretBearerList => "secret.bearer.list", - ControlRequest::SecretBearerChallenge { .. } => "secret.bearer.challenge", - ControlRequest::SecretBearerProve { .. } => "secret.bearer.prove", - ControlRequest::SecretBearerVerify { .. } => "secret.bearer.verify", - ControlRequest::SecretBearerRevoke { .. } => "secret.bearer.revoke", - ControlRequest::ModuleStub { .. } => "module.stub", - _ => "control.request", - } -} - -fn node_error_code(error: &NodeError) -> &'static str { - match error { - NodeError::Unauthorized(_) => "unauthorized", - NodeError::PeerNotFound(_) => "peer_not_found", - NodeError::IrohEndpointUnavailable => "iroh_endpoint_unavailable", - NodeError::KvNotFound(_) => "kv_not_found", - NodeError::DbNotFound(_) => "db_not_found", - NodeError::DocumentNotFound(_) => "document_not_found", - NodeError::Cas(_) => "cas_error", - NodeError::Store(_) => "store_error", - NodeError::Config(_) => "config_error", - NodeError::Codec(_) => "codec_error", - _ => "node_error", - } -} - async fn handle_stream(node: LocalNode, stream: UnixStream) -> Result<(), NodeError> { let mut reader = BufReader::new(stream); let mut line = String::new(); @@ -1161,7 +681,7 @@ async fn handle_stream(node: LocalNode, stream: UnixStream) -> Result<(), NodeEr resource = "resource:ssh-proxy:local", capability = "ssh_proxy.connect", stream = "ssh-proxy", - error_code = node_error_code(error), + error_code = local_control::node_error_code(error), %error, "streaming control request failed" ), @@ -1202,7 +722,7 @@ async fn handle_stream(node: LocalNode, stream: UnixStream) -> Result<(), NodeEr resource = %resource, capability = "pipe.forward", stream = "pipe-tcp", - error_code = node_error_code(error), + error_code = local_control::node_error_code(error), %error, "streaming control request failed" ), @@ -1243,7 +763,7 @@ async fn handle_stream(node: LocalNode, stream: UnixStream) -> Result<(), NodeEr resource = %resource, capability = "pipe.forward", stream = "pipe-unix", - error_code = node_error_code(error), + error_code = local_control::node_error_code(error), %error, "streaming control request failed" ), @@ -11889,7 +11409,7 @@ mod tests { #[test] fn control_trace_fields_exclude_bearer_secrets() { - let trace = control_request_trace_fields(&ControlRequest::KvSync { + let trace = local_control::control_request_trace_fields(&ControlRequest::KvSync { node: "node:peer".to_owned(), name: "prefs".to_owned(), bearer_secret: Some("private-token-value".to_owned()), @@ -11905,15 +11425,15 @@ mod tests { #[test] fn node_error_codes_are_stable_for_common_tracing_failures() { assert_eq!( - node_error_code(&NodeError::Unauthorized("missing grant".to_owned())), + local_control::node_error_code(&NodeError::Unauthorized("missing grant".to_owned())), "unauthorized" ); assert_eq!( - node_error_code(&NodeError::PeerNotFound("node:missing".to_owned())), + local_control::node_error_code(&NodeError::PeerNotFound("node:missing".to_owned())), "peer_not_found" ); assert_eq!( - node_error_code(&NodeError::IrohEndpointUnavailable), + local_control::node_error_code(&NodeError::IrohEndpointUnavailable), "iroh_endpoint_unavailable" ); } diff --git a/crates/geth-node/src/local_control.rs b/crates/geth-node/src/local_control.rs new file mode 100644 index 0000000..1f854dd --- /dev/null +++ b/crates/geth-node/src/local_control.rs @@ -0,0 +1,485 @@ +//! Local control request routing and trace classification. + +use super::*; + +pub async fn handle_request_async( + node: &LocalNode, + request: ControlRequest, +) -> Result { + let trace = control_request_trace_fields(&request); + tracing::info!( + command = trace.command, + peer_node = trace.peer_node.as_deref().unwrap_or(""), + resource = trace.resource.as_deref().unwrap_or(""), + capability = trace.capability.as_deref().unwrap_or(""), + stream = trace.stream.as_deref().unwrap_or(""), + "control request started" + ); + let result = match request { + ControlRequest::PeerCardExport { out } => export_peer_card(node, out, true).await, + ControlRequest::PeerPing { node: peer_node } => peer_ping(node, &peer_node).await, + ControlRequest::PeerAuthCheck { + node: peer_node, + resource, + capability, + } => peer_auth_check(node, &peer_node, resource, capability).await, + ControlRequest::KeychainSync { node: peer_node } => { + keychain_sync_from_peer(node, &peer_node).await + } + ControlRequest::AuthSync { node: peer_node } => auth_sync_from_peer(node, &peer_node).await, + ControlRequest::SyncNow { node: peer_node } => sync_now(node, peer_node).await, + ControlRequest::NodeEnrollSubmit { + owner_node, + request_id, + path, + } => node_enrollment_submit_to_peer(node, &owner_node, request_id, path).await, + ControlRequest::NodeEnrollSync { owner_node } => { + node_enrollment_sync_from_owner(node, &owner_node).await + } + ControlRequest::CasFetch { + node: peer_node, + hash, + bearer_secret, + } => cas_fetch_from_peer(node, &peer_node, hash, bearer_secret).await, + ControlRequest::CasRootSync { + node: peer_node, + name, + bearer_secret, + } => cas_root_sync_from_peer(node, &peer_node, &name, bearer_secret).await, + ControlRequest::SshCertSync { + node: peer_node, + bearer_secret, + } => ssh_cert_sync_from_peer(node, &peer_node, bearer_secret).await, + ControlRequest::SshRevocationSync { + node: peer_node, + bearer_secret, + } => ssh_revocation_sync_from_peer(node, &peer_node, bearer_secret).await, + ControlRequest::KvSync { + node: peer_node, + name, + bearer_secret, + } => kv_sync_from_peer(node, &peer_node, &name, bearer_secret).await, + ControlRequest::DbSync { + node: peer_node, + name, + limit, + bearer_secret, + } => db_sync_from_peer(node, &peer_node, &name, limit, bearer_secret).await, + ControlRequest::PubsubPub { + topic, + message, + node: Some(peer_node), + bearer_secret, + } => pubsub_publish_to_peer(node, &peer_node, topic, message, bearer_secret).await, + ControlRequest::PubsubSub { + topic, + node: Some(peer_node), + bearer_secret, + } => pubsub_subscribe_from_peer(node, &peer_node, topic, bearer_secret).await, + ControlRequest::PipeListen { + name, + node: Some(peer_node), + bearer_secret, + } => pipe_listen_on_peer(node, &peer_node, name, bearer_secret).await, + ControlRequest::PipeConnect { + target, + node: Some(peer_node), + bearer_secret, + } => pipe_connect_to_peer(node, &peer_node, target, bearer_secret).await, + ControlRequest::PipeSend { + target, + data_base64, + node: Some(peer_node), + bearer_secret, + } => pipe_send_to_peer(node, &peer_node, target, data_base64, bearer_secret).await, + ControlRequest::PipeTcpForward { .. } + | ControlRequest::PipeTcpStream { .. } + | ControlRequest::PipeUnixForward { .. } + | ControlRequest::PipeUnixStream { .. } => Err(NodeError::IrohPeer( + "pipe forwarding is a streaming control command; run without --json/--jsonl".to_owned(), + )), + ControlRequest::SshProxyConnect { + node: peer_node, + bearer_secret, + } => ssh_proxy_connect_to_peer(node, &peer_node, bearer_secret).await, + ControlRequest::SshAdminShell { + node: peer_node, + command, + bearer_secret, + } => ssh_admin_shell_to_peer(node, &peer_node, command, bearer_secret).await, + ControlRequest::DocumentSync { + node: peer_node, + name, + bearer_secret, + } => document_sync_from_peer(node, &peer_node, &name, bearer_secret).await, + ControlRequest::OverlaySend { + name, + node: peer_node, + packet_base64, + bearer_secret, + } => overlay_send_to_peer(node, &peer_node, name, packet_base64, bearer_secret).await, + ControlRequest::OverlayUp { + name, + bearer_secret, + mtu, + } => overlay_runtime_up(node, name, bearer_secret, mtu).await, + ControlRequest::OverlayDown { name } => overlay_runtime_down(node, &name), + other => { + let response = handle_request(node, other)?; + match &response { + ControlResponse::CasAdded { hash, .. } => { + mirror_local_blob_to_iroh_blobs(node, hash).await?; + } + ControlResponse::CasPrivateAdded { encrypted_hash, .. } => { + mirror_local_blob_to_iroh_blobs(node, encrypted_hash).await?; + } + ControlResponse::KvCreated { kv } => { + mirror_kv_store_to_iroh_docs(node, kv.name.as_str()).await?; + } + ControlResponse::KvSet { entry } => { + let name = entry + .store + .as_str() + .strip_prefix("kv:") + .unwrap_or(entry.store.as_str()); + mirror_kv_store_to_iroh_docs(node, name).await?; + } + ControlResponse::PubsubPublished { message } => { + broadcast_pubsub_gossip( + node, + message.topic.as_str(), + message.message.as_str(), + Vec::new(), + ) + .await?; + } + _ => {} + } + Ok(response) + } + }; + match &result { + Ok(_) => tracing::info!( + command = trace.command, + peer_node = trace.peer_node.as_deref().unwrap_or(""), + resource = trace.resource.as_deref().unwrap_or(""), + capability = trace.capability.as_deref().unwrap_or(""), + stream = trace.stream.as_deref().unwrap_or(""), + "control request completed" + ), + Err(error) => tracing::warn!( + command = trace.command, + peer_node = trace.peer_node.as_deref().unwrap_or(""), + resource = trace.resource.as_deref().unwrap_or(""), + capability = trace.capability.as_deref().unwrap_or(""), + stream = trace.stream.as_deref().unwrap_or(""), + error_code = node_error_code(error), + %error, + "control request failed" + ), + } + result +} + +#[derive(Debug, PartialEq, Eq)] +pub(crate) struct ControlTraceFields { + pub(crate) command: &'static str, + pub(crate) peer_node: Option, + pub(crate) resource: Option, + pub(crate) capability: Option, + pub(crate) stream: Option, +} + +pub(crate) fn control_request_trace_fields(request: &ControlRequest) -> ControlTraceFields { + let mut trace = ControlTraceFields { + command: control_request_command(request), + peer_node: None, + resource: None, + capability: None, + stream: None, + }; + match request { + ControlRequest::PeerPing { node } + | ControlRequest::KeychainSync { node } + | ControlRequest::AuthSync { node } => { + trace.peer_node = Some(node.clone()); + } + ControlRequest::PeerAuthCheck { + node, + resource, + capability, + } => { + trace.peer_node = Some(node.clone()); + trace.resource = Some(resource.clone()); + trace.capability = Some(capability.clone()); + } + ControlRequest::CasRootSync { node, name, .. } => { + trace.peer_node = Some(node.clone()); + trace.resource = Some(format!("resource:cas-tree:{name}")); + trace.capability = Some("cas.fetch".to_owned()); + trace.stream = Some(format!("cas-tree:{name}")); + } + ControlRequest::SyncNow { node } => { + trace.peer_node = node.clone(); + trace.stream = Some("all".to_owned()); + } + ControlRequest::NodeGrant { + node, + resource, + capability, + .. + } => { + trace.peer_node = Some(node.clone()); + trace.resource = Some(resource.clone()); + trace.capability = Some(capability.clone()); + } + ControlRequest::NodeRevokeGrant { resource, .. } => { + trace.resource = Some(resource.clone()); + } + ControlRequest::AuthExplain { + resource, + capability, + .. + } + | ControlRequest::AuthGrant { + resource, + capability, + .. + } => { + trace.resource = Some(resource.clone()); + trace.capability = Some(capability.clone()); + } + ControlRequest::AuthRevoke { resource, .. } + | ControlRequest::SecretCreate { resource } + | ControlRequest::SecretRotate { resource } + | ControlRequest::SecretBearerCreate { resource, .. } + | ControlRequest::SecretBearerChallenge { resource, .. } + | ControlRequest::SecretBearerProve { resource, .. } + | ControlRequest::SecretBearerVerify { resource, .. } + | ControlRequest::SecretBearerRevoke { resource, .. } + | ControlRequest::CasAddPrivate { resource, .. } + | ControlRequest::CasGetPrivate { resource, .. } => { + trace.resource = Some(resource.clone()); + } + ControlRequest::CasFetch { node, .. } => { + trace.peer_node = Some(node.clone()); + trace.resource = Some("resource:cas:local".to_owned()); + trace.capability = Some("cas.fetch".to_owned()); + } + ControlRequest::SshCertSync { node, .. } => { + trace.peer_node = Some(node.clone()); + trace.resource = Some("resource:ssh:certs".to_owned()); + trace.capability = Some("ssh_cert.sync".to_owned()); + trace.stream = Some("ssh-certs".to_owned()); + } + ControlRequest::SshRevocationSync { node, .. } => { + trace.peer_node = Some(node.clone()); + trace.resource = Some("resource:ssh:revocations".to_owned()); + trace.capability = Some("ssh_revocation.sync".to_owned()); + trace.stream = Some("ssh-revocations".to_owned()); + } + ControlRequest::KvSet { name, key, .. } => { + trace.resource = Some(format!("resource:kv:{name}")); + trace.capability = Some(format!("kv.write_key:{key}")); + trace.stream = Some(format!("kv:{name}")); + } + ControlRequest::KvSync { node, name, .. } => { + trace.peer_node = Some(node.clone()); + trace.resource = Some(format!("resource:kv:{name}")); + trace.capability = Some("kv.read".to_owned()); + trace.stream = Some(format!("kv:{name}")); + } + ControlRequest::KvCreate { name } | ControlRequest::KvGet { name, .. } => { + trace.resource = Some(format!("resource:kv:{name}")); + trace.stream = Some(format!("kv:{name}")); + } + ControlRequest::DbSync { node, name, .. } => { + trace.peer_node = Some(node.clone()); + trace.resource = Some(format!("resource:db:{name}")); + trace.capability = Some("db.sync".to_owned()); + trace.stream = Some(format!("db:{name}")); + } + ControlRequest::DbAdd { name, .. } + | ControlRequest::DbStatus { name } + | ControlRequest::DbChanges { name, .. } => { + trace.resource = Some(format!("resource:db:{name}")); + trace.stream = Some(format!("db:{name}")); + } + ControlRequest::DocumentSync { node, name, .. } => { + trace.peer_node = Some(node.clone()); + trace.resource = Some(format!("resource:document:{name}")); + trace.capability = Some("document.read".to_owned()); + trace.stream = Some(format!("document:{name}")); + } + ControlRequest::DocumentCreate { name } + | ControlRequest::DocumentStatus { name } + | ControlRequest::DocumentSet { name, .. } + | ControlRequest::DocumentGet { name } => { + trace.resource = Some(format!("resource:document:{name}")); + trace.stream = Some(format!("document:{name}")); + } + ControlRequest::PubsubPub { topic, node, .. } => { + trace.peer_node = node.clone(); + trace.resource = Some(format!("resource:pubsub:{topic}")); + trace.capability = Some("pubsub.publish".to_owned()); + trace.stream = Some(format!("pubsub:{topic}")); + } + ControlRequest::PubsubSub { topic, node, .. } => { + trace.peer_node = node.clone(); + trace.resource = Some(format!("resource:pubsub:{topic}")); + trace.capability = Some("pubsub.subscribe".to_owned()); + trace.stream = Some(format!("pubsub:{topic}")); + } + ControlRequest::PipeListen { name, node, .. } => { + trace.peer_node = node.clone(); + trace.resource = Some(format!("resource:pipe:{name}")); + trace.capability = Some("pipe.listen".to_owned()); + trace.stream = Some(format!("pipe:{name}")); + } + ControlRequest::PipeConnect { target, node, .. } + | ControlRequest::PipeSend { target, node, .. } => { + trace.peer_node = node.clone(); + trace.resource = Some(format!("resource:pipe:{target}")); + trace.capability = Some("pipe.connect".to_owned()); + trace.stream = Some(format!("pipe:{target}")); + } + ControlRequest::PipeTcpForward { + node, target_addr, .. + } + | ControlRequest::PipeTcpStream { + node, target_addr, .. + } => { + trace.peer_node = Some(node.clone()); + trace.resource = Some(format!("resource:pipe-tcp:{target_addr}")); + trace.capability = Some("pipe.forward".to_owned()); + trace.stream = Some("pipe-tcp".to_owned()); + } + ControlRequest::PipeUnixForward { + node, target_path, .. + } + | ControlRequest::PipeUnixStream { + node, target_path, .. + } => { + trace.peer_node = Some(node.clone()); + trace.resource = Some(format!("resource:pipe-unix:{}", target_path.display())); + trace.capability = Some("pipe.forward".to_owned()); + trace.stream = Some("pipe-unix".to_owned()); + } + ControlRequest::SshProxyConnect { node, .. } + | ControlRequest::SshProxyStream { node, .. } => { + trace.peer_node = Some(node.clone()); + trace.resource = Some("resource:ssh-proxy:local".to_owned()); + trace.capability = Some("ssh_proxy.connect".to_owned()); + trace.stream = Some("ssh-proxy".to_owned()); + } + ControlRequest::SshAdminShell { node, .. } => { + trace.peer_node = Some(node.clone()); + trace.resource = Some("resource:ssh-proxy:local".to_owned()); + trace.capability = Some("ssh_proxy.admin_shell".to_owned()); + trace.stream = Some("ssh-admin-shell".to_owned()); + } + ControlRequest::OverlaySend { name, node, .. } => { + trace.peer_node = Some(node.clone()); + trace.resource = Some(format!("resource:overlay:{name}")); + trace.capability = Some("overlay.route".to_owned()); + trace.stream = Some(format!("overlay:{name}")); + } + ControlRequest::OverlayUp { name, .. } => { + trace.resource = Some(format!("resource:overlay:{name}")); + trace.capability = Some("overlay.route".to_owned()); + trace.stream = Some(format!("overlay:{name}")); + } + _ => {} + } + trace +} + +fn control_request_command(request: &ControlRequest) -> &'static str { + match request { + ControlRequest::Status => "status", + ControlRequest::NodeId => "node-id", + ControlRequest::PeerCardExport { .. } => "peer.export", + ControlRequest::PeerCardImport { .. } => "peer.import", + ControlRequest::PeerCardList => "peer.list", + ControlRequest::PeerPing { .. } => "peer.ping", + ControlRequest::PeerAuthCheck { .. } => "peer.auth-check", + ControlRequest::ResourceList => "resource.list", + ControlRequest::ResourceCreate { .. } => "resource.create", + ControlRequest::KeychainSync { .. } => "keychain.sync", + ControlRequest::AuthSync { .. } => "auth.sync", + ControlRequest::SyncStatus => "sync.status", + ControlRequest::SyncNow { .. } => "sync.now", + ControlRequest::AuthExplain { .. } => "auth.explain", + ControlRequest::AuthGrant { .. } => "auth.grant", + ControlRequest::AuthRevoke { .. } => "auth.revoke", + ControlRequest::CasFetch { .. } => "cas.fetch", + ControlRequest::CasRootSync { .. } => "cas.root.sync", + ControlRequest::CasAddPrivate { .. } => "cas.add-private", + ControlRequest::CasGetPrivate { .. } => "cas.get-private", + ControlRequest::SshCertSync { .. } => "ssh.cert.sync", + ControlRequest::SshRevocationSync { .. } => "ssh.revocation.sync", + ControlRequest::KvCreate { .. } => "kv.create", + ControlRequest::KvSet { .. } => "kv.set", + ControlRequest::KvGet { .. } => "kv.get", + ControlRequest::KvSync { .. } => "kv.sync", + ControlRequest::DbAdd { .. } => "db.add", + ControlRequest::DbStatus { .. } => "db.status", + ControlRequest::DbChanges { .. } => "db.changes", + ControlRequest::DbSync { .. } => "db.sync", + ControlRequest::DocumentCreate { .. } => "document.create", + ControlRequest::DocumentStatus { .. } => "document.status", + ControlRequest::DocumentSet { .. } => "document.set", + ControlRequest::DocumentGet { .. } => "document.get", + ControlRequest::DocumentSync { .. } => "document.sync", + ControlRequest::PubsubPub { .. } => "pubsub.pub", + ControlRequest::PubsubSub { .. } => "pubsub.sub", + ControlRequest::PipeListen { .. } => "pipe.listen", + ControlRequest::PipeConnect { .. } => "pipe.connect", + ControlRequest::PipeSend { .. } => "pipe.send", + ControlRequest::PipeRecv { .. } => "pipe.recv", + ControlRequest::PipeTcpForward { .. } => "pipe.forward-tcp", + ControlRequest::PipeTcpStream { .. } => "pipe.tcp-stream", + ControlRequest::PipeUnixForward { .. } => "pipe.forward-unix", + ControlRequest::PipeUnixStream { .. } => "pipe.unix-stream", + ControlRequest::SshProxyConnect { .. } => "ssh.proxy", + ControlRequest::SshProxyStream { .. } => "ssh.proxy-stream", + ControlRequest::SshAdminShell { .. } => "ssh.admin-shell", + ControlRequest::OverlayStatus => "overlay.status", + ControlRequest::OverlayPlan { .. } => "overlay.plan", + ControlRequest::OverlayJoin { .. } => "overlay.join", + ControlRequest::OverlayLeave { .. } => "overlay.leave", + ControlRequest::OverlayInterfacePlan { .. } => "overlay.interface-plan", + ControlRequest::OverlayUp { .. } => "overlay.up", + ControlRequest::OverlayDown { .. } => "overlay.down", + ControlRequest::OverlayPeers { .. } => "overlay.peers", + ControlRequest::OverlaySend { .. } => "overlay.send", + ControlRequest::OverlayRecv { .. } => "overlay.recv", + ControlRequest::SecretStatus => "secret.status", + ControlRequest::SecretCreate { .. } => "secret.create", + ControlRequest::SecretRotate { .. } => "secret.rotate", + ControlRequest::SecretBearerCreate { .. } => "secret.bearer.create", + ControlRequest::SecretBearerList => "secret.bearer.list", + ControlRequest::SecretBearerChallenge { .. } => "secret.bearer.challenge", + ControlRequest::SecretBearerProve { .. } => "secret.bearer.prove", + ControlRequest::SecretBearerVerify { .. } => "secret.bearer.verify", + ControlRequest::SecretBearerRevoke { .. } => "secret.bearer.revoke", + ControlRequest::ModuleStub { .. } => "module.stub", + _ => "control.request", + } +} + +pub(crate) fn node_error_code(error: &NodeError) -> &'static str { + match error { + NodeError::Unauthorized(_) => "unauthorized", + NodeError::PeerNotFound(_) => "peer_not_found", + NodeError::IrohEndpointUnavailable => "iroh_endpoint_unavailable", + NodeError::KvNotFound(_) => "kv_not_found", + NodeError::DbNotFound(_) => "db_not_found", + NodeError::DocumentNotFound(_) => "document_not_found", + NodeError::Cas(_) => "cas_error", + NodeError::Store(_) => "store_error", + NodeError::Config(_) => "config_error", + NodeError::Codec(_) => "codec_error", + _ => "node_error", + } +} diff --git a/docs/architecture.md b/docs/architecture.md index 7ec2633..99a18df 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -9,8 +9,10 @@ requests. Within `geth-node`, daemon lifecycle code is separated from feature handlers: `daemon.rs` owns `geth daemon run` startup, local socket binding, shutdown signal handling, Iroh endpoint startup, the Iroh accept loop, and background -live-sync task spawning. Runtime registries for pubsub, pipes, and overlays live -behind narrow mutex-protected structs in `runtime.rs`. Local control and +live-sync task spawning. `local_control.rs` owns async local `ControlRequest` +routing and safe trace-field classification before delegating to feature +handlers. Runtime registries for pubsub, pipes, and overlays live behind narrow +mutex-protected structs in `runtime.rs`. Command-family handler modules and protected peer-control feature dispatch remain separate refactor targets. The local metadata store is SQLite product state. `geth-store` tracks a numeric diff --git a/docs/production-readiness-roadmap.md b/docs/production-readiness-roadmap.md index 0e767e4..0031608 100644 --- a/docs/production-readiness-roadmap.md +++ b/docs/production-readiness-roadmap.md @@ -45,12 +45,12 @@ behavior. ownership and locking rules. - `[x]` Existing daemon startup and status tests pass unchanged. -- `[ ]` Extract local control routing. +- `[~]` Extract local control routing. Acceptance criteria: - - `[ ]` Local `ControlRequest` dispatch is a routing layer, not the home of + - `[x]` Local `ControlRequest` dispatch is a routing layer, not the home of every feature implementation. - `[ ]` Each command family has a small handler module or function group. - - `[ ]` Local-only behavior remains covered by existing integration tests. + - `[x]` Local-only behavior remains covered by existing integration tests. - `[ ]` Extract protected peer-control routing. Acceptance criteria: