reconciliation_commit.rs (31706B)
1 //! Atomic durable reconciliation-source result commit. 2 3 use core::{cmp::Ordering, fmt}; 4 use std::error::Error; 5 6 use radroots_event::id::TradeId; 7 use radroots_service_sqlite::{ 8 ServiceSqliteTransaction, ServiceSqliteTransactionError, ServiceSqliteTransactionErrorKind, 9 }; 10 use sha2::{Digest, Sha256}; 11 use sqlx::Row; 12 13 use crate::{ 14 RhiReconciliationAttemptPlan, RhiReconciliationAttemptRepository, RhiReconciliationLease, 15 RhiReconciliationSourceCursorEvidence, RhiReconciliationSourceReplay, RhiStateHostMode, 16 RhiTradeSourceCursor, 17 reconciliation_job::{LeaseValidationError, validate_exact_lease}, 18 reconciliation_manifest::{RhiCommittedManifestMaterial, committed_manifest_material}, 19 reconciliation_replay::{ 20 RhiReconciliationReplayCommitFact, RhiReconciliationReplayCommitParts, 21 committed_cursor_evidence, 22 }, 23 source_ingest::{ 24 SourceOperationError, advance_dirty_generation, compare_cursor, read_checkpoint, 25 read_dirty, write_checkpoint, 26 }, 27 state_trade::{PersistenceOperationError, RhiTradeSourceObservation, persist}, 28 }; 29 30 /// Exact version of the atomic reconciliation-source commit contract. 31 pub const RHI_RECONCILIATION_COMMIT_CONTRACT_VERSION: u32 = 1; 32 33 const SOURCE_INVENTORY_DIGEST_DOMAIN: &[u8] = b"radroots.rhi.reconciliation_source_inventory.v1\0"; 34 35 const COUNT_ATTEMPT_SQL: &str = 36 "SELECT COUNT(*) AS row_count FROM evidence_reconciliations WHERE attempt_id = ?"; 37 const MATCH_ATTEMPT_SQL: &str = r#"SELECT COUNT(*) AS row_count 38 FROM evidence_reconciliations 39 WHERE attempt_id = ? AND job_id = ? AND trade_id = ? AND input_generation = ? 40 AND evidence_policy_sha256 = ? AND attempt_started_unix_ms = ? 41 AND deadline_unix_ms = ? AND source_count = ?"#; 42 const COUNT_SOURCE_RESULTS_SQL: &str = r#"SELECT COUNT(*) AS row_count 43 FROM evidence_reconciliation_sources WHERE attempt_id = ?"#; 44 const MATCH_SOURCE_RESULT_SQL: &str = r#"SELECT COUNT(*) AS row_count 45 FROM evidence_reconciliation_sources 46 WHERE attempt_id = ? AND request_id = ? AND source_ordinal = ? 47 AND source_id = ? AND trade_id = ? AND required = ? AND selector_sha256 = ? 48 AND replay_id = ? AND completion = ? AND started_unix_ms = ? AND finished_unix_ms = ? 49 AND accepted_event_count = ? AND accepted_event_bytes = ? 50 AND accepted_inventory_sha256 = ? 51 AND duplicate_observation_count = ? AND first_observed_unix_s IS ? 52 AND prior_cursor_created_at_unix_s IS ? AND prior_cursor_event_id IS ? 53 AND overlap_seconds = ? AND inclusive_since_unix_s = ? 54 AND candidate_created_at_unix_s IS ? AND candidate_event_id IS ? 55 AND checkpoint_advanced = ?"#; 56 const INSERT_ATTEMPT_SQL: &str = r#"INSERT INTO evidence_reconciliations ( 57 attempt_id, job_id, trade_id, input_generation, evidence_policy_sha256, 58 attempt_started_unix_ms, deadline_unix_ms, source_count 59 ) VALUES (?, ?, ?, ?, ?, ?, ?, ?)"#; 60 const INSERT_SOURCE_RESULT_SQL: &str = r#"INSERT INTO evidence_reconciliation_sources ( 61 attempt_id, request_id, source_ordinal, source_id, trade_id, required, 62 selector_sha256, replay_id, completion, started_unix_ms, finished_unix_ms, 63 accepted_event_count, accepted_event_bytes, accepted_inventory_sha256, 64 duplicate_observation_count, 65 first_observed_unix_s, prior_cursor_created_at_unix_s, prior_cursor_event_id, 66 overlap_seconds, inclusive_since_unix_s, candidate_created_at_unix_s, 67 candidate_event_id, checkpoint_advanced 68 ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)"#; 69 70 /// Stable source-free failure classification for one atomic result commit. 71 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 72 pub enum RhiReconciliationCommitErrorKind { 73 InvalidMode, 74 InvalidInput, 75 LeaseLost, 76 GenerationConflict, 77 Conflict, 78 Storage, 79 CommitOutcomeUnknown, 80 } 81 82 impl RhiReconciliationCommitErrorKind { 83 /// Returns the stable machine-readable failure code. 84 #[must_use] 85 pub const fn code(self) -> &'static str { 86 match self { 87 Self::InvalidMode => "reconciliation_commit_mode_invalid", 88 Self::InvalidInput => "reconciliation_commit_input_invalid", 89 Self::LeaseLost => "reconciliation_commit_lease_lost", 90 Self::GenerationConflict => "reconciliation_commit_generation_conflict", 91 Self::Conflict => "reconciliation_commit_conflict", 92 Self::Storage => "reconciliation_commit_storage_failed", 93 Self::CommitOutcomeUnknown => "reconciliation_commit_outcome_unknown", 94 } 95 } 96 } 97 98 /// Redacted source-free atomic reconciliation commit failure. 99 #[derive(Clone, Copy, PartialEq, Eq)] 100 pub struct RhiReconciliationCommitError { 101 kind: RhiReconciliationCommitErrorKind, 102 } 103 104 impl RhiReconciliationCommitError { 105 /// Returns the stable failure class. 106 #[must_use] 107 pub const fn kind(self) -> RhiReconciliationCommitErrorKind { 108 self.kind 109 } 110 111 /// Returns the stable machine-readable failure code. 112 #[must_use] 113 pub const fn code(self) -> &'static str { 114 self.kind.code() 115 } 116 } 117 118 impl fmt::Display for RhiReconciliationCommitError { 119 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 120 formatter.write_str(match self.kind { 121 RhiReconciliationCommitErrorKind::InvalidMode => { 122 "RHI reconciliation commit requires writable state" 123 } 124 RhiReconciliationCommitErrorKind::InvalidInput => { 125 "RHI reconciliation commit input is invalid" 126 } 127 RhiReconciliationCommitErrorKind::LeaseLost => { 128 "RHI reconciliation lease is no longer authoritative" 129 } 130 RhiReconciliationCommitErrorKind::GenerationConflict => { 131 "RHI reconciliation generation or cursor changed" 132 } 133 RhiReconciliationCommitErrorKind::Conflict => { 134 "RHI reconciliation evidence conflicts with durable state" 135 } 136 RhiReconciliationCommitErrorKind::Storage => { 137 "RHI reconciliation commit transaction failed" 138 } 139 RhiReconciliationCommitErrorKind::CommitOutcomeUnknown => { 140 "RHI reconciliation commit outcome is unknown" 141 } 142 }) 143 } 144 } 145 146 impl fmt::Debug for RhiReconciliationCommitError { 147 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 148 formatter 149 .debug_struct("RhiReconciliationCommitError") 150 .field("kind", &self.kind) 151 .finish() 152 } 153 } 154 155 impl Error for RhiReconciliationCommitError {} 156 157 /// Durable result of one exact bounded reconciliation-attempt commit. 158 pub struct RhiReconciliationSourceCommitOutcome { 159 created: bool, 160 source_result_count: u32, 161 checkpoint_advance_count: u32, 162 dirty_generation_advanced: bool, 163 committed_cursors: Box<[RhiReconciliationSourceCursorEvidence]>, 164 pub(crate) manifest_material: RhiCommittedManifestMaterial, 165 } 166 167 impl RhiReconciliationSourceCommitOutcome { 168 /// Reports whether this call created the immutable attempt inventory. 169 #[must_use] 170 pub const fn created(&self) -> bool { 171 self.created 172 } 173 174 /// Returns the exact committed source-result count. 175 #[must_use] 176 pub const fn source_result_count(&self) -> u32 { 177 self.source_result_count 178 } 179 180 /// Returns the exact cursor checkpoint-advance count. 181 #[must_use] 182 pub const fn checkpoint_advance_count(&self) -> u32 { 183 self.checkpoint_advance_count 184 } 185 186 /// Reports whether newly durable mutation/event evidence advanced dirty state once. 187 #[must_use] 188 pub const fn dirty_generation_advanced(&self) -> bool { 189 self.dirty_generation_advanced 190 } 191 192 /// Returns sealed cursor evidence minted only after durable commit confirmation. 193 #[must_use] 194 pub fn committed_cursors(&self) -> &[RhiReconciliationSourceCursorEvidence] { 195 &self.committed_cursors 196 } 197 } 198 199 impl fmt::Debug for RhiReconciliationSourceCommitOutcome { 200 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 201 formatter 202 .debug_struct("RhiReconciliationSourceCommitOutcome") 203 .field("created", &self.created) 204 .field("source_result_count", &self.source_result_count) 205 .field("checkpoint_advance_count", &self.checkpoint_advance_count) 206 .field("dirty_generation_advanced", &self.dirty_generation_advanced) 207 .field("committed_cursors", &self.committed_cursors.len()) 208 .field("manifest_material", &"[sealed]") 209 .finish() 210 } 211 } 212 213 impl RhiReconciliationAttemptRepository<'_> { 214 /// Atomically commits one exact bounded source-result inventory. 215 pub async fn commit_source_replays<I>( 216 &self, 217 lease: RhiReconciliationLease, 218 plan: RhiReconciliationAttemptPlan, 219 replays: I, 220 ) -> Result<RhiReconciliationSourceCommitOutcome, RhiReconciliationCommitError> 221 where 222 I: IntoIterator<Item = RhiReconciliationSourceReplay>, 223 { 224 if self.host().mode() != RhiStateHostMode::ReadWriteExisting { 225 return Err(failure(RhiReconciliationCommitErrorKind::InvalidMode)); 226 } 227 let replays = replays 228 .into_iter() 229 .take(plan.requests().len().saturating_add(1)) 230 .collect::<Vec<_>>(); 231 if replays.len() != plan.requests().len() 232 || plan.job_id() != lease.job().id() 233 || plan.input_generation() != lease.job().input_generation() 234 || plan.evidence_policy_digest() != lease.job().evidence_policy_digest() 235 || plan.deadline() > lease.lease_expires() 236 { 237 return Err(failure(RhiReconciliationCommitErrorKind::InvalidInput)); 238 } 239 let parts = replays 240 .into_iter() 241 .map(RhiReconciliationSourceReplay::into_commit_parts) 242 .collect::<Vec<_>>(); 243 if parts.iter().zip(plan.requests()).any(|(replay, request)| { 244 replay.request_id != request.id() 245 || replay.result.request_id() != request.id() 246 || replay.source_id.as_ref() != request.source_id() 247 || replay.trade_id != request.trade_id() 248 || replay.required != request.required() 249 || replay.trade_id != lease.job().trade_id() 250 || replay.policy_digest != *plan.evidence_policy_digest().as_bytes() 251 || replay.selector_digest != *request.selector_digest().as_bytes() 252 || replay.result.finished_at() >= lease.lease_expires() 253 || replay.facts.len() != replay.result.accepted_event_count() as usize 254 || replay 255 .eligible_cursor 256 .is_some_and(|_| replay.result.finished_at().get() / 1_000 == 0) 257 }) { 258 return Err(failure(RhiReconciliationCommitErrorKind::InvalidInput)); 259 } 260 let manifest_material = committed_manifest_material(&plan, &parts) 261 .map_err(|()| failure(RhiReconciliationCommitErrorKind::InvalidInput))?; 262 263 let raw = self 264 .host() 265 .sqlite_host() 266 .transaction(move |transaction| { 267 Box::pin(async move { commit(transaction, lease, plan, parts).await }) 268 }) 269 .await 270 .map_err(map_transaction_error)?; 271 let committed_cursors = raw 272 .cursor_scopes 273 .into_vec() 274 .into_iter() 275 .map(|scope| { 276 committed_cursor_evidence( 277 scope.source_id, 278 scope.trade_id, 279 scope.policy_digest, 280 scope.selector_digest, 281 scope.cursor, 282 ) 283 }) 284 .collect::<Vec<_>>() 285 .into_boxed_slice(); 286 Ok(RhiReconciliationSourceCommitOutcome { 287 created: raw.created, 288 source_result_count: raw.source_result_count, 289 checkpoint_advance_count: raw.checkpoint_advance_count, 290 dirty_generation_advanced: raw.dirty_generation_advanced, 291 committed_cursors, 292 manifest_material, 293 }) 294 } 295 } 296 297 struct RawCommitOutcome { 298 created: bool, 299 source_result_count: u32, 300 checkpoint_advance_count: u32, 301 dirty_generation_advanced: bool, 302 cursor_scopes: Box<[CommittedCursorScope]>, 303 } 304 305 struct CommittedCursorScope { 306 source_id: Box<str>, 307 trade_id: TradeId, 308 policy_digest: [u8; 32], 309 selector_digest: [u8; 32], 310 cursor: RhiTradeSourceCursor, 311 } 312 313 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 314 enum CommitOperationError { 315 LeaseLost, 316 GenerationConflict, 317 Conflict, 318 Storage, 319 } 320 321 async fn commit( 322 transaction: &mut ServiceSqliteTransaction<'_>, 323 lease: RhiReconciliationLease, 324 plan: RhiReconciliationAttemptPlan, 325 parts: Vec<RhiReconciliationReplayCommitParts>, 326 ) -> Result<RawCommitOutcome, CommitOperationError> { 327 if reconcile_existing(transaction, &plan, &parts).await? { 328 return raw_outcome(false, false, &parts); 329 } 330 validate_exact_lease(transaction, lease) 331 .await 332 .map_err(|error| match error { 333 LeaseValidationError::LeaseLost => CommitOperationError::LeaseLost, 334 LeaseValidationError::Storage => CommitOperationError::Storage, 335 })?; 336 let policy = plan.evidence_policy_digest(); 337 let dirty = read_dirty(transaction, lease.job().trade_id()) 338 .await 339 .map_err(map_source_error)? 340 .filter(|state| state.generation.get() == plan.input_generation() && state.policy == policy) 341 .ok_or(CommitOperationError::GenerationConflict)?; 342 let mut checkpoints = Vec::with_capacity(parts.len()); 343 for part in &parts { 344 let checkpoint = 345 read_checkpoint(transaction, part.source_id.as_ref(), policy, part.trade_id) 346 .await 347 .map_err(map_source_error)?; 348 if checkpoint.map(|value| value.cursor) != part.prior_cursor { 349 return Err(CommitOperationError::GenerationConflict); 350 } 351 checkpoints.push(checkpoint); 352 } 353 354 insert_attempt(transaction, &plan).await?; 355 let mut new_relevant_evidence = false; 356 for part in &parts { 357 for fact in &part.facts { 358 let observation = RhiTradeSourceObservation::from_persistence_parts( 359 part.source_id.clone(), 360 policy, 361 &fact.record, 362 fact.observed_at, 363 ); 364 let outcome = persist(transaction, &fact.record, &observation) 365 .await 366 .map_err(map_persistence_error)?; 367 new_relevant_evidence |= outcome.mutation_inserted() || outcome.signed_event_inserted(); 368 } 369 } 370 371 if new_relevant_evidence { 372 let updated_at = parts 373 .iter() 374 .map(|part| part.result.finished_at().get() / 1_000) 375 .max() 376 .ok_or(CommitOperationError::Storage)?; 377 advance_dirty_generation( 378 transaction, 379 lease.job().trade_id(), 380 policy, 381 updated_at, 382 Some(dirty), 383 ) 384 .await 385 .map_err(map_source_error)?; 386 } 387 388 for ((ordinal, part), checkpoint) in parts.iter().enumerate().zip(checkpoints) { 389 if let Some(candidate) = part.eligible_cursor { 390 write_checkpoint( 391 transaction, 392 part.source_id.as_ref(), 393 policy, 394 part.trade_id, 395 checkpoint, 396 candidate, 397 part.result.finished_at().get() / 1_000, 398 ) 399 .await 400 .map_err(map_source_error)?; 401 } 402 insert_source_result(transaction, plan.id().as_bytes(), ordinal, part).await?; 403 } 404 raw_outcome(true, new_relevant_evidence, &parts) 405 } 406 407 fn raw_outcome( 408 created: bool, 409 dirty_generation_advanced: bool, 410 parts: &[RhiReconciliationReplayCommitParts], 411 ) -> Result<RawCommitOutcome, CommitOperationError> { 412 let cursor_scopes = parts 413 .iter() 414 .filter_map(|part| { 415 part.eligible_cursor.map(|cursor| CommittedCursorScope { 416 source_id: part.source_id.clone(), 417 trade_id: part.trade_id, 418 policy_digest: part.policy_digest, 419 selector_digest: part.selector_digest, 420 cursor, 421 }) 422 }) 423 .collect::<Vec<_>>() 424 .into_boxed_slice(); 425 Ok(RawCommitOutcome { 426 created, 427 source_result_count: u32::try_from(parts.len()) 428 .map_err(|_| CommitOperationError::Storage)?, 429 checkpoint_advance_count: u32::try_from(cursor_scopes.len()) 430 .map_err(|_| CommitOperationError::Storage)?, 431 dirty_generation_advanced, 432 cursor_scopes, 433 }) 434 } 435 436 async fn reconcile_existing( 437 transaction: &mut ServiceSqliteTransaction<'_>, 438 plan: &RhiReconciliationAttemptPlan, 439 parts: &[RhiReconciliationReplayCommitParts], 440 ) -> Result<bool, CommitOperationError> { 441 let count = count_query(transaction, COUNT_ATTEMPT_SQL, plan.id().as_bytes()).await?; 442 if count == 0 { 443 return Ok(false); 444 } 445 if count != 1 || match_attempt(transaction, plan).await? != 1 { 446 return Err(CommitOperationError::Conflict); 447 } 448 let source_count = 449 count_query(transaction, COUNT_SOURCE_RESULTS_SQL, plan.id().as_bytes()).await?; 450 if source_count != parts.len() as u64 { 451 return Err(CommitOperationError::Conflict); 452 } 453 for (ordinal, part) in parts.iter().enumerate() { 454 if match_source_result(transaction, plan.id().as_bytes(), ordinal, part).await? != 1 { 455 return Err(CommitOperationError::Conflict); 456 } 457 if let Some(candidate) = part.eligible_cursor { 458 let checkpoint = read_checkpoint( 459 transaction, 460 part.source_id.as_ref(), 461 plan.evidence_policy_digest(), 462 part.trade_id, 463 ) 464 .await 465 .map_err(map_source_error)? 466 .ok_or(CommitOperationError::Conflict)?; 467 if compare_cursor(checkpoint.cursor, candidate) == Ordering::Less { 468 return Err(CommitOperationError::Conflict); 469 } 470 } 471 } 472 Ok(true) 473 } 474 475 async fn count_query( 476 transaction: &mut ServiceSqliteTransaction<'_>, 477 sql: &'static str, 478 id: &[u8; 32], 479 ) -> Result<u64, CommitOperationError> { 480 sqlx::query(sql) 481 .bind(id.as_slice()) 482 .fetch_one(&mut *transaction) 483 .await 484 .map_err(|_| CommitOperationError::Storage)? 485 .try_get::<i64, _>("row_count") 486 .ok() 487 .and_then(|value| u64::try_from(value).ok()) 488 .ok_or(CommitOperationError::Storage) 489 } 490 491 async fn match_attempt( 492 transaction: &mut ServiceSqliteTransaction<'_>, 493 plan: &RhiReconciliationAttemptPlan, 494 ) -> Result<u64, CommitOperationError> { 495 let trade = plan 496 .requests() 497 .first() 498 .ok_or(CommitOperationError::Storage)? 499 .trade_id(); 500 count_row( 501 sqlx::query(MATCH_ATTEMPT_SQL) 502 .bind(plan.id().as_bytes().as_slice()) 503 .bind(plan.job_id().as_bytes().as_slice()) 504 .bind(trade.as_bytes().as_slice()) 505 .bind(i64_value(plan.input_generation())?) 506 .bind(plan.evidence_policy_digest().as_bytes().as_slice()) 507 .bind(i64_value(plan.attempt_started_at().get())?) 508 .bind(i64_value(plan.deadline().get())?) 509 .bind(i64::try_from(plan.requests().len()).map_err(|_| CommitOperationError::Storage)?), 510 transaction, 511 ) 512 .await 513 } 514 515 async fn match_source_result( 516 transaction: &mut ServiceSqliteTransaction<'_>, 517 attempt_id: &[u8; 32], 518 ordinal: usize, 519 part: &RhiReconciliationReplayCommitParts, 520 ) -> Result<u64, CommitOperationError> { 521 let prior_created = part 522 .prior_cursor 523 .map(|cursor| i64_value(cursor.created_at_unix_seconds())) 524 .transpose()?; 525 let prior_event = part.prior_cursor.map(|cursor| cursor.event_id()); 526 let candidate_created = part 527 .cursor_candidate 528 .map(|cursor| i64_value(cursor.created_at_unix_seconds())) 529 .transpose()?; 530 let candidate_event = part.cursor_candidate.map(|cursor| cursor.event_id()); 531 let inventory_digest = accepted_inventory_digest(part)?; 532 count_row( 533 sqlx::query(MATCH_SOURCE_RESULT_SQL) 534 .bind(attempt_id.as_slice()) 535 .bind(part.request_id.as_bytes().as_slice()) 536 .bind(i64::try_from(ordinal).map_err(|_| CommitOperationError::Storage)?) 537 .bind(part.source_id.as_ref()) 538 .bind(part.trade_id.as_bytes().as_slice()) 539 .bind(i64::from(part.required)) 540 .bind(part.selector_digest.as_slice()) 541 .bind(part.replay_id.as_bytes().as_slice()) 542 .bind(part.result.outcome().code()) 543 .bind(i64_value(part.result.started_at().get())?) 544 .bind(i64_value(part.result.finished_at().get())?) 545 .bind(i64::from(part.result.accepted_event_count())) 546 .bind(i64_value(part.result.accepted_event_bytes())?) 547 .bind(inventory_digest.as_slice()) 548 .bind(i64::from(part.duplicate_observations)) 549 .bind( 550 part.first_observed_at 551 .map(|value| i64_value(value.get())) 552 .transpose()?, 553 ) 554 .bind(prior_created) 555 .bind(prior_event.as_ref().map(<[u8; 32]>::as_slice)) 556 .bind(i64_value(part.overlap_seconds)?) 557 .bind(i64_value(part.since_unix_seconds)?) 558 .bind(candidate_created) 559 .bind(candidate_event.as_ref().map(<[u8; 32]>::as_slice)) 560 .bind(i64::from(part.eligible_cursor.is_some())), 561 transaction, 562 ) 563 .await 564 } 565 566 async fn insert_attempt( 567 transaction: &mut ServiceSqliteTransaction<'_>, 568 plan: &RhiReconciliationAttemptPlan, 569 ) -> Result<(), CommitOperationError> { 570 let trade = plan 571 .requests() 572 .first() 573 .ok_or(CommitOperationError::Storage)? 574 .trade_id(); 575 let result = sqlx::query(INSERT_ATTEMPT_SQL) 576 .bind(plan.id().as_bytes().as_slice()) 577 .bind(plan.job_id().as_bytes().as_slice()) 578 .bind(trade.as_bytes().as_slice()) 579 .bind(i64_value(plan.input_generation())?) 580 .bind(plan.evidence_policy_digest().as_bytes().as_slice()) 581 .bind(i64_value(plan.attempt_started_at().get())?) 582 .bind(i64_value(plan.deadline().get())?) 583 .bind(i64::try_from(plan.requests().len()).map_err(|_| CommitOperationError::Storage)?) 584 .execute(&mut *transaction) 585 .await 586 .map_err(|_| CommitOperationError::Storage)?; 587 if result.rows_affected() == 1 { 588 Ok(()) 589 } else { 590 Err(CommitOperationError::Storage) 591 } 592 } 593 594 async fn insert_source_result( 595 transaction: &mut ServiceSqliteTransaction<'_>, 596 attempt_id: &[u8; 32], 597 ordinal: usize, 598 part: &RhiReconciliationReplayCommitParts, 599 ) -> Result<(), CommitOperationError> { 600 let prior_created = part 601 .prior_cursor 602 .map(|cursor| i64_value(cursor.created_at_unix_seconds())) 603 .transpose()?; 604 let prior_event = part.prior_cursor.map(|cursor| cursor.event_id()); 605 let candidate_created = part 606 .cursor_candidate 607 .map(|cursor| i64_value(cursor.created_at_unix_seconds())) 608 .transpose()?; 609 let candidate_event = part.cursor_candidate.map(|cursor| cursor.event_id()); 610 let inventory_digest = accepted_inventory_digest(part)?; 611 let result = sqlx::query(INSERT_SOURCE_RESULT_SQL) 612 .bind(attempt_id.as_slice()) 613 .bind(part.request_id.as_bytes().as_slice()) 614 .bind(i64::try_from(ordinal).map_err(|_| CommitOperationError::Storage)?) 615 .bind(part.source_id.as_ref()) 616 .bind(part.trade_id.as_bytes().as_slice()) 617 .bind(i64::from(part.required)) 618 .bind(part.selector_digest.as_slice()) 619 .bind(part.replay_id.as_bytes().as_slice()) 620 .bind(part.result.outcome().code()) 621 .bind(i64_value(part.result.started_at().get())?) 622 .bind(i64_value(part.result.finished_at().get())?) 623 .bind(i64::from(part.result.accepted_event_count())) 624 .bind(i64_value(part.result.accepted_event_bytes())?) 625 .bind(inventory_digest.as_slice()) 626 .bind(i64::from(part.duplicate_observations)) 627 .bind( 628 part.first_observed_at 629 .map(|value| i64_value(value.get())) 630 .transpose()?, 631 ) 632 .bind(prior_created) 633 .bind(prior_event.as_ref().map(<[u8; 32]>::as_slice)) 634 .bind(i64_value(part.overlap_seconds)?) 635 .bind(i64_value(part.since_unix_seconds)?) 636 .bind(candidate_created) 637 .bind(candidate_event.as_ref().map(<[u8; 32]>::as_slice)) 638 .bind(i64::from(part.eligible_cursor.is_some())) 639 .execute(&mut *transaction) 640 .await 641 .map_err(|_| CommitOperationError::Storage)?; 642 if result.rows_affected() == 1 { 643 Ok(()) 644 } else { 645 Err(CommitOperationError::Storage) 646 } 647 } 648 649 fn accepted_inventory_digest( 650 part: &RhiReconciliationReplayCommitParts, 651 ) -> Result<[u8; 32], CommitOperationError> { 652 accepted_inventory_digest_for_facts(&part.facts) 653 } 654 655 pub(crate) fn committed_inventory_digest( 656 part: &RhiReconciliationReplayCommitParts, 657 ) -> Option<[u8; 32]> { 658 accepted_inventory_digest(part).ok() 659 } 660 661 fn accepted_inventory_digest_for_facts( 662 facts: &[RhiReconciliationReplayCommitFact], 663 ) -> Result<[u8; 32], CommitOperationError> { 664 let mut digest = Sha256::new(); 665 digest.update(SOURCE_INVENTORY_DIGEST_DOMAIN); 666 digest.update( 667 u32::try_from(facts.len()) 668 .map_err(|_| CommitOperationError::Storage)? 669 .to_be_bytes(), 670 ); 671 for fact in facts { 672 let record = &fact.record; 673 digest.update(record.mutation_id); 674 digest.update(record.trade_id); 675 update_framed(&mut digest, record.contract_id.as_bytes())?; 676 digest.update(record.schema_version.to_be_bytes()); 677 digest.update(record.event_id); 678 digest.update(record.event_signature); 679 digest.update(record.author_pubkey); 680 digest.update(record.event_kind.to_be_bytes()); 681 digest.update(record.authored_at_unix_s.to_be_bytes()); 682 update_framed(&mut digest, &record.canonical_content)?; 683 update_framed(&mut digest, &record.canonical_event_json)?; 684 digest.update(fact.observed_at.get().to_be_bytes()); 685 } 686 Ok(digest.finalize().into()) 687 } 688 689 fn update_framed(digest: &mut Sha256, bytes: &[u8]) -> Result<(), CommitOperationError> { 690 digest.update( 691 u64::try_from(bytes.len()) 692 .map_err(|_| CommitOperationError::Storage)? 693 .to_be_bytes(), 694 ); 695 digest.update(bytes); 696 Ok(()) 697 } 698 699 async fn count_row<'query>( 700 query: sqlx::query::Query<'query, sqlx::Sqlite, sqlx::sqlite::SqliteArguments>, 701 transaction: &mut ServiceSqliteTransaction<'_>, 702 ) -> Result<u64, CommitOperationError> { 703 query 704 .fetch_one(&mut *transaction) 705 .await 706 .map_err(|_| CommitOperationError::Storage)? 707 .try_get::<i64, _>("row_count") 708 .ok() 709 .and_then(|value| u64::try_from(value).ok()) 710 .ok_or(CommitOperationError::Storage) 711 } 712 713 fn i64_value(value: u64) -> Result<i64, CommitOperationError> { 714 i64::try_from(value).map_err(|_| CommitOperationError::Storage) 715 } 716 717 fn map_persistence_error(error: PersistenceOperationError) -> CommitOperationError { 718 match error { 719 PersistenceOperationError::MutationConflict 720 | PersistenceOperationError::SignedEventConflict => CommitOperationError::Conflict, 721 PersistenceOperationError::Storage => CommitOperationError::Storage, 722 } 723 } 724 725 fn map_source_error(error: SourceOperationError) -> CommitOperationError { 726 match error { 727 SourceOperationError::GenerationConflict => CommitOperationError::GenerationConflict, 728 SourceOperationError::Persistence(error) => map_persistence_error(error), 729 SourceOperationError::Storage => CommitOperationError::Storage, 730 } 731 } 732 733 fn map_transaction_error( 734 error: ServiceSqliteTransactionError<CommitOperationError>, 735 ) -> RhiReconciliationCommitError { 736 if error.kind() == ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown { 737 return failure(RhiReconciliationCommitErrorKind::CommitOutcomeUnknown); 738 } 739 failure(match error.operation_error().copied() { 740 Some(CommitOperationError::LeaseLost) => RhiReconciliationCommitErrorKind::LeaseLost, 741 Some(CommitOperationError::GenerationConflict) => { 742 RhiReconciliationCommitErrorKind::GenerationConflict 743 } 744 Some(CommitOperationError::Conflict) => RhiReconciliationCommitErrorKind::Conflict, 745 Some(CommitOperationError::Storage) | None => RhiReconciliationCommitErrorKind::Storage, 746 }) 747 } 748 749 const fn failure(kind: RhiReconciliationCommitErrorKind) -> RhiReconciliationCommitError { 750 RhiReconciliationCommitError { kind } 751 } 752 753 #[cfg(test)] 754 mod tests { 755 use super::*; 756 757 fn fact(observed_at: u64, content: &[u8]) -> RhiReconciliationReplayCommitFact { 758 RhiReconciliationReplayCommitFact { 759 record: crate::state_trade::PersistenceRecord { 760 mutation_id: [0x11; 32], 761 trade_id: [0x22; 16], 762 contract_id: "radroots.trade.proposal.v1", 763 schema_version: 1, 764 event_id: [0x33; 32], 765 event_signature: [0x44; 64], 766 author_pubkey: [0x55; 32], 767 event_kind: 3_470, 768 authored_at_unix_s: 1_784_347_200, 769 canonical_content: content.into(), 770 canonical_event_json: br#"{"id":"example"}"#.as_slice().into(), 771 }, 772 observed_at: crate::RhiTradeMutationObservedAtUnixSeconds::new(observed_at) 773 .expect("observation"), 774 } 775 } 776 777 #[test] 778 fn source_inventory_digest_is_exact_framed_and_provenance_sensitive() { 779 assert_eq!( 780 accepted_inventory_digest_for_facts(&[]).expect("empty digest"), 781 [ 782 0xc6, 0x77, 0x8e, 0xb5, 0x38, 0xe6, 0x1a, 0x6d, 0xa4, 0xe5, 0xba, 0xe1, 0x9c, 0xb1, 783 0xb3, 0xf9, 0xcd, 0x14, 0x39, 0x8b, 0x7d, 0xc2, 0xfb, 0x26, 0x3c, 0xa0, 0x86, 0xf2, 784 0xf5, 0xf3, 0x0c, 0xdb, 785 ] 786 ); 787 let first = accepted_inventory_digest_for_facts(&[fact(1_784_347_201, b"alpha")]) 788 .expect("first digest"); 789 let changed_provenance = 790 accepted_inventory_digest_for_facts(&[fact(1_784_347_202, b"alpha")]) 791 .expect("provenance digest"); 792 let changed_content = accepted_inventory_digest_for_facts(&[fact(1_784_347_201, b"bravo")]) 793 .expect("content digest"); 794 assert_ne!(first, changed_provenance); 795 assert_ne!(first, changed_content); 796 } 797 798 #[test] 799 fn errors_are_closed_redacted_and_source_free() { 800 for kind in [ 801 RhiReconciliationCommitErrorKind::InvalidMode, 802 RhiReconciliationCommitErrorKind::InvalidInput, 803 RhiReconciliationCommitErrorKind::LeaseLost, 804 RhiReconciliationCommitErrorKind::GenerationConflict, 805 RhiReconciliationCommitErrorKind::Conflict, 806 RhiReconciliationCommitErrorKind::Storage, 807 RhiReconciliationCommitErrorKind::CommitOutcomeUnknown, 808 ] { 809 let error = failure(kind); 810 assert_eq!(error.kind(), kind); 811 assert!(Error::source(&error).is_none()); 812 assert!(!format!("{error} {error:?}").contains("trade-primary")); 813 } 814 } 815 }