rhi

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

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 }