reconciliation_finalization_commit.rs (43919B)
1 //! Atomic durable finalization of one independently verified attestation. 2 3 use core::fmt; 4 use std::error::Error; 5 6 use radroots_service_sqlite::{ 7 ServiceSqliteTransaction, ServiceSqliteTransactionError, ServiceSqliteTransactionErrorKind, 8 }; 9 use radroots_trade::evidence::{ 10 RadrootsTradeEvidenceOutcomeV1, RadrootsTradeEvidenceSourceCompletionV1, 11 RadrootsTradeEvidenceSourceRequirementV1, 12 }; 13 use sha2::{Digest, Sha256}; 14 15 use crate::{ 16 RhiPublicationAuthority, RhiPublicationMode, RhiPublicationOutboxId, 17 RhiReconciliationAttemptRepository, RhiReconciliationLease, RhiReconciliationOutcome, 18 RhiReconciliationUnixMilliseconds, RhiSignedEvidenceAttestation, RhiStateHostMode, 19 reconciliation_finalization::{FinalizationIdentity, validate_finalization_identity}, 20 }; 21 22 /// Exact version of the atomic reconciliation-finalization commit contract. 23 pub const RHI_RECONCILIATION_FINALIZATION_COMMIT_CONTRACT_VERSION: u32 = 1; 24 25 const OUTBOX_ID_DOMAIN: &[u8] = b"radroots.rhi.publication_outbox.v1\0"; 26 27 const INSERT_MANIFEST_SQL: &str = r#"INSERT INTO evidence_manifests ( 28 manifest_sha256, attempt_id, trade_id, trade_generation, 29 evidence_policy_sha256, canonical_manifest, observed_at_unix_s, 30 source_count, observation_count 31 ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)"#; 32 const INSERT_PROJECTION_SQL: &str = r#"INSERT INTO trade_projections ( 33 projection_sha256, manifest_sha256, shared_projection_sha256, 34 reducer_contract, reducer_contract_version, issue_count 35 ) VALUES (?, ?, ?, 'radroots.trade.reducer.v1', 1, ?)"#; 36 const INSERT_REPORT_SQL: &str = r#"INSERT INTO attestation_reports ( 37 statement_sha256, manifest_sha256, projection_sha256, trade_id, 38 claim_mutation_id, issuer_public_key, outcome, canonical_report, 39 observed_at_unix_s, supersedes_statement_sha256, supersedes_event_id 40 ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)"#; 41 const INSERT_SIGNED_EVENT_SQL: &str = r#"INSERT INTO signed_attestation_events ( 42 event_id, statement_sha256, event_sha256, issuer_public_key, 43 authored_at_unix_s, canonical_event_json 44 ) VALUES (?, ?, ?, ?, ?, ?)"#; 45 const INSERT_OUTBOX_SQL: &str = r#"INSERT INTO publication_outbox ( 46 outbox_id, event_id, event_sha256, publication_authority_sha256, 47 target_set_sha256, target_count, required_target_count, max_attempts, 48 initial_backoff_ms, maximum_backoff_ms, attempt_deadline_ms, 49 state, revision, next_attempt_unix_ms, lease_owner, 50 lease_expires_unix_ms, created_at_unix_ms, updated_at_unix_ms 51 ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 52 'pending', 1, ?, NULL, NULL, ?, ?)"#; 53 const INSERT_TARGET_SQL: &str = r#"INSERT INTO publication_targets ( 54 outbox_id, target_ordinal, relay_id, required, state, revision, 55 attempt_count, next_attempt_unix_ms, last_attempt_id, updated_at_unix_ms 56 ) VALUES (?, ?, ?, ?, 'pending', 1, 0, ?, NULL, ?)"#; 57 const COMPLETE_JOB_SQL: &str = r#"UPDATE reconciliation_jobs 58 SET state = 'completed', revision = revision + 1, next_attempt_unix_ms = NULL, 59 lease_owner = NULL, lease_expires_unix_ms = NULL, updated_at_unix_ms = ? 60 WHERE job_id = ? AND state = 'leased' AND revision = ? 61 AND lease_owner = ? AND lease_expires_unix_ms = ? 62 AND updated_at_unix_ms <= ?"#; 63 64 /// Stable source-free atomic-finalization failure class. 65 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 66 pub enum RhiReconciliationFinalizationCommitErrorKind { 67 InvalidMode, 68 InvalidInput, 69 LeaseLost, 70 GenerationConflict, 71 AttemptUnavailable, 72 PublicationQueueFull, 73 Conflict, 74 Storage, 75 CommitOutcomeUnknown, 76 } 77 78 impl RhiReconciliationFinalizationCommitErrorKind { 79 /// Returns the stable machine-readable failure code. 80 #[must_use] 81 pub const fn code(self) -> &'static str { 82 match self { 83 Self::InvalidMode => "reconciliation_finalization_commit_mode_invalid", 84 Self::InvalidInput => "reconciliation_finalization_commit_input_invalid", 85 Self::LeaseLost => "reconciliation_finalization_commit_lease_lost", 86 Self::GenerationConflict => "reconciliation_finalization_commit_generation_conflict", 87 Self::AttemptUnavailable => "reconciliation_finalization_commit_attempt_unavailable", 88 Self::PublicationQueueFull => "reconciliation_publication_queue_full", 89 Self::Conflict => "reconciliation_finalization_commit_conflict", 90 Self::Storage => "reconciliation_finalization_commit_storage_failed", 91 Self::CommitOutcomeUnknown => "reconciliation_finalization_commit_outcome_unknown", 92 } 93 } 94 } 95 96 /// Redacted source-free atomic-finalization failure. 97 #[derive(Clone, Copy, PartialEq, Eq)] 98 pub struct RhiReconciliationFinalizationCommitError { 99 kind: RhiReconciliationFinalizationCommitErrorKind, 100 } 101 102 impl RhiReconciliationFinalizationCommitError { 103 /// Returns the stable failure class. 104 #[must_use] 105 pub const fn kind(self) -> RhiReconciliationFinalizationCommitErrorKind { 106 self.kind 107 } 108 109 /// Returns the stable machine-readable failure code. 110 #[must_use] 111 pub const fn code(self) -> &'static str { 112 self.kind.code() 113 } 114 } 115 116 impl fmt::Display for RhiReconciliationFinalizationCommitError { 117 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 118 formatter.write_str(match self.kind { 119 RhiReconciliationFinalizationCommitErrorKind::InvalidMode => { 120 "RHI reconciliation finalization requires writable state" 121 } 122 RhiReconciliationFinalizationCommitErrorKind::InvalidInput => { 123 "RHI reconciliation finalization commit input is invalid" 124 } 125 RhiReconciliationFinalizationCommitErrorKind::LeaseLost => { 126 "RHI reconciliation finalization lease is no longer authoritative" 127 } 128 RhiReconciliationFinalizationCommitErrorKind::GenerationConflict => { 129 "RHI reconciliation finalization generation changed" 130 } 131 RhiReconciliationFinalizationCommitErrorKind::AttemptUnavailable => { 132 "RHI reconciliation finalization attempt is unavailable" 133 } 134 RhiReconciliationFinalizationCommitErrorKind::PublicationQueueFull => { 135 "RHI publication queue capacity is exhausted" 136 } 137 RhiReconciliationFinalizationCommitErrorKind::Conflict => { 138 "RHI reconciliation finalization conflicts with durable state" 139 } 140 RhiReconciliationFinalizationCommitErrorKind::Storage => { 141 "RHI reconciliation finalization transaction failed" 142 } 143 RhiReconciliationFinalizationCommitErrorKind::CommitOutcomeUnknown => { 144 "RHI reconciliation finalization outcome is unknown" 145 } 146 }) 147 } 148 } 149 150 impl fmt::Debug for RhiReconciliationFinalizationCommitError { 151 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 152 formatter 153 .debug_struct("RhiReconciliationFinalizationCommitError") 154 .field("kind", &self.kind) 155 .finish() 156 } 157 } 158 159 impl Error for RhiReconciliationFinalizationCommitError {} 160 161 /// Durable result of one exact atomic finalization commit. 162 #[derive(Clone, Copy, PartialEq, Eq)] 163 pub struct RhiReconciliationFinalizationCommitOutcome { 164 created: bool, 165 publication_mode: RhiPublicationMode, 166 target_count: u8, 167 outbox_id: Option<RhiPublicationOutboxId>, 168 } 169 170 impl RhiReconciliationFinalizationCommitOutcome { 171 /// Reports whether this call created the immutable finalization inventory. 172 #[must_use] 173 pub const fn created(self) -> bool { 174 self.created 175 } 176 177 /// Returns the explicit configured publication mode committed with the result. 178 #[must_use] 179 pub const fn publication_mode(self) -> RhiPublicationMode { 180 self.publication_mode 181 } 182 183 /// Returns the exact immutable target count, or zero when publication is disabled. 184 #[must_use] 185 pub const fn target_count(self) -> u8 { 186 self.target_count 187 } 188 189 /// Returns the immutable outbox identity when publication is required. 190 #[must_use] 191 pub const fn outbox_id(self) -> Option<RhiPublicationOutboxId> { 192 self.outbox_id 193 } 194 } 195 196 impl fmt::Debug for RhiReconciliationFinalizationCommitOutcome { 197 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 198 formatter 199 .debug_struct("RhiReconciliationFinalizationCommitOutcome") 200 .field("created", &self.created) 201 .field("publication_mode", &self.publication_mode) 202 .field("target_count", &self.target_count) 203 .field("outbox_id", &self.outbox_id.map(|_| "[redacted]")) 204 .finish() 205 } 206 } 207 208 impl RhiReconciliationAttemptRepository<'_> { 209 /// Atomically persists one sealed signed result and finalizes its exact job. 210 /// 211 /// The attestation is borrowed so a caller can reconcile a lost commit 212 /// acknowledgement by retrying the same sealed value. Exact prior success 213 /// is recognized before the consumed lease is checked. No relay, source, 214 /// clock, entropy, filesystem, or task operation occurs in this boundary. 215 pub async fn commit_finalization( 216 &self, 217 attestation: &RhiSignedEvidenceAttestation, 218 publication: &RhiPublicationAuthority, 219 now: RhiReconciliationUnixMilliseconds, 220 ) -> Result<RhiReconciliationFinalizationCommitOutcome, RhiReconciliationFinalizationCommitError> 221 { 222 if self.host().mode() != RhiStateHostMode::ReadWriteExisting { 223 return Err(failure( 224 RhiReconciliationFinalizationCommitErrorKind::InvalidMode, 225 )); 226 } 227 if publication.configuration_sha256() 228 != self.host().metadata().configuration_digest().as_bytes() 229 { 230 return Err(failure( 231 RhiReconciliationFinalizationCommitErrorKind::InvalidInput, 232 )); 233 } 234 let record = FinalizationRecord::from_inputs(attestation, publication, now) 235 .ok_or_else(|| failure(RhiReconciliationFinalizationCommitErrorKind::InvalidInput))?; 236 self.host() 237 .sqlite_host() 238 .transaction(move |transaction| { 239 Box::pin(async move { commit(transaction, &record).await }) 240 }) 241 .await 242 .map_err(map_transaction_error) 243 } 244 } 245 246 struct FinalizationTarget { 247 ordinal: u8, 248 relay_id: Box<str>, 249 required: bool, 250 } 251 252 struct FinalizationSource { 253 source_id: Box<str>, 254 required: bool, 255 completion: RadrootsTradeEvidenceSourceCompletionV1, 256 admitted_event_count: u32, 257 } 258 259 struct RequiredPublication { 260 outbox_id: [u8; 32], 261 authority_sha256: [u8; 32], 262 target_set_sha256: [u8; 32], 263 queue_capacity: u32, 264 maximum_attempts: u16, 265 initial_backoff_ms: u64, 266 maximum_backoff_ms: u64, 267 attempt_deadline_ms: u64, 268 target_count: u8, 269 targets: Box<[FinalizationTarget]>, 270 } 271 272 enum FinalizationPublication { 273 Disabled, 274 Required(RequiredPublication), 275 } 276 277 struct FinalizationRecord { 278 lease: RhiReconciliationLease, 279 identity: FinalizationIdentity, 280 now: RhiReconciliationUnixMilliseconds, 281 manifest_sha256: [u8; 32], 282 canonical_manifest: Box<[u8]>, 283 observed_at_unix_s: u64, 284 source_count: u8, 285 sources: Box<[FinalizationSource]>, 286 observation_count: u32, 287 projection_sha256: [u8; 32], 288 shared_projection_sha256: [u8; 32], 289 issue_count: u32, 290 statement_sha256: [u8; 32], 291 claim_mutation_id: [u8; 32], 292 issuer_public_key: [u8; 32], 293 outcome: &'static str, 294 canonical_report: Box<[u8]>, 295 supersession: Option<([u8; 32], [u8; 32])>, 296 event_id: [u8; 32], 297 event_sha256: [u8; 32], 298 authored_at_unix_s: u64, 299 canonical_event_json: Box<[u8]>, 300 publication: FinalizationPublication, 301 } 302 303 impl FinalizationRecord { 304 fn from_inputs( 305 attestation: &RhiSignedEvidenceAttestation, 306 publication: &RhiPublicationAuthority, 307 now: RhiReconciliationUnixMilliseconds, 308 ) -> Option<Self> { 309 let fence = attestation.fence(); 310 let (lease, identity) = fence.validation_parts(); 311 let evaluation = fence.evaluation(); 312 let projection = evaluation.projection(); 313 let manifest = projection.manifest(); 314 let projection_sha256 = projection.digest()?; 315 let shared_projection_sha256 = projection.shared_projection_digest()?; 316 let source_count = u8::try_from(manifest.source_count()).ok()?; 317 let sources = manifest 318 .inner() 319 .sources() 320 .iter() 321 .map(|source| { 322 let result = source.result(); 323 FinalizationSource { 324 source_id: source.source_id().as_str().into(), 325 required: matches!( 326 result.requirement(), 327 RadrootsTradeEvidenceSourceRequirementV1::Required 328 ), 329 completion: result.completion(), 330 admitted_event_count: result.admitted_event_count(), 331 } 332 }) 333 .collect::<Vec<_>>() 334 .into_boxed_slice(); 335 let observation_count = u32::try_from(manifest.observation_count()).ok()?; 336 let issue_count = u32::try_from(projection.issue_count()).ok()?; 337 let observed_at_unix_s = manifest.observed_at_unix_seconds(); 338 if attestation.report_observed_at_unix_seconds() != observed_at_unix_s 339 || [ 340 now.get(), 341 observed_at_unix_s, 342 attestation.created_at_unix_seconds(), 343 manifest.trade_generation(), 344 ] 345 .iter() 346 .any(|value| i64::try_from(*value).is_err()) 347 { 348 return None; 349 } 350 let publication = match publication.mode() { 351 RhiPublicationMode::Disabled => { 352 if !publication.targets().is_empty() || publication.retry_policy().is_some() { 353 return None; 354 } 355 FinalizationPublication::Disabled 356 } 357 RhiPublicationMode::Required => { 358 let retry = publication.retry_policy()?; 359 let targets = publication 360 .targets() 361 .iter() 362 .map(|target| FinalizationTarget { 363 ordinal: target.ordinal(), 364 relay_id: target.relay_id().into(), 365 required: target.required(), 366 }) 367 .collect::<Vec<_>>(); 368 if targets.is_empty() 369 || targets.len() > crate::RHI_PUBLICATION_MAX_TARGETS 370 || targets 371 .iter() 372 .enumerate() 373 .any(|(ordinal, target)| usize::from(target.ordinal) != ordinal) 374 { 375 return None; 376 } 377 let target_count = u8::try_from(targets.len()).ok()?; 378 FinalizationPublication::Required(RequiredPublication { 379 outbox_id: outbox_id( 380 attestation.event_id(), 381 publication.authority_sha256(), 382 publication.target_set_sha256(), 383 ), 384 authority_sha256: *publication.authority_sha256(), 385 target_set_sha256: *publication.target_set_sha256(), 386 queue_capacity: publication.queue_capacity(), 387 maximum_attempts: retry.maximum_attempts(), 388 initial_backoff_ms: retry.initial_backoff_milliseconds(), 389 maximum_backoff_ms: retry.maximum_backoff_milliseconds(), 390 attempt_deadline_ms: retry.attempt_deadline_milliseconds(), 391 target_count, 392 targets: targets.into_boxed_slice(), 393 }) 394 } 395 }; 396 Some(Self { 397 lease, 398 identity, 399 now, 400 manifest_sha256: manifest.digest(), 401 canonical_manifest: manifest.canonical_bytes().into(), 402 observed_at_unix_s, 403 source_count, 404 sources, 405 observation_count, 406 projection_sha256, 407 shared_projection_sha256, 408 issue_count, 409 statement_sha256: attestation.statement_digest(), 410 claim_mutation_id: attestation.claim_mutation_id_bytes(), 411 issuer_public_key: attestation.issuer_public_key_bytes(), 412 outcome: outcome_code(attestation.outcome()), 413 canonical_report: attestation.canonical_report_bytes().into(), 414 supersession: attestation.supersession_bytes(), 415 event_id: *attestation.event_id(), 416 event_sha256: *attestation.signed_event_sha256(), 417 authored_at_unix_s: attestation.created_at_unix_seconds(), 418 canonical_event_json: attestation.signed_event_bytes().into(), 419 publication, 420 }) 421 } 422 423 fn outcome(&self, created: bool) -> RhiReconciliationFinalizationCommitOutcome { 424 let (publication_mode, target_count, outbox_id) = match &self.publication { 425 FinalizationPublication::Disabled => (RhiPublicationMode::Disabled, 0, None), 426 FinalizationPublication::Required(required) => ( 427 RhiPublicationMode::Required, 428 required.target_count, 429 Some(RhiPublicationOutboxId::from_committed_bytes( 430 required.outbox_id, 431 )), 432 ), 433 }; 434 RhiReconciliationFinalizationCommitOutcome { 435 created, 436 publication_mode, 437 target_count, 438 outbox_id, 439 } 440 } 441 } 442 443 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 444 enum OperationError { 445 LeaseLost, 446 GenerationConflict, 447 AttemptUnavailable, 448 QueueFull, 449 Conflict, 450 Storage, 451 } 452 453 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 454 enum ExistingState { 455 Absent, 456 Exact, 457 } 458 459 async fn commit( 460 transaction: &mut ServiceSqliteTransaction<'_>, 461 record: &FinalizationRecord, 462 ) -> Result<RhiReconciliationFinalizationCommitOutcome, OperationError> { 463 match reconcile_existing(transaction, record).await? { 464 ExistingState::Exact => return Ok(record.outcome(false)), 465 ExistingState::Absent => {} 466 } 467 validate_finalization_identity(transaction, record.lease, record.identity, record.now) 468 .await 469 .map_err(|error| match error { 470 crate::reconciliation_finalization::FinalizationOperationError::LeaseLost => { 471 OperationError::LeaseLost 472 } 473 crate::reconciliation_finalization::FinalizationOperationError::GenerationConflict => { 474 OperationError::GenerationConflict 475 } 476 crate::reconciliation_finalization::FinalizationOperationError::AttemptUnavailable => { 477 OperationError::AttemptUnavailable 478 } 479 crate::reconciliation_finalization::FinalizationOperationError::Storage => { 480 OperationError::Storage 481 } 482 })?; 483 validate_source_inventory(transaction, record).await?; 484 validate_advanced_checkpoints(transaction, record).await?; 485 validate_supersession(transaction, record).await?; 486 if let FinalizationPublication::Required(required) = &record.publication { 487 let active: i64 = 488 sqlx::query_scalar("SELECT COUNT(*) FROM publication_outbox WHERE state != 'complete'") 489 .fetch_one(&mut *transaction) 490 .await 491 .map_err(|_| OperationError::Storage)?; 492 let active = u64::try_from(active).map_err(|_| OperationError::Storage)?; 493 if active >= u64::from(required.queue_capacity) { 494 return Err(OperationError::QueueFull); 495 } 496 } 497 498 insert_finalization(transaction, record).await?; 499 let lease_job = record.lease.job(); 500 let result = sqlx::query(COMPLETE_JOB_SQL) 501 .bind(i64_value(record.now.get())?) 502 .bind(record.identity.job_id.as_bytes().as_slice()) 503 .bind(i64_value(lease_job.revision())?) 504 .bind(record.lease.owner_bytes().as_slice()) 505 .bind(i64_value(record.lease.lease_expires().get())?) 506 .bind(i64_value(record.now.get())?) 507 .execute(&mut *transaction) 508 .await 509 .map_err(|_| OperationError::Storage)?; 510 if result.rows_affected() != 1 { 511 return Err(OperationError::LeaseLost); 512 } 513 Ok(record.outcome(true)) 514 } 515 516 async fn validate_source_inventory( 517 transaction: &mut ServiceSqliteTransaction<'_>, 518 record: &FinalizationRecord, 519 ) -> Result<(), OperationError> { 520 let attempt_count: i64 = sqlx::query_scalar( 521 "SELECT COUNT(*) FROM evidence_reconciliations WHERE attempt_id = ? AND source_count = ?", 522 ) 523 .bind(record.identity.attempt_id.as_bytes().as_slice()) 524 .bind(i64::from(record.source_count)) 525 .fetch_one(&mut *transaction) 526 .await 527 .map_err(|_| OperationError::Storage)?; 528 if attempt_count != 1 { 529 return Err(OperationError::AttemptUnavailable); 530 } 531 let count: i64 = sqlx::query_scalar( 532 "SELECT COUNT(*) FROM evidence_reconciliation_sources WHERE attempt_id = ?", 533 ) 534 .bind(record.identity.attempt_id.as_bytes().as_slice()) 535 .fetch_one(&mut *transaction) 536 .await 537 .map_err(|_| OperationError::Storage)?; 538 if count != i64::from(record.source_count) { 539 return Err(OperationError::AttemptUnavailable); 540 } 541 for (ordinal, source) in record.sources.iter().enumerate() { 542 let completion = match source.completion { 543 RadrootsTradeEvidenceSourceCompletionV1::Complete => "complete", 544 RadrootsTradeEvidenceSourceCompletionV1::Unsupported => "unsupported", 545 RadrootsTradeEvidenceSourceCompletionV1::Incomplete => "incomplete", 546 }; 547 let source_match: i64 = sqlx::query_scalar( 548 r#"SELECT COUNT(*) FROM evidence_reconciliation_sources 549 WHERE attempt_id = ? AND source_ordinal = ? AND source_id = ? 550 AND required = ? AND accepted_event_count = ? 551 AND CASE ? 552 WHEN 'complete' THEN completion = 'complete' 553 WHEN 'unsupported' THEN completion = 'unsupported' 554 ELSE completion IN ( 555 'incomplete_timeout', 'incomplete_unavailable', 556 'incomplete_resource_limit', 'incomplete_unknown' 557 ) 558 END"#, 559 ) 560 .bind(record.identity.attempt_id.as_bytes().as_slice()) 561 .bind(i64::try_from(ordinal).map_err(|_| OperationError::Storage)?) 562 .bind(source.source_id.as_ref()) 563 .bind(i64::from(source.required)) 564 .bind(i64::from(source.admitted_event_count)) 565 .bind(completion) 566 .fetch_one(&mut *transaction) 567 .await 568 .map_err(|_| OperationError::Storage)?; 569 if source_match != 1 { 570 return Err(OperationError::AttemptUnavailable); 571 } 572 } 573 Ok(()) 574 } 575 576 async fn validate_advanced_checkpoints( 577 transaction: &mut ServiceSqliteTransaction<'_>, 578 record: &FinalizationRecord, 579 ) -> Result<(), OperationError> { 580 let bad_checkpoint: i64 = sqlx::query_scalar( 581 r#"SELECT COUNT(*) 582 FROM evidence_reconciliation_sources AS source 583 LEFT JOIN relay_checkpoints AS checkpoint 584 ON checkpoint.source_id = source.source_id 585 AND checkpoint.evidence_policy_sha256 = ? 586 AND checkpoint.trade_id = source.trade_id 587 WHERE source.attempt_id = ? AND source.checkpoint_advanced = 1 588 AND (checkpoint.cursor_created_at_unix_s IS NULL 589 OR checkpoint.cursor_event_id IS NULL 590 OR checkpoint.cursor_created_at_unix_s != source.candidate_created_at_unix_s 591 OR checkpoint.cursor_event_id != source.candidate_event_id)"#, 592 ) 593 .bind(record.identity.policy_digest.as_bytes().as_slice()) 594 .bind(record.identity.attempt_id.as_bytes().as_slice()) 595 .fetch_one(&mut *transaction) 596 .await 597 .map_err(|_| OperationError::Storage)?; 598 if bad_checkpoint == 0 { 599 Ok(()) 600 } else { 601 Err(OperationError::GenerationConflict) 602 } 603 } 604 605 async fn insert_finalization( 606 transaction: &mut ServiceSqliteTransaction<'_>, 607 record: &FinalizationRecord, 608 ) -> Result<(), OperationError> { 609 execute_one( 610 sqlx::query(INSERT_MANIFEST_SQL) 611 .bind(record.manifest_sha256.as_slice()) 612 .bind(record.identity.attempt_id.as_bytes().as_slice()) 613 .bind(record.identity.trade_id.as_bytes().as_slice()) 614 .bind(i64_value(record.identity.generation)?) 615 .bind(record.identity.policy_digest.as_bytes().as_slice()) 616 .bind(record.canonical_manifest.as_ref()) 617 .bind(i64_value(record.observed_at_unix_s)?) 618 .bind(i64::from(record.source_count)) 619 .bind(i64::from(record.observation_count)), 620 transaction, 621 ) 622 .await?; 623 execute_one( 624 sqlx::query(INSERT_PROJECTION_SQL) 625 .bind(record.projection_sha256.as_slice()) 626 .bind(record.manifest_sha256.as_slice()) 627 .bind(record.shared_projection_sha256.as_slice()) 628 .bind(i64::from(record.issue_count)), 629 transaction, 630 ) 631 .await?; 632 let supersedes_statement = record.supersession.as_ref().map(|value| value.0.as_slice()); 633 let supersedes_event = record.supersession.as_ref().map(|value| value.1.as_slice()); 634 execute_one( 635 sqlx::query(INSERT_REPORT_SQL) 636 .bind(record.statement_sha256.as_slice()) 637 .bind(record.manifest_sha256.as_slice()) 638 .bind(record.projection_sha256.as_slice()) 639 .bind(record.identity.trade_id.as_bytes().as_slice()) 640 .bind(record.claim_mutation_id.as_slice()) 641 .bind(record.issuer_public_key.as_slice()) 642 .bind(record.outcome) 643 .bind(record.canonical_report.as_ref()) 644 .bind(i64_value(record.observed_at_unix_s)?) 645 .bind(supersedes_statement) 646 .bind(supersedes_event), 647 transaction, 648 ) 649 .await?; 650 execute_one( 651 sqlx::query(INSERT_SIGNED_EVENT_SQL) 652 .bind(record.event_id.as_slice()) 653 .bind(record.statement_sha256.as_slice()) 654 .bind(record.event_sha256.as_slice()) 655 .bind(record.issuer_public_key.as_slice()) 656 .bind(i64_value(record.authored_at_unix_s)?) 657 .bind(record.canonical_event_json.as_ref()), 658 transaction, 659 ) 660 .await?; 661 if let FinalizationPublication::Required(required) = &record.publication { 662 let required_count = required 663 .targets 664 .iter() 665 .filter(|target| target.required) 666 .count(); 667 execute_one( 668 sqlx::query(INSERT_OUTBOX_SQL) 669 .bind(required.outbox_id.as_slice()) 670 .bind(record.event_id.as_slice()) 671 .bind(record.event_sha256.as_slice()) 672 .bind(required.authority_sha256.as_slice()) 673 .bind(required.target_set_sha256.as_slice()) 674 .bind(i64::try_from(required.targets.len()).map_err(|_| OperationError::Storage)?) 675 .bind(i64::try_from(required_count).map_err(|_| OperationError::Storage)?) 676 .bind(i64::from(required.maximum_attempts)) 677 .bind(i64_value(required.initial_backoff_ms)?) 678 .bind(i64_value(required.maximum_backoff_ms)?) 679 .bind(i64_value(required.attempt_deadline_ms)?) 680 .bind(i64_value(record.now.get())?) 681 .bind(i64_value(record.now.get())?) 682 .bind(i64_value(record.now.get())?), 683 transaction, 684 ) 685 .await?; 686 for target in &required.targets { 687 execute_one( 688 sqlx::query(INSERT_TARGET_SQL) 689 .bind(required.outbox_id.as_slice()) 690 .bind(i64::from(target.ordinal)) 691 .bind(target.relay_id.as_ref()) 692 .bind(i64::from(target.required)) 693 .bind(i64_value(record.now.get())?) 694 .bind(i64_value(record.now.get())?), 695 transaction, 696 ) 697 .await?; 698 } 699 } 700 Ok(()) 701 } 702 703 async fn reconcile_existing( 704 transaction: &mut ServiceSqliteTransaction<'_>, 705 record: &FinalizationRecord, 706 ) -> Result<ExistingState, OperationError> { 707 let footprint: i64 = sqlx::query_scalar( 708 r#"SELECT 709 (SELECT COUNT(*) FROM evidence_manifests 710 WHERE manifest_sha256 = ? OR attempt_id = ?) 711 + (SELECT COUNT(*) FROM trade_projections 712 WHERE projection_sha256 = ? OR manifest_sha256 = ?) 713 + (SELECT COUNT(*) FROM attestation_reports WHERE statement_sha256 = ?) 714 + (SELECT COUNT(*) FROM signed_attestation_events 715 WHERE event_id = ? OR statement_sha256 = ? OR event_sha256 = ?) 716 + (SELECT COUNT(*) FROM publication_outbox WHERE event_id = ?) 717 + (SELECT COUNT(*) FROM reconciliation_jobs 718 WHERE job_id = ? AND state = 'completed')"#, 719 ) 720 .bind(record.manifest_sha256.as_slice()) 721 .bind(record.identity.attempt_id.as_bytes().as_slice()) 722 .bind(record.projection_sha256.as_slice()) 723 .bind(record.manifest_sha256.as_slice()) 724 .bind(record.statement_sha256.as_slice()) 725 .bind(record.event_id.as_slice()) 726 .bind(record.statement_sha256.as_slice()) 727 .bind(record.event_sha256.as_slice()) 728 .bind(record.event_id.as_slice()) 729 .bind(record.identity.job_id.as_bytes().as_slice()) 730 .fetch_one(&mut *transaction) 731 .await 732 .map_err(|_| OperationError::Storage)?; 733 if footprint == 0 { 734 return Ok(ExistingState::Absent); 735 } 736 if manifest_matches(transaction, record).await? 737 && projection_matches(transaction, record).await? 738 && report_matches(transaction, record).await? 739 && event_matches(transaction, record).await? 740 && publication_matches(transaction, record).await? 741 && completed_job_matches(transaction, record).await? 742 { 743 validate_source_inventory(transaction, record).await?; 744 validate_supersession(transaction, record).await?; 745 Ok(ExistingState::Exact) 746 } else { 747 Err(OperationError::Conflict) 748 } 749 } 750 751 async fn validate_supersession( 752 transaction: &mut ServiceSqliteTransaction<'_>, 753 record: &FinalizationRecord, 754 ) -> Result<(), OperationError> { 755 let Some((statement_sha256, event_id)) = record.supersession else { 756 return Ok(()); 757 }; 758 let count: i64 = sqlx::query_scalar( 759 r#"SELECT COUNT(*) 760 FROM attestation_reports AS report 761 JOIN signed_attestation_events AS event 762 ON event.statement_sha256 = report.statement_sha256 763 WHERE report.statement_sha256 = ? AND event.event_id = ?"#, 764 ) 765 .bind(statement_sha256.as_slice()) 766 .bind(event_id.as_slice()) 767 .fetch_one(&mut *transaction) 768 .await 769 .map_err(|_| OperationError::Storage)?; 770 if count == 1 { 771 Ok(()) 772 } else { 773 Err(OperationError::Conflict) 774 } 775 } 776 777 async fn manifest_matches( 778 transaction: &mut ServiceSqliteTransaction<'_>, 779 record: &FinalizationRecord, 780 ) -> Result<bool, OperationError> { 781 match_count( 782 sqlx::query_scalar( 783 r#"SELECT COUNT(*) FROM evidence_manifests 784 WHERE manifest_sha256 = ? AND attempt_id = ? AND trade_id = ? 785 AND trade_generation = ? AND evidence_policy_sha256 = ? 786 AND canonical_manifest = ? AND observed_at_unix_s = ? 787 AND source_count = ? AND observation_count = ?"#, 788 ) 789 .bind(record.manifest_sha256.as_slice()) 790 .bind(record.identity.attempt_id.as_bytes().as_slice()) 791 .bind(record.identity.trade_id.as_bytes().as_slice()) 792 .bind(i64_value(record.identity.generation)?) 793 .bind(record.identity.policy_digest.as_bytes().as_slice()) 794 .bind(record.canonical_manifest.as_ref()) 795 .bind(i64_value(record.observed_at_unix_s)?) 796 .bind(i64::from(record.source_count)) 797 .bind(i64::from(record.observation_count)), 798 transaction, 799 ) 800 .await 801 } 802 803 async fn projection_matches( 804 transaction: &mut ServiceSqliteTransaction<'_>, 805 record: &FinalizationRecord, 806 ) -> Result<bool, OperationError> { 807 match_count( 808 sqlx::query_scalar( 809 r#"SELECT COUNT(*) FROM trade_projections 810 WHERE projection_sha256 = ? AND manifest_sha256 = ? 811 AND shared_projection_sha256 = ? 812 AND reducer_contract = 'radroots.trade.reducer.v1' 813 AND reducer_contract_version = 1 AND issue_count = ?"#, 814 ) 815 .bind(record.projection_sha256.as_slice()) 816 .bind(record.manifest_sha256.as_slice()) 817 .bind(record.shared_projection_sha256.as_slice()) 818 .bind(i64::from(record.issue_count)), 819 transaction, 820 ) 821 .await 822 } 823 824 async fn report_matches( 825 transaction: &mut ServiceSqliteTransaction<'_>, 826 record: &FinalizationRecord, 827 ) -> Result<bool, OperationError> { 828 let supersedes_statement = record.supersession.as_ref().map(|value| value.0.as_slice()); 829 let supersedes_event = record.supersession.as_ref().map(|value| value.1.as_slice()); 830 match_count( 831 sqlx::query_scalar( 832 r#"SELECT COUNT(*) FROM attestation_reports 833 WHERE statement_sha256 = ? AND manifest_sha256 = ? AND projection_sha256 = ? 834 AND trade_id = ? AND claim_mutation_id = ? AND issuer_public_key = ? 835 AND outcome = ? AND canonical_report = ? AND observed_at_unix_s = ? 836 AND supersedes_statement_sha256 IS ? AND supersedes_event_id IS ?"#, 837 ) 838 .bind(record.statement_sha256.as_slice()) 839 .bind(record.manifest_sha256.as_slice()) 840 .bind(record.projection_sha256.as_slice()) 841 .bind(record.identity.trade_id.as_bytes().as_slice()) 842 .bind(record.claim_mutation_id.as_slice()) 843 .bind(record.issuer_public_key.as_slice()) 844 .bind(record.outcome) 845 .bind(record.canonical_report.as_ref()) 846 .bind(i64_value(record.observed_at_unix_s)?) 847 .bind(supersedes_statement) 848 .bind(supersedes_event), 849 transaction, 850 ) 851 .await 852 } 853 854 async fn event_matches( 855 transaction: &mut ServiceSqliteTransaction<'_>, 856 record: &FinalizationRecord, 857 ) -> Result<bool, OperationError> { 858 match_count( 859 sqlx::query_scalar( 860 r#"SELECT COUNT(*) FROM signed_attestation_events 861 WHERE event_id = ? AND statement_sha256 = ? AND event_sha256 = ? 862 AND issuer_public_key = ? AND authored_at_unix_s = ? 863 AND canonical_event_json = ?"#, 864 ) 865 .bind(record.event_id.as_slice()) 866 .bind(record.statement_sha256.as_slice()) 867 .bind(record.event_sha256.as_slice()) 868 .bind(record.issuer_public_key.as_slice()) 869 .bind(i64_value(record.authored_at_unix_s)?) 870 .bind(record.canonical_event_json.as_ref()), 871 transaction, 872 ) 873 .await 874 } 875 876 async fn publication_matches( 877 transaction: &mut ServiceSqliteTransaction<'_>, 878 record: &FinalizationRecord, 879 ) -> Result<bool, OperationError> { 880 match &record.publication { 881 FinalizationPublication::Disabled => { 882 let count: i64 = 883 sqlx::query_scalar("SELECT COUNT(*) FROM publication_outbox WHERE event_id = ?") 884 .bind(record.event_id.as_slice()) 885 .fetch_one(&mut *transaction) 886 .await 887 .map_err(|_| OperationError::Storage)?; 888 Ok(count == 0) 889 } 890 FinalizationPublication::Required(required) => { 891 let required_count = required 892 .targets 893 .iter() 894 .filter(|target| target.required) 895 .count(); 896 let outbox = match_count( 897 sqlx::query_scalar( 898 r#"SELECT COUNT(*) FROM publication_outbox 899 WHERE outbox_id = ? AND event_id = ? AND event_sha256 = ? 900 AND publication_authority_sha256 = ? AND target_set_sha256 = ? 901 AND target_count = ? AND required_target_count = ? AND max_attempts = ? 902 AND initial_backoff_ms = ? AND maximum_backoff_ms = ? 903 AND attempt_deadline_ms = ?"#, 904 ) 905 .bind(required.outbox_id.as_slice()) 906 .bind(record.event_id.as_slice()) 907 .bind(record.event_sha256.as_slice()) 908 .bind(required.authority_sha256.as_slice()) 909 .bind(required.target_set_sha256.as_slice()) 910 .bind(i64::try_from(required.targets.len()).map_err(|_| OperationError::Storage)?) 911 .bind(i64::try_from(required_count).map_err(|_| OperationError::Storage)?) 912 .bind(i64::from(required.maximum_attempts)) 913 .bind(i64_value(required.initial_backoff_ms)?) 914 .bind(i64_value(required.maximum_backoff_ms)?) 915 .bind(i64_value(required.attempt_deadline_ms)?), 916 transaction, 917 ) 918 .await?; 919 if !outbox { 920 return Ok(false); 921 } 922 let count: i64 = 923 sqlx::query_scalar("SELECT COUNT(*) FROM publication_targets WHERE outbox_id = ?") 924 .bind(required.outbox_id.as_slice()) 925 .fetch_one(&mut *transaction) 926 .await 927 .map_err(|_| OperationError::Storage)?; 928 if count 929 != i64::try_from(required.targets.len()).map_err(|_| OperationError::Storage)? 930 { 931 return Ok(false); 932 } 933 for target in &required.targets { 934 if !match_count( 935 sqlx::query_scalar( 936 r#"SELECT COUNT(*) FROM publication_targets 937 WHERE outbox_id = ? AND target_ordinal = ? AND relay_id = ? AND required = ?"#, 938 ) 939 .bind(required.outbox_id.as_slice()) 940 .bind(i64::from(target.ordinal)) 941 .bind(target.relay_id.as_ref()) 942 .bind(i64::from(target.required)), 943 transaction, 944 ) 945 .await? 946 { 947 return Ok(false); 948 } 949 } 950 Ok(true) 951 } 952 } 953 } 954 955 async fn completed_job_matches( 956 transaction: &mut ServiceSqliteTransaction<'_>, 957 record: &FinalizationRecord, 958 ) -> Result<bool, OperationError> { 959 let job = record.lease.job(); 960 match_count( 961 sqlx::query_scalar( 962 r#"SELECT COUNT(*) FROM reconciliation_jobs 963 WHERE job_id = ? AND trade_id = ? AND input_generation = ? 964 AND evidence_policy_sha256 = ? AND state = 'completed' 965 AND revision = ? AND attempt_count = ? AND failure_count = ? 966 AND max_attempts = ? AND next_attempt_unix_ms IS NULL 967 AND lease_owner IS NULL AND lease_expires_unix_ms IS NULL"#, 968 ) 969 .bind(record.identity.job_id.as_bytes().as_slice()) 970 .bind(record.identity.trade_id.as_bytes().as_slice()) 971 .bind(i64_value(record.identity.generation)?) 972 .bind(record.identity.policy_digest.as_bytes().as_slice()) 973 .bind(i64_value( 974 job.revision() 975 .checked_add(1) 976 .ok_or(OperationError::Storage)?, 977 )?) 978 .bind(i64::from(job.attempt_count())) 979 .bind(i64::from(job.failure_count())) 980 .bind(i64::from(job.policy().max_attempts())), 981 transaction, 982 ) 983 .await 984 } 985 986 async fn execute_one<'query>( 987 query: sqlx::query::Query<'query, sqlx::Sqlite, sqlx::sqlite::SqliteArguments>, 988 transaction: &mut ServiceSqliteTransaction<'_>, 989 ) -> Result<(), OperationError> { 990 let result = query 991 .execute(&mut *transaction) 992 .await 993 .map_err(|_| OperationError::Storage)?; 994 if result.rows_affected() == 1 { 995 Ok(()) 996 } else { 997 Err(OperationError::Storage) 998 } 999 } 1000 1001 async fn match_count<'query>( 1002 query: sqlx::query::QueryScalar<'query, sqlx::Sqlite, i64, sqlx::sqlite::SqliteArguments>, 1003 transaction: &mut ServiceSqliteTransaction<'_>, 1004 ) -> Result<bool, OperationError> { 1005 query 1006 .fetch_one(&mut *transaction) 1007 .await 1008 .map(|count| count == 1) 1009 .map_err(|_| OperationError::Storage) 1010 } 1011 1012 const fn outcome_code(outcome: RhiReconciliationOutcome) -> &'static str { 1013 match outcome { 1014 RadrootsTradeEvidenceOutcomeV1::Valid => "valid", 1015 RadrootsTradeEvidenceOutcomeV1::Invalid => "invalid", 1016 RadrootsTradeEvidenceOutcomeV1::Indeterminate => "indeterminate", 1017 } 1018 } 1019 1020 fn outbox_id( 1021 event_id: &[u8; 32], 1022 authority_sha256: &[u8; 32], 1023 target_set_sha256: &[u8; 32], 1024 ) -> [u8; 32] { 1025 let mut digest = Sha256::new(); 1026 digest.update(OUTBOX_ID_DOMAIN); 1027 digest.update(event_id); 1028 digest.update(authority_sha256); 1029 digest.update(target_set_sha256); 1030 digest.finalize().into() 1031 } 1032 1033 fn i64_value(value: u64) -> Result<i64, OperationError> { 1034 i64::try_from(value).map_err(|_| OperationError::Storage) 1035 } 1036 1037 fn map_transaction_error( 1038 error: ServiceSqliteTransactionError<OperationError>, 1039 ) -> RhiReconciliationFinalizationCommitError { 1040 if error.kind() == ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown { 1041 return failure(RhiReconciliationFinalizationCommitErrorKind::CommitOutcomeUnknown); 1042 } 1043 failure(match error.operation_error().copied() { 1044 Some(OperationError::LeaseLost) => RhiReconciliationFinalizationCommitErrorKind::LeaseLost, 1045 Some(OperationError::GenerationConflict) => { 1046 RhiReconciliationFinalizationCommitErrorKind::GenerationConflict 1047 } 1048 Some(OperationError::AttemptUnavailable) => { 1049 RhiReconciliationFinalizationCommitErrorKind::AttemptUnavailable 1050 } 1051 Some(OperationError::QueueFull) => { 1052 RhiReconciliationFinalizationCommitErrorKind::PublicationQueueFull 1053 } 1054 Some(OperationError::Conflict) => RhiReconciliationFinalizationCommitErrorKind::Conflict, 1055 Some(OperationError::Storage) | None => { 1056 RhiReconciliationFinalizationCommitErrorKind::Storage 1057 } 1058 }) 1059 } 1060 1061 const fn failure( 1062 kind: RhiReconciliationFinalizationCommitErrorKind, 1063 ) -> RhiReconciliationFinalizationCommitError { 1064 RhiReconciliationFinalizationCommitError { kind } 1065 } 1066 1067 #[cfg(test)] 1068 mod tests { 1069 use super::*; 1070 1071 #[test] 1072 fn outbox_identity_is_exact_and_domain_separated() { 1073 assert_eq!( 1074 outbox_id(&[0x11; 32], &[0x22; 32], &[0x33; 32]), 1075 [ 1076 0xf7, 0x80, 0x86, 0x80, 0xd5, 0xa6, 0x84, 0x1f, 0x2e, 0xcd, 0xaf, 0xb0, 0xc3, 0x1c, 1077 0x71, 0xb7, 0x4f, 0x0f, 0xf4, 0x02, 0x7f, 0x85, 0xb2, 0x49, 0x1f, 0x00, 0xc8, 0xb9, 1078 0xa4, 0xf5, 0x9e, 0x37, 1079 ] 1080 ); 1081 assert_ne!( 1082 outbox_id(&[0x11; 32], &[0x22; 32], &[0x33; 32]), 1083 outbox_id(&[0x11; 32], &[0x22; 32], &[0x34; 32]) 1084 ); 1085 } 1086 1087 #[test] 1088 fn errors_are_closed_source_free_and_redacted() { 1089 for kind in [ 1090 RhiReconciliationFinalizationCommitErrorKind::InvalidMode, 1091 RhiReconciliationFinalizationCommitErrorKind::InvalidInput, 1092 RhiReconciliationFinalizationCommitErrorKind::LeaseLost, 1093 RhiReconciliationFinalizationCommitErrorKind::GenerationConflict, 1094 RhiReconciliationFinalizationCommitErrorKind::AttemptUnavailable, 1095 RhiReconciliationFinalizationCommitErrorKind::PublicationQueueFull, 1096 RhiReconciliationFinalizationCommitErrorKind::Conflict, 1097 RhiReconciliationFinalizationCommitErrorKind::Storage, 1098 RhiReconciliationFinalizationCommitErrorKind::CommitOutcomeUnknown, 1099 ] { 1100 let error = failure(kind); 1101 assert_eq!(error.kind(), kind); 1102 assert!(error.code().starts_with("reconciliation_")); 1103 assert!(Error::source(&error).is_none()); 1104 let rendered = format!("{error} {error:?}"); 1105 assert!(!rendered.contains("relay-primary")); 1106 assert!(!rendered.contains("11111111")); 1107 assert!(!rendered.contains("SELECT")); 1108 } 1109 } 1110 }