pub mod service; use geth_auth::AuthExplanation; use geth_cas::{LocalCas, hash_path}; use geth_config::{GethConfig, GethPaths, RelayMode}; use geth_control::{ CasBlob, ControlRequest, ControlResponse, KeychainStatusResponse, NodeIdResponse, StatusResponse, }; use geth_crypto::AgentKey; use geth_iroh::{EndpointStatus, GethIrohConfig, GethIrohEndpoint, GethRelayMode}; use geth_resource::ResourceDescriptor; use geth_ssh_identity::{ SshCertApproval, SshCertKind, SshCertRequest, SshCertRequestStatus, SshCertificateRecord, SshRevocationEntry, SshRevocationKind, build_ssh_cert_sign_command, cert_request_id, certificate_id, revocation_id, ssh_public_key_fingerprint, }; use geth_store::{ Store, StoredResource, StoredSshCertRequest, StoredSshCertificate, StoredSshRevocation, }; use geth_types::{ NodeId, ResourceId, ResourceKind, ResourceName, SshCertId, SshCertRequestId, UnixMillis, }; use std::path::Path; use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; use tokio::net::{UnixListener, UnixStream}; #[derive(Debug, thiserror::Error)] pub enum NodeError { #[error("config error: {0}")] Config(#[from] geth_config::ConfigError), #[error("crypto error: {0}")] Crypto(#[from] geth_crypto::CryptoError), #[error("store error: {0}")] Store(#[from] geth_store::StoreError), #[error("cas error: {0}")] Cas(#[from] geth_cas::CasError), #[error("control error: {0}")] Control(#[from] geth_control::ControlError), #[error("json error: {0}")] Json(#[from] serde_json::Error), #[error("io error: {0}")] Io(#[from] std::io::Error), #[error("invalid resource kind: {0}")] InvalidResourceKind(String), #[error("invalid ssh certificate kind: {0}")] InvalidSshCertKind(String), #[error("invalid ssh certificate request status: {0}")] InvalidSshCertStatus(String), #[error("invalid ssh revocation kind: {0}")] InvalidSshRevocationKind(String), #[error("ssh certificate request not found: {0}")] SshCertRequestNotFound(String), #[error("ssh certificate request must include at least one principal")] MissingSshCertPrincipal, #[error("ssh certificate flow error: {0}")] SshCertFlow(#[from] geth_ssh_identity::SshCertFlowError), } #[derive(Clone, Debug)] pub struct LocalNode { pub paths: GethPaths, pub agent_id: String, pub node_id: String, pub iroh_status: EndpointStatus, } pub fn init_node(paths: &GethPaths) -> Result { paths.ensure_base_dirs()?; if !paths.config_file().exists() { std::fs::write(paths.config_file(), GethConfig::default_toml())?; } let key = AgentKey::load_or_create(&paths.agent_key())?; let agent_id = key.agent_id().to_string(); let node_id = stable_node_id(&agent_id); let store = Store::open(&paths.metadata_db())?; store.upsert_agent(&agent_id, &key.public_key_hex())?; store.upsert_node(&node_id, "local", &agent_id)?; store.insert_resource(&StoredResource { resource_id: "resource:cas:local".to_owned(), kind: ResourceKind::Cas.to_string(), name: "local-cas".to_owned(), status: "active".to_owned(), })?; Ok(LocalNode { paths: paths.clone(), agent_id, node_id, iroh_status: EndpointStatus::scaffolded(), }) } pub fn open_node(paths: &GethPaths) -> Result { init_node(paths) } pub async fn run_daemon(paths: GethPaths) -> Result<(), NodeError> { let mut node = init_node(&paths)?; let _iroh_endpoint = start_daemon_iroh_endpoint(&mut node).await?; if Path::new(&paths.socket_path()).exists() { std::fs::remove_file(paths.socket_path())?; } let listener = UnixListener::bind(paths.socket_path())?; tracing::info!(socket = %paths.socket_path().display(), "geth daemon listening"); loop { let (stream, _) = listener.accept().await?; let node = node.clone(); tokio::spawn(async move { if let Err(error) = handle_stream(node, stream).await { tracing::warn!(%error, "control request failed"); } }); } } pub async fn send_control( paths: &GethPaths, request: ControlRequest, ) -> Result { let mut stream = UnixStream::connect(paths.socket_path()).await?; stream .write_all(geth_control::encode_request(&request)?.as_bytes()) .await?; stream.shutdown().await?; let mut reader = BufReader::new(stream); let mut line = String::new(); reader.read_line(&mut line).await?; Ok(geth_control::decode_response(&line)?) } async fn handle_stream(node: LocalNode, stream: UnixStream) -> Result<(), NodeError> { let mut reader = BufReader::new(stream); let mut line = String::new(); reader.read_line(&mut line).await?; let request = geth_control::decode_request(&line)?; let response = match handle_request(&node, request) { Ok(response) => response, Err(error) => ControlResponse::Error { message: error.to_string(), }, }; let mut stream = reader.into_inner(); stream .write_all(geth_control::encode_response(&response)?.as_bytes()) .await?; Ok(()) } pub fn handle_request( node: &LocalNode, request: ControlRequest, ) -> Result { let store = Store::open(&node.paths.metadata_db())?; match request { ControlRequest::Status => Ok(ControlResponse::Status(StatusResponse { home: node.paths.home().to_path_buf(), socket: node.paths.socket_path(), agent_id: node.agent_id.clone(), node_id: node.node_id.clone(), iroh_enabled: node.iroh_status.enabled, endpoint_id: node.iroh_status.endpoint_id.clone(), iroh_relay_mode: node.iroh_status.relay_mode.clone(), iroh: node.iroh_status.note.clone(), })), ControlRequest::NodeId => Ok(ControlResponse::NodeId(NodeIdResponse { agent_id: node.agent_id.clone(), node_id: node.node_id.clone(), endpoint_id: node.iroh_status.endpoint_id.clone(), })), ControlRequest::ResourceList => Ok(ControlResponse::ResourceList { resources: store .list_resources()? .into_iter() .map(stored_resource_to_descriptor) .collect::, _>>()?, }), ControlRequest::ResourceCreate { kind, name } => { let kind = kind .parse::() .map_err(|_| NodeError::InvalidResourceKind(kind.clone()))?; let id = format!("resource:{}:{}", kind, name); let stored = StoredResource { resource_id: id, kind: kind.to_string(), name, status: "active".to_owned(), }; store.insert_resource(&stored)?; Ok(ControlResponse::ResourceCreated { resource: stored_resource_to_descriptor(stored)?, }) } ControlRequest::CasAdd { path } => { let cas = LocalCas::new(node.paths.cas_dir()); let info = cas.add_path(&path)?; store.record_cas_object( info.hash.as_str(), info.size_bytes, &info.path.to_string_lossy(), )?; Ok(ControlResponse::CasAdded { hash: info.hash, size_bytes: info.size_bytes, }) } ControlRequest::CasGet { hash, out } => { let cas = LocalCas::new(node.paths.cas_dir()); let size_bytes = cas.get_to_path(&hash, &out)?; Ok(ControlResponse::CasGot { hash, out, size_bytes, }) } ControlRequest::CasHash { path } => Ok(ControlResponse::CasHash { hash: hash_path(&path)?, }), ControlRequest::CasHas { hash } => { let cas = LocalCas::new(node.paths.cas_dir()); let present = cas.has(&hash)?; Ok(ControlResponse::CasHas { hash, present }) } ControlRequest::CasList => { let cas = LocalCas::new(node.paths.cas_dir()); Ok(ControlResponse::CasList { blobs: cas .list()? .into_iter() .map(|blob| CasBlob { hash: blob.hash, size_bytes: blob.size_bytes, }) .collect(), }) } ControlRequest::KeychainStatus => { Ok(ControlResponse::KeychainStatus(KeychainStatusResponse { initialized: false, admin_keys: 0, users: 0, devices: 0, nodes: 1, })) } ControlRequest::AuthExplain { subject, resource, capability, } => { if store.get_peer_card(&subject)?.is_some() { Ok(ControlResponse::AuthExplain( AuthExplanation::discovered_candidate(subject, resource, capability), )) } else { Ok(ControlResponse::AuthExplain(AuthExplanation::stub( subject, resource, capability, ))) } } ControlRequest::SshCertRequest { public_key_path, cert_kind, principals, requested_validity, renewal_of, reason, } => { if principals.is_empty() { return Err(NodeError::MissingSshCertPrincipal); } let cert_kind = cert_kind .parse::() .map_err(|_| NodeError::InvalidSshCertKind(cert_kind.clone()))?; let public_key = std::fs::read_to_string(&public_key_path)?; let created_at = UnixMillis(geth_store::now_ms()); let request = SshCertRequest { id: cert_request_id( &NodeId::new(node.node_id.clone()), &public_key, &principals, created_at, ), requester_node: NodeId::new(node.node_id.clone()), public_key_fingerprint: ssh_public_key_fingerprint(&public_key), public_key, cert_kind, principals, requested_validity, renewal_of: renewal_of.map(SshCertId::new), reason, status: SshCertRequestStatus::Pending, created_at, }; store.insert_ssh_cert_request(&stored_from_ssh_cert_request(&request))?; Ok(ControlResponse::SshCertRequested { request }) } ControlRequest::SshCertRequests => Ok(ControlResponse::SshCertRequests { requests: store .list_ssh_cert_requests()? .into_iter() .map(ssh_cert_request_from_stored) .collect::, _>>()?, }), ControlRequest::SshCertApprove { request_id, ca_key_path, valid_for, serial, out, } => { let stored = store .get_ssh_cert_request(&request_id)? .ok_or_else(|| NodeError::SshCertRequestNotFound(request_id.clone()))?; let mut request = ssh_cert_request_from_stored(stored)?; request.status = SshCertRequestStatus::Approved; store.update_ssh_cert_request_status(request.id.as_str(), request.status.as_str())?; let public_key_path = out.clone().unwrap_or_else(|| { node.paths .home() .join("ssh-cert-requests") .join(format!("{}.pub", request.id)) }); if let Some(parent) = public_key_path.parent() { std::fs::create_dir_all(parent)?; } std::fs::write(&public_key_path, &request.public_key)?; let valid_for = valid_for .or_else(|| request.requested_validity.clone()) .unwrap_or_else(|| "+52w".to_owned()); let signing_command = build_ssh_cert_sign_command( &request, &ca_key_path, &public_key_path, &valid_for, serial, )?; let approval = SshCertApproval { request_id: request.id, approved_by_node: NodeId::new(node.node_id.clone()), ca_key_path: ca_key_path.display().to_string(), key_id: request_id, valid_for, serial, output_path: Some(expected_openssh_cert_path(&public_key_path)), signing_command, note: "request approved; run the signing command on the CA/YubiKey machine, then import the resulting -cert.pub file".to_owned(), }; Ok(ControlResponse::SshCertApproved { approval }) } ControlRequest::SshCertImport { request_id, cert_path, } => { let certificate = std::fs::read_to_string(&cert_path)?; let record = SshCertificateRecord { id: certificate_id(&certificate), request_id: SshCertRequestId::new(request_id.clone()), certificate_fingerprint: ssh_public_key_fingerprint(&certificate), certificate, imported_at: UnixMillis(geth_store::now_ms()), }; store.insert_ssh_certificate(&stored_from_ssh_certificate(&record))?; store.update_ssh_cert_request_status( &request_id, SshCertRequestStatus::Signed.as_str(), )?; Ok(ControlResponse::SshCertImported { certificate: record, }) } ControlRequest::SshCertList => Ok(ControlResponse::SshCertList { requests: store .list_ssh_cert_requests()? .into_iter() .map(ssh_cert_request_from_stored) .collect::, _>>()?, certificates: store .list_ssh_certificates()? .into_iter() .map(ssh_certificate_from_stored) .collect(), }), ControlRequest::SshRevocationAdd { kind, target, reason, } => { let kind = kind .parse::() .map_err(|_| NodeError::InvalidSshRevocationKind(kind.clone()))?; let created_at = UnixMillis(geth_store::now_ms()); let revocation = SshRevocationEntry { id: revocation_id(&kind, &target, created_at), kind, target, reason, created_at, published: true, }; store.insert_ssh_revocation(&stored_from_ssh_revocation(&revocation))?; Ok(ControlResponse::SshRevocationAdded { revocation }) } ControlRequest::SshRevocationList => Ok(ControlResponse::SshRevocationList { revocations: store .list_ssh_revocations()? .into_iter() .map(ssh_revocation_from_stored) .collect::, _>>()?, }), ControlRequest::SshRevocationExport { out } => { let revocations = store .list_ssh_revocations()? .into_iter() .map(ssh_revocation_from_stored) .collect::, _>>()?; if let Some(parent) = out.parent() { std::fs::create_dir_all(parent)?; } let mut body = String::new(); for revocation in &revocations { body.push_str(&serde_json::to_string(revocation)?); body.push('\n'); } std::fs::write(&out, body)?; Ok(ControlResponse::SshRevocationExported { out, count: revocations.len(), }) } ControlRequest::ModuleStub { module, command } => { Ok(ControlResponse::NotImplemented { module, command }) } } } fn stored_resource_to_descriptor(stored: StoredResource) -> Result { let kind = stored .kind .parse::() .map_err(|_| NodeError::InvalidResourceKind(stored.kind.clone()))?; Ok(ResourceDescriptor::local( ResourceId::new(stored.resource_id), kind, ResourceName::new(stored.name), )) } fn stable_node_id(agent_id: &str) -> String { format!("node:{agent_id}") } async fn start_daemon_iroh_endpoint( node: &mut LocalNode, ) -> Result, NodeError> { let node_config = GethConfig::load(&node.paths.config_file())?; let relay_mode = node_config.iroh.relay_mode; let iroh_relay_mode = config_relay_mode_to_iroh(relay_mode); let config = GethIrohConfig::local_with_relay(node.paths.iroh_key(), iroh_relay_mode); match geth_iroh::start_endpoint(&config).await { Ok(endpoint) => { let status = endpoint.status(); if let Some(endpoint_id) = &status.endpoint_id { let store = Store::open(&node.paths.metadata_db())?; store.upsert_node_endpoint(endpoint_id, &node.node_id, &node.agent_id, "iroh")?; } node.iroh_status = status; Ok(Some(endpoint)) } Err(error) => { node.iroh_status = EndpointStatus { enabled: false, endpoint_id: None, relay_mode: relay_mode.as_str().to_owned(), note: format!("Iroh endpoint failed to start: {error}"), }; Ok(None) } } } fn config_relay_mode_to_iroh(mode: RelayMode) -> GethRelayMode { match mode { RelayMode::Disabled => GethRelayMode::Disabled, RelayMode::Default => GethRelayMode::Default, RelayMode::Staging => GethRelayMode::Staging, } } fn expected_openssh_cert_path(public_key_path: &Path) -> String { let text = public_key_path.display().to_string(); if let Some(prefix) = text.strip_suffix(".pub") { format!("{prefix}-cert.pub") } else { format!("{text}-cert.pub") } } fn stored_from_ssh_cert_request(request: &SshCertRequest) -> StoredSshCertRequest { StoredSshCertRequest { request_id: request.id.to_string(), requester_node: request.requester_node.to_string(), public_key: request.public_key.clone(), public_key_fingerprint: request.public_key_fingerprint.clone(), cert_kind: request.cert_kind.to_string(), principals: request.principals.clone(), requested_validity: request.requested_validity.clone(), renewal_of: request.renewal_of.as_ref().map(ToString::to_string), reason: request.reason.clone(), status: request.status.to_string(), created_at_ms: request.created_at.0, } } fn ssh_cert_request_from_stored(stored: StoredSshCertRequest) -> Result { let cert_kind = stored .cert_kind .parse::() .map_err(|_| NodeError::InvalidSshCertKind(stored.cert_kind.clone()))?; let status = stored .status .parse::() .map_err(|_| NodeError::InvalidSshCertStatus(stored.status.clone()))?; Ok(SshCertRequest { id: SshCertRequestId::new(stored.request_id), requester_node: NodeId::new(stored.requester_node), public_key: stored.public_key, public_key_fingerprint: stored.public_key_fingerprint, cert_kind, principals: stored.principals, requested_validity: stored.requested_validity, renewal_of: stored.renewal_of.map(SshCertId::new), reason: stored.reason, status, created_at: UnixMillis(stored.created_at_ms), }) } fn stored_from_ssh_certificate(certificate: &SshCertificateRecord) -> StoredSshCertificate { StoredSshCertificate { cert_id: certificate.id.to_string(), request_id: certificate.request_id.to_string(), certificate: certificate.certificate.clone(), certificate_fingerprint: certificate.certificate_fingerprint.clone(), imported_at_ms: certificate.imported_at.0, } } fn ssh_certificate_from_stored(stored: StoredSshCertificate) -> SshCertificateRecord { SshCertificateRecord { id: SshCertId::new(stored.cert_id), request_id: SshCertRequestId::new(stored.request_id), certificate: stored.certificate, certificate_fingerprint: stored.certificate_fingerprint, imported_at: UnixMillis(stored.imported_at_ms), } } fn stored_from_ssh_revocation(revocation: &SshRevocationEntry) -> StoredSshRevocation { StoredSshRevocation { revocation_id: revocation.id.to_string(), kind: revocation.kind.to_string(), target: revocation.target.clone(), reason: revocation.reason.clone(), created_at_ms: revocation.created_at.0, published: revocation.published, } } fn ssh_revocation_from_stored( stored: StoredSshRevocation, ) -> Result { let kind = stored .kind .parse::() .map_err(|_| NodeError::InvalidSshRevocationKind(stored.kind.clone()))?; Ok(SshRevocationEntry { id: geth_types::SshRevocationId::new(stored.revocation_id), kind, target: stored.target, reason: stored.reason, created_at: UnixMillis(stored.created_at_ms), published: stored.published, }) }