refactor: extract local control routing
This commit is contained in:
parent
d5a91c2396
commit
f88bc851a4
4 changed files with 501 additions and 494 deletions
|
|
@ -1,6 +1,7 @@
|
||||||
pub mod backup;
|
pub mod backup;
|
||||||
mod daemon;
|
mod daemon;
|
||||||
pub mod doctor;
|
pub mod doctor;
|
||||||
|
mod local_control;
|
||||||
mod peer_client;
|
mod peer_client;
|
||||||
mod runtime;
|
mod runtime;
|
||||||
pub mod service;
|
pub mod service;
|
||||||
|
|
@ -64,6 +65,7 @@ use geth_types::{
|
||||||
};
|
};
|
||||||
use iroh::protocol::ProtocolHandler;
|
use iroh::protocol::ProtocolHandler;
|
||||||
use iroh_docs::api::protocol::{AddrInfoOptions, ShareMode};
|
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 peer_client::{request_overlay_wire, request_peer_control, request_pipe_wire};
|
||||||
use runtime::{
|
use runtime::{
|
||||||
NodeRuntime, OverlayTunCounters, OverlayTunRuntime, PipeRuntime, PubsubGossipTopicRuntime,
|
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<ControlResponse, NodeError> {
|
|
||||||
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<String>,
|
|
||||||
resource: Option<String>,
|
|
||||||
capability: Option<String>,
|
|
||||||
stream: Option<String>,
|
|
||||||
}
|
|
||||||
|
|
||||||
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> {
|
async fn handle_stream(node: LocalNode, stream: UnixStream) -> Result<(), NodeError> {
|
||||||
let mut reader = BufReader::new(stream);
|
let mut reader = BufReader::new(stream);
|
||||||
let mut line = String::new();
|
let mut line = String::new();
|
||||||
|
|
@ -1161,7 +681,7 @@ async fn handle_stream(node: LocalNode, stream: UnixStream) -> Result<(), NodeEr
|
||||||
resource = "resource:ssh-proxy:local",
|
resource = "resource:ssh-proxy:local",
|
||||||
capability = "ssh_proxy.connect",
|
capability = "ssh_proxy.connect",
|
||||||
stream = "ssh-proxy",
|
stream = "ssh-proxy",
|
||||||
error_code = node_error_code(error),
|
error_code = local_control::node_error_code(error),
|
||||||
%error,
|
%error,
|
||||||
"streaming control request failed"
|
"streaming control request failed"
|
||||||
),
|
),
|
||||||
|
|
@ -1202,7 +722,7 @@ async fn handle_stream(node: LocalNode, stream: UnixStream) -> Result<(), NodeEr
|
||||||
resource = %resource,
|
resource = %resource,
|
||||||
capability = "pipe.forward",
|
capability = "pipe.forward",
|
||||||
stream = "pipe-tcp",
|
stream = "pipe-tcp",
|
||||||
error_code = node_error_code(error),
|
error_code = local_control::node_error_code(error),
|
||||||
%error,
|
%error,
|
||||||
"streaming control request failed"
|
"streaming control request failed"
|
||||||
),
|
),
|
||||||
|
|
@ -1243,7 +763,7 @@ async fn handle_stream(node: LocalNode, stream: UnixStream) -> Result<(), NodeEr
|
||||||
resource = %resource,
|
resource = %resource,
|
||||||
capability = "pipe.forward",
|
capability = "pipe.forward",
|
||||||
stream = "pipe-unix",
|
stream = "pipe-unix",
|
||||||
error_code = node_error_code(error),
|
error_code = local_control::node_error_code(error),
|
||||||
%error,
|
%error,
|
||||||
"streaming control request failed"
|
"streaming control request failed"
|
||||||
),
|
),
|
||||||
|
|
@ -11889,7 +11409,7 @@ mod tests {
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn control_trace_fields_exclude_bearer_secrets() {
|
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(),
|
node: "node:peer".to_owned(),
|
||||||
name: "prefs".to_owned(),
|
name: "prefs".to_owned(),
|
||||||
bearer_secret: Some("private-token-value".to_owned()),
|
bearer_secret: Some("private-token-value".to_owned()),
|
||||||
|
|
@ -11905,15 +11425,15 @@ mod tests {
|
||||||
#[test]
|
#[test]
|
||||||
fn node_error_codes_are_stable_for_common_tracing_failures() {
|
fn node_error_codes_are_stable_for_common_tracing_failures() {
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
node_error_code(&NodeError::Unauthorized("missing grant".to_owned())),
|
local_control::node_error_code(&NodeError::Unauthorized("missing grant".to_owned())),
|
||||||
"unauthorized"
|
"unauthorized"
|
||||||
);
|
);
|
||||||
assert_eq!(
|
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"
|
"peer_not_found"
|
||||||
);
|
);
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
node_error_code(&NodeError::IrohEndpointUnavailable),
|
local_control::node_error_code(&NodeError::IrohEndpointUnavailable),
|
||||||
"iroh_endpoint_unavailable"
|
"iroh_endpoint_unavailable"
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
|
||||||
485
crates/geth-node/src/local_control.rs
Normal file
485
crates/geth-node/src/local_control.rs
Normal file
|
|
@ -0,0 +1,485 @@
|
||||||
|
//! Local control request routing and trace classification.
|
||||||
|
|
||||||
|
use super::*;
|
||||||
|
|
||||||
|
pub async fn handle_request_async(
|
||||||
|
node: &LocalNode,
|
||||||
|
request: ControlRequest,
|
||||||
|
) -> Result<ControlResponse, NodeError> {
|
||||||
|
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<String>,
|
||||||
|
pub(crate) resource: Option<String>,
|
||||||
|
pub(crate) capability: Option<String>,
|
||||||
|
pub(crate) stream: Option<String>,
|
||||||
|
}
|
||||||
|
|
||||||
|
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",
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -9,8 +9,10 @@ requests.
|
||||||
Within `geth-node`, daemon lifecycle code is separated from feature handlers:
|
Within `geth-node`, daemon lifecycle code is separated from feature handlers:
|
||||||
`daemon.rs` owns `geth daemon run` startup, local socket binding, shutdown
|
`daemon.rs` owns `geth daemon run` startup, local socket binding, shutdown
|
||||||
signal handling, Iroh endpoint startup, the Iroh accept loop, and background
|
signal handling, Iroh endpoint startup, the Iroh accept loop, and background
|
||||||
live-sync task spawning. Runtime registries for pubsub, pipes, and overlays live
|
live-sync task spawning. `local_control.rs` owns async local `ControlRequest`
|
||||||
behind narrow mutex-protected structs in `runtime.rs`. Local control and
|
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.
|
protected peer-control feature dispatch remain separate refactor targets.
|
||||||
|
|
||||||
The local metadata store is SQLite product state. `geth-store` tracks a numeric
|
The local metadata store is SQLite product state. `geth-store` tracks a numeric
|
||||||
|
|
|
||||||
|
|
@ -45,12 +45,12 @@ behavior.
|
||||||
ownership and locking rules.
|
ownership and locking rules.
|
||||||
- `[x]` Existing daemon startup and status tests pass unchanged.
|
- `[x]` Existing daemon startup and status tests pass unchanged.
|
||||||
|
|
||||||
- `[ ]` Extract local control routing.
|
- `[~]` Extract local control routing.
|
||||||
Acceptance criteria:
|
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.
|
every feature implementation.
|
||||||
- `[ ]` Each command family has a small handler module or function group.
|
- `[ ]` 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.
|
- `[ ]` Extract protected peer-control routing.
|
||||||
Acceptance criteria:
|
Acceptance criteria:
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue