diff --git a/crates/geth-node/src/lib.rs b/crates/geth-node/src/lib.rs index 6e70e5b..194898a 100644 --- a/crates/geth-node/src/lib.rs +++ b/crates/geth-node/src/lib.rs @@ -11622,10 +11622,24 @@ fn stored_from_ssh_cert_request(request: &SshCertRequest) -> StoredSshCertReques } } +#[cfg(test)] fn insert_ssh_cert_request_if_not_conflicting( store: &Store, request: &SshCertRequest, ) -> Result { + match stage_ssh_cert_request_if_not_conflicting(store, request)? { + Some(stored) => { + store.insert_ssh_cert_request(&stored)?; + Ok(true) + } + None => Ok(false), + } +} + +fn stage_ssh_cert_request_if_not_conflicting( + store: &Store, + request: &SshCertRequest, +) -> Result, NodeError> { let stored = stored_from_ssh_cert_request(request); if let Some(existing) = store.get_ssh_cert_request(&stored.request_id)? { if existing != stored { @@ -11634,11 +11648,10 @@ fn insert_ssh_cert_request_if_not_conflicting( "rejected conflicting SSH certificate request during sync" ); } - return Ok(false); + return Ok(None); } verify_ssh_cert_request_provenance(request)?; - store.insert_ssh_cert_request(&stored)?; - Ok(true) + Ok(Some(stored)) } fn ssh_cert_request_from_stored(stored: StoredSshCertRequest) -> Result { @@ -11685,10 +11698,24 @@ fn stored_from_ssh_certificate(certificate: &SshCertificateRecord) -> StoredSshC } } +#[cfg(test)] fn insert_ssh_certificate_if_not_conflicting( store: &Store, certificate: &SshCertificateRecord, ) -> Result { + match stage_ssh_certificate_if_not_conflicting(store, certificate)? { + Some(stored) => { + store.insert_ssh_certificate(&stored)?; + Ok(true) + } + None => Ok(false), + } +} + +fn stage_ssh_certificate_if_not_conflicting( + store: &Store, + certificate: &SshCertificateRecord, +) -> Result, NodeError> { let stored = stored_from_ssh_certificate(certificate); if let Some(existing) = store .list_ssh_certificates()? @@ -11701,11 +11728,10 @@ fn insert_ssh_certificate_if_not_conflicting( "rejected conflicting SSH certificate during sync" ); } - return Ok(false); + return Ok(None); } verify_ssh_certificate_provenance(certificate)?; - store.insert_ssh_certificate(&stored)?; - Ok(true) + Ok(Some(stored)) } fn ssh_certificate_from_stored( @@ -11741,10 +11767,24 @@ fn stored_from_ssh_revocation(revocation: &SshRevocationEntry) -> StoredSshRevoc } } +#[cfg(test)] fn insert_ssh_revocation_if_not_conflicting( store: &Store, revocation: &SshRevocationEntry, ) -> Result { + match stage_ssh_revocation_if_not_conflicting(store, revocation)? { + Some(stored) => { + store.insert_ssh_revocation(&stored)?; + Ok(true) + } + None => Ok(false), + } +} + +fn stage_ssh_revocation_if_not_conflicting( + store: &Store, + revocation: &SshRevocationEntry, +) -> Result, NodeError> { let stored = stored_from_ssh_revocation(revocation); if let Some(existing) = store .list_ssh_revocations()? @@ -11757,10 +11797,69 @@ fn insert_ssh_revocation_if_not_conflicting( "rejected conflicting SSH revocation during sync" ); } - return Ok(false); + return Ok(None); } verify_ssh_revocation_provenance(revocation)?; - store.insert_ssh_revocation(&stored)?; + Ok(Some(stored)) +} + +fn stage_unique_ssh_cert_request( + staged: &mut Vec, + stored: StoredSshCertRequest, +) -> Result { + if let Some(existing) = staged + .iter() + .find(|existing| existing.request_id == stored.request_id) + { + if existing == &stored { + return Ok(false); + } + return Err(NodeError::IrohPeer(format!( + "conflicting SSH certificate request in sync batch: {}", + stored.request_id + ))); + } + staged.push(stored); + Ok(true) +} + +fn stage_unique_ssh_certificate( + staged: &mut Vec, + stored: StoredSshCertificate, +) -> Result { + if let Some(existing) = staged + .iter() + .find(|existing| existing.cert_id == stored.cert_id) + { + if existing == &stored { + return Ok(false); + } + return Err(NodeError::IrohPeer(format!( + "conflicting SSH certificate in sync batch: {}", + stored.cert_id + ))); + } + staged.push(stored); + Ok(true) +} + +fn stage_unique_ssh_revocation( + staged: &mut Vec, + stored: StoredSshRevocation, +) -> Result { + if let Some(existing) = staged + .iter() + .find(|existing| existing.revocation_id == stored.revocation_id) + { + if existing == &stored { + return Ok(false); + } + return Err(NodeError::IrohPeer(format!( + "conflicting SSH revocation in sync batch: {}", + stored.revocation_id + ))); + } + staged.push(stored); Ok(true) } @@ -11809,25 +11908,39 @@ fn import_ssh_distribution_log_entries( let mut requests_imported = 0; let mut certificates_imported = 0; let mut revocations_imported = 0; + let mut requests = Vec::new(); + let mut certificates = Vec::new(); + let mut revocations = Vec::new(); for entry in entries { match entry { SshDistributionLogEntry::CertRequest { request, .. } => { - if insert_ssh_cert_request_if_not_conflicting(store, request)? { + if let Some(stored) = stage_ssh_cert_request_if_not_conflicting(store, request)? { + if !stage_unique_ssh_cert_request(&mut requests, stored)? { + continue; + } requests_imported += 1; } } SshDistributionLogEntry::Certificate { certificate, .. } => { - if insert_ssh_certificate_if_not_conflicting(store, certificate)? { + if let Some(stored) = stage_ssh_certificate_if_not_conflicting(store, certificate)? + { + if !stage_unique_ssh_certificate(&mut certificates, stored)? { + continue; + } certificates_imported += 1; } } SshDistributionLogEntry::Revocation { revocation, .. } => { - if insert_ssh_revocation_if_not_conflicting(store, revocation)? { + if let Some(stored) = stage_ssh_revocation_if_not_conflicting(store, revocation)? { + if !stage_unique_ssh_revocation(&mut revocations, stored)? { + continue; + } revocations_imported += 1; } } } } + store.insert_ssh_distribution_records(&requests, &certificates, &revocations)?; Ok(( requests_imported, certificates_imported, @@ -12263,6 +12376,37 @@ mod tests { Some("original") ); + let mut certificate = SshCertificateRecord { + id: SshCertId::new("ssh-cert:1"), + request_id: request.id.clone(), + certificate: "ssh-ed25519-cert-v01@openssh.com AAAA".to_owned(), + certificate_fingerprint: "SHA256:cert-one".to_owned(), + imported_at: UnixMillis(2), + provenance: None, + }; + add_test_ssh_certificate_provenance(&mut certificate); + assert!( + insert_ssh_certificate_if_not_conflicting(&store, &certificate) + .expect("insert certificate") + ); + let mut conflicting_certificate = certificate.clone(); + conflicting_certificate.certificate = "ssh-ed25519-cert-v01@openssh.com BBBB".to_owned(); + conflicting_certificate.provenance = None; + add_test_ssh_certificate_provenance(&mut conflicting_certificate); + assert!( + !insert_ssh_certificate_if_not_conflicting(&store, &conflicting_certificate) + .expect("reject conflicting certificate") + ); + assert_eq!( + store + .list_ssh_certificates() + .expect("list certificates") + .first() + .expect("certificate") + .certificate, + "ssh-ed25519-cert-v01@openssh.com AAAA" + ); + let mut revocation = SshRevocationEntry { id: geth_types::SshRevocationId::new("ssh-revocation:1"), kind: SshRevocationKind::KeyId, @@ -12298,6 +12442,46 @@ mod tests { ); } + #[test] + fn ssh_sync_import_rejects_batch_conflicts_without_partial_write() { + let store = Store::open_memory().expect("open"); + let mut request = SshCertRequest { + id: SshCertRequestId::new("ssh-cert-request:batch-conflict"), + requester_node: NodeId::new("node:right"), + public_key: "ssh-ed25519 AAAA".to_owned(), + public_key_fingerprint: "SHA256:one".to_owned(), + cert_kind: SshCertKind::User, + principals: vec!["eric".to_owned()], + requested_validity: None, + renewal_of: None, + reason: Some("original".to_owned()), + status: SshCertRequestStatus::Pending, + created_at: UnixMillis(1), + provenance: None, + }; + add_test_ssh_cert_request_provenance(&mut request); + let mut conflicting_request = request.clone(); + conflicting_request.reason = Some("conflict".to_owned()); + conflicting_request.provenance = None; + add_test_ssh_cert_request_provenance(&mut conflicting_request); + + let result = import_ssh_distribution_log_entries( + &store, + &[ + SshDistributionLogEntry::cert_request(request.clone()), + SshDistributionLogEntry::cert_request(conflicting_request), + ], + ); + + assert!(result.is_err()); + assert_eq!( + store + .get_ssh_cert_request(request.id.as_str()) + .expect("get request"), + None + ); + } + fn add_test_ssh_cert_request_provenance(request: &mut SshCertRequest) { let key = AgentKey::generate(); let signature = key @@ -12316,6 +12500,24 @@ mod tests { }); } + fn add_test_ssh_certificate_provenance(certificate: &mut SshCertificateRecord) { + let key = AgentKey::generate(); + let signature = key + .sign_canonical( + SSH_CERT_ISSUANCE_NAMESPACE, + &ssh_certificate_signing_payload(certificate), + ) + .expect("sign certificate"); + certificate.provenance = Some(SshRecordProvenance { + namespace: SSH_CERT_ISSUANCE_NAMESPACE.to_owned(), + signer_node: NodeId::new("node:right"), + signer_agent: key.agent_id().to_string(), + signer_public_key: key.public_key_hex(), + signature_hex: hex::encode(signature), + signed_at: UnixMillis(2), + }); + } + fn add_test_ssh_revocation_provenance(revocation: &mut SshRevocationEntry) { let key = AgentKey::generate(); let signature = key diff --git a/crates/geth-store/src/lib.rs b/crates/geth-store/src/lib.rs index 47413c2..e4b4113 100644 --- a/crates/geth-store/src/lib.rs +++ b/crates/geth-store/src/lib.rs @@ -1353,7 +1353,34 @@ impl Store { &self, request: &StoredSshCertRequest, ) -> Result<(), StoreError> { - self.conn.execute( + Self::insert_ssh_cert_request_tx(&self.conn, request) + } + + pub fn insert_ssh_distribution_records( + &self, + requests: &[StoredSshCertRequest], + certificates: &[StoredSshCertificate], + revocations: &[StoredSshRevocation], + ) -> Result<(), StoreError> { + let tx = self.conn.unchecked_transaction()?; + for request in requests { + Self::insert_ssh_cert_request_tx(&tx, request)?; + } + for certificate in certificates { + Self::insert_ssh_certificate_tx(&tx, certificate)?; + } + for revocation in revocations { + Self::insert_ssh_revocation_tx(&tx, revocation)?; + } + tx.commit()?; + Ok(()) + } + + fn insert_ssh_cert_request_tx( + conn: &rusqlite::Connection, + request: &StoredSshCertRequest, + ) -> Result<(), StoreError> { + conn.execute( r#"INSERT OR REPLACE INTO ssh_cert_requests( request_id, requester_node, public_key, public_key_fingerprint, cert_kind, principals_json, requested_validity, renewal_of, reason, status, created_at_ms, @@ -1438,7 +1465,14 @@ impl Store { &self, certificate: &StoredSshCertificate, ) -> Result<(), StoreError> { - self.conn.execute( + Self::insert_ssh_certificate_tx(&self.conn, certificate) + } + + fn insert_ssh_certificate_tx( + conn: &rusqlite::Connection, + certificate: &StoredSshCertificate, + ) -> Result<(), StoreError> { + conn.execute( r#"INSERT OR REPLACE INTO ssh_certificates( cert_id, request_id, certificate, certificate_fingerprint, imported_at_ms, provenance_json @@ -1502,7 +1536,14 @@ impl Store { &self, revocation: &StoredSshRevocation, ) -> Result<(), StoreError> { - self.conn.execute( + Self::insert_ssh_revocation_tx(&self.conn, revocation) + } + + fn insert_ssh_revocation_tx( + conn: &rusqlite::Connection, + revocation: &StoredSshRevocation, + ) -> Result<(), StoreError> { + conn.execute( r#"INSERT OR REPLACE INTO ssh_revocations( revocation_id, kind, target, reason, created_at_ms, published, provenance_json ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)"#, @@ -1993,6 +2034,65 @@ mod tests { ); } + #[test] + fn ssh_distribution_records_batch_commits_multiple_records() { + let store = Store::open_memory().expect("open"); + let request = StoredSshCertRequest { + request_id: "ssh-cert-request:batch".to_owned(), + requester_node: "node:laptop".to_owned(), + public_key: "ssh-ed25519 AAAA test".to_owned(), + public_key_fingerprint: "ssh:blake3:test".to_owned(), + cert_kind: "user".to_owned(), + principals: vec!["eric".to_owned()], + requested_validity: Some("+52w".to_owned()), + renewal_of: None, + reason: Some("renewal".to_owned()), + status: "pending".to_owned(), + created_at_ms: 1, + provenance_json: Some(r#"{"test":true}"#.to_owned()), + }; + let certificate = StoredSshCertificate { + cert_id: "ssh-cert:batch".to_owned(), + request_id: request.request_id.clone(), + certificate: "ssh-ed25519-cert-v01@openssh.com AAAA test".to_owned(), + certificate_fingerprint: "ssh:blake3:cert".to_owned(), + imported_at_ms: 2, + provenance_json: Some(r#"{"test":true}"#.to_owned()), + }; + let revocation = StoredSshRevocation { + revocation_id: "ssh-revocation:batch".to_owned(), + kind: "public-key".to_owned(), + target: "ssh:blake3:test".to_owned(), + reason: Some("lost key".to_owned()), + created_at_ms: 3, + published: true, + provenance_json: Some(r#"{"test":true}"#.to_owned()), + }; + + store + .insert_ssh_distribution_records( + std::slice::from_ref(&request), + std::slice::from_ref(&certificate), + std::slice::from_ref(&revocation), + ) + .expect("batch insert SSH records"); + + assert_eq!( + store + .get_ssh_cert_request("ssh-cert-request:batch") + .expect("get request"), + Some(request) + ); + assert_eq!( + store.list_ssh_certificates().expect("list certificates"), + vec![certificate] + ); + assert_eq!( + store.list_ssh_revocations().expect("list revocations"), + vec![revocation] + ); + } + #[test] fn peer_card_roundtrip() { let store = Store::open_memory().expect("open"); diff --git a/docs/production-readiness-roadmap.md b/docs/production-readiness-roadmap.md index 92a8bec..8915bb3 100644 --- a/docs/production-readiness-roadmap.md +++ b/docs/production-readiness-roadmap.md @@ -72,7 +72,7 @@ behavior. points. - `[ ]` No generic `geth-common` crate is introduced. -- `[~]` Extract live-sync engine. +- `[x]` Extract live-sync engine. Acceptance criteria: - `[x]` Sync cursor keys, cursor persistence, stream health recording, and local sync-status reduction live in a sync-focused module. @@ -127,12 +127,14 @@ it. - `[x]` Fresh database creation and repeated opens are idempotent. - `[x]` Old schema fixtures migrate to the current schema in tests. -- `[~]` Make multi-table writes transactional. +- `[x]` Make multi-table writes transactional. Acceptance criteria: - `[x]` Signed operation and signature imports commit atomically. - `[x]` Keychain/auth sync imports stage accepted records and commit accepted op/signature groups through store batch transactions. - - `[ ]` Multi-record sync imports cannot leave partial state after a local + - `[x]` SSH distribution sync imports stage accepted records and commit + requests, certificates, and revocations through one store transaction. + - `[x]` Multi-record sync imports cannot leave partial state after a local error. - `[x]` Tests cover failure injection for at least one multi-table path.