fix: batch ssh sync imports

This commit is contained in:
Eric Wendland 2026-07-05 23:19:05 +02:00
commit 7c0399b2ef
3 changed files with 321 additions and 17 deletions

View file

@ -11622,10 +11622,24 @@ fn stored_from_ssh_cert_request(request: &SshCertRequest) -> StoredSshCertReques
} }
} }
#[cfg(test)]
fn insert_ssh_cert_request_if_not_conflicting( fn insert_ssh_cert_request_if_not_conflicting(
store: &Store, store: &Store,
request: &SshCertRequest, request: &SshCertRequest,
) -> Result<bool, NodeError> { ) -> Result<bool, NodeError> {
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<Option<StoredSshCertRequest>, NodeError> {
let stored = stored_from_ssh_cert_request(request); let stored = stored_from_ssh_cert_request(request);
if let Some(existing) = store.get_ssh_cert_request(&stored.request_id)? { if let Some(existing) = store.get_ssh_cert_request(&stored.request_id)? {
if existing != stored { if existing != stored {
@ -11634,11 +11648,10 @@ fn insert_ssh_cert_request_if_not_conflicting(
"rejected conflicting SSH certificate request during sync" "rejected conflicting SSH certificate request during sync"
); );
} }
return Ok(false); return Ok(None);
} }
verify_ssh_cert_request_provenance(request)?; verify_ssh_cert_request_provenance(request)?;
store.insert_ssh_cert_request(&stored)?; Ok(Some(stored))
Ok(true)
} }
fn ssh_cert_request_from_stored(stored: StoredSshCertRequest) -> Result<SshCertRequest, NodeError> { fn ssh_cert_request_from_stored(stored: StoredSshCertRequest) -> Result<SshCertRequest, NodeError> {
@ -11685,10 +11698,24 @@ fn stored_from_ssh_certificate(certificate: &SshCertificateRecord) -> StoredSshC
} }
} }
#[cfg(test)]
fn insert_ssh_certificate_if_not_conflicting( fn insert_ssh_certificate_if_not_conflicting(
store: &Store, store: &Store,
certificate: &SshCertificateRecord, certificate: &SshCertificateRecord,
) -> Result<bool, NodeError> { ) -> Result<bool, NodeError> {
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<Option<StoredSshCertificate>, NodeError> {
let stored = stored_from_ssh_certificate(certificate); let stored = stored_from_ssh_certificate(certificate);
if let Some(existing) = store if let Some(existing) = store
.list_ssh_certificates()? .list_ssh_certificates()?
@ -11701,11 +11728,10 @@ fn insert_ssh_certificate_if_not_conflicting(
"rejected conflicting SSH certificate during sync" "rejected conflicting SSH certificate during sync"
); );
} }
return Ok(false); return Ok(None);
} }
verify_ssh_certificate_provenance(certificate)?; verify_ssh_certificate_provenance(certificate)?;
store.insert_ssh_certificate(&stored)?; Ok(Some(stored))
Ok(true)
} }
fn ssh_certificate_from_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( fn insert_ssh_revocation_if_not_conflicting(
store: &Store, store: &Store,
revocation: &SshRevocationEntry, revocation: &SshRevocationEntry,
) -> Result<bool, NodeError> { ) -> Result<bool, NodeError> {
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<Option<StoredSshRevocation>, NodeError> {
let stored = stored_from_ssh_revocation(revocation); let stored = stored_from_ssh_revocation(revocation);
if let Some(existing) = store if let Some(existing) = store
.list_ssh_revocations()? .list_ssh_revocations()?
@ -11757,10 +11797,69 @@ fn insert_ssh_revocation_if_not_conflicting(
"rejected conflicting SSH revocation during sync" "rejected conflicting SSH revocation during sync"
); );
} }
return Ok(false); return Ok(None);
} }
verify_ssh_revocation_provenance(revocation)?; verify_ssh_revocation_provenance(revocation)?;
store.insert_ssh_revocation(&stored)?; Ok(Some(stored))
}
fn stage_unique_ssh_cert_request(
staged: &mut Vec<StoredSshCertRequest>,
stored: StoredSshCertRequest,
) -> Result<bool, NodeError> {
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<StoredSshCertificate>,
stored: StoredSshCertificate,
) -> Result<bool, NodeError> {
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<StoredSshRevocation>,
stored: StoredSshRevocation,
) -> Result<bool, NodeError> {
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) Ok(true)
} }
@ -11809,25 +11908,39 @@ fn import_ssh_distribution_log_entries(
let mut requests_imported = 0; let mut requests_imported = 0;
let mut certificates_imported = 0; let mut certificates_imported = 0;
let mut revocations_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 { for entry in entries {
match entry { match entry {
SshDistributionLogEntry::CertRequest { request, .. } => { 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; requests_imported += 1;
} }
} }
SshDistributionLogEntry::Certificate { certificate, .. } => { 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; certificates_imported += 1;
} }
} }
SshDistributionLogEntry::Revocation { revocation, .. } => { 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; revocations_imported += 1;
} }
} }
} }
} }
store.insert_ssh_distribution_records(&requests, &certificates, &revocations)?;
Ok(( Ok((
requests_imported, requests_imported,
certificates_imported, certificates_imported,
@ -12263,6 +12376,37 @@ mod tests {
Some("original") 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 { let mut revocation = SshRevocationEntry {
id: geth_types::SshRevocationId::new("ssh-revocation:1"), id: geth_types::SshRevocationId::new("ssh-revocation:1"),
kind: SshRevocationKind::KeyId, 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) { fn add_test_ssh_cert_request_provenance(request: &mut SshCertRequest) {
let key = AgentKey::generate(); let key = AgentKey::generate();
let signature = key 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) { fn add_test_ssh_revocation_provenance(revocation: &mut SshRevocationEntry) {
let key = AgentKey::generate(); let key = AgentKey::generate();
let signature = key let signature = key

View file

@ -1353,7 +1353,34 @@ impl Store {
&self, &self,
request: &StoredSshCertRequest, request: &StoredSshCertRequest,
) -> Result<(), StoreError> { ) -> 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( r#"INSERT OR REPLACE INTO ssh_cert_requests(
request_id, requester_node, public_key, public_key_fingerprint, cert_kind, request_id, requester_node, public_key, public_key_fingerprint, cert_kind,
principals_json, requested_validity, renewal_of, reason, status, created_at_ms, principals_json, requested_validity, renewal_of, reason, status, created_at_ms,
@ -1438,7 +1465,14 @@ impl Store {
&self, &self,
certificate: &StoredSshCertificate, certificate: &StoredSshCertificate,
) -> Result<(), StoreError> { ) -> 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( r#"INSERT OR REPLACE INTO ssh_certificates(
cert_id, request_id, certificate, certificate_fingerprint, imported_at_ms, cert_id, request_id, certificate, certificate_fingerprint, imported_at_ms,
provenance_json provenance_json
@ -1502,7 +1536,14 @@ impl Store {
&self, &self,
revocation: &StoredSshRevocation, revocation: &StoredSshRevocation,
) -> Result<(), StoreError> { ) -> 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( r#"INSERT OR REPLACE INTO ssh_revocations(
revocation_id, kind, target, reason, created_at_ms, published, provenance_json revocation_id, kind, target, reason, created_at_ms, published, provenance_json
) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)"#, ) 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] #[test]
fn peer_card_roundtrip() { fn peer_card_roundtrip() {
let store = Store::open_memory().expect("open"); let store = Store::open_memory().expect("open");

View file

@ -72,7 +72,7 @@ behavior.
points. points.
- `[ ]` No generic `geth-common` crate is introduced. - `[ ]` No generic `geth-common` crate is introduced.
- `[~]` Extract live-sync engine. - `[x]` Extract live-sync engine.
Acceptance criteria: Acceptance criteria:
- `[x]` Sync cursor keys, cursor persistence, stream health recording, and - `[x]` Sync cursor keys, cursor persistence, stream health recording, and
local sync-status reduction live in a sync-focused module. 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]` Fresh database creation and repeated opens are idempotent.
- `[x]` Old schema fixtures migrate to the current schema in tests. - `[x]` Old schema fixtures migrate to the current schema in tests.
- `[~]` Make multi-table writes transactional. - `[x]` Make multi-table writes transactional.
Acceptance criteria: Acceptance criteria:
- `[x]` Signed operation and signature imports commit atomically. - `[x]` Signed operation and signature imports commit atomically.
- `[x]` Keychain/auth sync imports stage accepted records and commit accepted - `[x]` Keychain/auth sync imports stage accepted records and commit accepted
op/signature groups through store batch transactions. 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. error.
- `[x]` Tests cover failure injection for at least one multi-table path. - `[x]` Tests cover failure injection for at least one multi-table path.