rhi

Coordinated trade for connected markets
git clone https://radroots.dev/git/rhi.git
Log | Files | Refs | README | LICENSE

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 }