pub mod service; use geth_auth::{AuthExplanation, AuthOp, AuthOpKind}; use geth_cas::{ BlobInfoSummary, FileConflict, FileConflictKind, FileConflictResolution, FileConflictStatus, FileRoot, FileRootScan, LocalCas, hash_path, }; use geth_config::{GethConfig, GethPaths, RelayMode}; use geth_control::{ CasBlob, ControlRequest, ControlResponse, KeychainStatusResponse, NodeIdResponse, StatusResponse, }; use geth_crypto::AgentKey; use geth_db::DbResource; use geth_document::{DocumentResource, DocumentState}; use geth_iroh::{EndpointStatus, GethIrohConfig, GethIrohEndpoint, GethRelayMode}; use geth_keychain::{KeychainOp, KeychainOpKind}; use geth_kv::{KvEntry, KvResource}; use geth_pipe::{PipeConnection, PipeListener}; use geth_pubsub::PubsubMessage; use geth_resource::ResourceDescriptor; use geth_secrets::{BearerAccess, ResourceMasterSecret}; use geth_ssh_identity::{ SshCertApproval, SshCertKind, SshCertRequest, SshCertRequestStatus, SshCertificateRecord, SshRevocationEntry, SshRevocationExportFormat, SshRevocationKind, build_ssh_cert_sign_command, cert_request_id, certificate_id, openssh_krl_spec, revocation_id, ssh_public_key_fingerprint, }; use geth_store::{ Store, StoredAuthOp, StoredDbResource, StoredDocumentResource, StoredFileConflict, StoredFileRoot, StoredKeychainOp, StoredKvEntry, StoredKvStore, StoredResource, StoredResourceSecret, StoredSshCertRequest, StoredSshCertificate, StoredSshRevocation, }; use geth_types::{ AuthOpId, Capability, KeyId, NodeId, PrincipalId, ResourceId, ResourceKind, ResourceName, SshCertId, SshCertRequestId, UnixMillis, }; use std::collections::{BTreeMap, VecDeque}; use std::path::Path; use std::sync::{Arc, Mutex}; 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("db error: {0}")] Db(#[from] geth_db::DbError), #[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 db resource name: {0}")] InvalidDbName(String), #[error("db path does not exist or is not a file: {0}")] InvalidDbPath(String), #[error("db resource not found: {0}")] DbNotFound(String), #[error("invalid kv store name: {0}")] InvalidKvName(String), #[error("invalid kv key: {0}")] InvalidKvKey(String), #[error("kv store not found: {0}")] KvNotFound(String), #[error("invalid document name: {0}")] InvalidDocumentName(String), #[error("document not found: {0}")] DocumentNotFound(String), #[error("document error: {0}")] Document(#[from] geth_document::DocumentError), #[error("resource not found: {0}")] ResourceNotFound(String), #[error("secrets error: {0}")] Secrets(#[from] geth_secrets::SecretsError), #[error("pubsub error: {0}")] Pubsub(#[from] geth_pubsub::PubsubError), #[error("pipe error: {0}")] Pipe(#[from] geth_pipe::PipeError), #[error("runtime state lock poisoned")] RuntimeLockPoisoned, #[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), #[error("ssh identity error: {0}")] SshIdentity(#[from] geth_ssh_identity::SshIdentityError), } #[derive(Clone, Debug)] pub struct LocalNode { pub paths: GethPaths, pub agent_id: String, pub node_id: String, pub iroh_status: EndpointStatus, runtime: Arc, } #[derive(Debug)] struct NodeRuntime { pubsub: Mutex, pipes: Mutex, } #[derive(Debug, Default)] struct PubsubRuntime { messages: VecDeque, } #[derive(Debug, Default)] struct PipeRuntime { listeners: BTreeMap, connections: VecDeque, } const PUBSUB_RING_LIMIT: usize = 256; const PIPE_CONNECTION_RING_LIMIT: usize = 256; 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(), runtime: Arc::new(NodeRuntime { pubsub: Mutex::new(PubsubRuntime::default()), pipes: Mutex::new(PipeRuntime::default()), }), }) } 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_local_discovery: node.iroh_status.local_discovery, 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::CasPin { hash } => { let cas = LocalCas::new(node.paths.cas_dir()); if !cas.has(&hash)? { return Err(NodeError::Cas(geth_cas::CasError::NotFound( hash.to_string(), ))); } store.pin_cas_object(hash.as_str())?; Ok(ControlResponse::CasPinned { hash, pinned: true }) } ControlRequest::CasUnpin { hash } => { geth_cas::validate_hash(&hash)?; store.unpin_cas_object(hash.as_str())?; Ok(ControlResponse::CasPinned { hash, pinned: false, }) } ControlRequest::CasCleanup { dry_run } => { let cas = LocalCas::new(node.paths.cas_dir()); let mut removed = Vec::new(); let mut retained_pinned = Vec::new(); for blob in cas.list()? { if store.is_cas_object_pinned(blob.hash.as_str())? { retained_pinned.push(blob.hash); continue; } if !dry_run && cas.remove(&blob.hash)? { store.delete_cas_object(blob.hash.as_str())?; } removed.push(blob.hash); } Ok(ControlResponse::CasCleanup { removed, retained_pinned, dry_run, }) } ControlRequest::CasList => { let cas = LocalCas::new(node.paths.cas_dir()); let blobs = cas .list()? .into_iter() .map(|blob| { Ok(CasBlob { pinned: store.is_cas_object_pinned(blob.hash.as_str())?, hash: blob.hash, size_bytes: blob.size_bytes, }) }) .collect::, NodeError>>()?; Ok(ControlResponse::CasList { blobs }) } ControlRequest::CasRootAdd { name, path } => { geth_cas::validate_file_root_name(&name)?; if !path.is_dir() { return Err(NodeError::Cas(geth_cas::CasError::TreeRootNotDirectory( path.display().to_string(), ))); } let path = std::fs::canonicalize(path)?; let resource_id = format!("resource:cas-tree:{name}"); store.insert_resource(&StoredResource { resource_id: resource_id.clone(), kind: ResourceKind::Cas.to_string(), name: format!("file-root:{name}"), status: "active".to_owned(), })?; let stored = StoredFileRoot { root_id: format!("file-root:{name}"), resource_id, name, path: path.display().to_string(), latest_tree_hash: None, latest_tree_json: None, updated_at_ms: geth_store::now_ms(), }; store.upsert_file_root(&stored)?; Ok(ControlResponse::CasRootAdded { root: file_root_from_stored(&stored), }) } ControlRequest::CasRootList => Ok(ControlResponse::CasRootList { roots: store .list_file_roots()? .iter() .map(file_root_from_stored) .collect(), }), ControlRequest::CasRootScan { name } => { geth_cas::validate_file_root_name(&name)?; let mut stored = store .get_file_root_by_name(&name)? .ok_or_else(|| NodeError::ResourceNotFound(format!("file-root:{name}")))?; let previous_tree = stored .latest_tree_json .as_deref() .map(serde_json::from_str) .transpose()?; let cas = LocalCas::new(node.paths.cas_dir()); let scanned = cas.add_tree_path(Path::new(&stored.path))?; store.record_cas_object( scanned.object.hash.as_str(), scanned.object.size_bytes, &scanned.object.path.to_string_lossy(), )?; let changes = geth_cas::diff_tree_objects(previous_tree.as_ref(), &scanned.tree); stored.latest_tree_hash = Some(scanned.object.hash.to_string()); stored.latest_tree_json = Some(serde_json::to_string(&scanned.tree)?); stored.updated_at_ms = geth_store::now_ms(); store.upsert_file_root(&stored)?; Ok(ControlResponse::CasRootScanned { scan: FileRootScan { root: file_root_from_stored(&stored), tree: BlobInfoSummary { hash: scanned.object.hash, size_bytes: scanned.object.size_bytes, }, changes, note: geth_cas::file_root_scan_note().to_owned(), }, }) } ControlRequest::CasConflictRecord { root, path, kind, detail, base_tree, local_tree, remote_tree, } => { geth_cas::validate_file_root_name(&root)?; if let Some(hash) = &base_tree { geth_cas::validate_hash(hash)?; } if let Some(hash) = &local_tree { geth_cas::validate_hash(hash)?; } if let Some(hash) = &remote_tree { geth_cas::validate_hash(hash)?; } let root_record = store .get_file_root_by_name(&root)? .ok_or_else(|| NodeError::ResourceNotFound(format!("file-root:{root}")))?; let kind = FileConflictKind::parse(&kind)?; let created_at = geth_store::now_ms(); let conflict_id = generated_file_conflict_id(&root, &path, kind.as_str(), created_at); let conflict = StoredFileConflict { conflict_id, root_name: root, resource_id: root_record.resource_id, path, kind: kind.as_str().to_owned(), status: FileConflictStatus::Open.as_str().to_owned(), base_tree_hash: base_tree.map(|hash| hash.to_string()), local_tree_hash: local_tree.map(|hash| hash.to_string()), remote_tree_hash: remote_tree.map(|hash| hash.to_string()), detail, resolution: None, resolution_note: None, created_at_ms: created_at, resolved_at_ms: None, }; store.upsert_file_conflict(&conflict)?; Ok(ControlResponse::CasConflictRecorded { conflict: file_conflict_from_stored(conflict)?, }) } ControlRequest::CasConflictList { root } => { if let Some(root) = root.as_deref() { geth_cas::validate_file_root_name(root)?; } let conflicts = store .list_file_conflicts(root.as_deref())? .into_iter() .map(file_conflict_from_stored) .collect::, _>>()?; Ok(ControlResponse::CasConflictList { conflicts }) } ControlRequest::CasConflictResolve { conflict_id, resolution, note, } => { let resolution = FileConflictResolution::parse(&resolution)?; let mut conflict = store .get_file_conflict(&conflict_id)? .ok_or_else(|| NodeError::ResourceNotFound(conflict_id.clone()))?; conflict.status = FileConflictStatus::Resolved.as_str().to_owned(); conflict.resolution = Some(resolution.as_str().to_owned()); conflict.resolution_note = note; conflict.resolved_at_ms = Some(geth_store::now_ms()); store.upsert_file_conflict(&conflict)?; Ok(ControlResponse::CasConflictResolved { conflict: file_conflict_from_stored(conflict)?, }) } ControlRequest::KeychainInit { admin_key_path } => { let mut ops = Vec::new(); let created_at = UnixMillis(geth_store::now_ms()); let init = KeychainOp { id: generated_keychain_op_id("keychain-init", "local", created_at), created_at, kind: KeychainOpKind::KeychainInit, }; store_keychain_op(&store, &init)?; ops.push(init); if let Some(admin_key_path) = admin_key_path { let public_key = std::fs::read_to_string(admin_key_path)?; let created_at = UnixMillis(geth_store::now_ms()); let admin_key = KeyId::new(ssh_public_key_fingerprint(&public_key)); let op = KeychainOp { id: generated_keychain_op_id("admin-key-add", admin_key.as_str(), created_at), created_at, kind: KeychainOpKind::AdminKeyAdd { key: admin_key }, }; store_keychain_op(&store, &op)?; ops.push(op); } Ok(ControlResponse::KeychainInitialized { ops }) } ControlRequest::KeychainStatus => { let view = geth_keychain::reduce_keychain_ops(&load_keychain_ops(&store)?); Ok(ControlResponse::KeychainStatus(KeychainStatusResponse { initialized: view.initialized, admin_keys: view.admin_keys.len(), users: view.users.len(), devices: view.devices.len(), nodes: view.nodes.len(), })) } ControlRequest::SecretStatus => Ok(ControlResponse::SecretStatus { secrets: store .list_resource_secrets()? .into_iter() .map(resource_secret_from_stored) .collect(), }), ControlRequest::SecretCreate { resource } => { ensure_resource_exists(&store, &resource)?; let secret = create_resource_secret(&store, &resource, 1)?; Ok(ControlResponse::SecretCreated { secret }) } ControlRequest::SecretRotate { resource } => { ensure_resource_exists(&store, &resource)?; let next_epoch = store .latest_resource_secret(&resource)? .map(|secret| secret.epoch + 1) .unwrap_or(1); let secret = create_resource_secret(&store, &resource, next_epoch)?; Ok(ControlResponse::SecretCreated { secret }) } ControlRequest::SecretBearerCreate { resource, capabilities, expires_at_ms, } => { ensure_resource_exists(&store, &resource)?; let capabilities = capabilities .into_iter() .map(Capability::new) .collect::>(); geth_secrets::validate_bearer_capabilities(&capabilities)?; let created_at = UnixMillis(geth_store::now_ms()); let secret = geth_types::SecretId::new(format!( "bearer:{}", geth_crypto::blake3_hex( format!( "{resource}\0{}\0{}", capabilities .iter() .map(ToString::to_string) .collect::>() .join(","), created_at.0 ) .as_bytes() ) )); let access = BearerAccess::resource_scoped( secret.clone(), ResourceId::new(resource.clone()), capabilities.clone(), ); let op = AuthOp { id: generated_auth_op_id("bearer-create", &resource, secret.as_str(), created_at), resource: ResourceId::new(resource), created_at, kind: AuthOpKind::BearerAccessCreate { secret, capabilities, expires_at: expires_at_ms.map(UnixMillis), }, }; store_auth_op(&store, &op)?; Ok(ControlResponse::SecretBearerCreated { access: BearerAccess { expires_at: expires_at_ms.map(UnixMillis), ..access }, }) } ControlRequest::SecretBearerList => Ok(ControlResponse::SecretBearerList { access: load_bearer_access(&store)?, }), ControlRequest::SecretBearerRevoke { resource, secret } => { ensure_resource_exists(&store, &resource)?; let created_at = UnixMillis(geth_store::now_ms()); let op = AuthOp { id: generated_auth_op_id("bearer-revoke", &resource, &secret, created_at), resource: ResourceId::new(resource.clone()), created_at, kind: AuthOpKind::BearerAccessRevoke { secret: secret.clone().into(), }, }; store_auth_op(&store, &op)?; Ok(ControlResponse::SecretBearerRevoked { resource, secret }) } ControlRequest::AuthExplain { subject, resource, capability, } => { let ops = load_auth_ops_for_resource(&store, &resource)?; let discovered = store.get_peer_card(&subject)?.is_some(); if ops.is_empty() { if discovered { Ok(ControlResponse::AuthExplain( AuthExplanation::discovered_candidate(subject, resource, capability), )) } else { Ok(ControlResponse::AuthExplain(AuthExplanation::stub( subject, resource, capability, ))) } } else { let mut explanation = geth_auth::explain_auth_ops( &ops, PrincipalId::new(subject.clone()), ResourceId::new(resource.clone()), Capability::new(capability.clone()), ); if discovered && !explanation.allowed { explanation.reason = format!( "subject is a discovered peer candidate only; discovery does not grant trust or authorization; {}", explanation.reason ); } Ok(ControlResponse::AuthExplain(explanation)) } } ControlRequest::AuthGrant { subject, resource, capability, grant_id, } => { let created_at = UnixMillis(geth_store::now_ms()); let grant_id = grant_id.unwrap_or_else(|| generated_grant_id(&subject, &resource, &capability)); let op = AuthOp { id: generated_auth_op_id("grant-create", &resource, &grant_id, created_at), resource: ResourceId::new(resource), created_at, kind: AuthOpKind::GrantCreate { grant_id, principal: PrincipalId::new(subject), capabilities: vec![Capability::new(capability)], }, }; store_auth_op(&store, &op)?; Ok(ControlResponse::AuthOpRecorded { op }) } ControlRequest::AuthRevoke { resource, grant_id } => { let created_at = UnixMillis(geth_store::now_ms()); let op = AuthOp { id: generated_auth_op_id("grant-revoke", &resource, &grant_id, created_at), resource: ResourceId::new(resource), created_at, kind: AuthOpKind::GrantRevoke { grant_id }, }; store_auth_op(&store, &op)?; Ok(ControlResponse::AuthOpRecorded { op }) } 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, format } => { let revocations = store .list_ssh_revocations()? .into_iter() .map(ssh_revocation_from_stored) .collect::, _>>()?; let format = format .parse::() .map_err(NodeError::SshIdentity)?; if let Some(parent) = out.parent() { std::fs::create_dir_all(parent)?; } let (body, note) = match format { SshRevocationExportFormat::Jsonl => { let mut body = String::new(); for revocation in &revocations { body.push_str(&serde_json::to_string(revocation)?); body.push('\n'); } ( body, "JSONL geth revocation metadata; not an OpenSSH KRL binary".to_owned(), ) } SshRevocationExportFormat::OpenSshKrlSpec => ( openssh_krl_spec(&revocations)?, "OpenSSH KRL specification; generate a binary KRL with ssh-keygen -k -f [-s ] ".to_owned(), ), }; std::fs::write(&out, body)?; Ok(ControlResponse::SshRevocationExported { out, format: format.to_string(), count: revocations.len(), note, }) } ControlRequest::DbAdd { name, path } => { geth_db::validate_db_name(&name).map_err(|_| NodeError::InvalidDbName(name.clone()))?; if !path.is_file() { return Err(NodeError::InvalidDbPath(path.display().to_string())); } let path = std::fs::canonicalize(path)?; let schema_metadata = geth_db::schema_metadata(&path)?; let crsqlite_changes = geth_db::crsqlite_change_metadata(&path)?; let resource_id = format!("resource:db:{name}"); let db_id = format!("db:{name}"); let resource = StoredResource { resource_id: resource_id.clone(), kind: ResourceKind::Db.to_string(), name: name.clone(), status: "active".to_owned(), }; store.insert_resource(&resource)?; let stored = StoredDbResource { db_id, resource_id, name, path: path.display().to_string(), }; store.insert_db_resource(&stored)?; Ok(ControlResponse::DbAdded { db: db_resource_from_stored_with_schema(&stored, schema_metadata, crsqlite_changes), }) } ControlRequest::DbStatus { name } => { geth_db::validate_db_name(&name).map_err(|_| NodeError::InvalidDbName(name.clone()))?; let stored = store .get_db_resource_by_name(&name)? .ok_or_else(|| NodeError::DbNotFound(name.clone()))?; Ok(ControlResponse::DbStatus { db: db_resource_from_stored(&stored)?, }) } ControlRequest::DbChanges { name, after_db_version, limit, } => { geth_db::validate_db_name(&name).map_err(|_| NodeError::InvalidDbName(name.clone()))?; let stored = store .get_db_resource_by_name(&name)? .ok_or_else(|| NodeError::DbNotFound(name.clone()))?; let batch = geth_db::extract_crsqlite_changes( Path::new(&stored.path), after_db_version, limit, )?; Ok(ControlResponse::DbChanges { db: db_resource_from_stored(&stored)?, batch, }) } ControlRequest::KvCreate { name } => { geth_kv::validate_kv_name(&name).map_err(|_| NodeError::InvalidKvName(name.clone()))?; let resource_id = format!("resource:kv:{name}"); let kv_id = format!("kv:{name}"); let resource = StoredResource { resource_id: resource_id.clone(), kind: ResourceKind::Kv.to_string(), name: name.clone(), status: "active".to_owned(), }; store.insert_resource(&resource)?; let stored = StoredKvStore { kv_id, resource_id, name, }; store.insert_kv_store(&stored)?; Ok(ControlResponse::KvCreated { kv: kv_resource_from_stored(&stored), }) } ControlRequest::KvSet { name, key, value } => { geth_kv::validate_kv_name(&name).map_err(|_| NodeError::InvalidKvName(name.clone()))?; geth_kv::validate_kv_key(&key).map_err(|_| NodeError::InvalidKvKey(key.clone()))?; let kv = store .get_kv_store_by_name(&name)? .ok_or_else(|| NodeError::KvNotFound(name.clone()))?; let stored = StoredKvEntry { kv_id: kv.kv_id, key, value, updated_at_ms: geth_store::now_ms(), }; store.set_kv_entry(&stored)?; Ok(ControlResponse::KvSet { entry: kv_entry_from_stored(stored), }) } ControlRequest::KvGet { name, key } => { geth_kv::validate_kv_name(&name).map_err(|_| NodeError::InvalidKvName(name.clone()))?; geth_kv::validate_kv_key(&key).map_err(|_| NodeError::InvalidKvKey(key.clone()))?; let kv = store .get_kv_store_by_name(&name)? .ok_or_else(|| NodeError::KvNotFound(name.clone()))?; let entry = store .get_kv_entry(&kv.kv_id, &key)? .map(kv_entry_from_stored); Ok(ControlResponse::KvGet { entry }) } ControlRequest::DocumentCreate { name } => { geth_document::validate_document_name(&name) .map_err(|_| NodeError::InvalidDocumentName(name.clone()))?; let resource_id = format!("resource:document:{name}"); let document_id = format!("document:{name}"); let resource = StoredResource { resource_id: resource_id.clone(), kind: ResourceKind::Document.to_string(), name: name.clone(), status: "active".to_owned(), }; store.insert_resource(&resource)?; let stored = StoredDocumentResource { document_id, resource_id, name, state_json: "{}".to_owned(), updated_at_ms: geth_store::now_ms(), }; store.insert_document_resource(&stored)?; Ok(ControlResponse::DocumentCreated { document: document_resource_from_stored(&stored), }) } ControlRequest::DocumentStatus { name } => { geth_document::validate_document_name(&name) .map_err(|_| NodeError::InvalidDocumentName(name.clone()))?; let stored = store .get_document_resource_by_name(&name)? .ok_or_else(|| NodeError::DocumentNotFound(name.clone()))?; Ok(ControlResponse::DocumentStatus { document: document_resource_from_stored(&stored), }) } ControlRequest::DocumentSet { name, state_json } => { geth_document::validate_document_name(&name) .map_err(|_| NodeError::InvalidDocumentName(name.clone()))?; let mut stored = store .get_document_resource_by_name(&name)? .ok_or_else(|| NodeError::DocumentNotFound(name.clone()))?; stored.state_json = geth_document::normalize_document_state(&state_json)?; stored.updated_at_ms = geth_store::now_ms(); store.insert_document_resource(&stored)?; Ok(ControlResponse::DocumentSet { state: document_state_from_stored(&stored), }) } ControlRequest::DocumentGet { name } => { geth_document::validate_document_name(&name) .map_err(|_| NodeError::InvalidDocumentName(name.clone()))?; let stored = store .get_document_resource_by_name(&name)? .ok_or_else(|| NodeError::DocumentNotFound(name.clone()))?; Ok(ControlResponse::DocumentGet { state: document_state_from_stored(&stored), }) } ControlRequest::PubsubPub { topic, message } => { geth_pubsub::validate_topic(&topic)?; geth_pubsub::validate_message(&message)?; let message = PubsubMessage { topic: topic.clone().into(), message, published_at: UnixMillis(geth_store::now_ms()), }; let mut runtime = node .runtime .pubsub .lock() .map_err(|_| NodeError::RuntimeLockPoisoned)?; runtime.messages.push_back(message.clone()); while runtime.messages.len() > PUBSUB_RING_LIMIT { runtime.messages.pop_front(); } Ok(ControlResponse::PubsubPublished { message }) } ControlRequest::PubsubSub { topic } => { geth_pubsub::validate_topic(&topic)?; let runtime = node .runtime .pubsub .lock() .map_err(|_| NodeError::RuntimeLockPoisoned)?; let messages = runtime .messages .iter() .filter(|message| message.topic.as_str() == topic) .cloned() .collect(); Ok(ControlResponse::PubsubMessages { topic, messages, note: geth_pubsub::pubsub_storage_warning().to_owned(), }) } ControlRequest::PipeListen { name } => { geth_pipe::validate_pipe_name(&name)?; let listener = PipeListener { id: format!("pipe:{name}").into(), name: name.clone(), listened_at: UnixMillis(geth_store::now_ms()), note: geth_pipe::local_pipe_runtime_note().to_owned(), }; let mut runtime = node .runtime .pipes .lock() .map_err(|_| NodeError::RuntimeLockPoisoned)?; runtime.listeners.insert(name, listener.clone()); Ok(ControlResponse::PipeListening { listener }) } ControlRequest::PipeConnect { target } => { geth_pipe::validate_pipe_name(&target)?; let mut runtime = node .runtime .pipes .lock() .map_err(|_| NodeError::RuntimeLockPoisoned)?; let local_listener_found = runtime.listeners.contains_key(&target); let connection = PipeConnection { target, connected_at: UnixMillis(geth_store::now_ms()), local_listener_found, note: geth_pipe::local_pipe_runtime_note().to_owned(), }; runtime.connections.push_back(connection.clone()); while runtime.connections.len() > PIPE_CONNECTION_RING_LIMIT { runtime.connections.pop_front(); } Ok(ControlResponse::PipeConnected { connection }) } 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 db_resource_from_stored(stored: &StoredDbResource) -> Result { let path = Path::new(&stored.path); let metadata = path.metadata().ok(); let (schema_metadata, crsqlite_changes) = if metadata .as_ref() .is_some_and(std::fs::Metadata::is_file) { ( geth_db::schema_metadata(path).unwrap_or_else(|error| format!("unavailable:{error}")), geth_db::crsqlite_change_metadata(path) .unwrap_or_else(|error| geth_db::CrSqliteChangeMetadata::error(error.to_string())), ) } else { ( "unavailable:path-missing".to_owned(), geth_db::CrSqliteChangeMetadata::error("path-missing"), ) }; Ok(db_resource_from_stored_with_schema( stored, schema_metadata, crsqlite_changes, )) } fn db_resource_from_stored_with_schema( stored: &StoredDbResource, schema_metadata: String, crsqlite_changes: geth_db::CrSqliteChangeMetadata, ) -> DbResource { let path = Path::new(&stored.path); let metadata = path.metadata().ok(); DbResource { id: stored.db_id.clone().into(), resource: stored.resource_id.clone().into(), name: stored.name.clone(), path: stored.path.clone(), path_exists: metadata.as_ref().is_some_and(std::fs::Metadata::is_file), size_bytes: metadata.map(|metadata| metadata.len()), schema_metadata, crsqlite_changes, sync_status: "local-only".to_owned(), } } fn kv_resource_from_stored(stored: &StoredKvStore) -> KvResource { KvResource { id: stored.kv_id.clone().into(), resource: stored.resource_id.clone().into(), name: stored.name.clone(), sync_status: "local-only".to_owned(), } } fn kv_entry_from_stored(stored: StoredKvEntry) -> KvEntry { KvEntry { store: stored.kv_id.into(), key: stored.key, value: stored.value, } } fn document_resource_from_stored(stored: &StoredDocumentResource) -> DocumentResource { DocumentResource { id: stored.document_id.clone().into(), resource: stored.resource_id.clone().into(), name: stored.name.clone(), sync_status: "local-only".to_owned(), state_bytes: stored.state_json.len() as u64, } } fn document_state_from_stored(stored: &StoredDocumentResource) -> DocumentState { DocumentState { document: document_resource_from_stored(stored), state_json: stored.state_json.clone(), updated_at: UnixMillis(stored.updated_at_ms), } } fn file_root_from_stored(stored: &StoredFileRoot) -> FileRoot { FileRoot { id: stored.root_id.clone(), resource: stored.resource_id.clone(), name: stored.name.clone(), path: stored.path.clone(), latest_tree: stored.latest_tree_hash.clone().map(Into::into), updated_at_ms: stored.updated_at_ms, } } fn file_conflict_from_stored(stored: StoredFileConflict) -> Result { Ok(FileConflict { id: stored.conflict_id, root: stored.root_name, resource: stored.resource_id, path: stored.path, kind: FileConflictKind::parse(&stored.kind)?, status: FileConflictStatus::parse(&stored.status)?, base_tree: stored.base_tree_hash.map(Into::into), local_tree: stored.local_tree_hash.map(Into::into), remote_tree: stored.remote_tree_hash.map(Into::into), detail: stored.detail, resolution: stored .resolution .as_deref() .map(FileConflictResolution::parse) .transpose()?, resolution_note: stored.resolution_note, created_at_ms: stored.created_at_ms, resolved_at_ms: stored.resolved_at_ms, }) } fn ensure_resource_exists(store: &Store, resource_id: &str) -> Result<(), NodeError> { if store .list_resources()? .into_iter() .any(|resource| resource.resource_id == resource_id) { Ok(()) } else { Err(NodeError::ResourceNotFound(resource_id.to_owned())) } } fn create_resource_secret( store: &Store, resource_id: &str, epoch: u64, ) -> Result { let created_at = UnixMillis(geth_store::now_ms()); let secret_id = format!( "secret:{}", geth_crypto::blake3_hex(format!("{resource_id}\0{epoch}\0{}", created_at.0).as_bytes()) ); let stored = StoredResourceSecret { secret_id, resource_id: resource_id.to_owned(), epoch, status: "active".to_owned(), created_at_ms: created_at.0, }; store.insert_resource_secret(&stored)?; Ok(resource_secret_from_stored(stored)) } fn resource_secret_from_stored(stored: StoredResourceSecret) -> ResourceMasterSecret { ResourceMasterSecret { id: stored.secret_id.into(), resource: stored.resource_id.into(), epoch: stored.epoch, created_at: UnixMillis(stored.created_at_ms), } } fn load_bearer_access(store: &Store) -> Result, NodeError> { let ops = store .list_auth_ops()? .into_iter() .map(|stored| serde_json::from_str(&stored.op_json).map_err(NodeError::from)) .collect::, NodeError>>()?; let view = geth_auth::reduce_auth_ops(&ops); Ok(view .bearer_access .into_values() .map(|record| BearerAccess { secret: record.secret, resource: record.resource, capabilities: record.capabilities, expires_at: record.expires_at, may_delegate: false, }) .collect()) } fn store_auth_op(store: &Store, op: &AuthOp) -> Result<(), NodeError> { store.insert_auth_op(&StoredAuthOp { op_id: op.id.to_string(), resource_id: op.resource.to_string(), op_json: serde_json::to_string(op)?, created_at_ms: op.created_at.0, })?; Ok(()) } fn load_auth_ops_for_resource(store: &Store, resource: &str) -> Result, NodeError> { store .list_auth_ops_for_resource(resource)? .into_iter() .map(|stored| serde_json::from_str(&stored.op_json).map_err(NodeError::from)) .collect() } fn store_keychain_op(store: &Store, op: &KeychainOp) -> Result<(), NodeError> { store.insert_keychain_op(&StoredKeychainOp { op_id: op.id.to_string(), op_json: serde_json::to_string(op)?, created_at_ms: op.created_at.0, })?; Ok(()) } fn load_keychain_ops(store: &Store) -> Result, NodeError> { store .list_keychain_ops()? .into_iter() .map(|stored| serde_json::from_str(&stored.op_json).map_err(NodeError::from)) .collect() } fn generated_grant_id(subject: &str, resource: &str, capability: &str) -> String { format!( "grant:{}", geth_crypto::blake3_hex(format!("{subject}\0{resource}\0{capability}").as_bytes()) ) } fn generated_file_conflict_id(root: &str, path: &str, kind: &str, created_at_ms: i64) -> String { format!( "file-conflict:{}", geth_crypto::blake3_hex(format!("{created_at_ms}\0{root}\0{path}\0{kind}").as_bytes()) ) } fn generated_auth_op_id( kind: &str, resource: &str, stable_id: &str, created_at: UnixMillis, ) -> AuthOpId { AuthOpId::new(format!( "auth-op:{}", geth_crypto::blake3_hex( format!("{}\0{kind}\0{resource}\0{stable_id}", created_at.0).as_bytes() ) )) } fn generated_keychain_op_id(kind: &str, stable_id: &str, created_at: UnixMillis) -> AuthOpId { AuthOpId::new(format!( "keychain-op:{}", geth_crypto::blake3_hex(format!("{}\0{kind}\0{stable_id}", created_at.0).as_bytes()) )) } 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.clone(); let relay_mode_label = relay_mode.label(); let local_discovery = node_config.iroh.local_discovery; let iroh_relay_mode = config_relay_mode_to_iroh(&relay_mode, &node_config.iroh.relay_maps); let mut config = GethIrohConfig::local_with_relay(node.paths.iroh_key(), iroh_relay_mode); config.local_discovery = local_discovery; 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_label, local_discovery, note: format!("Iroh endpoint failed to start: {error}"), }; Ok(None) } } } fn config_relay_mode_to_iroh( mode: &RelayMode, relay_maps: &std::collections::BTreeMap, ) -> GethRelayMode { match mode { RelayMode::Disabled => GethRelayMode::Disabled, RelayMode::Default => GethRelayMode::Default, RelayMode::Staging => GethRelayMode::Staging, RelayMode::Custom { map } => { let relay_map = relay_maps .get(map) .expect("custom relay map was validated during config load"); GethRelayMode::Custom { name: map.clone(), relay_urls: relay_map.urls.clone(), } } } } 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, }) }