reconciliation_job.rs (44693B)
1 //! Bounded durable reconciliation-job scheduling and lease state machine. 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 use sha2::{Digest, Sha256}; 11 use sqlx::Row; 12 13 use crate::{ 14 RhiConfigDocumentV1, RhiEvidencePolicyDigest, RhiReconciliationJobRepository, RhiStateHostMode, 15 }; 16 17 /// Exact version of the durable reconciliation-job contract. 18 pub const RHI_RECONCILIATION_JOB_CONTRACT_VERSION: u32 = 1; 19 20 /// Absolute hard ceiling for active durable reconciliation jobs. 21 pub const RHI_RECONCILIATION_JOB_MAX_ACTIVE: u32 = 65_536; 22 23 const JOB_ID_DOMAIN: &[u8] = b"radroots.rhi.reconciliation_job.v1\0"; 24 const LEASE_OWNER_BYTES: usize = 16; 25 const MAX_UNIX_MILLISECONDS: u64 = i64::MAX as u64; 26 const MAX_ATTEMPTS: u16 = 100; 27 const MAX_LEASE_MILLISECONDS: u64 = 300_000; 28 const MAX_RENEWAL_MILLISECONDS: u64 = 150_000; 29 const MAX_INITIAL_BACKOFF_MILLISECONDS: u64 = 60_000; 30 const MAX_BACKOFF_MILLISECONDS: u64 = 3_600_000; 31 32 const READ_DIRTY_SQL: &str = r#"SELECT generation, 33 length(evidence_policy_sha256) AS evidence_policy_bytes, 34 substr(evidence_policy_sha256, 1, 33) AS evidence_policy_sha256 35 FROM trade_dirty_generations 36 WHERE trade_id = ? 37 LIMIT 1"#; 38 const READ_JOB_SQL: &str = r#"SELECT 39 length(job_id) AS job_id_bytes, substr(job_id, 1, 33) AS job_id, 40 length(trade_id) AS trade_id_bytes, substr(trade_id, 1, 17) AS trade_id, 41 input_generation, 42 length(evidence_policy_sha256) AS evidence_policy_bytes, 43 substr(evidence_policy_sha256, 1, 33) AS evidence_policy_sha256, 44 length(CAST(state AS BLOB)) AS state_bytes, substr(state, 1, 12) AS state, 45 revision, attempt_count, failure_count, max_attempts, 46 lease_duration_ms, lease_renewal_ms, initial_backoff_ms, maximum_backoff_ms, 47 next_attempt_unix_ms, 48 CASE WHEN lease_owner IS NULL THEN NULL ELSE length(lease_owner) END AS lease_owner_bytes, 49 CASE WHEN lease_owner IS NULL THEN NULL ELSE substr(lease_owner, 1, 17) END AS lease_owner, 50 lease_expires_unix_ms, created_at_unix_ms, updated_at_unix_ms 51 FROM reconciliation_jobs 52 WHERE job_id = ? 53 LIMIT 1"#; 54 const READ_ACTIVE_JOB_SQL: &str = r#"SELECT 55 length(job_id) AS job_id_bytes, substr(job_id, 1, 33) AS job_id, 56 length(trade_id) AS trade_id_bytes, substr(trade_id, 1, 17) AS trade_id, 57 input_generation, 58 length(evidence_policy_sha256) AS evidence_policy_bytes, 59 substr(evidence_policy_sha256, 1, 33) AS evidence_policy_sha256, 60 length(CAST(state AS BLOB)) AS state_bytes, substr(state, 1, 12) AS state, 61 revision, attempt_count, failure_count, max_attempts, 62 lease_duration_ms, lease_renewal_ms, initial_backoff_ms, maximum_backoff_ms, 63 next_attempt_unix_ms, 64 CASE WHEN lease_owner IS NULL THEN NULL ELSE length(lease_owner) END AS lease_owner_bytes, 65 CASE WHEN lease_owner IS NULL THEN NULL ELSE substr(lease_owner, 1, 17) END AS lease_owner, 66 lease_expires_unix_ms, created_at_unix_ms, updated_at_unix_ms 67 FROM reconciliation_jobs 68 WHERE trade_id = ? AND state IN ('ready', 'leased') 69 LIMIT 2"#; 70 const ACTIVE_JOB_COUNT_SQL: &str = r#"SELECT COUNT(*) AS active_count 71 FROM reconciliation_jobs 72 WHERE state IN ('ready', 'leased')"#; 73 const SUPERSEDE_JOB_SQL: &str = r#"UPDATE reconciliation_jobs 74 SET state = 'superseded', revision = revision + 1, 75 next_attempt_unix_ms = NULL, lease_owner = NULL, lease_expires_unix_ms = NULL, 76 updated_at_unix_ms = ? 77 WHERE job_id = ? AND revision = ? AND state IN ('ready', 'leased')"#; 78 const INSERT_JOB_SQL: &str = r#"INSERT INTO reconciliation_jobs ( 79 job_id, trade_id, input_generation, evidence_policy_sha256, state, revision, 80 attempt_count, failure_count, max_attempts, lease_duration_ms, 81 lease_renewal_ms, initial_backoff_ms, maximum_backoff_ms, 82 next_attempt_unix_ms, lease_owner, lease_expires_unix_ms, 83 created_at_unix_ms, updated_at_unix_ms 84 ) VALUES (?, ?, ?, ?, 'ready', 1, 0, 0, ?, ?, ?, ?, ?, ?, NULL, NULL, ?, ?)"#; 85 const EXHAUST_EXPIRED_SQL: &str = r#"UPDATE reconciliation_jobs 86 SET state = 'exhausted', revision = revision + 1, 87 failure_count = attempt_count, lease_owner = NULL, lease_expires_unix_ms = NULL, 88 updated_at_unix_ms = ? 89 WHERE job_id IN ( 90 SELECT job_id FROM reconciliation_jobs 91 WHERE state = 'leased' AND lease_expires_unix_ms <= ? 92 AND attempt_count >= max_attempts AND updated_at_unix_ms <= ? 93 ORDER BY lease_expires_unix_ms, created_at_unix_ms, job_id 94 LIMIT 65536 95 )"#; 96 const READ_CLAIMABLE_SQL: &str = r#"SELECT 97 length(job_id) AS job_id_bytes, substr(job_id, 1, 33) AS job_id, 98 length(trade_id) AS trade_id_bytes, substr(trade_id, 1, 17) AS trade_id, 99 input_generation, 100 length(evidence_policy_sha256) AS evidence_policy_bytes, 101 substr(evidence_policy_sha256, 1, 33) AS evidence_policy_sha256, 102 length(CAST(state AS BLOB)) AS state_bytes, substr(state, 1, 12) AS state, 103 revision, attempt_count, failure_count, max_attempts, 104 lease_duration_ms, lease_renewal_ms, initial_backoff_ms, maximum_backoff_ms, 105 next_attempt_unix_ms, 106 CASE WHEN lease_owner IS NULL THEN NULL ELSE length(lease_owner) END AS lease_owner_bytes, 107 CASE WHEN lease_owner IS NULL THEN NULL ELSE substr(lease_owner, 1, 17) END AS lease_owner, 108 lease_expires_unix_ms, created_at_unix_ms, updated_at_unix_ms 109 FROM reconciliation_jobs 110 WHERE attempt_count < max_attempts AND updated_at_unix_ms <= ? AND ( 111 (state = 'ready' AND next_attempt_unix_ms <= ?) 112 OR (state = 'leased' AND lease_expires_unix_ms <= ?) 113 ) 114 ORDER BY 115 CASE state WHEN 'ready' THEN next_attempt_unix_ms ELSE lease_expires_unix_ms END, 116 created_at_unix_ms, job_id 117 LIMIT 1"#; 118 const CLAIM_JOB_SQL: &str = r#"UPDATE reconciliation_jobs 119 SET state = 'leased', revision = revision + 1, attempt_count = attempt_count + 1, 120 next_attempt_unix_ms = NULL, lease_owner = ?, lease_expires_unix_ms = ?, 121 updated_at_unix_ms = ? 122 WHERE job_id = ? AND revision = ? AND attempt_count < max_attempts 123 AND updated_at_unix_ms <= ? AND ( 124 (state = 'ready' AND next_attempt_unix_ms <= ?) 125 OR (state = 'leased' AND lease_expires_unix_ms <= ?) 126 )"#; 127 const RENEW_JOB_SQL: &str = r#"UPDATE reconciliation_jobs 128 SET revision = revision + 1, lease_expires_unix_ms = ?, updated_at_unix_ms = ? 129 WHERE job_id = ? AND revision = ? AND state = 'leased' 130 AND lease_owner = ? AND lease_expires_unix_ms = ? 131 AND lease_expires_unix_ms > ? AND updated_at_unix_ms <= ?"#; 132 const RETRY_JOB_SQL: &str = r#"UPDATE reconciliation_jobs 133 SET state = 'ready', revision = revision + 1, failure_count = failure_count + 1, 134 next_attempt_unix_ms = ?, lease_owner = NULL, lease_expires_unix_ms = NULL, 135 updated_at_unix_ms = ? 136 WHERE job_id = ? AND revision = ? AND state = 'leased' 137 AND lease_owner = ? AND lease_expires_unix_ms = ? 138 AND lease_expires_unix_ms > ? AND attempt_count < max_attempts"#; 139 const EXHAUST_JOB_SQL: &str = r#"UPDATE reconciliation_jobs 140 SET state = 'exhausted', revision = revision + 1, failure_count = failure_count + 1, 141 next_attempt_unix_ms = NULL, lease_owner = NULL, lease_expires_unix_ms = NULL, 142 updated_at_unix_ms = ? 143 WHERE job_id = ? AND revision = ? AND state = 'leased' 144 AND lease_owner = ? AND lease_expires_unix_ms = ? 145 AND lease_expires_unix_ms > ? AND attempt_count >= max_attempts"#; 146 147 /// Stable durable reconciliation-job lifecycle state. 148 #[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)] 149 pub enum RhiReconciliationJobState { 150 Ready, 151 Leased, 152 Exhausted, 153 Superseded, 154 Completed, 155 } 156 157 impl RhiReconciliationJobState { 158 /// Returns the exact machine-contract spelling. 159 #[must_use] 160 pub const fn code(self) -> &'static str { 161 match self { 162 Self::Ready => "ready", 163 Self::Leased => "leased", 164 Self::Exhausted => "exhausted", 165 Self::Superseded => "superseded", 166 Self::Completed => "completed", 167 } 168 } 169 } 170 171 /// Stable source-free durable job failure class. 172 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 173 pub enum RhiReconciliationJobErrorKind { 174 InvalidMode, 175 InvalidInput, 176 QueueFull, 177 DirtyGenerationConflict, 178 LeaseLost, 179 NotReady, 180 Storage, 181 CommitOutcomeUnknown, 182 } 183 184 impl RhiReconciliationJobErrorKind { 185 /// Returns the stable machine-readable failure code. 186 #[must_use] 187 pub const fn code(self) -> &'static str { 188 match self { 189 Self::InvalidMode => "reconciliation_job_mode_invalid", 190 Self::InvalidInput => "reconciliation_job_input_invalid", 191 Self::QueueFull => "reconciliation_queue_full", 192 Self::DirtyGenerationConflict => "reconciliation_dirty_generation_conflict", 193 Self::LeaseLost => "reconciliation_lease_lost", 194 Self::NotReady => "reconciliation_job_not_ready", 195 Self::Storage => "reconciliation_job_storage_failed", 196 Self::CommitOutcomeUnknown => "reconciliation_job_commit_outcome_unknown", 197 } 198 } 199 } 200 201 /// Redacted source-free durable job failure. 202 #[derive(Clone, Copy, PartialEq, Eq)] 203 pub struct RhiReconciliationJobError { 204 kind: RhiReconciliationJobErrorKind, 205 } 206 207 impl RhiReconciliationJobError { 208 const fn new(kind: RhiReconciliationJobErrorKind) -> Self { 209 Self { kind } 210 } 211 212 /// Returns the stable failure class. 213 #[must_use] 214 pub const fn kind(self) -> RhiReconciliationJobErrorKind { 215 self.kind 216 } 217 218 /// Returns the stable machine-readable failure code. 219 #[must_use] 220 pub const fn code(self) -> &'static str { 221 self.kind.code() 222 } 223 } 224 225 impl fmt::Display for RhiReconciliationJobError { 226 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 227 formatter.write_str(match self.kind { 228 RhiReconciliationJobErrorKind::InvalidMode => { 229 "RHI reconciliation jobs require writable state" 230 } 231 RhiReconciliationJobErrorKind::InvalidInput => { 232 "RHI reconciliation job input is invalid" 233 } 234 RhiReconciliationJobErrorKind::QueueFull => { 235 "RHI reconciliation job capacity is exhausted" 236 } 237 RhiReconciliationJobErrorKind::DirtyGenerationConflict => { 238 "RHI reconciliation dirty generation changed" 239 } 240 RhiReconciliationJobErrorKind::LeaseLost => { 241 "RHI reconciliation lease is no longer authoritative" 242 } 243 RhiReconciliationJobErrorKind::NotReady => "RHI reconciliation job is not ready", 244 RhiReconciliationJobErrorKind::Storage => "RHI reconciliation job transaction failed", 245 RhiReconciliationJobErrorKind::CommitOutcomeUnknown => { 246 "RHI reconciliation job commit outcome is unknown" 247 } 248 }) 249 } 250 } 251 252 impl fmt::Debug for RhiReconciliationJobError { 253 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 254 formatter 255 .debug_struct("RhiReconciliationJobError") 256 .field("kind", &self.kind) 257 .finish() 258 } 259 } 260 261 impl Error for RhiReconciliationJobError {} 262 263 /// Validated numeric scheduling authority copied into each durable job. 264 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 265 pub struct RhiReconciliationJobPolicy { 266 queue_capacity: u32, 267 lease_duration_ms: u64, 268 lease_renewal_ms: u64, 269 max_attempts: u16, 270 initial_backoff_ms: u64, 271 maximum_backoff_ms: u64, 272 } 273 274 impl RhiReconciliationJobPolicy { 275 /// Constructs an explicit bounded policy with no ambient defaults. 276 pub fn new( 277 queue_capacity: u32, 278 lease_duration_ms: u64, 279 lease_renewal_ms: u64, 280 max_attempts: u16, 281 initial_backoff_ms: u64, 282 maximum_backoff_ms: u64, 283 ) -> Result<Self, RhiReconciliationJobError> { 284 if queue_capacity == 0 285 || queue_capacity > RHI_RECONCILIATION_JOB_MAX_ACTIVE 286 || !(1_000..=MAX_LEASE_MILLISECONDS).contains(&lease_duration_ms) 287 || !(100..=MAX_RENEWAL_MILLISECONDS).contains(&lease_renewal_ms) 288 || lease_renewal_ms >= lease_duration_ms 289 || max_attempts == 0 290 || max_attempts > MAX_ATTEMPTS 291 || !(1..=MAX_INITIAL_BACKOFF_MILLISECONDS).contains(&initial_backoff_ms) 292 || !(1..=MAX_BACKOFF_MILLISECONDS).contains(&maximum_backoff_ms) 293 || initial_backoff_ms > maximum_backoff_ms 294 { 295 return Err(RhiReconciliationJobError::new( 296 RhiReconciliationJobErrorKind::InvalidInput, 297 )); 298 } 299 Ok(Self { 300 queue_capacity, 301 lease_duration_ms, 302 lease_renewal_ms, 303 max_attempts, 304 initial_backoff_ms, 305 maximum_backoff_ms, 306 }) 307 } 308 309 /// Extracts the exact admitted reconciliation policy from one immutable configuration. 310 pub fn from_configuration( 311 configuration: &RhiConfigDocumentV1, 312 ) -> Result<Self, RhiReconciliationJobError> { 313 let value = |pointer: &str| { 314 configuration 315 .normalized() 316 .pointer(pointer) 317 .and_then(serde_json::Value::as_u64) 318 .ok_or_else(|| { 319 RhiReconciliationJobError::new(RhiReconciliationJobErrorKind::InvalidInput) 320 }) 321 }; 322 Self::new( 323 u32::try_from(value("/reconciliation/queue_capacity")?).map_err(|_| { 324 RhiReconciliationJobError::new(RhiReconciliationJobErrorKind::InvalidInput) 325 })?, 326 value("/reconciliation/lease_ms")?, 327 value("/reconciliation/lease_renewal_ms")?, 328 u16::try_from(value("/reconciliation/max_attempts")?).map_err(|_| { 329 RhiReconciliationJobError::new(RhiReconciliationJobErrorKind::InvalidInput) 330 })?, 331 value("/reconciliation/initial_backoff_ms")?, 332 value("/reconciliation/maximum_backoff_ms")?, 333 ) 334 } 335 336 #[must_use] 337 pub const fn queue_capacity(self) -> u32 { 338 self.queue_capacity 339 } 340 341 #[must_use] 342 pub const fn lease_duration_milliseconds(self) -> u64 { 343 self.lease_duration_ms 344 } 345 346 #[must_use] 347 pub const fn lease_renewal_milliseconds(self) -> u64 { 348 self.lease_renewal_ms 349 } 350 351 #[must_use] 352 pub const fn max_attempts(self) -> u16 { 353 self.max_attempts 354 } 355 356 #[must_use] 357 pub const fn initial_backoff_milliseconds(self) -> u64 { 358 self.initial_backoff_ms 359 } 360 361 #[must_use] 362 pub const fn maximum_backoff_milliseconds(self) -> u64 { 363 self.maximum_backoff_ms 364 } 365 } 366 367 /// Injected wall-clock millisecond used only for durable scheduling evidence. 368 #[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)] 369 pub struct RhiReconciliationUnixMilliseconds(u64); 370 371 impl RhiReconciliationUnixMilliseconds { 372 pub fn new(value: u64) -> Result<Self, RhiReconciliationJobError> { 373 if value > MAX_UNIX_MILLISECONDS { 374 return Err(RhiReconciliationJobError::new( 375 RhiReconciliationJobErrorKind::InvalidInput, 376 )); 377 } 378 Ok(Self(value)) 379 } 380 381 #[must_use] 382 pub const fn get(self) -> u64 { 383 self.0 384 } 385 } 386 387 /// Injected full-jitter delay for one failed reconciliation attempt. 388 #[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)] 389 pub struct RhiReconciliationRetryDelayMilliseconds(u64); 390 391 impl RhiReconciliationRetryDelayMilliseconds { 392 pub fn new(value: u64) -> Result<Self, RhiReconciliationJobError> { 393 if value > MAX_BACKOFF_MILLISECONDS { 394 return Err(RhiReconciliationJobError::new( 395 RhiReconciliationJobErrorKind::InvalidInput, 396 )); 397 } 398 Ok(Self(value)) 399 } 400 401 #[must_use] 402 pub const fn get(self) -> u64 { 403 self.0 404 } 405 } 406 407 /// Stable process-local owner token for compare-and-swap leases. 408 #[derive(Clone, Copy, PartialEq, Eq, Hash)] 409 pub struct RhiReconciliationLeaseOwner([u8; LEASE_OWNER_BYTES]); 410 411 impl RhiReconciliationLeaseOwner { 412 pub fn from_bytes(bytes: [u8; LEASE_OWNER_BYTES]) -> Result<Self, RhiReconciliationJobError> { 413 if bytes.iter().all(|byte| *byte == 0) { 414 return Err(RhiReconciliationJobError::new( 415 RhiReconciliationJobErrorKind::InvalidInput, 416 )); 417 } 418 Ok(Self(bytes)) 419 } 420 } 421 422 impl fmt::Debug for RhiReconciliationLeaseOwner { 423 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 424 formatter.write_str("RhiReconciliationLeaseOwner([redacted])") 425 } 426 } 427 428 /// Deterministic identity of one trade generation and evidence policy. 429 #[derive(Clone, Copy, PartialEq, Eq, Hash)] 430 pub struct RhiReconciliationJobId([u8; 32]); 431 432 impl RhiReconciliationJobId { 433 #[must_use] 434 pub const fn as_bytes(&self) -> &[u8; 32] { 435 &self.0 436 } 437 } 438 439 impl fmt::Debug for RhiReconciliationJobId { 440 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 441 formatter.write_str("RhiReconciliationJobId([redacted])") 442 } 443 } 444 445 /// Immutable decoded view of one retained durable job row. 446 #[derive(Clone, Copy, PartialEq, Eq)] 447 pub struct RhiReconciliationJob { 448 id: RhiReconciliationJobId, 449 trade_id: TradeId, 450 input_generation: u64, 451 policy_digest: RhiEvidencePolicyDigest, 452 state: RhiReconciliationJobState, 453 revision: u64, 454 attempt_count: u16, 455 failure_count: u16, 456 policy: RhiReconciliationJobPolicy, 457 next_attempt: Option<RhiReconciliationUnixMilliseconds>, 458 lease_owner: Option<RhiReconciliationLeaseOwner>, 459 lease_expires: Option<RhiReconciliationUnixMilliseconds>, 460 created_at: RhiReconciliationUnixMilliseconds, 461 updated_at: RhiReconciliationUnixMilliseconds, 462 } 463 464 impl RhiReconciliationJob { 465 pub(crate) const fn policy(self) -> RhiReconciliationJobPolicy { 466 self.policy 467 } 468 469 pub(crate) const fn attempt_policy_matches(self, expected: RhiReconciliationJobPolicy) -> bool { 470 self.policy.lease_duration_ms == expected.lease_duration_ms 471 && self.policy.lease_renewal_ms == expected.lease_renewal_ms 472 && self.policy.max_attempts == expected.max_attempts 473 && self.policy.initial_backoff_ms == expected.initial_backoff_ms 474 && self.policy.maximum_backoff_ms == expected.maximum_backoff_ms 475 } 476 477 #[must_use] 478 pub const fn id(self) -> RhiReconciliationJobId { 479 self.id 480 } 481 482 #[must_use] 483 pub const fn trade_id(self) -> TradeId { 484 self.trade_id 485 } 486 487 #[must_use] 488 pub const fn input_generation(self) -> u64 { 489 self.input_generation 490 } 491 492 #[must_use] 493 pub const fn evidence_policy_digest(self) -> RhiEvidencePolicyDigest { 494 self.policy_digest 495 } 496 497 #[must_use] 498 pub const fn state(self) -> RhiReconciliationJobState { 499 self.state 500 } 501 502 #[must_use] 503 pub const fn revision(self) -> u64 { 504 self.revision 505 } 506 507 #[must_use] 508 pub const fn attempt_count(self) -> u16 { 509 self.attempt_count 510 } 511 512 #[must_use] 513 pub const fn failure_count(self) -> u16 { 514 self.failure_count 515 } 516 517 #[must_use] 518 pub const fn next_attempt(self) -> Option<RhiReconciliationUnixMilliseconds> { 519 self.next_attempt 520 } 521 522 #[must_use] 523 pub const fn created_at(self) -> RhiReconciliationUnixMilliseconds { 524 self.created_at 525 } 526 527 #[must_use] 528 pub const fn updated_at(self) -> RhiReconciliationUnixMilliseconds { 529 self.updated_at 530 } 531 } 532 533 impl fmt::Debug for RhiReconciliationJob { 534 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 535 formatter 536 .debug_struct("RhiReconciliationJob") 537 .field("identity", &"[redacted]") 538 .field("state", &self.state) 539 .field("revision", &self.revision) 540 .field("attempt_count", &self.attempt_count) 541 .field("failure_count", &self.failure_count) 542 .field("next_attempt", &self.next_attempt) 543 .field("lease", &self.lease_owner.map(|_| "[redacted]")) 544 .field("lease_expires", &self.lease_expires) 545 .finish() 546 } 547 } 548 549 /// Idempotent schedule result for one exact dirty generation. 550 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 551 pub struct RhiReconciliationScheduleOutcome { 552 job: RhiReconciliationJob, 553 created: bool, 554 } 555 556 impl RhiReconciliationScheduleOutcome { 557 #[must_use] 558 pub const fn job(self) -> RhiReconciliationJob { 559 self.job 560 } 561 562 #[must_use] 563 pub const fn created(self) -> bool { 564 self.created 565 } 566 } 567 568 /// Non-forgeable compare-and-swap lease returned by a successful claim. 569 /// 570 /// ```compile_fail 571 /// use rhi::RhiReconciliationLease; 572 /// 573 /// let _ = RhiReconciliationLease { job: todo!(), owner: todo!(), lease_expires: todo!() }; 574 /// ``` 575 #[derive(Clone, Copy, PartialEq, Eq)] 576 pub struct RhiReconciliationLease { 577 job: RhiReconciliationJob, 578 owner: RhiReconciliationLeaseOwner, 579 lease_expires: RhiReconciliationUnixMilliseconds, 580 } 581 582 impl RhiReconciliationLease { 583 #[must_use] 584 pub const fn job(self) -> RhiReconciliationJob { 585 self.job 586 } 587 588 #[must_use] 589 pub const fn lease_expires(self) -> RhiReconciliationUnixMilliseconds { 590 self.lease_expires 591 } 592 593 pub(crate) const fn owner_bytes(self) -> [u8; LEASE_OWNER_BYTES] { 594 self.owner.0 595 } 596 597 #[must_use] 598 pub fn renewal_due(self) -> RhiReconciliationUnixMilliseconds { 599 RhiReconciliationUnixMilliseconds( 600 self.lease_expires() 601 .get() 602 .saturating_sub(self.job.policy.lease_renewal_ms), 603 ) 604 } 605 606 #[must_use] 607 pub fn retry_delay_upper_bound(self) -> u64 { 608 retry_delay_upper_bound(self.job.policy, self.job.failure_count) 609 } 610 } 611 612 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 613 pub(crate) enum LeaseValidationError { 614 LeaseLost, 615 Storage, 616 } 617 618 pub(crate) async fn validate_exact_lease( 619 transaction: &mut ServiceSqliteTransaction<'_>, 620 lease: RhiReconciliationLease, 621 ) -> Result<(), LeaseValidationError> { 622 let current = read_job(transaction, lease.job.id) 623 .await 624 .map_err(|_| LeaseValidationError::Storage)?; 625 if current == Some(lease.job) 626 && lease.job.state == RhiReconciliationJobState::Leased 627 && lease.job.lease_expires == Some(lease.lease_expires) 628 { 629 Ok(()) 630 } else { 631 Err(LeaseValidationError::LeaseLost) 632 } 633 } 634 635 impl fmt::Debug for RhiReconciliationLease { 636 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 637 formatter 638 .debug_struct("RhiReconciliationLease") 639 .field("job", &self.job) 640 .field("owner", &"[redacted]") 641 .finish() 642 } 643 } 644 645 impl RhiReconciliationJobRepository<'_> { 646 /// Idempotently schedules the current durable dirty generation for one trade. 647 pub async fn schedule_trade( 648 &self, 649 trade_id: TradeId, 650 policy: RhiReconciliationJobPolicy, 651 now: RhiReconciliationUnixMilliseconds, 652 ) -> Result<RhiReconciliationScheduleOutcome, RhiReconciliationJobError> { 653 require_writable(self)?; 654 self.host() 655 .sqlite_host() 656 .transaction(move |transaction| { 657 Box::pin(async move { schedule(transaction, trade_id, policy, now).await }) 658 }) 659 .await 660 .map_err(map_transaction_error) 661 } 662 663 /// Claims the oldest eligible job or reclaims one expired lease. 664 pub async fn claim_next( 665 &self, 666 owner: RhiReconciliationLeaseOwner, 667 now: RhiReconciliationUnixMilliseconds, 668 ) -> Result<Option<RhiReconciliationLease>, RhiReconciliationJobError> { 669 require_writable(self)?; 670 self.host() 671 .sqlite_host() 672 .transaction(move |transaction| { 673 Box::pin(async move { claim_next(transaction, owner, now).await }) 674 }) 675 .await 676 .map_err(map_transaction_error) 677 } 678 679 /// Renews an unexpired lease only after its configured renewal point. 680 pub async fn renew( 681 &self, 682 lease: RhiReconciliationLease, 683 now: RhiReconciliationUnixMilliseconds, 684 ) -> Result<RhiReconciliationLease, RhiReconciliationJobError> { 685 require_writable(self)?; 686 if now >= lease.lease_expires() { 687 return Err(RhiReconciliationJobError::new( 688 RhiReconciliationJobErrorKind::LeaseLost, 689 )); 690 } 691 if now < lease.renewal_due() { 692 return Err(RhiReconciliationJobError::new( 693 RhiReconciliationJobErrorKind::NotReady, 694 )); 695 } 696 self.host() 697 .sqlite_host() 698 .transaction(move |transaction| { 699 Box::pin(async move { renew(transaction, lease, now).await }) 700 }) 701 .await 702 .map_err(map_transaction_error) 703 } 704 705 /// Records one failed leased attempt and schedules bounded full-jitter retry. 706 pub async fn record_failure( 707 &self, 708 lease: RhiReconciliationLease, 709 now: RhiReconciliationUnixMilliseconds, 710 delay: RhiReconciliationRetryDelayMilliseconds, 711 ) -> Result<RhiReconciliationJob, RhiReconciliationJobError> { 712 require_writable(self)?; 713 if now >= lease.lease_expires() { 714 return Err(RhiReconciliationJobError::new( 715 RhiReconciliationJobErrorKind::LeaseLost, 716 )); 717 } 718 if delay.get() > lease.retry_delay_upper_bound() { 719 return Err(RhiReconciliationJobError::new( 720 RhiReconciliationJobErrorKind::InvalidInput, 721 )); 722 } 723 self.host() 724 .sqlite_host() 725 .transaction(move |transaction| { 726 Box::pin(async move { record_failure(transaction, lease, now, delay).await }) 727 }) 728 .await 729 .map_err(map_transaction_error) 730 } 731 } 732 733 fn require_writable( 734 repository: &RhiReconciliationJobRepository<'_>, 735 ) -> Result<(), RhiReconciliationJobError> { 736 if repository.host().mode() == RhiStateHostMode::ReadWriteExisting { 737 Ok(()) 738 } else { 739 Err(RhiReconciliationJobError::new( 740 RhiReconciliationJobErrorKind::InvalidMode, 741 )) 742 } 743 } 744 745 async fn schedule( 746 transaction: &mut ServiceSqliteTransaction<'_>, 747 trade_id: TradeId, 748 policy: RhiReconciliationJobPolicy, 749 now: RhiReconciliationUnixMilliseconds, 750 ) -> Result<RhiReconciliationScheduleOutcome, OperationError> { 751 let dirty = read_dirty(transaction, trade_id).await?; 752 let id = job_id(trade_id, dirty.generation, dirty.policy); 753 if let Some(job) = read_job(transaction, id).await? { 754 if job.trade_id != trade_id 755 || job.input_generation != dirty.generation 756 || job.policy_digest != dirty.policy 757 { 758 return Err(OperationError::Storage); 759 } 760 return Ok(RhiReconciliationScheduleOutcome { 761 job, 762 created: false, 763 }); 764 } 765 766 if let Some(active) = read_active_job(transaction, trade_id).await? { 767 if active.input_generation >= dirty.generation { 768 return Err(OperationError::DirtyGenerationConflict); 769 } 770 let result = sqlx::query(SUPERSEDE_JOB_SQL) 771 .bind(i64_value(now.get())?) 772 .bind(active.id.as_bytes().as_slice()) 773 .bind(i64_value(active.revision)?) 774 .execute(&mut *transaction) 775 .await 776 .map_err(|_| OperationError::Storage)?; 777 if result.rows_affected() != 1 { 778 return Err(OperationError::DirtyGenerationConflict); 779 } 780 } 781 782 let active = sqlx::query(ACTIVE_JOB_COUNT_SQL) 783 .fetch_one(&mut *transaction) 784 .await 785 .map_err(|_| OperationError::Storage)? 786 .try_get::<i64, _>("active_count") 787 .ok() 788 .and_then(|value| u32::try_from(value).ok()) 789 .ok_or(OperationError::Storage)?; 790 if active >= policy.queue_capacity { 791 return Err(OperationError::QueueFull); 792 } 793 794 let now = i64_value(now.get())?; 795 let result = sqlx::query(INSERT_JOB_SQL) 796 .bind(id.as_bytes().as_slice()) 797 .bind(trade_id.as_bytes().as_slice()) 798 .bind(i64_value(dirty.generation)?) 799 .bind(dirty.policy.as_bytes().as_slice()) 800 .bind(i64::from(policy.max_attempts)) 801 .bind(i64_value(policy.lease_duration_ms)?) 802 .bind(i64_value(policy.lease_renewal_ms)?) 803 .bind(i64_value(policy.initial_backoff_ms)?) 804 .bind(i64_value(policy.maximum_backoff_ms)?) 805 .bind(now) 806 .bind(now) 807 .bind(now) 808 .execute(&mut *transaction) 809 .await 810 .map_err(|_| OperationError::Storage)?; 811 if result.rows_affected() != 1 { 812 return Err(OperationError::Storage); 813 } 814 let job = read_job(transaction, id) 815 .await? 816 .ok_or(OperationError::Storage)?; 817 Ok(RhiReconciliationScheduleOutcome { job, created: true }) 818 } 819 820 async fn claim_next( 821 transaction: &mut ServiceSqliteTransaction<'_>, 822 owner: RhiReconciliationLeaseOwner, 823 now: RhiReconciliationUnixMilliseconds, 824 ) -> Result<Option<RhiReconciliationLease>, OperationError> { 825 let now_value = i64_value(now.get())?; 826 sqlx::query(EXHAUST_EXPIRED_SQL) 827 .bind(now_value) 828 .bind(now_value) 829 .bind(now_value) 830 .execute(&mut *transaction) 831 .await 832 .map_err(|_| OperationError::Storage)?; 833 let Some(row) = sqlx::query(READ_CLAIMABLE_SQL) 834 .bind(now_value) 835 .bind(now_value) 836 .bind(now_value) 837 .fetch_optional(&mut *transaction) 838 .await 839 .map_err(|_| OperationError::Storage)? 840 else { 841 return Ok(None); 842 }; 843 let candidate = decode_job(row)?; 844 let expires = now 845 .get() 846 .checked_add(candidate.policy.lease_duration_ms) 847 .filter(|value| *value <= MAX_UNIX_MILLISECONDS) 848 .ok_or(OperationError::InvalidInput)?; 849 let result = sqlx::query(CLAIM_JOB_SQL) 850 .bind(owner.0.as_slice()) 851 .bind(i64_value(expires)?) 852 .bind(now_value) 853 .bind(candidate.id.as_bytes().as_slice()) 854 .bind(i64_value(candidate.revision)?) 855 .bind(now_value) 856 .bind(now_value) 857 .bind(now_value) 858 .execute(&mut *transaction) 859 .await 860 .map_err(|_| OperationError::Storage)?; 861 if result.rows_affected() != 1 { 862 return Err(OperationError::LeaseLost); 863 } 864 let job = read_job(transaction, candidate.id) 865 .await? 866 .filter(|job| job.state == RhiReconciliationJobState::Leased) 867 .ok_or(OperationError::Storage)?; 868 let lease_expires = job.lease_expires.ok_or(OperationError::Storage)?; 869 Ok(Some(RhiReconciliationLease { 870 job, 871 owner, 872 lease_expires, 873 })) 874 } 875 876 async fn renew( 877 transaction: &mut ServiceSqliteTransaction<'_>, 878 lease: RhiReconciliationLease, 879 now: RhiReconciliationUnixMilliseconds, 880 ) -> Result<RhiReconciliationLease, OperationError> { 881 let previous_expiry = lease.lease_expires().get(); 882 let expires = now 883 .get() 884 .checked_add(lease.job.policy.lease_duration_ms) 885 .filter(|value| *value > previous_expiry && *value <= MAX_UNIX_MILLISECONDS) 886 .ok_or(OperationError::NotReady)?; 887 let result = sqlx::query(RENEW_JOB_SQL) 888 .bind(i64_value(expires)?) 889 .bind(i64_value(now.get())?) 890 .bind(lease.job.id.as_bytes().as_slice()) 891 .bind(i64_value(lease.job.revision)?) 892 .bind(lease.owner.0.as_slice()) 893 .bind(i64_value(previous_expiry)?) 894 .bind(i64_value(now.get())?) 895 .bind(i64_value(now.get())?) 896 .execute(&mut *transaction) 897 .await 898 .map_err(|_| OperationError::Storage)?; 899 if result.rows_affected() != 1 { 900 return Err(OperationError::LeaseLost); 901 } 902 let job = read_job(transaction, lease.job.id) 903 .await? 904 .filter(|job| job.state == RhiReconciliationJobState::Leased) 905 .ok_or(OperationError::Storage)?; 906 Ok(RhiReconciliationLease { 907 job, 908 owner: lease.owner, 909 lease_expires: job.lease_expires.ok_or(OperationError::Storage)?, 910 }) 911 } 912 913 async fn record_failure( 914 transaction: &mut ServiceSqliteTransaction<'_>, 915 lease: RhiReconciliationLease, 916 now: RhiReconciliationUnixMilliseconds, 917 delay: RhiReconciliationRetryDelayMilliseconds, 918 ) -> Result<RhiReconciliationJob, OperationError> { 919 let now_value = i64_value(now.get())?; 920 let query = if lease.job.attempt_count >= lease.job.policy.max_attempts { 921 sqlx::query(EXHAUST_JOB_SQL) 922 .bind(now_value) 923 .bind(lease.job.id.as_bytes().as_slice()) 924 .bind(i64_value(lease.job.revision)?) 925 .bind(lease.owner.0.as_slice()) 926 .bind(i64_value(lease.lease_expires().get())?) 927 .bind(now_value) 928 } else { 929 let next = now 930 .get() 931 .checked_add(delay.get()) 932 .filter(|value| *value <= MAX_UNIX_MILLISECONDS) 933 .ok_or(OperationError::InvalidInput)?; 934 sqlx::query(RETRY_JOB_SQL) 935 .bind(i64_value(next)?) 936 .bind(now_value) 937 .bind(lease.job.id.as_bytes().as_slice()) 938 .bind(i64_value(lease.job.revision)?) 939 .bind(lease.owner.0.as_slice()) 940 .bind(i64_value(lease.lease_expires().get())?) 941 .bind(now_value) 942 }; 943 let result = query 944 .execute(&mut *transaction) 945 .await 946 .map_err(|_| OperationError::Storage)?; 947 if result.rows_affected() != 1 { 948 return Err(OperationError::LeaseLost); 949 } 950 read_job(transaction, lease.job.id) 951 .await? 952 .ok_or(OperationError::Storage) 953 } 954 955 async fn read_dirty( 956 transaction: &mut ServiceSqliteTransaction<'_>, 957 trade_id: TradeId, 958 ) -> Result<DirtyGeneration, OperationError> { 959 let row = sqlx::query(READ_DIRTY_SQL) 960 .bind(trade_id.as_bytes().as_slice()) 961 .fetch_optional(&mut *transaction) 962 .await 963 .map_err(|_| OperationError::Storage)? 964 .ok_or(OperationError::DirtyGenerationConflict)?; 965 Ok(DirtyGeneration { 966 generation: positive_u64(&row, "generation")?, 967 policy: RhiEvidencePolicyDigest::from_bytes(exact_bytes::<32>( 968 &row, 969 "evidence_policy_sha256", 970 "evidence_policy_bytes", 971 )?), 972 }) 973 } 974 975 async fn read_job( 976 transaction: &mut ServiceSqliteTransaction<'_>, 977 id: RhiReconciliationJobId, 978 ) -> Result<Option<RhiReconciliationJob>, OperationError> { 979 sqlx::query(READ_JOB_SQL) 980 .bind(id.as_bytes().as_slice()) 981 .fetch_optional(&mut *transaction) 982 .await 983 .map_err(|_| OperationError::Storage)? 984 .map(decode_job) 985 .transpose() 986 } 987 988 async fn read_active_job( 989 transaction: &mut ServiceSqliteTransaction<'_>, 990 trade_id: TradeId, 991 ) -> Result<Option<RhiReconciliationJob>, OperationError> { 992 let rows = sqlx::query(READ_ACTIVE_JOB_SQL) 993 .bind(trade_id.as_bytes().as_slice()) 994 .fetch_all(&mut *transaction) 995 .await 996 .map_err(|_| OperationError::Storage)?; 997 match rows.len() { 998 0 => Ok(None), 999 1 => Ok(Some(decode_job( 1000 rows.into_iter().next().ok_or(OperationError::Storage)?, 1001 )?)), 1002 _ => Err(OperationError::Storage), 1003 } 1004 } 1005 1006 fn decode_job(row: sqlx::sqlite::SqliteRow) -> Result<RhiReconciliationJob, OperationError> { 1007 let id = RhiReconciliationJobId(exact_bytes::<32>(&row, "job_id", "job_id_bytes")?); 1008 let trade_id = TradeId::from_bytes(exact_bytes::<16>(&row, "trade_id", "trade_id_bytes")?); 1009 let input_generation = positive_u64(&row, "input_generation")?; 1010 let policy_digest = RhiEvidencePolicyDigest::from_bytes(exact_bytes::<32>( 1011 &row, 1012 "evidence_policy_sha256", 1013 "evidence_policy_bytes", 1014 )?); 1015 let state = match bounded_text(&row, "state", "state_bytes", 10)?.as_str() { 1016 "ready" => RhiReconciliationJobState::Ready, 1017 "leased" => RhiReconciliationJobState::Leased, 1018 "exhausted" => RhiReconciliationJobState::Exhausted, 1019 "superseded" => RhiReconciliationJobState::Superseded, 1020 "completed" => RhiReconciliationJobState::Completed, 1021 _ => return Err(OperationError::Storage), 1022 }; 1023 let revision = positive_u64(&row, "revision")?; 1024 let attempt_count = bounded_u16(&row, "attempt_count", 0, MAX_ATTEMPTS)?; 1025 let failure_count = bounded_u16(&row, "failure_count", 0, attempt_count)?; 1026 let max_attempts = bounded_u16(&row, "max_attempts", 1, MAX_ATTEMPTS)?; 1027 let lease_duration_ms = bounded_u64(&row, "lease_duration_ms", 1_000, MAX_LEASE_MILLISECONDS)?; 1028 let lease_renewal_ms = bounded_u64(&row, "lease_renewal_ms", 100, MAX_RENEWAL_MILLISECONDS)?; 1029 let initial_backoff_ms = bounded_u64( 1030 &row, 1031 "initial_backoff_ms", 1032 1, 1033 MAX_INITIAL_BACKOFF_MILLISECONDS, 1034 )?; 1035 let maximum_backoff_ms = bounded_u64(&row, "maximum_backoff_ms", 1, MAX_BACKOFF_MILLISECONDS)?; 1036 let policy = RhiReconciliationJobPolicy::new( 1037 RHI_RECONCILIATION_JOB_MAX_ACTIVE, 1038 lease_duration_ms, 1039 lease_renewal_ms, 1040 max_attempts, 1041 initial_backoff_ms, 1042 maximum_backoff_ms, 1043 ) 1044 .map_err(|_| OperationError::Storage)?; 1045 let next_attempt = optional_millis(&row, "next_attempt_unix_ms")?; 1046 let lease_expires = optional_millis(&row, "lease_expires_unix_ms")?; 1047 let lease_owner = optional_owner(&row)?; 1048 let created_at = 1049 RhiReconciliationUnixMilliseconds(nonnegative_u64(&row, "created_at_unix_ms")?); 1050 let updated_at = 1051 RhiReconciliationUnixMilliseconds(nonnegative_u64(&row, "updated_at_unix_ms")?); 1052 if updated_at < created_at 1053 || failure_count > attempt_count 1054 || job_id(trade_id, input_generation, policy_digest) != id 1055 || !valid_state_fields( 1056 state, 1057 attempt_count, 1058 max_attempts, 1059 next_attempt, 1060 lease_owner, 1061 lease_expires, 1062 ) 1063 { 1064 return Err(OperationError::Storage); 1065 } 1066 Ok(RhiReconciliationJob { 1067 id, 1068 trade_id, 1069 input_generation, 1070 policy_digest, 1071 state, 1072 revision, 1073 attempt_count, 1074 failure_count, 1075 policy, 1076 next_attempt, 1077 lease_owner, 1078 lease_expires, 1079 created_at, 1080 updated_at, 1081 }) 1082 } 1083 1084 fn valid_state_fields( 1085 state: RhiReconciliationJobState, 1086 attempt_count: u16, 1087 max_attempts: u16, 1088 next_attempt: Option<RhiReconciliationUnixMilliseconds>, 1089 owner: Option<RhiReconciliationLeaseOwner>, 1090 expires: Option<RhiReconciliationUnixMilliseconds>, 1091 ) -> bool { 1092 match state { 1093 RhiReconciliationJobState::Ready => { 1094 attempt_count < max_attempts 1095 && next_attempt.is_some() 1096 && owner.is_none() 1097 && expires.is_none() 1098 } 1099 RhiReconciliationJobState::Leased => { 1100 attempt_count > 0 1101 && attempt_count <= max_attempts 1102 && next_attempt.is_none() 1103 && owner.is_some() 1104 && expires.is_some() 1105 } 1106 RhiReconciliationJobState::Exhausted 1107 | RhiReconciliationJobState::Superseded 1108 | RhiReconciliationJobState::Completed => { 1109 next_attempt.is_none() && owner.is_none() && expires.is_none() 1110 } 1111 } 1112 } 1113 1114 fn job_id( 1115 trade_id: TradeId, 1116 generation: u64, 1117 policy: RhiEvidencePolicyDigest, 1118 ) -> RhiReconciliationJobId { 1119 let mut hasher = Sha256::new(); 1120 hasher.update(JOB_ID_DOMAIN); 1121 hasher.update(trade_id.as_bytes()); 1122 hasher.update(generation.to_be_bytes()); 1123 hasher.update(policy.as_bytes()); 1124 RhiReconciliationJobId(hasher.finalize().into()) 1125 } 1126 1127 fn retry_delay_upper_bound(policy: RhiReconciliationJobPolicy, failure_count: u16) -> u64 { 1128 let exponent = u32::from(failure_count.min(63)); 1129 policy 1130 .initial_backoff_ms 1131 .saturating_mul(1_u64.checked_shl(exponent).unwrap_or(u64::MAX)) 1132 .min(policy.maximum_backoff_ms) 1133 } 1134 1135 fn exact_bytes<const N: usize>( 1136 row: &sqlx::sqlite::SqliteRow, 1137 field: &str, 1138 length_field: &str, 1139 ) -> Result<[u8; N], OperationError> { 1140 if row.try_get::<i64, _>(length_field).ok() != i64::try_from(N).ok() { 1141 return Err(OperationError::Storage); 1142 } 1143 row.try_get::<Vec<u8>, _>(field) 1144 .map_err(|_| OperationError::Storage)? 1145 .try_into() 1146 .map_err(|_| OperationError::Storage) 1147 } 1148 1149 fn optional_owner( 1150 row: &sqlx::sqlite::SqliteRow, 1151 ) -> Result<Option<RhiReconciliationLeaseOwner>, OperationError> { 1152 let length = row 1153 .try_get::<Option<i64>, _>("lease_owner_bytes") 1154 .map_err(|_| OperationError::Storage)?; 1155 match length { 1156 None => Ok(None), 1157 Some(value) if value == LEASE_OWNER_BYTES as i64 => { 1158 let bytes: [u8; LEASE_OWNER_BYTES] = row 1159 .try_get::<Vec<u8>, _>("lease_owner") 1160 .map_err(|_| OperationError::Storage)? 1161 .try_into() 1162 .map_err(|_| OperationError::Storage)?; 1163 RhiReconciliationLeaseOwner::from_bytes(bytes) 1164 .map(Some) 1165 .map_err(|_| OperationError::Storage) 1166 } 1167 Some(_) => Err(OperationError::Storage), 1168 } 1169 } 1170 1171 fn bounded_text( 1172 row: &sqlx::sqlite::SqliteRow, 1173 field: &str, 1174 length_field: &str, 1175 maximum: usize, 1176 ) -> Result<String, OperationError> { 1177 let length = row 1178 .try_get::<i64, _>(length_field) 1179 .ok() 1180 .and_then(|value| usize::try_from(value).ok()) 1181 .filter(|value| *value > 0 && *value <= maximum) 1182 .ok_or(OperationError::Storage)?; 1183 let value = row 1184 .try_get::<String, _>(field) 1185 .map_err(|_| OperationError::Storage)?; 1186 if value.len() == length { 1187 Ok(value) 1188 } else { 1189 Err(OperationError::Storage) 1190 } 1191 } 1192 1193 fn bounded_u16( 1194 row: &sqlx::sqlite::SqliteRow, 1195 field: &str, 1196 minimum: u16, 1197 maximum: u16, 1198 ) -> Result<u16, OperationError> { 1199 row.try_get::<i64, _>(field) 1200 .ok() 1201 .and_then(|value| u16::try_from(value).ok()) 1202 .filter(|value| (*value >= minimum) && (*value <= maximum)) 1203 .ok_or(OperationError::Storage) 1204 } 1205 1206 fn bounded_u64( 1207 row: &sqlx::sqlite::SqliteRow, 1208 field: &str, 1209 minimum: u64, 1210 maximum: u64, 1211 ) -> Result<u64, OperationError> { 1212 row.try_get::<i64, _>(field) 1213 .ok() 1214 .and_then(|value| u64::try_from(value).ok()) 1215 .filter(|value| (*value >= minimum) && (*value <= maximum)) 1216 .ok_or(OperationError::Storage) 1217 } 1218 1219 fn positive_u64(row: &sqlx::sqlite::SqliteRow, field: &str) -> Result<u64, OperationError> { 1220 bounded_u64(row, field, 1, MAX_UNIX_MILLISECONDS) 1221 } 1222 1223 fn nonnegative_u64(row: &sqlx::sqlite::SqliteRow, field: &str) -> Result<u64, OperationError> { 1224 bounded_u64(row, field, 0, MAX_UNIX_MILLISECONDS) 1225 } 1226 1227 fn optional_millis( 1228 row: &sqlx::sqlite::SqliteRow, 1229 field: &str, 1230 ) -> Result<Option<RhiReconciliationUnixMilliseconds>, OperationError> { 1231 row.try_get::<Option<i64>, _>(field) 1232 .map_err(|_| OperationError::Storage)? 1233 .map(|value| { 1234 u64::try_from(value) 1235 .map(RhiReconciliationUnixMilliseconds) 1236 .map_err(|_| OperationError::Storage) 1237 }) 1238 .transpose() 1239 } 1240 1241 fn i64_value(value: u64) -> Result<i64, OperationError> { 1242 i64::try_from(value).map_err(|_| OperationError::InvalidInput) 1243 } 1244 1245 #[derive(Clone, Copy)] 1246 struct DirtyGeneration { 1247 generation: u64, 1248 policy: RhiEvidencePolicyDigest, 1249 } 1250 1251 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 1252 enum OperationError { 1253 InvalidInput, 1254 QueueFull, 1255 DirtyGenerationConflict, 1256 LeaseLost, 1257 NotReady, 1258 Storage, 1259 } 1260 1261 fn map_transaction_error( 1262 error: ServiceSqliteTransactionError<OperationError>, 1263 ) -> RhiReconciliationJobError { 1264 if error.kind() == ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown { 1265 return RhiReconciliationJobError::new(RhiReconciliationJobErrorKind::CommitOutcomeUnknown); 1266 } 1267 RhiReconciliationJobError::new(match error.operation_error().copied() { 1268 Some(OperationError::InvalidInput) => RhiReconciliationJobErrorKind::InvalidInput, 1269 Some(OperationError::QueueFull) => RhiReconciliationJobErrorKind::QueueFull, 1270 Some(OperationError::DirtyGenerationConflict) => { 1271 RhiReconciliationJobErrorKind::DirtyGenerationConflict 1272 } 1273 Some(OperationError::LeaseLost) => RhiReconciliationJobErrorKind::LeaseLost, 1274 Some(OperationError::NotReady) => RhiReconciliationJobErrorKind::NotReady, 1275 Some(OperationError::Storage) | None => RhiReconciliationJobErrorKind::Storage, 1276 }) 1277 }