Add smooth node enrollment flow

This commit is contained in:
Eric Wendland 2026-05-21 18:01:38 +02:00
commit 27a79768e4
12 changed files with 1662 additions and 22 deletions

View file

@ -1,7 +1,7 @@
pub mod service;
use base64::Engine;
use geth_auth::{AuthExplanation, AuthOp, AuthOpKind};
use geth_auth::{AUTH_SIGNATURE_NAMESPACE, AuthExplanation, AuthOp, AuthOpKind, AuthOpSignature};
use geth_cas::{
BlobInfoSummary, CasTreeEntryKind, CasTreeObject, FileConflict, FileConflictKind,
FileConflictResolution, FileConflictStatus, FileRoot, FileRootScan, LocalCas, hash_path,
@ -20,7 +20,11 @@ use geth_discovery::{
};
use geth_document::{DocumentResource, DocumentState};
use geth_iroh::{EndpointStatus, GethIrohConfig, GethIrohEndpoint, GethRelayMode};
use geth_keychain::{KeychainOp, KeychainOpKind, KeychainOpSignature};
use geth_keychain::{
KeychainOp, KeychainOpKind, KeychainOpSignature, NODE_ENROLLMENT_REQUEST_NAMESPACE,
NodeEnrollmentCapability, NodeEnrollmentProvenance, NodeEnrollmentRequest,
NodeEnrollmentStatus, node_enrollment_request_signing_payload,
};
use geth_kv::{KvEntry, KvResource, KvSyncEntry};
use geth_pipe::{PipeConnection, PipeListener, PipeMessage};
use geth_pubsub::PubsubMessage;
@ -37,10 +41,10 @@ use geth_ssh_identity::{
};
use geth_ssh_proxy::SshProxyConnection;
use geth_store::{
Store, StoredAuthOp, StoredDbResource, StoredDocumentResource, StoredFileConflict,
StoredFileRoot, StoredKeychainOp, StoredKeychainSignature, StoredKvEntry, StoredKvStore,
StoredModuleState, StoredPeerCard, StoredResource, StoredResourceSecret, StoredSshCertRequest,
StoredSshCertificate, StoredSshRevocation,
Store, StoredAuthOp, StoredAuthSignature, StoredDbResource, StoredDocumentResource,
StoredFileConflict, StoredFileRoot, StoredKeychainOp, StoredKeychainSignature, StoredKvEntry,
StoredKvStore, StoredModuleState, StoredNodeEnrollmentRequest, StoredPeerCard, StoredResource,
StoredResourceSecret, StoredSshCertRequest, StoredSshCertificate, StoredSshRevocation,
};
use geth_types::{
AuthOpId, BlobHash, Capability, DeviceId, KeyId, NodeId, PrincipalId, ResourceId, ResourceKind,
@ -85,6 +89,12 @@ pub enum NodeError {
SigningKeyRequired(String),
#[error("invalid init capability grant, expected <resource>=<capability>: {0}")]
InvalidInitGrant(String),
#[error("invalid enrollment capability, expected <resource>=<capability>: {0}")]
InvalidEnrollmentCapability(String),
#[error("node enrollment request not found: {0}")]
NodeEnrollmentRequestNotFound(String),
#[error("node enrollment request has invalid provenance: {0}")]
InvalidNodeEnrollment(String),
#[error("invalid db resource name: {0}")]
InvalidDbName(String),
#[error("db path does not exist or is not a file: {0}")]
@ -641,6 +651,15 @@ pub async fn handle_request_async(
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::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,
@ -1020,6 +1039,8 @@ async fn peer_ping(node: &LocalNode, peer_node: &str) -> Result<ControlResponse,
)),
PeerControlResponse::SyncStatus { .. }
| PeerControlResponse::KeychainSynced { .. }
| PeerControlResponse::AuthSynced { .. }
| PeerControlResponse::NodeEnrollmentSubmitted { .. }
| PeerControlResponse::CasFetched { .. }
| PeerControlResponse::CasRootSynced { .. }
| PeerControlResponse::SshCertSynced { .. }
@ -1139,6 +1160,8 @@ async fn peer_auth_check(
)),
PeerControlResponse::SyncStatus { .. }
| PeerControlResponse::KeychainSynced { .. }
| PeerControlResponse::AuthSynced { .. }
| PeerControlResponse::NodeEnrollmentSubmitted { .. }
| PeerControlResponse::CasFetched { .. }
| PeerControlResponse::CasRootSynced { .. }
| PeerControlResponse::SshCertSynced { .. }
@ -1291,6 +1314,8 @@ async fn cas_fetch_from_peer(
| PeerControlResponse::AuthChecked { .. }
| PeerControlResponse::SyncStatus { .. }
| PeerControlResponse::KeychainSynced { .. }
| PeerControlResponse::AuthSynced { .. }
| PeerControlResponse::NodeEnrollmentSubmitted { .. }
| PeerControlResponse::CasRootSynced { .. }
| PeerControlResponse::SshCertSynced { .. }
| PeerControlResponse::SshRevocationSynced { .. }
@ -2651,6 +2676,191 @@ async fn keychain_sync_from_peer(
}
}
async fn auth_sync_from_peer(
node: &LocalNode,
peer_node: &str,
) -> Result<ControlResponse, NodeError> {
let response = request_peer_control(node, peer_node, "auth-sync", |peer_card, nonce| {
PeerControlRequest::AuthSync { peer_card, nonce }
})
.await?;
match response {
PeerControlResponse::AuthSynced {
node_id,
agent_id,
endpoint_id,
ops,
signatures,
note,
..
} => {
let store = Store::open(&node.paths.metadata_db())?;
let trusted_admins =
geth_keychain::reduce_keychain_ops(&load_keychain_ops(&store)?).admin_keys;
let mut ops_imported = 0;
let mut signatures_imported = 0;
let mut invalid_ops_rejected = 0;
for op in ops {
let op_signatures = signatures
.iter()
.filter(|signature| signature.op_id == op.id)
.cloned()
.collect::<Vec<_>>();
if op_signatures.is_empty() {
invalid_ops_rejected += 1;
continue;
}
let valid_signatures = op_signatures
.iter()
.filter(|signature| {
if !trusted_admins.contains(&signature.signer)
|| !auth_signature_uses_claimed_key(signature)
{
return false;
}
let stored = stored_auth_signature_from_signature(signature);
verify_auth_signature_with_ssh(node, &op, &stored).unwrap_or(false)
})
.cloned()
.collect::<Vec<_>>();
if valid_signatures.is_empty() {
invalid_ops_rejected += 1;
continue;
}
let existing = load_auth_ops(&store)?
.into_iter()
.find(|existing| existing.id == op.id);
if existing.as_ref().is_some_and(|existing| existing != &op) {
invalid_ops_rejected += 1;
continue;
}
store_auth_op(&store, &op)?;
ops_imported += usize::from(existing.is_none());
for signature in valid_signatures {
store.insert_auth_signature(&StoredAuthSignature {
op_id: signature.op_id.to_string(),
signer: signature.signer.to_string(),
signer_public_key: signature.signer_public_key.clone(),
namespace: signature.namespace.clone(),
signature: signature.signature.clone(),
created_at_ms: signature.created_at.0,
})?;
signatures_imported += 1;
}
}
Ok(ControlResponse::AuthSynced {
peer_node_id: node_id,
peer_agent_id: agent_id,
endpoint_id,
ops_imported,
signatures_imported,
invalid_ops_rejected,
note: format!("{note}; imported only auth ops signed by trusted admin SSH keys"),
})
}
PeerControlResponse::Error { message } => Err(NodeError::IrohPeer(message)),
_ => Err(NodeError::IrohPeer(
"peer returned wrong response type to auth sync".to_owned(),
)),
}
}
async fn node_enrollment_submit_to_peer(
node: &LocalNode,
owner_node: &str,
request_id: Option<String>,
path: Option<PathBuf>,
) -> Result<ControlResponse, NodeError> {
let store = Store::open(&node.paths.metadata_db())?;
let request = if let Some(path) = path {
let request = read_node_enrollment_request_file(&path)?;
insert_node_enrollment_request_if_not_conflicting(&store, &request)?;
request
} else {
let request_id = request_id
.or_else(|| {
latest_pending_node_enrollment_request(&store).map(|request| request.id.to_string())
})
.ok_or_else(|| NodeError::NodeEnrollmentRequestNotFound("latest-pending".to_owned()))?;
load_node_enrollment_request(&store, &request_id)?
};
let response = request_peer_control(
node,
owner_node,
"node-enrollment-submit",
|peer_card, nonce| PeerControlRequest::NodeEnrollmentSubmit {
peer_card,
request: request.clone(),
nonce,
},
)
.await?;
match response {
PeerControlResponse::NodeEnrollmentSubmitted {
node_id,
accepted,
note,
..
} => Ok(ControlResponse::NodeEnrollmentSubmitted {
request_id: request.id.to_string(),
owner_node_id: node_id,
accepted,
note,
}),
PeerControlResponse::Error { message } => Err(NodeError::IrohPeer(message)),
_ => Err(NodeError::IrohPeer(
"peer returned wrong response type to node enrollment submit".to_owned(),
)),
}
}
async fn node_enrollment_sync_from_owner(
node: &LocalNode,
owner_node: &str,
) -> Result<ControlResponse, NodeError> {
let keychain = keychain_sync_from_peer(node, owner_node).await?;
let auth = auth_sync_from_peer(node, owner_node).await?;
let (keychain_ops_imported, keychain_signatures_imported, mut invalid_ops_rejected) =
match keychain {
ControlResponse::KeychainSynced {
ops_imported,
signatures_imported,
invalid_ops_rejected,
..
} => (ops_imported, signatures_imported, invalid_ops_rejected),
other => {
return Err(NodeError::IrohPeer(format!(
"unexpected keychain sync response during enrollment sync: {other:?}"
)));
}
};
let (auth_ops_imported, auth_signatures_imported, auth_invalid) = match auth {
ControlResponse::AuthSynced {
ops_imported,
signatures_imported,
invalid_ops_rejected,
..
} => (ops_imported, signatures_imported, invalid_ops_rejected),
other => {
return Err(NodeError::IrohPeer(format!(
"unexpected auth sync response during enrollment sync: {other:?}"
)));
}
};
invalid_ops_rejected += auth_invalid;
Ok(ControlResponse::NodeEnrollmentSynced {
owner_node: owner_node.to_owned(),
keychain_ops_imported,
keychain_signatures_imported,
auth_ops_imported,
auth_signatures_imported,
invalid_ops_rejected,
note:
"enrollment sync pulled signed keychain and signed auth operations from the owner node"
.to_owned(),
})
}
async fn db_sync_from_peer(
node: &LocalNode,
peer_node: &str,
@ -3102,6 +3312,14 @@ async fn request_peer_control(
nonce: response_nonce,
..
}
| PeerControlResponse::AuthSynced {
nonce: response_nonce,
..
}
| PeerControlResponse::NodeEnrollmentSubmitted {
nonce: response_nonce,
..
}
| PeerControlResponse::SshRevocationSynced {
nonce: response_nonce,
..
@ -3280,6 +3498,12 @@ async fn run_live_sync_once(node: &LocalNode) -> Result<(), NodeError> {
None
}
};
if let Err(error) = keychain_sync_from_peer(node, &peer.peer_id).await {
tracing::debug!(peer = %peer.peer_id, %error, "live keychain sync failed");
}
if let Err(error) = auth_sync_from_peer(node, &peer.peer_id).await {
tracing::debug!(peer = %peer.peer_id, %error, "live auth sync failed");
}
if should_live_sync_stream(
&store,
&peer.peer_id,
@ -3457,6 +3681,61 @@ async fn handle_iroh_control_connection(
note: "keychain sync returns signed operation-log data; receiver must verify OpenSSH signatures before import".to_owned(),
}
}
PeerControlRequest::AuthSync { peer_card, nonce } => {
peer_card.validate_candidate()?;
ensure_peer_card_matches_endpoint(&peer_card, &remote_endpoint_id)?;
let discovered = DiscoveredPeer::candidate(
peer_card.clone(),
UnixMillis(geth_store::now_ms()),
DiscoverySource::PeerExchange,
)?;
let store = Store::open(&node.paths.metadata_db())?;
store.upsert_peer_card(&StoredPeerCard {
peer_id: peer_card.node_id.to_string(),
card_json: serde_json::to_string(&peer_card)?,
updated_at_ms: discovered.discovered_at.0,
})?;
PeerControlResponse::AuthSynced {
node_id: node.node_id.clone(),
agent_id: node.agent_id.clone(),
endpoint_id: node.iroh_status.endpoint_id.clone().unwrap_or_default(),
remote_endpoint_id,
ops: load_auth_ops(&store)?,
signatures: load_auth_signatures(&store)?,
nonce,
note: "auth sync returns signed resource operation-log data; receiver must verify admin signatures before import".to_owned(),
}
}
PeerControlRequest::NodeEnrollmentSubmit {
peer_card,
request,
nonce,
} => {
peer_card.validate_candidate()?;
ensure_peer_card_matches_endpoint(&peer_card, &remote_endpoint_id)?;
let discovered = DiscoveredPeer::candidate(
peer_card.clone(),
UnixMillis(geth_store::now_ms()),
DiscoverySource::PeerExchange,
)?;
let store = Store::open(&node.paths.metadata_db())?;
store.upsert_peer_card(&StoredPeerCard {
peer_id: peer_card.node_id.to_string(),
card_json: serde_json::to_string(&peer_card)?,
updated_at_ms: discovered.discovered_at.0,
})?;
let accepted = insert_node_enrollment_request_if_not_conflicting(&store, &request)?;
PeerControlResponse::NodeEnrollmentSubmitted {
node_id: node.node_id.clone(),
agent_id: node.agent_id.clone(),
endpoint_id: node.iroh_status.endpoint_id.clone().unwrap_or_default(),
remote_endpoint_id,
request_id: request.id.to_string(),
accepted,
nonce,
note: "stored signed node enrollment request for owner review".to_owned(),
}
}
PeerControlRequest::AuthCheck {
peer_card,
resource,
@ -5388,6 +5667,7 @@ pub fn handle_request(
let op = record_node_grant(&store, target.as_str(), &resource, &capability, grant_id)?;
Ok(ControlResponse::NodeGrantUpdated {
op,
signatures: Vec::new(),
note: "recorded resource-scoped node capability grant; auth explain can show the grant path".to_owned(),
})
}
@ -5402,9 +5682,66 @@ pub fn handle_request(
store_auth_op(&store, &op)?;
Ok(ControlResponse::NodeGrantUpdated {
op,
signatures: Vec::new(),
note: "recorded capability grant revocation".to_owned(),
})
}
ControlRequest::NodeEnrollRequest {
node_name,
capabilities,
reason,
out,
} => {
let mut request =
create_node_enrollment_request(node, node_name, capabilities, reason)?;
sign_node_enrollment_request(node, &mut request)?;
store_node_enrollment_request(&store, &request)?;
if let Some(path) = out.as_ref() {
write_node_enrollment_request_file(path, &request)?;
}
Ok(ControlResponse::NodeEnrollmentRequested {
request,
out,
note: "created signed node enrollment request; submit it to an owner node or move the JSON file to the owner machine".to_owned(),
})
}
ControlRequest::NodeEnrollImport { path } => {
let request = read_node_enrollment_request_file(&path)?;
insert_node_enrollment_request_if_not_conflicting(&store, &request)?;
Ok(ControlResponse::NodeEnrollmentImported {
request,
note: "imported signed node enrollment request for owner review".to_owned(),
})
}
ControlRequest::NodeEnrollList { status } => {
let requests = load_node_enrollment_requests(&store)?
.into_iter()
.filter(|request| {
status
.as_ref()
.is_none_or(|status| request.status.as_str() == status)
})
.collect();
Ok(ControlResponse::NodeEnrollmentList {
requests,
note: "node enrollment requests are signed by the requesting agent; approve on an owner machine with an admin signing key".to_owned(),
})
}
ControlRequest::NodeEnrollApprove {
request_id,
signing_key_path,
admin_key_path,
node_name,
capabilities,
} => approve_node_enrollment(
&store,
node,
&request_id,
&signing_key_path,
admin_key_path.as_deref(),
node_name,
capabilities,
),
ControlRequest::SecretStatus => Ok(ControlResponse::SecretStatus {
secrets: store
.list_resource_secrets()?
@ -6305,7 +6642,10 @@ pub fn handle_request(
ControlRequest::SshProxyConnect { .. }
| ControlRequest::SshProxyStream { .. }
| ControlRequest::SshAdminShell { .. }
| ControlRequest::KeychainSync { .. } => Err(NodeError::IrohEndpointUnavailable),
| ControlRequest::KeychainSync { .. }
| ControlRequest::AuthSync { .. }
| ControlRequest::NodeEnrollSubmit { .. }
| ControlRequest::NodeEnrollSync { .. } => Err(NodeError::IrohEndpointUnavailable),
ControlRequest::ModuleStub { module, command } => {
Ok(ControlResponse::NotImplemented { module, command })
}
@ -6968,6 +7308,78 @@ fn store_auth_op(store: &Store, op: &AuthOp) -> Result<(), NodeError> {
Ok(())
}
fn store_and_sign_auth_ops(
store: &Store,
node: &LocalNode,
ops: &[AuthOp],
signing_key_path: &Path,
admin_key_path: Option<&Path>,
) -> Result<Vec<AuthOpSignature>, NodeError> {
for op in ops {
store_auth_op(store, op)?;
}
let (signer, signer_public_key) = keychain_signer_from_paths(signing_key_path, admin_key_path)?;
ops.iter()
.map(|op| {
sign_auth_op_with_ssh(
store,
node,
op,
signing_key_path,
&signer,
&signer_public_key,
)
})
.collect()
}
fn sign_auth_op_with_ssh(
store: &Store,
node: &LocalNode,
op: &AuthOp,
signing_key_path: &Path,
signer: &KeyId,
signer_public_key: &str,
) -> Result<AuthOpSignature, NodeError> {
geth_ssh_identity::ensure_ssh_keygen_available()?;
let signature_dir = node.paths.home().join("auth-signatures");
std::fs::create_dir_all(&signature_dir)?;
let payload_path = signature_dir.join(format!(
"{}.payload",
geth_crypto::blake3_hex(op.id.as_str().as_bytes())
));
std::fs::write(&payload_path, geth_auth::auth_signing_payload(op)?)?;
let output =
geth_ssh_identity::sign_command(signing_key_path, AUTH_SIGNATURE_NAMESPACE, &payload_path)
.output()?;
if !output.status.success() {
return Err(geth_ssh_identity::SshIdentityError::SshKeygenFailed(
String::from_utf8_lossy(&output.stderr).trim().to_owned(),
)
.into());
}
let signature_path = Path::new(&format!("{}.sig", payload_path.display())).to_path_buf();
let signature_bytes = std::fs::read(signature_path)?;
let created_at = UnixMillis(geth_store::now_ms());
let signature = AuthOpSignature {
op_id: op.id.clone(),
signer: signer.clone(),
signer_public_key: signer_public_key.to_owned(),
namespace: AUTH_SIGNATURE_NAMESPACE.to_owned(),
signature: signature_bytes,
created_at,
};
store.insert_auth_signature(&StoredAuthSignature {
op_id: signature.op_id.to_string(),
signer: signature.signer.to_string(),
signer_public_key: signature.signer_public_key.clone(),
namespace: signature.namespace.clone(),
signature: signature.signature.clone(),
created_at_ms: signature.created_at.0,
})?;
Ok(signature)
}
fn load_auth_ops_for_resource(store: &Store, resource: &str) -> Result<Vec<AuthOp>, NodeError> {
store
.list_auth_ops_for_resource(resource)?
@ -6976,6 +7388,94 @@ fn load_auth_ops_for_resource(store: &Store, resource: &str) -> Result<Vec<AuthO
.collect()
}
fn load_auth_ops(store: &Store) -> Result<Vec<AuthOp>, NodeError> {
store
.list_auth_ops()?
.into_iter()
.map(|stored| serde_json::from_str(&stored.op_json).map_err(NodeError::from))
.collect()
}
fn load_auth_signatures(store: &Store) -> Result<Vec<AuthOpSignature>, NodeError> {
Ok(store
.list_auth_signatures()?
.into_iter()
.map(|stored| AuthOpSignature {
op_id: AuthOpId::new(stored.op_id),
signer: KeyId::new(stored.signer),
signer_public_key: stored.signer_public_key,
namespace: stored.namespace,
signature: stored.signature,
created_at: UnixMillis(stored.created_at_ms),
})
.collect())
}
fn stored_auth_signature_from_signature(signature: &AuthOpSignature) -> StoredAuthSignature {
StoredAuthSignature {
op_id: signature.op_id.to_string(),
signer: signature.signer.to_string(),
signer_public_key: signature.signer_public_key.clone(),
namespace: signature.namespace.clone(),
signature: signature.signature.clone(),
created_at_ms: signature.created_at.0,
}
}
fn auth_signature_uses_claimed_key(signature: &AuthOpSignature) -> bool {
KeyId::new(ssh_public_key_fingerprint(&signature.signer_public_key)) == signature.signer
}
fn verify_auth_signature_with_ssh(
node: &LocalNode,
op: &AuthOp,
signature: &StoredAuthSignature,
) -> Result<bool, NodeError> {
if signature.signer_public_key.trim().is_empty() {
return Ok(false);
}
if geth_ssh_identity::ensure_ssh_keygen_available().is_err() {
return Ok(false);
}
let verify_dir = node.paths.home().join("auth-signatures").join("verify");
std::fs::create_dir_all(&verify_dir)?;
let stable_id = geth_crypto::blake3_hex(
format!(
"{}\0{}\0{}\0{}",
op.id, signature.signer, signature.namespace, signature.created_at_ms
)
.as_bytes(),
);
let payload_path = verify_dir.join(format!("{stable_id}.payload"));
let signature_path = verify_dir.join(format!("{stable_id}.sig"));
let allowed_signers_path = verify_dir.join(format!("{stable_id}.allowed-signers"));
std::fs::write(&payload_path, geth_auth::auth_signing_payload(op)?)?;
std::fs::write(&signature_path, &signature.signature)?;
std::fs::write(
&allowed_signers_path,
format!(
"{} {}\n",
signature.signer,
signature.signer_public_key.trim()
),
)?;
let payload = std::fs::File::open(&payload_path)?;
let output = std::process::Command::new("ssh-keygen")
.arg("-Y")
.arg("verify")
.arg("-f")
.arg(&allowed_signers_path)
.arg("-I")
.arg(&signature.signer)
.arg("-n")
.arg(&signature.namespace)
.arg("-s")
.arg(&signature_path)
.stdin(std::process::Stdio::from(payload))
.output()?;
Ok(output.status.success())
}
fn store_keychain_op(store: &Store, op: &KeychainOp) -> Result<(), NodeError> {
store.insert_keychain_op(&StoredKeychainOp {
op_id: op.id.to_string(),
@ -7269,6 +7769,297 @@ fn stable_slug(value: &str) -> String {
}
}
fn parse_enrollment_capabilities(
capabilities: Vec<String>,
) -> Result<Vec<NodeEnrollmentCapability>, NodeError> {
capabilities
.into_iter()
.map(|value| {
let (resource, capability) = value
.split_once('=')
.ok_or_else(|| NodeError::InvalidEnrollmentCapability(value.clone()))?;
Ok(NodeEnrollmentCapability {
resource: ResourceId::new(resource.to_owned()),
capability: Capability::new(capability.to_owned()),
})
})
.collect()
}
fn create_node_enrollment_request(
node: &LocalNode,
node_name: String,
capabilities: Vec<String>,
reason: Option<String>,
) -> Result<NodeEnrollmentRequest, NodeError> {
let key = AgentKey::load(&node.paths.agent_key())?;
let created_at = UnixMillis(geth_store::now_ms());
let requested_capabilities = parse_enrollment_capabilities(capabilities)?;
let endpoint_id = node.iroh_status.endpoint_id.clone();
let stable = format!(
"{}\0{}\0{}\0{}",
node.node_id, node.agent_id, node_name, created_at.0
);
Ok(NodeEnrollmentRequest {
id: AuthOpId::new(format!(
"node-enrollment:{}",
geth_crypto::blake3_hex(stable.as_bytes())
)),
requester_node: NodeId::new(node.node_id.clone()),
requester_agent: node.agent_id.clone().into(),
requester_agent_public_key: key.public_key_hex(),
requested_node_name: node_name,
requested_capabilities,
endpoint_id,
reason,
status: NodeEnrollmentStatus::Pending,
created_at,
provenance: None,
})
}
fn sign_node_enrollment_request(
node: &LocalNode,
request: &mut NodeEnrollmentRequest,
) -> Result<(), NodeError> {
let key = AgentKey::load(&node.paths.agent_key())?;
let signature = key.sign_canonical(
NODE_ENROLLMENT_REQUEST_NAMESPACE,
&node_enrollment_request_signing_payload(request),
)?;
request.provenance = Some(NodeEnrollmentProvenance {
namespace: NODE_ENROLLMENT_REQUEST_NAMESPACE.to_owned(),
signer_node: NodeId::new(node.node_id.clone()),
signer_agent: node.agent_id.clone().into(),
signer_public_key: key.public_key_hex(),
signature_hex: hex::encode(signature),
signed_at: UnixMillis(geth_store::now_ms()),
});
Ok(())
}
fn verify_node_enrollment_request(request: &NodeEnrollmentRequest) -> Result<(), NodeError> {
let provenance = request
.provenance
.as_ref()
.ok_or_else(|| NodeError::InvalidNodeEnrollment("missing signed provenance".to_owned()))?;
if provenance.namespace != NODE_ENROLLMENT_REQUEST_NAMESPACE {
return Err(NodeError::InvalidNodeEnrollment(format!(
"invalid namespace {}",
provenance.namespace
)));
}
if provenance.signer_node != request.requester_node
|| provenance.signer_agent != request.requester_agent
|| provenance.signer_public_key != request.requester_agent_public_key
{
return Err(NodeError::InvalidNodeEnrollment(
"provenance does not match requesting node/agent".to_owned(),
));
}
verify_record_provenance(
&provenance.signer_public_key,
NODE_ENROLLMENT_REQUEST_NAMESPACE,
&node_enrollment_request_signing_payload(request),
&provenance.signature_hex,
)
}
fn store_node_enrollment_request(
store: &Store,
request: &NodeEnrollmentRequest,
) -> Result<(), NodeError> {
store.insert_node_enrollment_request(&StoredNodeEnrollmentRequest {
request_id: request.id.to_string(),
requester_node: request.requester_node.to_string(),
request_json: serde_json::to_string(request)?,
status: request.status.to_string(),
created_at_ms: request.created_at.0,
})?;
Ok(())
}
fn insert_node_enrollment_request_if_not_conflicting(
store: &Store,
request: &NodeEnrollmentRequest,
) -> Result<bool, NodeError> {
verify_node_enrollment_request(request)?;
let stored = StoredNodeEnrollmentRequest {
request_id: request.id.to_string(),
requester_node: request.requester_node.to_string(),
request_json: serde_json::to_string(request)?,
status: request.status.to_string(),
created_at_ms: request.created_at.0,
};
if let Some(existing) = store.get_node_enrollment_request(&stored.request_id)? {
return Ok(existing == stored);
}
store.insert_node_enrollment_request(&stored)?;
Ok(true)
}
fn load_node_enrollment_request(
store: &Store,
request_id: &str,
) -> Result<NodeEnrollmentRequest, NodeError> {
store
.get_node_enrollment_request(request_id)?
.ok_or_else(|| NodeError::NodeEnrollmentRequestNotFound(request_id.to_owned()))
.and_then(node_enrollment_request_from_stored)
}
fn load_node_enrollment_requests(store: &Store) -> Result<Vec<NodeEnrollmentRequest>, NodeError> {
store
.list_node_enrollment_requests()?
.into_iter()
.map(node_enrollment_request_from_stored)
.collect()
}
fn node_enrollment_request_from_stored(
stored: StoredNodeEnrollmentRequest,
) -> Result<NodeEnrollmentRequest, NodeError> {
let mut request: NodeEnrollmentRequest = serde_json::from_str(&stored.request_json)?;
request.status = stored
.status
.parse()
.map_err(|error: geth_keychain::KeychainError| {
NodeError::InvalidNodeEnrollment(error.to_string())
})?;
Ok(request)
}
fn latest_pending_node_enrollment_request(store: &Store) -> Option<NodeEnrollmentRequest> {
load_node_enrollment_requests(store)
.ok()
.and_then(|requests| {
requests
.into_iter()
.filter(|request| request.status == NodeEnrollmentStatus::Pending)
.max_by_key(|request| request.created_at.0)
})
}
fn write_node_enrollment_request_file(
path: &Path,
request: &NodeEnrollmentRequest,
) -> Result<(), NodeError> {
std::fs::write(path, serde_json::to_vec_pretty(request)?)?;
Ok(())
}
fn read_node_enrollment_request_file(path: &Path) -> Result<NodeEnrollmentRequest, NodeError> {
let request = serde_json::from_slice(&std::fs::read(path)?)?;
Ok(request)
}
fn approve_node_enrollment(
store: &Store,
node: &LocalNode,
request_id: &str,
signing_key_path: &Path,
admin_key_path: Option<&Path>,
node_name: Option<String>,
extra_capabilities: Vec<String>,
) -> Result<ControlResponse, NodeError> {
let mut request = load_node_enrollment_request(store, request_id)?;
verify_node_enrollment_request(&request)?;
let now = UnixMillis(geth_store::now_ms());
let owner = first_keychain_user(store).unwrap_or_else(|| UserId::new("user:owner"));
let device = DeviceId::new(format!(
"device:{}",
stable_slug(request.requester_node.as_str())
));
let node_id = request.requester_node.clone();
let name = node_name.unwrap_or_else(|| request.requested_node_name.clone());
let mut keychain_ops = vec![
KeychainOp {
id: generated_keychain_op_id("device-add", device.as_str(), now),
created_at: now,
kind: KeychainOpKind::DeviceAdd {
device: device.clone(),
user: owner,
},
},
KeychainOp {
id: generated_keychain_op_id("node-add", node_id.as_str(), now),
created_at: now,
kind: KeychainOpKind::NodeAdd {
node: node_id.clone(),
device,
name,
},
},
KeychainOp {
id: generated_keychain_op_id("agent-bind", request.requester_agent.as_str(), now),
created_at: now,
kind: KeychainOpKind::AgentBind {
agent: request.requester_agent.clone(),
node: node_id.clone(),
},
},
];
if let Some(endpoint) = request.endpoint_id.clone() {
keychain_ops.push(KeychainOp {
id: generated_keychain_op_id("node-endpoint-add", &endpoint, now),
created_at: now,
kind: KeychainOpKind::NodeEndpointAdd {
node: node_id.clone(),
endpoint,
},
});
}
let keychain_signatures = store_and_sign_keychain_ops(
store,
node,
&keychain_ops,
Some(signing_key_path),
admin_key_path,
)?;
let mut grant_caps = request.requested_capabilities.clone();
grant_caps.extend(parse_enrollment_capabilities(extra_capabilities)?);
let auth_ops = grant_caps
.into_iter()
.map(|grant| {
let resource = grant.resource.to_string();
let capability = grant.capability.to_string();
let grant_id = generated_grant_id(node_id.as_str(), &resource, &capability);
AuthOp {
id: generated_auth_op_id("grant-create", &resource, &grant_id, now),
resource: grant.resource,
created_at: now,
kind: AuthOpKind::GrantCreate {
grant_id,
principal: PrincipalId::new(node_id.to_string()),
capabilities: vec![grant.capability],
},
}
})
.collect::<Vec<_>>();
let auth_signatures =
store_and_sign_auth_ops(store, node, &auth_ops, signing_key_path, admin_key_path)?;
request.status = NodeEnrollmentStatus::Approved;
store_node_enrollment_request(store, &request)?;
Ok(ControlResponse::NodeEnrollmentApproved {
request,
keychain_ops,
keychain_signatures,
auth_ops,
auth_signatures,
note: "approved enrollment; requester can run keychain sync and auth sync from this owner node".to_owned(),
})
}
fn first_keychain_user(store: &Store) -> Option<UserId> {
geth_keychain::reduce_keychain_ops(&load_keychain_ops(store).ok()?)
.users
.keys()
.next()
.cloned()
}
fn generated_grant_id(subject: &str, resource: &str, capability: &str) -> String {
format!(
"grant:{}",