rhi

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

reconciliation_finalization.rs (14661B)


      1 //! Generation-fenced reconciliation-finalization preflight.
      2 
      3 use core::fmt;
      4 use std::error::Error;
      5 
      6 use radroots_event::id::TradeId;
      7 use radroots_service_sqlite::{
      8     ServiceSqliteTransaction, ServiceSqliteTransactionError, ServiceSqliteTransactionErrorKind,
      9 };
     10 
     11 use crate::{
     12     RhiEvidencePolicyDigest, RhiReconciliationAttemptId, RhiReconciliationAttemptRepository,
     13     RhiReconciliationEvaluation, RhiReconciliationJobId, RhiReconciliationJobState,
     14     RhiReconciliationLease, RhiReconciliationUnixMilliseconds, RhiStateHostMode,
     15     reconciliation_attempt::attempt_id,
     16     reconciliation_job::{LeaseValidationError, validate_exact_lease},
     17     source_ingest::{SourceOperationError, read_dirty},
     18 };
     19 
     20 /// Exact version of the reconciliation-finalization fence contract.
     21 pub const RHI_RECONCILIATION_FINALIZATION_CONTRACT_VERSION: u32 = 1;
     22 
     23 const MATCH_COMMITTED_ATTEMPT_SQL: &str = r#"SELECT COUNT(*)
     24 FROM evidence_reconciliations
     25 WHERE attempt_id = ? AND job_id = ? AND trade_id = ? AND input_generation = ?
     26     AND evidence_policy_sha256 = ?"#;
     27 
     28 /// Stable source-free failure class for finalization preflight.
     29 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
     30 pub enum RhiReconciliationFinalizationErrorKind {
     31     InvalidMode,
     32     InvalidInput,
     33     LeaseLost,
     34     GenerationConflict,
     35     AttemptUnavailable,
     36     ProjectionUnavailable,
     37     Storage,
     38     CommitOutcomeUnknown,
     39 }
     40 
     41 impl RhiReconciliationFinalizationErrorKind {
     42     /// Returns the stable machine-readable failure code.
     43     #[must_use]
     44     pub const fn code(self) -> &'static str {
     45         match self {
     46             Self::InvalidMode => "reconciliation_finalization_mode_invalid",
     47             Self::InvalidInput => "reconciliation_finalization_input_invalid",
     48             Self::LeaseLost => "reconciliation_finalization_lease_lost",
     49             Self::GenerationConflict => "reconciliation_finalization_generation_conflict",
     50             Self::AttemptUnavailable => "reconciliation_finalization_attempt_unavailable",
     51             Self::ProjectionUnavailable => "reconciliation_finalization_projection_unavailable",
     52             Self::Storage => "reconciliation_finalization_storage_failed",
     53             Self::CommitOutcomeUnknown => "reconciliation_finalization_outcome_unknown",
     54         }
     55     }
     56 }
     57 
     58 /// Redacted source-free finalization-preflight failure.
     59 #[derive(Clone, Copy, PartialEq, Eq)]
     60 pub struct RhiReconciliationFinalizationError {
     61     kind: RhiReconciliationFinalizationErrorKind,
     62 }
     63 
     64 impl RhiReconciliationFinalizationError {
     65     /// Returns the stable failure class.
     66     #[must_use]
     67     pub const fn kind(self) -> RhiReconciliationFinalizationErrorKind {
     68         self.kind
     69     }
     70 
     71     /// Returns the stable machine-readable failure code.
     72     #[must_use]
     73     pub const fn code(self) -> &'static str {
     74         self.kind.code()
     75     }
     76 }
     77 
     78 impl fmt::Display for RhiReconciliationFinalizationError {
     79     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
     80         formatter.write_str(match self.kind {
     81             RhiReconciliationFinalizationErrorKind::InvalidMode => {
     82                 "RHI reconciliation finalization requires writable state"
     83             }
     84             RhiReconciliationFinalizationErrorKind::InvalidInput => {
     85                 "RHI reconciliation finalization input is invalid"
     86             }
     87             RhiReconciliationFinalizationErrorKind::LeaseLost => {
     88                 "RHI reconciliation finalization lease is no longer authoritative"
     89             }
     90             RhiReconciliationFinalizationErrorKind::GenerationConflict => {
     91                 "RHI reconciliation finalization generation changed"
     92             }
     93             RhiReconciliationFinalizationErrorKind::AttemptUnavailable => {
     94                 "RHI reconciliation finalization attempt is unavailable"
     95             }
     96             RhiReconciliationFinalizationErrorKind::ProjectionUnavailable => {
     97                 "RHI reconciliation finalization projection is unavailable"
     98             }
     99             RhiReconciliationFinalizationErrorKind::Storage => {
    100                 "RHI reconciliation finalization preflight failed"
    101             }
    102             RhiReconciliationFinalizationErrorKind::CommitOutcomeUnknown => {
    103                 "RHI reconciliation finalization preflight outcome is unknown"
    104             }
    105         })
    106     }
    107 }
    108 
    109 impl fmt::Debug for RhiReconciliationFinalizationError {
    110     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    111         formatter
    112             .debug_struct("RhiReconciliationFinalizationError")
    113             .field("kind", &self.kind)
    114             .finish()
    115     }
    116 }
    117 
    118 impl Error for RhiReconciliationFinalizationError {}
    119 
    120 /// Sealed preflight capability for one exact evaluated reconciliation attempt.
    121 ///
    122 /// This value proves only that its lease, generation, policy, and committed
    123 /// attempt matched during one bounded read-only preflight transaction. The
    124 /// eventual Step199 writer must rerun the same validator inside its atomic
    125 /// transaction before any mutation.
    126 ///
    127 /// ```compile_fail
    128 /// use rhi::RhiReconciliationFinalizationFence;
    129 ///
    130 /// let _forged = RhiReconciliationFinalizationFence { evaluation: todo!() };
    131 /// ```
    132 pub struct RhiReconciliationFinalizationFence {
    133     lease: RhiReconciliationLease,
    134     identity: FinalizationIdentity,
    135     evaluation: RhiReconciliationEvaluation,
    136 }
    137 
    138 impl RhiReconciliationFinalizationFence {
    139     /// Returns the exact finalization-fence contract version.
    140     #[must_use]
    141     pub const fn contract_version(&self) -> u32 {
    142         RHI_RECONCILIATION_FINALIZATION_CONTRACT_VERSION
    143     }
    144 
    145     /// Returns the sealed Step194 evaluation retained by this preflight.
    146     #[must_use]
    147     pub const fn evaluation(&self) -> &RhiReconciliationEvaluation {
    148         &self.evaluation
    149     }
    150 }
    151 
    152 impl fmt::Debug for RhiReconciliationFinalizationFence {
    153     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    154         formatter
    155             .debug_struct("RhiReconciliationFinalizationFence")
    156             .field("coverage", &self.evaluation.coverage())
    157             .field("outcome", &self.evaluation.outcome())
    158             .finish_non_exhaustive()
    159     }
    160 }
    161 
    162 #[derive(Clone, Copy)]
    163 pub(crate) struct FinalizationIdentity {
    164     pub(crate) attempt_id: RhiReconciliationAttemptId,
    165     pub(crate) job_id: RhiReconciliationJobId,
    166     pub(crate) trade_id: TradeId,
    167     pub(crate) generation: u64,
    168     pub(crate) policy_digest: RhiEvidencePolicyDigest,
    169 }
    170 
    171 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
    172 pub(crate) enum FinalizationOperationError {
    173     LeaseLost,
    174     GenerationConflict,
    175     AttemptUnavailable,
    176     Storage,
    177 }
    178 
    179 impl RhiReconciliationAttemptRepository<'_> {
    180     /// Performs one bounded read-only preflight for later atomic finalization.
    181     ///
    182     /// The returned capability is not commit authority. Durable finalization
    183     /// must rerun the same exact checks inside its Step199 write transaction.
    184     pub async fn prepare_finalization(
    185         &self,
    186         lease: RhiReconciliationLease,
    187         evaluation: RhiReconciliationEvaluation,
    188         now: RhiReconciliationUnixMilliseconds,
    189     ) -> Result<RhiReconciliationFinalizationFence, RhiReconciliationFinalizationError> {
    190         if self.host().mode() != RhiStateHostMode::ReadWriteExisting {
    191             return Err(failure(RhiReconciliationFinalizationErrorKind::InvalidMode));
    192         }
    193         let identity = finalization_identity(lease, &evaluation, now)?;
    194         let fence = RhiReconciliationFinalizationFence {
    195             lease,
    196             identity,
    197             evaluation,
    198         };
    199         self.host()
    200             .sqlite_host()
    201             .transaction(move |transaction| {
    202                 Box::pin(async move {
    203                     validate_finalization_fence(transaction, &fence, now).await?;
    204                     Ok(fence)
    205                 })
    206             })
    207             .await
    208             .map_err(map_transaction_error)
    209     }
    210 }
    211 
    212 impl RhiReconciliationFinalizationFence {
    213     pub(crate) const fn validation_parts(&self) -> (RhiReconciliationLease, FinalizationIdentity) {
    214         (self.lease, self.identity)
    215     }
    216 }
    217 
    218 pub(crate) async fn validate_finalization_fence(
    219     transaction: &mut ServiceSqliteTransaction<'_>,
    220     fence: &RhiReconciliationFinalizationFence,
    221     now: RhiReconciliationUnixMilliseconds,
    222 ) -> Result<(), FinalizationOperationError> {
    223     validate_finalization_identity(transaction, fence.lease, fence.identity, now).await
    224 }
    225 
    226 fn finalization_identity(
    227     lease: RhiReconciliationLease,
    228     evaluation: &RhiReconciliationEvaluation,
    229     now: RhiReconciliationUnixMilliseconds,
    230 ) -> Result<FinalizationIdentity, RhiReconciliationFinalizationError> {
    231     let job = lease.job();
    232     let manifest = evaluation.projection().manifest();
    233     if job.state() != RhiReconciliationJobState::Leased
    234         || job.attempt_count() == 0
    235         || manifest.job_id() != job.id()
    236         || manifest.attempt_id() != attempt_id(job.id(), job.attempt_count())
    237         || manifest.trade_id() != &job.trade_id()
    238         || manifest.trade_generation() != job.input_generation()
    239         || manifest.inner().evidence_policy_digest().as_bytes()
    240             != job.evidence_policy_digest().as_bytes()
    241     {
    242         return Err(failure(
    243             RhiReconciliationFinalizationErrorKind::InvalidInput,
    244         ));
    245     }
    246     if now >= lease.lease_expires() {
    247         return Err(failure(RhiReconciliationFinalizationErrorKind::LeaseLost));
    248     }
    249     if evaluation.projection().digest().is_none() {
    250         return Err(failure(
    251             RhiReconciliationFinalizationErrorKind::ProjectionUnavailable,
    252         ));
    253     }
    254     Ok(FinalizationIdentity {
    255         attempt_id: manifest.attempt_id(),
    256         job_id: manifest.job_id(),
    257         trade_id: job.trade_id(),
    258         generation: job.input_generation(),
    259         policy_digest: job.evidence_policy_digest(),
    260     })
    261 }
    262 
    263 pub(crate) async fn validate_finalization_identity(
    264     transaction: &mut ServiceSqliteTransaction<'_>,
    265     lease: RhiReconciliationLease,
    266     identity: FinalizationIdentity,
    267     now: RhiReconciliationUnixMilliseconds,
    268 ) -> Result<(), FinalizationOperationError> {
    269     if now >= lease.lease_expires() {
    270         return Err(FinalizationOperationError::LeaseLost);
    271     }
    272     validate_exact_lease(transaction, lease)
    273         .await
    274         .map_err(|error| match error {
    275             LeaseValidationError::LeaseLost => FinalizationOperationError::LeaseLost,
    276             LeaseValidationError::Storage => FinalizationOperationError::Storage,
    277         })?;
    278     read_dirty(transaction, identity.trade_id)
    279         .await
    280         .map_err(|error| match error {
    281             SourceOperationError::GenerationConflict => {
    282                 FinalizationOperationError::GenerationConflict
    283             }
    284             SourceOperationError::Persistence(_) | SourceOperationError::Storage => {
    285                 FinalizationOperationError::Storage
    286             }
    287         })?
    288         .filter(|dirty| {
    289             dirty.generation.get() == identity.generation && dirty.policy == identity.policy_digest
    290         })
    291         .ok_or(FinalizationOperationError::GenerationConflict)?;
    292     let generation =
    293         i64::try_from(identity.generation).map_err(|_| FinalizationOperationError::Storage)?;
    294     let count = sqlx::query_scalar::<_, i64>(MATCH_COMMITTED_ATTEMPT_SQL)
    295         .bind(identity.attempt_id.as_bytes().as_slice())
    296         .bind(identity.job_id.as_bytes().as_slice())
    297         .bind(identity.trade_id.as_bytes().as_slice())
    298         .bind(generation)
    299         .bind(identity.policy_digest.as_bytes().as_slice())
    300         .fetch_one(&mut *transaction)
    301         .await
    302         .map_err(|_| FinalizationOperationError::Storage)?;
    303     if count == 1 {
    304         Ok(())
    305     } else {
    306         Err(FinalizationOperationError::AttemptUnavailable)
    307     }
    308 }
    309 
    310 fn map_transaction_error(
    311     error: ServiceSqliteTransactionError<FinalizationOperationError>,
    312 ) -> RhiReconciliationFinalizationError {
    313     if error.kind() == ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown {
    314         return failure(RhiReconciliationFinalizationErrorKind::CommitOutcomeUnknown);
    315     }
    316     failure(match error.operation_error().copied() {
    317         Some(FinalizationOperationError::LeaseLost) => {
    318             RhiReconciliationFinalizationErrorKind::LeaseLost
    319         }
    320         Some(FinalizationOperationError::GenerationConflict) => {
    321             RhiReconciliationFinalizationErrorKind::GenerationConflict
    322         }
    323         Some(FinalizationOperationError::AttemptUnavailable) => {
    324             RhiReconciliationFinalizationErrorKind::AttemptUnavailable
    325         }
    326         Some(FinalizationOperationError::Storage) | None => {
    327             RhiReconciliationFinalizationErrorKind::Storage
    328         }
    329     })
    330 }
    331 
    332 const fn failure(
    333     kind: RhiReconciliationFinalizationErrorKind,
    334 ) -> RhiReconciliationFinalizationError {
    335     RhiReconciliationFinalizationError { kind }
    336 }
    337 
    338 #[cfg(test)]
    339 mod tests {
    340     use super::*;
    341 
    342     #[test]
    343     fn error_inventory_is_exact_source_free_and_redacted() {
    344         let cases = [
    345             (
    346                 RhiReconciliationFinalizationErrorKind::InvalidMode,
    347                 "reconciliation_finalization_mode_invalid",
    348             ),
    349             (
    350                 RhiReconciliationFinalizationErrorKind::InvalidInput,
    351                 "reconciliation_finalization_input_invalid",
    352             ),
    353             (
    354                 RhiReconciliationFinalizationErrorKind::LeaseLost,
    355                 "reconciliation_finalization_lease_lost",
    356             ),
    357             (
    358                 RhiReconciliationFinalizationErrorKind::GenerationConflict,
    359                 "reconciliation_finalization_generation_conflict",
    360             ),
    361             (
    362                 RhiReconciliationFinalizationErrorKind::AttemptUnavailable,
    363                 "reconciliation_finalization_attempt_unavailable",
    364             ),
    365             (
    366                 RhiReconciliationFinalizationErrorKind::ProjectionUnavailable,
    367                 "reconciliation_finalization_projection_unavailable",
    368             ),
    369             (
    370                 RhiReconciliationFinalizationErrorKind::Storage,
    371                 "reconciliation_finalization_storage_failed",
    372             ),
    373             (
    374                 RhiReconciliationFinalizationErrorKind::CommitOutcomeUnknown,
    375                 "reconciliation_finalization_outcome_unknown",
    376             ),
    377         ];
    378         for (kind, code) in cases {
    379             let error = failure(kind);
    380             assert_eq!(error.kind(), kind);
    381             assert_eq!(error.code(), code);
    382             assert!(Error::source(&error).is_none());
    383             let rendered = format!("{error} {error:?}");
    384             assert!(!rendered.contains("11111111"));
    385             assert!(!rendered.contains("trade-primary"));
    386         }
    387     }
    388 }