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 }