rhi

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

publication_execution.rs (71086B)


      1 //! Durable exact-byte publication claims, outcomes, retry, and recovery.
      2 
      3 use core::fmt;
      4 use std::error::Error;
      5 
      6 use radroots_service_sqlite::{
      7     ServiceSqliteTransaction, ServiceSqliteTransactionError, ServiceSqliteTransactionErrorKind,
      8 };
      9 use radroots_transport::BoxFuture;
     10 use sha2::{Digest as _, Sha256};
     11 use sqlx::Row as _;
     12 
     13 use crate::{
     14     RHI_PUBLICATION_ATTEMPT_NUMBER_MAXIMUM, RHI_PUBLICATION_TARGET_ORDINAL_MAXIMUM,
     15     RhiCommittedPublication, RhiJitterBoundMilliseconds, RhiPublicationAttemptEvidence,
     16     RhiPublicationAttemptId, RhiPublicationAttemptOutcome, RhiPublicationAuthority,
     17     RhiPublicationMode, RhiPublicationOutboxId, RhiPublicationOutboxRepository,
     18     RhiPublicationTargetState, RhiPublicationUnixMilliseconds, RhiRuntimeAdapterErrorKind,
     19     RhiStateHostMode, RhiTimeEntropyAdapters,
     20     publication_attempt::derive_attempt_id,
     21     publication_submission::{ReadError, read_committed},
     22 };
     23 
     24 /// Exact version of the durable publication-execution contract.
     25 pub const RHI_PUBLICATION_EXECUTION_CONTRACT_VERSION: u32 = 1;
     26 
     27 const LEASE_OWNER_BYTES: usize = 16;
     28 const MAX_UNIX_MILLISECONDS: u64 = i64::MAX as u64;
     29 const TARGET_SET_DOMAIN: &[u8] = b"radroots.rhi.publication_target_set.v1\0";
     30 const READ_OUTBOX_SQL: &str = r#"SELECT
     31     CASE WHEN typeof(outbox_id) = 'blob' AND length(outbox_id) = 32
     32         THEN outbox_id ELSE NULL END AS outbox_id,
     33     CASE WHEN typeof(event_sha256) = 'blob' AND length(event_sha256) = 32
     34         THEN event_sha256 ELSE NULL END AS event_sha256,
     35     CASE WHEN typeof(publication_authority_sha256) = 'blob'
     36             AND length(publication_authority_sha256) = 32
     37         THEN publication_authority_sha256 ELSE NULL END AS publication_authority_sha256,
     38     CASE WHEN typeof(target_set_sha256) = 'blob' AND length(target_set_sha256) = 32
     39         THEN target_set_sha256 ELSE NULL END AS target_set_sha256,
     40     target_count, required_target_count, max_attempts,
     41     initial_backoff_ms, maximum_backoff_ms, attempt_deadline_ms,
     42     length(CAST(state AS BLOB)) AS state_bytes, substr(state, 1, 9) AS state,
     43     revision, next_attempt_unix_ms,
     44     CASE WHEN lease_owner IS NULL THEN NULL ELSE length(lease_owner) END AS lease_owner_bytes,
     45     CASE WHEN lease_owner IS NULL THEN NULL ELSE substr(lease_owner, 1, 17) END AS lease_owner,
     46     lease_expires_unix_ms, created_at_unix_ms, updated_at_unix_ms
     47 FROM publication_outbox
     48 WHERE outbox_id = ?
     49 LIMIT 2"#;
     50 
     51 const READ_CLAIMABLE_OUTBOX_SQL: &str = r#"SELECT
     52     CASE WHEN typeof(outbox_id) = 'blob' AND length(outbox_id) = 32
     53         THEN outbox_id ELSE NULL END AS outbox_id,
     54     CASE WHEN typeof(event_sha256) = 'blob' AND length(event_sha256) = 32
     55         THEN event_sha256 ELSE NULL END AS event_sha256,
     56     CASE WHEN typeof(publication_authority_sha256) = 'blob'
     57             AND length(publication_authority_sha256) = 32
     58         THEN publication_authority_sha256 ELSE NULL END AS publication_authority_sha256,
     59     CASE WHEN typeof(target_set_sha256) = 'blob' AND length(target_set_sha256) = 32
     60         THEN target_set_sha256 ELSE NULL END AS target_set_sha256,
     61     target_count, required_target_count, max_attempts,
     62     initial_backoff_ms, maximum_backoff_ms, attempt_deadline_ms,
     63     length(CAST(state AS BLOB)) AS state_bytes, substr(state, 1, 9) AS state,
     64     revision, next_attempt_unix_ms,
     65     NULL AS lease_owner_bytes, NULL AS lease_owner,
     66     lease_expires_unix_ms, created_at_unix_ms, updated_at_unix_ms
     67 FROM publication_outbox
     68 WHERE state = 'pending' AND next_attempt_unix_ms <= ?
     69 ORDER BY next_attempt_unix_ms, created_at_unix_ms, outbox_id
     70 LIMIT 1"#;
     71 
     72 const READ_EXPIRED_OUTBOX_SQL: &str = r#"SELECT
     73     CASE WHEN typeof(outbox_id) = 'blob' AND length(outbox_id) = 32
     74         THEN outbox_id ELSE NULL END AS outbox_id,
     75     CASE WHEN typeof(event_sha256) = 'blob' AND length(event_sha256) = 32
     76         THEN event_sha256 ELSE NULL END AS event_sha256,
     77     CASE WHEN typeof(publication_authority_sha256) = 'blob'
     78             AND length(publication_authority_sha256) = 32
     79         THEN publication_authority_sha256 ELSE NULL END AS publication_authority_sha256,
     80     CASE WHEN typeof(target_set_sha256) = 'blob' AND length(target_set_sha256) = 32
     81         THEN target_set_sha256 ELSE NULL END AS target_set_sha256,
     82     target_count, required_target_count, max_attempts,
     83     initial_backoff_ms, maximum_backoff_ms, attempt_deadline_ms,
     84     length(CAST(state AS BLOB)) AS state_bytes, substr(state, 1, 9) AS state,
     85     revision, next_attempt_unix_ms,
     86     length(lease_owner) AS lease_owner_bytes, substr(lease_owner, 1, 17) AS lease_owner,
     87     lease_expires_unix_ms, created_at_unix_ms, updated_at_unix_ms
     88 FROM publication_outbox
     89 WHERE state = 'leased' AND lease_expires_unix_ms <= ?
     90 ORDER BY lease_expires_unix_ms, created_at_unix_ms, outbox_id
     91 LIMIT 1"#;
     92 
     93 const READ_TARGETS_SQL: &str = r#"SELECT target_ordinal,
     94     length(CAST(relay_id AS BLOB)) AS relay_id_bytes, substr(relay_id, 1, 65) AS relay_id,
     95     required,
     96     length(CAST(state AS BLOB)) AS state_bytes, substr(state, 1, 14) AS state,
     97     revision, attempt_count, next_attempt_unix_ms,
     98     CASE WHEN last_attempt_id IS NULL THEN NULL ELSE length(last_attempt_id) END
     99         AS last_attempt_id_bytes,
    100     CASE WHEN last_attempt_id IS NULL THEN NULL ELSE substr(last_attempt_id, 1, 33) END
    101         AS last_attempt_id,
    102     updated_at_unix_ms
    103 FROM publication_targets
    104 WHERE outbox_id = ?
    105 ORDER BY target_ordinal
    106 LIMIT 33"#;
    107 
    108 const CLAIM_OUTBOX_SQL: &str = r#"UPDATE publication_outbox
    109 SET state = 'leased', revision = revision + 1,
    110     next_attempt_unix_ms = NULL, lease_owner = ?, lease_expires_unix_ms = ?,
    111     updated_at_unix_ms = ?
    112 WHERE outbox_id = ? AND revision = ? AND state = 'pending'
    113     AND next_attempt_unix_ms <= ? AND updated_at_unix_ms <= ?"#;
    114 
    115 const PREPARE_TARGET_SQL: &str = r#"UPDATE publication_targets
    116 SET state = 'submitted', revision = revision + 1,
    117     attempt_count = attempt_count + 1, next_attempt_unix_ms = NULL,
    118     last_attempt_id = ?, updated_at_unix_ms = ?
    119 WHERE outbox_id = ? AND target_ordinal = ? AND revision = ?
    120     AND state IN ('pending', 'failed', 'rate_limited', 'unknown')
    121     AND attempt_count < ? AND next_attempt_unix_ms <= ?
    122     AND updated_at_unix_ms <= ?"#;
    123 
    124 const INSERT_ATTEMPT_SQL: &str = r#"INSERT INTO publication_attempts (
    125     attempt_id, outbox_id, target_ordinal, attempt_number, event_sha256,
    126     lease_owner, started_at_unix_ms, finished_at_unix_ms, outcome, result_code
    127 ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)"#;
    128 
    129 const READ_ATTEMPT_SQL: &str = r#"SELECT
    130     CASE WHEN typeof(attempt_id) = 'blob' AND length(attempt_id) = 32
    131         THEN attempt_id ELSE NULL END AS attempt_id,
    132     CASE WHEN typeof(outbox_id) = 'blob' AND length(outbox_id) = 32
    133         THEN outbox_id ELSE NULL END AS outbox_id,
    134     target_ordinal, attempt_number,
    135     CASE WHEN typeof(event_sha256) = 'blob' AND length(event_sha256) = 32
    136         THEN event_sha256 ELSE NULL END AS event_sha256,
    137     CASE WHEN typeof(lease_owner) = 'blob' AND length(lease_owner) = 16
    138         THEN lease_owner ELSE NULL END AS lease_owner,
    139     started_at_unix_ms, finished_at_unix_ms,
    140     length(CAST(outcome AS BLOB)) AS outcome_bytes, substr(outcome, 1, 14) AS outcome,
    141     length(CAST(result_code AS BLOB)) AS result_code_bytes,
    142     substr(result_code, 1, 65) AS result_code
    143 FROM publication_attempts
    144 WHERE attempt_id = ?
    145 LIMIT 2"#;
    146 
    147 const UPDATE_TARGET_OUTCOME_SQL: &str = r#"UPDATE publication_targets
    148 SET state = ?, revision = revision + 1, next_attempt_unix_ms = ?,
    149     updated_at_unix_ms = ?
    150 WHERE outbox_id = ? AND target_ordinal = ? AND revision = ?
    151     AND state = 'submitted' AND attempt_count = ? AND last_attempt_id = ?"#;
    152 
    153 const UPDATE_OUTBOX_AFTER_ATTEMPT_SQL: &str = r#"UPDATE publication_outbox
    154 SET state = ?, revision = revision + 1, next_attempt_unix_ms = ?,
    155     lease_owner = NULL, lease_expires_unix_ms = NULL, updated_at_unix_ms = ?
    156 WHERE outbox_id = ? AND revision = ? AND state = 'leased'
    157     AND lease_owner = ? AND lease_expires_unix_ms = ?
    158     AND lease_expires_unix_ms > ? AND updated_at_unix_ms <= ?"#;
    159 
    160 const UPDATE_OUTBOX_RECOVERY_SQL: &str = r#"UPDATE publication_outbox
    161 SET state = ?, revision = revision + 1, next_attempt_unix_ms = ?,
    162     lease_owner = NULL, lease_expires_unix_ms = NULL, updated_at_unix_ms = ?
    163 WHERE outbox_id = ? AND revision = ? AND state = 'leased'
    164     AND lease_owner = ? AND lease_expires_unix_ms = ?
    165     AND lease_expires_unix_ms <= ? AND updated_at_unix_ms <= ?"#;
    166 
    167 /// Stable durable outbox lifecycle state.
    168 #[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
    169 pub enum RhiPublicationOutboxState {
    170     Pending,
    171     Leased,
    172     Complete,
    173     Blocked,
    174 }
    175 
    176 impl RhiPublicationOutboxState {
    177     /// Returns the exact machine-contract spelling.
    178     #[must_use]
    179     pub const fn code(self) -> &'static str {
    180         match self {
    181             Self::Pending => "pending",
    182             Self::Leased => "leased",
    183             Self::Complete => "complete",
    184             Self::Blocked => "blocked",
    185         }
    186     }
    187 }
    188 
    189 /// Stable source-free durable publication-execution failure class.
    190 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
    191 pub enum RhiPublicationExecutionErrorKind {
    192     InvalidMode,
    193     InvalidInput,
    194     NotReady,
    195     LeaseLost,
    196     Invariant,
    197     ClockUnavailable,
    198     EntropyUnavailable,
    199     Storage,
    200     CommitOutcomeUnknown,
    201 }
    202 
    203 impl RhiPublicationExecutionErrorKind {
    204     /// Returns the stable machine-readable failure code.
    205     #[must_use]
    206     pub const fn code(self) -> &'static str {
    207         match self {
    208             Self::InvalidMode => "publication_execution_mode_invalid",
    209             Self::InvalidInput => "publication_execution_input_invalid",
    210             Self::NotReady => "publication_execution_not_ready",
    211             Self::LeaseLost => "publication_execution_lease_lost",
    212             Self::Invariant => "publication_execution_invariant_failed",
    213             Self::ClockUnavailable => "publication_execution_clock_unavailable",
    214             Self::EntropyUnavailable => "publication_execution_entropy_unavailable",
    215             Self::Storage => "publication_execution_storage_failed",
    216             Self::CommitOutcomeUnknown => "publication_execution_commit_outcome_unknown",
    217         }
    218     }
    219 }
    220 
    221 /// Redacted source-free durable publication-execution failure.
    222 #[derive(Clone, Copy, PartialEq, Eq)]
    223 pub struct RhiPublicationExecutionError {
    224     kind: RhiPublicationExecutionErrorKind,
    225 }
    226 
    227 impl RhiPublicationExecutionError {
    228     const fn new(kind: RhiPublicationExecutionErrorKind) -> Self {
    229         Self { kind }
    230     }
    231 
    232     /// Returns the stable failure class.
    233     #[must_use]
    234     pub const fn kind(self) -> RhiPublicationExecutionErrorKind {
    235         self.kind
    236     }
    237 
    238     /// Returns the stable machine-readable failure code.
    239     #[must_use]
    240     pub const fn code(self) -> &'static str {
    241         self.kind.code()
    242     }
    243 }
    244 
    245 impl fmt::Display for RhiPublicationExecutionError {
    246     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    247         formatter.write_str(match self.kind {
    248             RhiPublicationExecutionErrorKind::InvalidMode => {
    249                 "RHI publication execution requires writable state"
    250             }
    251             RhiPublicationExecutionErrorKind::InvalidInput => {
    252                 "RHI publication execution input is invalid"
    253             }
    254             RhiPublicationExecutionErrorKind::NotReady => "RHI publication work is not ready",
    255             RhiPublicationExecutionErrorKind::LeaseLost => {
    256                 "RHI publication lease is no longer authoritative"
    257             }
    258             RhiPublicationExecutionErrorKind::Invariant => "RHI publication state invariant failed",
    259             RhiPublicationExecutionErrorKind::ClockUnavailable => {
    260                 "RHI publication clock is unavailable"
    261             }
    262             RhiPublicationExecutionErrorKind::EntropyUnavailable => {
    263                 "RHI publication retry entropy is unavailable"
    264             }
    265             RhiPublicationExecutionErrorKind::Storage => "RHI publication transaction failed",
    266             RhiPublicationExecutionErrorKind::CommitOutcomeUnknown => {
    267                 "RHI publication commit outcome is unknown"
    268             }
    269         })
    270     }
    271 }
    272 
    273 impl fmt::Debug for RhiPublicationExecutionError {
    274     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    275         formatter
    276             .debug_struct("RhiPublicationExecutionError")
    277             .field("kind", &self.kind)
    278             .finish()
    279     }
    280 }
    281 
    282 impl Error for RhiPublicationExecutionError {}
    283 
    284 /// Stable process-local owner token for one compare-and-swap publication lease.
    285 #[derive(Clone, Copy, PartialEq, Eq, Hash)]
    286 pub struct RhiPublicationLeaseOwner([u8; LEASE_OWNER_BYTES]);
    287 
    288 impl RhiPublicationLeaseOwner {
    289     /// Validates one injected nonzero lease-owner identity.
    290     pub fn from_bytes(
    291         bytes: [u8; LEASE_OWNER_BYTES],
    292     ) -> Result<Self, RhiPublicationExecutionError> {
    293         if bytes.iter().all(|byte| *byte == 0) {
    294             return Err(RhiPublicationExecutionError::new(
    295                 RhiPublicationExecutionErrorKind::InvalidInput,
    296             ));
    297         }
    298         Ok(Self(bytes))
    299     }
    300 }
    301 
    302 impl fmt::Debug for RhiPublicationLeaseOwner {
    303     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    304         formatter.write_str("RhiPublicationLeaseOwner([redacted])")
    305     }
    306 }
    307 
    308 /// Bounded caller-injected retry delay in whole milliseconds.
    309 #[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord)]
    310 pub struct RhiPublicationRetryDelayMilliseconds(u64);
    311 
    312 impl RhiPublicationRetryDelayMilliseconds {
    313     /// Validates a delay against the absolute publication backoff ceiling.
    314     pub fn new(value: u64) -> Result<Self, RhiPublicationExecutionError> {
    315         if value > 3_600_000 {
    316             return Err(RhiPublicationExecutionError::new(
    317                 RhiPublicationExecutionErrorKind::InvalidInput,
    318             ));
    319         }
    320         Ok(Self(value))
    321     }
    322 
    323     /// Returns the exact delay.
    324     #[must_use]
    325     pub const fn get(self) -> u64 {
    326         self.0
    327     }
    328 }
    329 
    330 /// Non-forgeable compare-and-swap authority for one claimed outbox.
    331 #[derive(Clone, Copy, PartialEq, Eq)]
    332 pub struct RhiPublicationLease {
    333     outbox: OutboxRecord,
    334     owner: RhiPublicationLeaseOwner,
    335     expires_at: RhiPublicationUnixMilliseconds,
    336 }
    337 
    338 impl RhiPublicationLease {
    339     /// Returns the exact claimed outbox identity.
    340     #[must_use]
    341     pub const fn outbox_id(self) -> RhiPublicationOutboxId {
    342         self.outbox.id
    343     }
    344 
    345     /// Returns the exact lease expiry.
    346     #[must_use]
    347     pub const fn expires_at(self) -> RhiPublicationUnixMilliseconds {
    348         self.expires_at
    349     }
    350 }
    351 
    352 impl fmt::Debug for RhiPublicationLease {
    353     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    354         formatter
    355             .debug_struct("RhiPublicationLease")
    356             .field("outbox", &"[redacted]")
    357             .field("revision", &self.outbox.revision)
    358             .field("expires_at", &self.expires_at)
    359             .finish()
    360     }
    361 }
    362 
    363 /// Sealed exact-byte remote submission prepared only after durable Submitted.
    364 #[must_use = "prepared publication must be executed exactly or left for unknown recovery"]
    365 pub struct RhiPreparedPublicationAttempt {
    366     lease: RhiPublicationLease,
    367     target: TargetRecord,
    368     attempt_id: RhiPublicationAttemptId,
    369     started_at: RhiPublicationUnixMilliseconds,
    370     deadline_at: RhiPublicationUnixMilliseconds,
    371     publication: RhiCommittedPublication,
    372 }
    373 
    374 impl RhiPreparedPublicationAttempt {
    375     /// Returns the exact attempt identity.
    376     #[must_use]
    377     pub const fn attempt_id(&self) -> RhiPublicationAttemptId {
    378         self.attempt_id
    379     }
    380 
    381     /// Returns the stable configured relay identity without an endpoint or secret.
    382     #[must_use]
    383     pub fn relay_id(&self) -> &str {
    384         &self.target.relay_id
    385     }
    386 
    387     /// Returns the exact committed signed-event bytes with no transformation.
    388     #[must_use]
    389     pub const fn exact_signed_event_bytes(&self) -> &[u8] {
    390         self.publication.exact_signed_event_bytes()
    391     }
    392 
    393     /// Returns the absolute attempt deadline.
    394     #[must_use]
    395     pub const fn deadline_at(&self) -> RhiPublicationUnixMilliseconds {
    396         self.deadline_at
    397     }
    398 
    399     /// Returns the one-based durable attempt number.
    400     #[must_use]
    401     pub const fn attempt_number(&self) -> u16 {
    402         self.target.attempt_count
    403     }
    404 
    405     fn retry_upper_bound(&self) -> u64 {
    406         retry_upper_bound(self.lease.outbox, self.target.attempt_count)
    407     }
    408 }
    409 
    410 impl fmt::Debug for RhiPreparedPublicationAttempt {
    411     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    412         formatter
    413             .debug_struct("RhiPreparedPublicationAttempt")
    414             .field("identity", &"[redacted]")
    415             .field("target_ordinal", &self.target.ordinal)
    416             .field("attempt_number", &self.target.attempt_count)
    417             .field("deadline_at", &self.deadline_at)
    418             .finish()
    419     }
    420 }
    421 
    422 /// Closed exact-byte transport boundary used by the durable executor.
    423 ///
    424 /// The adapter receives the original committed payload and must submit that
    425 /// byte slice unchanged. Dropping the future after durable preparation leaves
    426 /// Submitted evidence; expired-lease recovery records Unknown before retry.
    427 pub trait RhiExactPublicationSink: Send + Sync {
    428     /// Submits one exact prepared payload and returns only a closed observation.
    429     fn submit_exact<'a>(
    430         &'a self,
    431         attempt: &'a RhiPreparedPublicationAttempt,
    432     ) -> BoxFuture<'a, RhiPublicationAttemptOutcome>;
    433 }
    434 
    435 /// Confirmed durable result for one exact publication attempt.
    436 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
    437 pub struct RhiPublicationAttemptCommit {
    438     outbox_id: RhiPublicationOutboxId,
    439     attempt_id: RhiPublicationAttemptId,
    440     target_ordinal: u8,
    441     attempt_number: u16,
    442     outcome: RhiPublicationAttemptOutcome,
    443     target_state: RhiPublicationTargetState,
    444     outbox_state: RhiPublicationOutboxState,
    445 }
    446 
    447 impl RhiPublicationAttemptCommit {
    448     #[must_use]
    449     pub const fn outbox_id(self) -> RhiPublicationOutboxId {
    450         self.outbox_id
    451     }
    452 
    453     #[must_use]
    454     pub const fn attempt_id(self) -> RhiPublicationAttemptId {
    455         self.attempt_id
    456     }
    457 
    458     #[must_use]
    459     pub const fn target_ordinal(self) -> u8 {
    460         self.target_ordinal
    461     }
    462 
    463     #[must_use]
    464     pub const fn attempt_number(self) -> u16 {
    465         self.attempt_number
    466     }
    467 
    468     #[must_use]
    469     pub const fn outcome(self) -> RhiPublicationAttemptOutcome {
    470         self.outcome
    471     }
    472 
    473     #[must_use]
    474     pub const fn target_state(self) -> RhiPublicationTargetState {
    475         self.target_state
    476     }
    477 
    478     #[must_use]
    479     pub const fn outbox_state(self) -> RhiPublicationOutboxState {
    480         self.outbox_state
    481     }
    482 }
    483 
    484 impl RhiPublicationOutboxRepository<'_> {
    485     /// Claims the oldest due outbox using one injected owner and wall time.
    486     pub async fn claim_next_publication(
    487         &self,
    488         owner: RhiPublicationLeaseOwner,
    489         now: RhiPublicationUnixMilliseconds,
    490         authority: &RhiPublicationAuthority,
    491     ) -> Result<Option<RhiPublicationLease>, RhiPublicationExecutionError> {
    492         require_writable(self)?;
    493         let Some(authority) = AuthorityBinding::from_authority(authority) else {
    494             return Ok(None);
    495         };
    496         self.host()
    497             .sqlite_host()
    498             .transaction(move |transaction| {
    499                 Box::pin(async move { claim_next(transaction, owner, now, authority).await })
    500             })
    501             .await
    502             .map_err(map_transaction_error)
    503     }
    504 
    505     /// Persists Submitted before exposing exact bytes to remote I/O.
    506     pub async fn prepare_next_publication_target(
    507         &self,
    508         lease: RhiPublicationLease,
    509         started_at: RhiPublicationUnixMilliseconds,
    510     ) -> Result<RhiPreparedPublicationAttempt, RhiPublicationExecutionError> {
    511         require_writable(self)?;
    512         self.host()
    513             .sqlite_host()
    514             .transaction(move |transaction| {
    515                 Box::pin(async move { prepare_next(transaction, lease, started_at).await })
    516             })
    517             .await
    518             .map_err(map_transaction_error)
    519     }
    520 
    521     /// Commits one closed observed outcome by exact compare-and-swap.
    522     pub async fn record_publication_outcome(
    523         &self,
    524         prepared: &RhiPreparedPublicationAttempt,
    525         finished_at: RhiPublicationUnixMilliseconds,
    526         outcome: RhiPublicationAttemptOutcome,
    527         retry_delay: RhiPublicationRetryDelayMilliseconds,
    528     ) -> Result<RhiPublicationAttemptCommit, RhiPublicationExecutionError> {
    529         require_writable(self)?;
    530         let input = RecordInput::from_prepared(prepared, finished_at, outcome, retry_delay)?;
    531         self.host()
    532             .sqlite_host()
    533             .transaction(move |transaction| {
    534                 Box::pin(async move { record_outcome(transaction, input).await })
    535             })
    536             .await
    537             .map_err(map_transaction_error)
    538     }
    539 
    540     /// Recovers at most one expired lease and persists Unknown for Submitted work.
    541     pub async fn recover_one_expired_publication(
    542         &self,
    543         adapters: &RhiTimeEntropyAdapters,
    544         now: RhiPublicationUnixMilliseconds,
    545         authority: &RhiPublicationAuthority,
    546     ) -> Result<bool, RhiPublicationExecutionError> {
    547         require_writable(self)?;
    548         let Some(authority) = AuthorityBinding::from_authority(authority) else {
    549             return Ok(false);
    550         };
    551         let candidate = self
    552             .host()
    553             .sqlite_host()
    554             .transaction(move |transaction| {
    555                 Box::pin(async move { read_expired_candidate(transaction, now, authority).await })
    556             })
    557             .await
    558             .map_err(map_transaction_error)?;
    559         let Some(candidate) = candidate else {
    560             return Ok(false);
    561         };
    562         let cap = candidate.retry_upper_bound;
    563         let delay = sample_retry_delay(adapters, cap)?;
    564         self.host()
    565             .sqlite_host()
    566             .transaction(move |transaction| {
    567                 Box::pin(async move { recover_expired(transaction, candidate, now, delay).await })
    568             })
    569             .await
    570             .map_err(map_transaction_error)
    571     }
    572 
    573     /// Executes at most one exact-byte attempt through the injected sink.
    574     ///
    575     /// Cancellation before a claim has no effect. Cancellation after durable
    576     /// preparation leaves Submitted state; lease-expiry recovery records
    577     /// Unknown and schedules the exact same committed bytes without rebuilding.
    578     pub async fn execute_next_publication(
    579         &self,
    580         owner: RhiPublicationLeaseOwner,
    581         adapters: &RhiTimeEntropyAdapters,
    582         sink: &dyn RhiExactPublicationSink,
    583         authority: &RhiPublicationAuthority,
    584     ) -> Result<Option<RhiPublicationAttemptCommit>, RhiPublicationExecutionError> {
    585         let now = publication_now(adapters)?;
    586         self.recover_one_expired_publication(adapters, now, authority)
    587             .await?;
    588         let Some(lease) = self.claim_next_publication(owner, now, authority).await? else {
    589             return Ok(None);
    590         };
    591         let prepared = self.prepare_next_publication_target(lease, now).await?;
    592         let outcome = match sink.submit_exact(&prepared).await {
    593             RhiPublicationAttemptOutcome::Submitted => RhiPublicationAttemptOutcome::Unknown,
    594             outcome => outcome,
    595         };
    596         let finished_at = publication_now(adapters)?;
    597         let retry_delay = if retryable_outcome(outcome) {
    598             sample_retry_delay(adapters, prepared.retry_upper_bound())?
    599         } else {
    600             RhiPublicationRetryDelayMilliseconds(0)
    601         };
    602         self.record_publication_outcome(&prepared, finished_at, outcome, retry_delay)
    603             .await
    604             .map(Some)
    605     }
    606 }
    607 
    608 fn require_writable(
    609     repository: &RhiPublicationOutboxRepository<'_>,
    610 ) -> Result<(), RhiPublicationExecutionError> {
    611     if repository.host().mode() == RhiStateHostMode::ReadWriteExisting {
    612         Ok(())
    613     } else {
    614         Err(RhiPublicationExecutionError::new(
    615             RhiPublicationExecutionErrorKind::InvalidMode,
    616         ))
    617     }
    618 }
    619 
    620 fn publication_now(
    621     adapters: &RhiTimeEntropyAdapters,
    622 ) -> Result<RhiPublicationUnixMilliseconds, RhiPublicationExecutionError> {
    623     let value = adapters.now_utc_milliseconds().map_err(|error| {
    624         RhiPublicationExecutionError::new(match error.kind() {
    625             RhiRuntimeAdapterErrorKind::WallClockUnavailable => {
    626                 RhiPublicationExecutionErrorKind::ClockUnavailable
    627             }
    628             _ => RhiPublicationExecutionErrorKind::ClockUnavailable,
    629         })
    630     })?;
    631     RhiPublicationUnixMilliseconds::new(value).map_err(|_| {
    632         RhiPublicationExecutionError::new(RhiPublicationExecutionErrorKind::ClockUnavailable)
    633     })
    634 }
    635 
    636 fn sample_retry_delay(
    637     adapters: &RhiTimeEntropyAdapters,
    638     cap: u64,
    639 ) -> Result<RhiPublicationRetryDelayMilliseconds, RhiPublicationExecutionError> {
    640     if cap == 0 {
    641         return Ok(RhiPublicationRetryDelayMilliseconds(0));
    642     }
    643     let bound = RhiJitterBoundMilliseconds::new(cap).map_err(|_| {
    644         RhiPublicationExecutionError::new(RhiPublicationExecutionErrorKind::InvalidInput)
    645     })?;
    646     adapters
    647         .sample_full_jitter(bound)
    648         .map(|delay| RhiPublicationRetryDelayMilliseconds(delay.get()))
    649         .map_err(|_| {
    650             RhiPublicationExecutionError::new(RhiPublicationExecutionErrorKind::EntropyUnavailable)
    651         })
    652 }
    653 
    654 async fn claim_next(
    655     transaction: &mut ServiceSqliteTransaction<'_>,
    656     owner: RhiPublicationLeaseOwner,
    657     now: RhiPublicationUnixMilliseconds,
    658     authority: AuthorityBinding,
    659 ) -> Result<Option<RhiPublicationLease>, OperationError> {
    660     let rows = sqlx::query(READ_CLAIMABLE_OUTBOX_SQL)
    661         .bind(i64_value(now.get())?)
    662         .fetch_all(&mut *transaction)
    663         .await
    664         .map_err(|_| OperationError::Storage)?;
    665     let Some(row) = exactly_zero_or_one(rows)? else {
    666         return Ok(None);
    667     };
    668     let outbox = decode_outbox(row)?;
    669     validate_authority(outbox, authority)?;
    670     validate_target_inventory(transaction, outbox).await?;
    671     let expires = now
    672         .get()
    673         .checked_add(outbox.attempt_deadline_ms)
    674         .filter(|value| *value <= MAX_UNIX_MILLISECONDS)
    675         .ok_or(OperationError::InvalidInput)?;
    676     let result = sqlx::query(CLAIM_OUTBOX_SQL)
    677         .bind(owner.0.as_slice())
    678         .bind(i64_value(expires)?)
    679         .bind(i64_value(now.get())?)
    680         .bind(outbox.id.as_bytes().as_slice())
    681         .bind(i64_value(outbox.revision)?)
    682         .bind(i64_value(now.get())?)
    683         .bind(i64_value(now.get())?)
    684         .execute(&mut *transaction)
    685         .await
    686         .map_err(|_| OperationError::Storage)?;
    687     if result.rows_affected() != 1 {
    688         return Err(OperationError::LeaseLost);
    689     }
    690     let leased = read_outbox(transaction, outbox.id)
    691         .await?
    692         .filter(|record| {
    693             record.state == RhiPublicationOutboxState::Leased
    694                 && record.lease_owner == Some(owner)
    695                 && record.lease_expires == Some(RhiPublicationUnixMilliseconds(expires))
    696         })
    697         .ok_or(OperationError::Invariant)?;
    698     Ok(Some(RhiPublicationLease {
    699         outbox: leased,
    700         owner,
    701         expires_at: RhiPublicationUnixMilliseconds(expires),
    702     }))
    703 }
    704 
    705 async fn prepare_next(
    706     transaction: &mut ServiceSqliteTransaction<'_>,
    707     lease: RhiPublicationLease,
    708     started_at: RhiPublicationUnixMilliseconds,
    709 ) -> Result<RhiPreparedPublicationAttempt, OperationError> {
    710     validate_lease(transaction, lease, started_at, true).await?;
    711     let targets = read_targets(transaction, lease.outbox.id).await?;
    712     validate_targets(lease.outbox, &targets)?;
    713     let candidate = targets
    714         .into_iter()
    715         .find(|target| target_is_due(target, lease.outbox.max_attempts, started_at))
    716         .ok_or(OperationError::NotReady)?;
    717     let attempt_number = candidate
    718         .attempt_count
    719         .checked_add(1)
    720         .filter(|value| *value <= lease.outbox.max_attempts)
    721         .ok_or(OperationError::Invariant)?;
    722     let attempt_id = derive_attempt_id(
    723         lease.outbox.id,
    724         lease.outbox.event_sha256,
    725         candidate.ordinal,
    726         attempt_number,
    727     );
    728     let result = sqlx::query(PREPARE_TARGET_SQL)
    729         .bind(attempt_id.as_bytes().as_slice())
    730         .bind(i64_value(started_at.get())?)
    731         .bind(lease.outbox.id.as_bytes().as_slice())
    732         .bind(i64::from(candidate.ordinal))
    733         .bind(i64_value(candidate.revision)?)
    734         .bind(i64::from(lease.outbox.max_attempts))
    735         .bind(i64_value(started_at.get())?)
    736         .bind(i64_value(started_at.get())?)
    737         .execute(&mut *transaction)
    738         .await
    739         .map_err(|_| OperationError::Storage)?;
    740     if result.rows_affected() != 1 {
    741         return Err(OperationError::LeaseLost);
    742     }
    743     let target = read_targets(transaction, lease.outbox.id)
    744         .await?
    745         .into_iter()
    746         .find(|target| target.ordinal == candidate.ordinal)
    747         .filter(|target| {
    748             target.state == RhiPublicationTargetState::Submitted
    749                 && target.attempt_count == attempt_number
    750                 && target.last_attempt_id == Some(attempt_id)
    751         })
    752         .ok_or(OperationError::Invariant)?;
    753     let publication = read_committed(transaction, lease.outbox.id)
    754         .await
    755         .map_err(|error| match error {
    756             ReadError::Storage => OperationError::Storage,
    757             ReadError::NotFound | ReadError::Binding => OperationError::Invariant,
    758         })?;
    759     if publication.event_sha256() != &lease.outbox.event_sha256 {
    760         return Err(OperationError::Invariant);
    761     }
    762     Ok(RhiPreparedPublicationAttempt {
    763         lease,
    764         target,
    765         attempt_id,
    766         started_at,
    767         deadline_at: lease.expires_at,
    768         publication,
    769     })
    770 }
    771 
    772 async fn record_outcome(
    773     transaction: &mut ServiceSqliteTransaction<'_>,
    774     input: RecordInput,
    775 ) -> Result<RhiPublicationAttemptCommit, OperationError> {
    776     if let Some(existing) = read_attempt(transaction, input.attempt_id).await? {
    777         return reconcile_recorded(transaction, input, existing).await;
    778     }
    779     validate_lease(transaction, input.lease, input.finished_at, true).await?;
    780     let target = read_targets(transaction, input.lease.outbox.id)
    781         .await?
    782         .into_iter()
    783         .find(|target| target.ordinal == input.target_ordinal)
    784         .ok_or(OperationError::Invariant)?;
    785     if target.state != RhiPublicationTargetState::Submitted
    786         || target.revision != input.target_revision
    787         || target.attempt_count != input.attempt_number
    788         || target.last_attempt_id != Some(input.attempt_id)
    789         || target.updated_at != input.started_at
    790     {
    791         return Err(OperationError::LeaseLost);
    792     }
    793     insert_attempt(transaction, input).await?;
    794     let (next_state, next_attempt) = target_outcome_schedule(input)?;
    795     let result = sqlx::query(UPDATE_TARGET_OUTCOME_SQL)
    796         .bind(next_state.code())
    797         .bind(optional_i64(next_attempt)?)
    798         .bind(i64_value(input.finished_at.get())?)
    799         .bind(input.lease.outbox.id.as_bytes().as_slice())
    800         .bind(i64::from(input.target_ordinal))
    801         .bind(i64_value(input.target_revision)?)
    802         .bind(i64::from(input.attempt_number))
    803         .bind(input.attempt_id.as_bytes().as_slice())
    804         .execute(&mut *transaction)
    805         .await
    806         .map_err(|_| OperationError::Storage)?;
    807     if result.rows_affected() != 1 {
    808         return Err(OperationError::LeaseLost);
    809     }
    810     let targets = read_targets(transaction, input.lease.outbox.id).await?;
    811     let disposition = disposition(input.lease.outbox, &targets)?;
    812     update_outbox_after_attempt(transaction, input, disposition).await?;
    813     Ok(RhiPublicationAttemptCommit {
    814         outbox_id: input.lease.outbox.id,
    815         attempt_id: input.attempt_id,
    816         target_ordinal: input.target_ordinal,
    817         attempt_number: input.attempt_number,
    818         outcome: input.outcome,
    819         target_state: next_state,
    820         outbox_state: disposition.state,
    821     })
    822 }
    823 
    824 async fn read_expired_candidate(
    825     transaction: &mut ServiceSqliteTransaction<'_>,
    826     now: RhiPublicationUnixMilliseconds,
    827     authority: AuthorityBinding,
    828 ) -> Result<Option<RecoveryCandidate>, OperationError> {
    829     let rows = sqlx::query(READ_EXPIRED_OUTBOX_SQL)
    830         .bind(i64_value(now.get())?)
    831         .fetch_all(&mut *transaction)
    832         .await
    833         .map_err(|_| OperationError::Storage)?;
    834     let Some(row) = exactly_zero_or_one(rows)? else {
    835         return Ok(None);
    836     };
    837     let outbox = decode_outbox(row)?;
    838     validate_authority(outbox, authority)?;
    839     let targets = read_targets(transaction, outbox.id).await?;
    840     validate_targets(outbox, &targets)?;
    841     let submitted_attempts: Vec<_> = targets
    842         .iter()
    843         .filter(|target| target.state == RhiPublicationTargetState::Submitted)
    844         .map(|target| target.attempt_count)
    845         .collect();
    846     if submitted_attempts.len() > 1 {
    847         return Err(OperationError::Invariant);
    848     }
    849     let retry_upper_bound = submitted_attempts
    850         .first()
    851         .copied()
    852         .filter(|attempt| *attempt < outbox.max_attempts)
    853         .map_or(0, |attempt| retry_upper_bound(outbox, attempt));
    854     Ok(Some(RecoveryCandidate {
    855         outbox,
    856         retry_upper_bound,
    857     }))
    858 }
    859 
    860 async fn recover_expired(
    861     transaction: &mut ServiceSqliteTransaction<'_>,
    862     candidate: RecoveryCandidate,
    863     now: RhiPublicationUnixMilliseconds,
    864     delay: RhiPublicationRetryDelayMilliseconds,
    865 ) -> Result<bool, OperationError> {
    866     let current = read_outbox(transaction, candidate.outbox.id)
    867         .await?
    868         .filter(|current| *current == candidate.outbox)
    869         .ok_or(OperationError::LeaseLost)?;
    870     if current.state != RhiPublicationOutboxState::Leased
    871         || current.lease_expires.is_none_or(|expires| expires > now)
    872     {
    873         return Err(OperationError::LeaseLost);
    874     }
    875     let owner = current.lease_owner.ok_or(OperationError::Invariant)?;
    876     let mut targets = read_targets(transaction, current.id).await?;
    877     validate_targets(current, &targets)?;
    878     for target in targets
    879         .iter_mut()
    880         .filter(|target| target.state == RhiPublicationTargetState::Submitted)
    881     {
    882         let attempt_id = target.last_attempt_id.ok_or(OperationError::Invariant)?;
    883         let publication =
    884             read_committed(transaction, current.id)
    885                 .await
    886                 .map_err(|error| match error {
    887                     ReadError::Storage => OperationError::Storage,
    888                     ReadError::NotFound | ReadError::Binding => OperationError::Invariant,
    889                 })?;
    890         let evidence = RhiPublicationAttemptEvidence::new(
    891             &publication,
    892             u32::from(target.ordinal),
    893             target.attempt_count,
    894             target.updated_at,
    895             now,
    896             RhiPublicationAttemptOutcome::Unknown,
    897         )
    898         .map_err(|_| OperationError::Invariant)?;
    899         if attempt_id != evidence.id() || read_attempt(transaction, attempt_id).await?.is_some() {
    900             return Err(OperationError::Invariant);
    901         }
    902         let input = RecordInput {
    903             lease: RhiPublicationLease {
    904                 outbox: current,
    905                 owner,
    906                 expires_at: current.lease_expires.ok_or(OperationError::Invariant)?,
    907             },
    908             target_ordinal: target.ordinal,
    909             target_revision: target.revision,
    910             attempt_number: target.attempt_count,
    911             attempt_id,
    912             started_at: target.updated_at,
    913             finished_at: now,
    914             outcome: RhiPublicationAttemptOutcome::Unknown,
    915             retry_delay: delay,
    916         };
    917         insert_attempt(transaction, input).await?;
    918         let (state, next) = target_outcome_schedule(input)?;
    919         let result = sqlx::query(UPDATE_TARGET_OUTCOME_SQL)
    920             .bind(state.code())
    921             .bind(optional_i64(next)?)
    922             .bind(i64_value(now.get())?)
    923             .bind(current.id.as_bytes().as_slice())
    924             .bind(i64::from(target.ordinal))
    925             .bind(i64_value(target.revision)?)
    926             .bind(i64::from(target.attempt_count))
    927             .bind(attempt_id.as_bytes().as_slice())
    928             .execute(&mut *transaction)
    929             .await
    930             .map_err(|_| OperationError::Storage)?;
    931         if result.rows_affected() != 1 {
    932             return Err(OperationError::LeaseLost);
    933         }
    934     }
    935     targets = read_targets(transaction, current.id).await?;
    936     let disposition = disposition(current, &targets)?;
    937     let result = sqlx::query(UPDATE_OUTBOX_RECOVERY_SQL)
    938         .bind(disposition.state.code())
    939         .bind(optional_i64(disposition.next_attempt)?)
    940         .bind(i64_value(now.get())?)
    941         .bind(current.id.as_bytes().as_slice())
    942         .bind(i64_value(current.revision)?)
    943         .bind(owner.0.as_slice())
    944         .bind(i64_value(
    945             current
    946                 .lease_expires
    947                 .ok_or(OperationError::Invariant)?
    948                 .get(),
    949         )?)
    950         .bind(i64_value(now.get())?)
    951         .bind(i64_value(now.get())?)
    952         .execute(&mut *transaction)
    953         .await
    954         .map_err(|_| OperationError::Storage)?;
    955     if result.rows_affected() != 1 {
    956         return Err(OperationError::LeaseLost);
    957     }
    958     Ok(true)
    959 }
    960 
    961 async fn validate_lease(
    962     transaction: &mut ServiceSqliteTransaction<'_>,
    963     lease: RhiPublicationLease,
    964     now: RhiPublicationUnixMilliseconds,
    965     require_unexpired: bool,
    966 ) -> Result<(), OperationError> {
    967     let current = read_outbox(transaction, lease.outbox.id)
    968         .await?
    969         .filter(|current| *current == lease.outbox)
    970         .ok_or(OperationError::LeaseLost)?;
    971     if current.state != RhiPublicationOutboxState::Leased
    972         || current.lease_owner != Some(lease.owner)
    973         || current.lease_expires != Some(lease.expires_at)
    974         || (require_unexpired && lease.expires_at <= now)
    975     {
    976         return Err(OperationError::LeaseLost);
    977     }
    978     Ok(())
    979 }
    980 
    981 async fn insert_attempt(
    982     transaction: &mut ServiceSqliteTransaction<'_>,
    983     input: RecordInput,
    984 ) -> Result<(), OperationError> {
    985     let result = sqlx::query(INSERT_ATTEMPT_SQL)
    986         .bind(input.attempt_id.as_bytes().as_slice())
    987         .bind(input.lease.outbox.id.as_bytes().as_slice())
    988         .bind(i64::from(input.target_ordinal))
    989         .bind(i64::from(input.attempt_number))
    990         .bind(input.lease.outbox.event_sha256.as_slice())
    991         .bind(input.lease.owner.0.as_slice())
    992         .bind(i64_value(input.started_at.get())?)
    993         .bind(i64_value(input.finished_at.get())?)
    994         .bind(input.outcome.code())
    995         .bind(input.outcome.code())
    996         .execute(&mut *transaction)
    997         .await
    998         .map_err(|_| OperationError::Storage)?;
    999     if result.rows_affected() != 1 {
   1000         return Err(OperationError::Storage);
   1001     }
   1002     Ok(())
   1003 }
   1004 
   1005 async fn reconcile_recorded(
   1006     transaction: &mut ServiceSqliteTransaction<'_>,
   1007     input: RecordInput,
   1008     existing: AttemptRecord,
   1009 ) -> Result<RhiPublicationAttemptCommit, OperationError> {
   1010     if existing != AttemptRecord::from_input(input) {
   1011         return Err(OperationError::Invariant);
   1012     }
   1013     let outbox = read_outbox(transaction, input.lease.outbox.id)
   1014         .await?
   1015         .ok_or(OperationError::Invariant)?;
   1016     if !same_outbox_identity(outbox, input.lease.outbox)
   1017         || outbox.revision
   1018             != input
   1019                 .lease
   1020                 .outbox
   1021                 .revision
   1022                 .checked_add(1)
   1023                 .ok_or(OperationError::Invariant)?
   1024         || outbox.updated_at != input.finished_at
   1025         || outbox.lease_owner.is_some()
   1026         || outbox.lease_expires.is_some()
   1027     {
   1028         return Err(OperationError::Invariant);
   1029     }
   1030     let targets = read_targets(transaction, outbox.id).await?;
   1031     validate_targets(outbox, &targets)?;
   1032     let (expected_target_state, expected_target_schedule) = target_outcome_schedule(input)?;
   1033     let expected_target_revision = input
   1034         .target_revision
   1035         .checked_add(1)
   1036         .ok_or(OperationError::Invariant)?;
   1037     let target = targets
   1038         .iter()
   1039         .find(|target| target.ordinal == input.target_ordinal)
   1040         .filter(|target| {
   1041             target.revision == expected_target_revision
   1042                 && target.attempt_count == input.attempt_number
   1043                 && target.last_attempt_id == Some(input.attempt_id)
   1044                 && target.state == expected_target_state
   1045                 && target.next_attempt == expected_target_schedule
   1046                 && target.updated_at == input.finished_at
   1047         })
   1048         .ok_or(OperationError::Invariant)?;
   1049     let expected_disposition = disposition(outbox, &targets)?;
   1050     if outbox.state != expected_disposition.state
   1051         || outbox.next_attempt != expected_disposition.next_attempt
   1052     {
   1053         return Err(OperationError::Invariant);
   1054     }
   1055     Ok(RhiPublicationAttemptCommit {
   1056         outbox_id: outbox.id,
   1057         attempt_id: input.attempt_id,
   1058         target_ordinal: input.target_ordinal,
   1059         attempt_number: input.attempt_number,
   1060         outcome: input.outcome,
   1061         target_state: target.state,
   1062         outbox_state: outbox.state,
   1063     })
   1064 }
   1065 
   1066 fn same_outbox_identity(current: OutboxRecord, prior: OutboxRecord) -> bool {
   1067     current.id == prior.id
   1068         && current.event_sha256 == prior.event_sha256
   1069         && current.authority_sha256 == prior.authority_sha256
   1070         && current.target_set_sha256 == prior.target_set_sha256
   1071         && current.target_count == prior.target_count
   1072         && current.required_target_count == prior.required_target_count
   1073         && current.max_attempts == prior.max_attempts
   1074         && current.initial_backoff_ms == prior.initial_backoff_ms
   1075         && current.maximum_backoff_ms == prior.maximum_backoff_ms
   1076         && current.attempt_deadline_ms == prior.attempt_deadline_ms
   1077         && current.created_at == prior.created_at
   1078 }
   1079 
   1080 async fn update_outbox_after_attempt(
   1081     transaction: &mut ServiceSqliteTransaction<'_>,
   1082     input: RecordInput,
   1083     disposition: OutboxDisposition,
   1084 ) -> Result<(), OperationError> {
   1085     let result = sqlx::query(UPDATE_OUTBOX_AFTER_ATTEMPT_SQL)
   1086         .bind(disposition.state.code())
   1087         .bind(optional_i64(disposition.next_attempt)?)
   1088         .bind(i64_value(input.finished_at.get())?)
   1089         .bind(input.lease.outbox.id.as_bytes().as_slice())
   1090         .bind(i64_value(input.lease.outbox.revision)?)
   1091         .bind(input.lease.owner.0.as_slice())
   1092         .bind(i64_value(input.lease.expires_at.get())?)
   1093         .bind(i64_value(input.finished_at.get())?)
   1094         .bind(i64_value(input.finished_at.get())?)
   1095         .execute(&mut *transaction)
   1096         .await
   1097         .map_err(|_| OperationError::Storage)?;
   1098     if result.rows_affected() != 1 {
   1099         return Err(OperationError::LeaseLost);
   1100     }
   1101     Ok(())
   1102 }
   1103 
   1104 async fn read_outbox(
   1105     transaction: &mut ServiceSqliteTransaction<'_>,
   1106     id: RhiPublicationOutboxId,
   1107 ) -> Result<Option<OutboxRecord>, OperationError> {
   1108     let rows = sqlx::query(READ_OUTBOX_SQL)
   1109         .bind(id.as_bytes().as_slice())
   1110         .fetch_all(&mut *transaction)
   1111         .await
   1112         .map_err(|_| OperationError::Storage)?;
   1113     exactly_zero_or_one(rows)?.map(decode_outbox).transpose()
   1114 }
   1115 
   1116 async fn read_targets(
   1117     transaction: &mut ServiceSqliteTransaction<'_>,
   1118     id: RhiPublicationOutboxId,
   1119 ) -> Result<Vec<TargetRecord>, OperationError> {
   1120     sqlx::query(READ_TARGETS_SQL)
   1121         .bind(id.as_bytes().as_slice())
   1122         .fetch_all(&mut *transaction)
   1123         .await
   1124         .map_err(|_| OperationError::Storage)?
   1125         .into_iter()
   1126         .map(decode_target)
   1127         .collect()
   1128 }
   1129 
   1130 async fn validate_target_inventory(
   1131     transaction: &mut ServiceSqliteTransaction<'_>,
   1132     outbox: OutboxRecord,
   1133 ) -> Result<(), OperationError> {
   1134     let targets = read_targets(transaction, outbox.id).await?;
   1135     validate_targets(outbox, &targets)
   1136 }
   1137 
   1138 fn validate_targets(outbox: OutboxRecord, targets: &[TargetRecord]) -> Result<(), OperationError> {
   1139     if targets.len() != usize::from(outbox.target_count)
   1140         || targets.len()
   1141             > usize::try_from(RHI_PUBLICATION_TARGET_ORDINAL_MAXIMUM + 1)
   1142                 .map_err(|_| OperationError::Invariant)?
   1143         || targets
   1144             .iter()
   1145             .enumerate()
   1146             .any(|(index, target)| usize::from(target.ordinal) != index)
   1147         || targets.iter().filter(|target| target.required).count()
   1148             != usize::from(outbox.required_target_count)
   1149     {
   1150         return Err(OperationError::Invariant);
   1151     }
   1152     for (index, target) in targets.iter().enumerate() {
   1153         if targets[index + 1..]
   1154             .iter()
   1155             .any(|other| other.relay_id == target.relay_id)
   1156         {
   1157             return Err(OperationError::Invariant);
   1158         }
   1159     }
   1160     let mut digest = Sha256::new();
   1161     digest.update(TARGET_SET_DOMAIN);
   1162     digest.update(u32::from(outbox.target_count).to_be_bytes());
   1163     for target in targets {
   1164         digest.update(u32::from(target.ordinal).to_be_bytes());
   1165         digest.update(
   1166             u64::try_from(target.relay_id.len())
   1167                 .map_err(|_| OperationError::Invariant)?
   1168                 .to_be_bytes(),
   1169         );
   1170         digest.update(target.relay_id.as_bytes());
   1171         digest.update([u8::from(target.required)]);
   1172     }
   1173     let actual: [u8; 32] = digest.finalize().into();
   1174     if actual != outbox.target_set_sha256 {
   1175         return Err(OperationError::Invariant);
   1176     }
   1177     Ok(())
   1178 }
   1179 
   1180 fn validate_authority(
   1181     outbox: OutboxRecord,
   1182     authority: AuthorityBinding,
   1183 ) -> Result<(), OperationError> {
   1184     if outbox.authority_sha256 != authority.authority_sha256
   1185         || outbox.target_set_sha256 != authority.target_set_sha256
   1186         || outbox.target_count != authority.target_count
   1187         || outbox.required_target_count != authority.required_target_count
   1188         || outbox.max_attempts != authority.max_attempts
   1189         || outbox.initial_backoff_ms != authority.initial_backoff_ms
   1190         || outbox.maximum_backoff_ms != authority.maximum_backoff_ms
   1191         || outbox.attempt_deadline_ms != authority.attempt_deadline_ms
   1192     {
   1193         return Err(OperationError::Invariant);
   1194     }
   1195     Ok(())
   1196 }
   1197 
   1198 fn disposition(
   1199     outbox: OutboxRecord,
   1200     targets: &[TargetRecord],
   1201 ) -> Result<OutboxDisposition, OperationError> {
   1202     validate_targets(outbox, targets)?;
   1203     let required = targets.iter().filter(|target| target.required);
   1204     if required
   1205         .clone()
   1206         .all(|target| target.state == RhiPublicationTargetState::Accepted)
   1207     {
   1208         return Ok(OutboxDisposition {
   1209             state: RhiPublicationOutboxState::Complete,
   1210             next_attempt: None,
   1211         });
   1212     }
   1213     if required
   1214         .clone()
   1215         .any(|target| target_is_blocking(target, outbox.max_attempts))
   1216     {
   1217         return Ok(OutboxDisposition {
   1218             state: RhiPublicationOutboxState::Blocked,
   1219             next_attempt: None,
   1220         });
   1221     }
   1222     let next_attempt = required
   1223         .filter_map(|target| target.next_attempt)
   1224         .min()
   1225         .ok_or(OperationError::Invariant)?;
   1226     Ok(OutboxDisposition {
   1227         state: RhiPublicationOutboxState::Pending,
   1228         next_attempt: Some(next_attempt),
   1229     })
   1230 }
   1231 
   1232 fn target_is_blocking(target: &TargetRecord, max_attempts: u16) -> bool {
   1233     matches!(
   1234         target.state,
   1235         RhiPublicationTargetState::Rejected | RhiPublicationTargetState::AuthRequired
   1236     ) || (retryable_state(target.state)
   1237         && target.attempt_count >= max_attempts
   1238         && target.next_attempt.is_none())
   1239 }
   1240 
   1241 fn target_is_due(
   1242     target: &TargetRecord,
   1243     max_attempts: u16,
   1244     now: RhiPublicationUnixMilliseconds,
   1245 ) -> bool {
   1246     retryable_state(target.state)
   1247         && target.attempt_count < max_attempts
   1248         && target.next_attempt.is_some_and(|next| next <= now)
   1249 }
   1250 
   1251 const fn retryable_state(state: RhiPublicationTargetState) -> bool {
   1252     matches!(
   1253         state,
   1254         RhiPublicationTargetState::Pending
   1255             | RhiPublicationTargetState::Failed
   1256             | RhiPublicationTargetState::RateLimited
   1257             | RhiPublicationTargetState::Unknown
   1258     )
   1259 }
   1260 
   1261 const fn retryable_outcome(outcome: RhiPublicationAttemptOutcome) -> bool {
   1262     matches!(
   1263         outcome,
   1264         RhiPublicationAttemptOutcome::RateLimited
   1265             | RhiPublicationAttemptOutcome::Failed
   1266             | RhiPublicationAttemptOutcome::Unknown
   1267     )
   1268 }
   1269 
   1270 fn target_outcome_schedule(
   1271     input: RecordInput,
   1272 ) -> Result<
   1273     (
   1274         RhiPublicationTargetState,
   1275         Option<RhiPublicationUnixMilliseconds>,
   1276     ),
   1277     OperationError,
   1278 > {
   1279     let state = input.outcome.target_state();
   1280     if input.outcome == RhiPublicationAttemptOutcome::Submitted {
   1281         return Err(OperationError::InvalidInput);
   1282     }
   1283     let next = if retryable_outcome(input.outcome)
   1284         && input.attempt_number < input.lease.outbox.max_attempts
   1285     {
   1286         if input.retry_delay.get() > retry_upper_bound(input.lease.outbox, input.attempt_number) {
   1287             return Err(OperationError::InvalidInput);
   1288         }
   1289         Some(RhiPublicationUnixMilliseconds(
   1290             input
   1291                 .finished_at
   1292                 .get()
   1293                 .checked_add(input.retry_delay.get())
   1294                 .filter(|value| *value <= MAX_UNIX_MILLISECONDS)
   1295                 .ok_or(OperationError::InvalidInput)?,
   1296         ))
   1297     } else {
   1298         if input.retry_delay.get() != 0 {
   1299             return Err(OperationError::InvalidInput);
   1300         }
   1301         None
   1302     };
   1303     Ok((state, next))
   1304 }
   1305 
   1306 fn retry_upper_bound(outbox: OutboxRecord, attempt_number: u16) -> u64 {
   1307     let mut bound = outbox.initial_backoff_ms;
   1308     for _ in 1..attempt_number {
   1309         bound = bound.saturating_mul(2).min(outbox.maximum_backoff_ms);
   1310     }
   1311     bound.min(outbox.maximum_backoff_ms)
   1312 }
   1313 
   1314 fn decode_outbox(row: sqlx::sqlite::SqliteRow) -> Result<OutboxRecord, OperationError> {
   1315     let id = RhiPublicationOutboxId::from_committed_bytes(blob::<32>(&row, "outbox_id")?);
   1316     let state = decode_outbox_state(bounded_text(&row, "state", "state_bytes", 8)?.as_str())?;
   1317     let target_count = bounded_u16(&row, "target_count", 1, 32)? as u8;
   1318     let required_target_count =
   1319         bounded_u16(&row, "required_target_count", 0, target_count.into())? as u8;
   1320     let max_attempts = bounded_u16(
   1321         &row,
   1322         "max_attempts",
   1323         1,
   1324         RHI_PUBLICATION_ATTEMPT_NUMBER_MAXIMUM,
   1325     )?;
   1326     let initial_backoff_ms = bounded_u64(&row, "initial_backoff_ms", 1, 60_000)?;
   1327     let maximum_backoff_ms = bounded_u64(&row, "maximum_backoff_ms", 1, 3_600_000)?;
   1328     if initial_backoff_ms > maximum_backoff_ms {
   1329         return Err(OperationError::Invariant);
   1330     }
   1331     let attempt_deadline_ms = bounded_u64(&row, "attempt_deadline_ms", 100, 30_000)?;
   1332     let lease_owner = optional_owner(&row)?;
   1333     let lease_expires = optional_millis(&row, "lease_expires_unix_ms")?;
   1334     let next_attempt = optional_millis(&row, "next_attempt_unix_ms")?;
   1335     if !valid_outbox_shape(state, next_attempt, lease_owner, lease_expires) {
   1336         return Err(OperationError::Invariant);
   1337     }
   1338     Ok(OutboxRecord {
   1339         id,
   1340         event_sha256: blob::<32>(&row, "event_sha256")?,
   1341         authority_sha256: blob::<32>(&row, "publication_authority_sha256")?,
   1342         target_set_sha256: blob::<32>(&row, "target_set_sha256")?,
   1343         target_count,
   1344         required_target_count,
   1345         max_attempts,
   1346         initial_backoff_ms,
   1347         maximum_backoff_ms,
   1348         attempt_deadline_ms,
   1349         state,
   1350         revision: positive_u64(&row, "revision")?,
   1351         next_attempt,
   1352         lease_owner,
   1353         lease_expires,
   1354         created_at: nonnegative_millis(&row, "created_at_unix_ms")?,
   1355         updated_at: nonnegative_millis(&row, "updated_at_unix_ms")?,
   1356     })
   1357 }
   1358 
   1359 fn decode_target(row: sqlx::sqlite::SqliteRow) -> Result<TargetRecord, OperationError> {
   1360     let ordinal = bounded_u16(
   1361         &row,
   1362         "target_ordinal",
   1363         0,
   1364         u16::try_from(RHI_PUBLICATION_TARGET_ORDINAL_MAXIMUM)
   1365             .map_err(|_| OperationError::Invariant)?,
   1366     )? as u8;
   1367     let relay_id = bounded_text(&row, "relay_id", "relay_id_bytes", 64)?;
   1368     if relay_id.is_empty()
   1369         || !relay_id.bytes().enumerate().all(|(index, byte)| {
   1370             (byte.is_ascii_lowercase() || byte.is_ascii_digit() || matches!(byte, b'_' | b'-'))
   1371                 && (index != 0 || byte.is_ascii_lowercase())
   1372         })
   1373     {
   1374         return Err(OperationError::Invariant);
   1375     }
   1376     let required = match row.try_get::<i64, _>("required") {
   1377         Ok(0) => false,
   1378         Ok(1) => true,
   1379         _ => return Err(OperationError::Invariant),
   1380     };
   1381     let state = decode_target_state(bounded_text(&row, "state", "state_bytes", 13)?.as_str())?;
   1382     let attempt_count = bounded_u16(
   1383         &row,
   1384         "attempt_count",
   1385         0,
   1386         RHI_PUBLICATION_ATTEMPT_NUMBER_MAXIMUM,
   1387     )?;
   1388     let last_attempt_id = optional_attempt_id(&row)?;
   1389     if (attempt_count == 0) != last_attempt_id.is_none() {
   1390         return Err(OperationError::Invariant);
   1391     }
   1392     Ok(TargetRecord {
   1393         ordinal,
   1394         relay_id: relay_id.into_boxed_str(),
   1395         required,
   1396         state,
   1397         revision: positive_u64(&row, "revision")?,
   1398         attempt_count,
   1399         next_attempt: optional_millis(&row, "next_attempt_unix_ms")?,
   1400         last_attempt_id,
   1401         updated_at: nonnegative_millis(&row, "updated_at_unix_ms")?,
   1402     })
   1403 }
   1404 
   1405 async fn read_attempt(
   1406     transaction: &mut ServiceSqliteTransaction<'_>,
   1407     id: RhiPublicationAttemptId,
   1408 ) -> Result<Option<AttemptRecord>, OperationError> {
   1409     let rows = sqlx::query(READ_ATTEMPT_SQL)
   1410         .bind(id.as_bytes().as_slice())
   1411         .fetch_all(&mut *transaction)
   1412         .await
   1413         .map_err(|_| OperationError::Storage)?;
   1414     exactly_zero_or_one(rows)?.map(decode_attempt).transpose()
   1415 }
   1416 
   1417 fn decode_attempt(row: sqlx::sqlite::SqliteRow) -> Result<AttemptRecord, OperationError> {
   1418     let outcome =
   1419         decode_attempt_outcome(bounded_text(&row, "outcome", "outcome_bytes", 13)?.as_str())?;
   1420     let result_code = bounded_text(&row, "result_code", "result_code_bytes", 64)?;
   1421     if result_code != outcome.code() {
   1422         return Err(OperationError::Invariant);
   1423     }
   1424     Ok(AttemptRecord {
   1425         attempt_id: attempt_id_from_durable_bytes(blob::<32>(&row, "attempt_id")?),
   1426         outbox_id: RhiPublicationOutboxId::from_committed_bytes(blob::<32>(&row, "outbox_id")?),
   1427         target_ordinal: bounded_u16(&row, "target_ordinal", 0, 31)? as u8,
   1428         attempt_number: bounded_u16(&row, "attempt_number", 1, 100)?,
   1429         event_sha256: blob::<32>(&row, "event_sha256")?,
   1430         lease_owner: RhiPublicationLeaseOwner(blob::<16>(&row, "lease_owner")?),
   1431         started_at: nonnegative_millis(&row, "started_at_unix_ms")?,
   1432         finished_at: nonnegative_millis(&row, "finished_at_unix_ms")?,
   1433         outcome,
   1434     })
   1435 }
   1436 
   1437 fn decode_outbox_state(value: &str) -> Result<RhiPublicationOutboxState, OperationError> {
   1438     match value {
   1439         "pending" => Ok(RhiPublicationOutboxState::Pending),
   1440         "leased" => Ok(RhiPublicationOutboxState::Leased),
   1441         "complete" => Ok(RhiPublicationOutboxState::Complete),
   1442         "blocked" => Ok(RhiPublicationOutboxState::Blocked),
   1443         _ => Err(OperationError::Invariant),
   1444     }
   1445 }
   1446 
   1447 fn decode_target_state(value: &str) -> Result<RhiPublicationTargetState, OperationError> {
   1448     match value {
   1449         "pending" => Ok(RhiPublicationTargetState::Pending),
   1450         "submitted" => Ok(RhiPublicationTargetState::Submitted),
   1451         "accepted" => Ok(RhiPublicationTargetState::Accepted),
   1452         "rejected" => Ok(RhiPublicationTargetState::Rejected),
   1453         "rate_limited" => Ok(RhiPublicationTargetState::RateLimited),
   1454         "auth_required" => Ok(RhiPublicationTargetState::AuthRequired),
   1455         "failed" => Ok(RhiPublicationTargetState::Failed),
   1456         "unknown" => Ok(RhiPublicationTargetState::Unknown),
   1457         _ => Err(OperationError::Invariant),
   1458     }
   1459 }
   1460 
   1461 fn decode_attempt_outcome(value: &str) -> Result<RhiPublicationAttemptOutcome, OperationError> {
   1462     match value {
   1463         "submitted" => Ok(RhiPublicationAttemptOutcome::Submitted),
   1464         "accepted" => Ok(RhiPublicationAttemptOutcome::Accepted),
   1465         "rejected" => Ok(RhiPublicationAttemptOutcome::Rejected),
   1466         "rate_limited" => Ok(RhiPublicationAttemptOutcome::RateLimited),
   1467         "auth_required" => Ok(RhiPublicationAttemptOutcome::AuthRequired),
   1468         "failed" => Ok(RhiPublicationAttemptOutcome::Failed),
   1469         "unknown" => Ok(RhiPublicationAttemptOutcome::Unknown),
   1470         _ => Err(OperationError::Invariant),
   1471     }
   1472 }
   1473 
   1474 fn valid_outbox_shape(
   1475     state: RhiPublicationOutboxState,
   1476     next: Option<RhiPublicationUnixMilliseconds>,
   1477     owner: Option<RhiPublicationLeaseOwner>,
   1478     expires: Option<RhiPublicationUnixMilliseconds>,
   1479 ) -> bool {
   1480     match state {
   1481         RhiPublicationOutboxState::Pending => {
   1482             next.is_some() && owner.is_none() && expires.is_none()
   1483         }
   1484         RhiPublicationOutboxState::Leased => {
   1485             next.is_none() && owner.is_some() && expires.is_some_and(|value| value.get() > 0)
   1486         }
   1487         RhiPublicationOutboxState::Complete | RhiPublicationOutboxState::Blocked => {
   1488             next.is_none() && owner.is_none() && expires.is_none()
   1489         }
   1490     }
   1491 }
   1492 
   1493 fn exactly_zero_or_one(
   1494     mut rows: Vec<sqlx::sqlite::SqliteRow>,
   1495 ) -> Result<Option<sqlx::sqlite::SqliteRow>, OperationError> {
   1496     match rows.len() {
   1497         0 => Ok(None),
   1498         1 => Ok(rows.pop()),
   1499         _ => Err(OperationError::Invariant),
   1500     }
   1501 }
   1502 
   1503 fn blob<const N: usize>(
   1504     row: &sqlx::sqlite::SqliteRow,
   1505     column: &str,
   1506 ) -> Result<[u8; N], OperationError> {
   1507     row.try_get::<Option<Vec<u8>>, _>(column)
   1508         .map_err(|_| OperationError::Invariant)?
   1509         .ok_or(OperationError::Invariant)?
   1510         .try_into()
   1511         .map_err(|_| OperationError::Invariant)
   1512 }
   1513 
   1514 fn bounded_text(
   1515     row: &sqlx::sqlite::SqliteRow,
   1516     column: &str,
   1517     length_column: &str,
   1518     maximum: usize,
   1519 ) -> Result<String, OperationError> {
   1520     let length = row
   1521         .try_get::<i64, _>(length_column)
   1522         .ok()
   1523         .and_then(|value| usize::try_from(value).ok())
   1524         .filter(|value| *value <= maximum)
   1525         .ok_or(OperationError::Invariant)?;
   1526     let value = row
   1527         .try_get::<String, _>(column)
   1528         .map_err(|_| OperationError::Invariant)?;
   1529     if value.len() != length {
   1530         return Err(OperationError::Invariant);
   1531     }
   1532     Ok(value)
   1533 }
   1534 
   1535 fn bounded_u16(
   1536     row: &sqlx::sqlite::SqliteRow,
   1537     column: &str,
   1538     minimum: u16,
   1539     maximum: u16,
   1540 ) -> Result<u16, OperationError> {
   1541     row.try_get::<i64, _>(column)
   1542         .ok()
   1543         .and_then(|value| u16::try_from(value).ok())
   1544         .filter(|value| *value >= minimum && *value <= maximum)
   1545         .ok_or(OperationError::Invariant)
   1546 }
   1547 
   1548 fn bounded_u64(
   1549     row: &sqlx::sqlite::SqliteRow,
   1550     column: &str,
   1551     minimum: u64,
   1552     maximum: u64,
   1553 ) -> Result<u64, OperationError> {
   1554     row.try_get::<i64, _>(column)
   1555         .ok()
   1556         .and_then(|value| u64::try_from(value).ok())
   1557         .filter(|value| *value >= minimum && *value <= maximum)
   1558         .ok_or(OperationError::Invariant)
   1559 }
   1560 
   1561 fn positive_u64(row: &sqlx::sqlite::SqliteRow, column: &str) -> Result<u64, OperationError> {
   1562     bounded_u64(row, column, 1, MAX_UNIX_MILLISECONDS)
   1563 }
   1564 
   1565 fn nonnegative_millis(
   1566     row: &sqlx::sqlite::SqliteRow,
   1567     column: &str,
   1568 ) -> Result<RhiPublicationUnixMilliseconds, OperationError> {
   1569     bounded_u64(row, column, 0, MAX_UNIX_MILLISECONDS).map(RhiPublicationUnixMilliseconds)
   1570 }
   1571 
   1572 fn optional_millis(
   1573     row: &sqlx::sqlite::SqliteRow,
   1574     column: &str,
   1575 ) -> Result<Option<RhiPublicationUnixMilliseconds>, OperationError> {
   1576     row.try_get::<Option<i64>, _>(column)
   1577         .map_err(|_| OperationError::Invariant)?
   1578         .map(|value| {
   1579             u64::try_from(value)
   1580                 .ok()
   1581                 .filter(|value| *value <= MAX_UNIX_MILLISECONDS)
   1582                 .map(RhiPublicationUnixMilliseconds)
   1583                 .ok_or(OperationError::Invariant)
   1584         })
   1585         .transpose()
   1586 }
   1587 
   1588 fn optional_owner(
   1589     row: &sqlx::sqlite::SqliteRow,
   1590 ) -> Result<Option<RhiPublicationLeaseOwner>, OperationError> {
   1591     let length = row
   1592         .try_get::<Option<i64>, _>("lease_owner_bytes")
   1593         .map_err(|_| OperationError::Invariant)?;
   1594     let bytes = row
   1595         .try_get::<Option<Vec<u8>>, _>("lease_owner")
   1596         .map_err(|_| OperationError::Invariant)?;
   1597     match (length, bytes) {
   1598         (None, None) => Ok(None),
   1599         (Some(length), Some(bytes))
   1600             if length == LEASE_OWNER_BYTES as i64 && bytes.len() == LEASE_OWNER_BYTES =>
   1601         {
   1602             let bytes: [u8; LEASE_OWNER_BYTES] =
   1603                 bytes.try_into().map_err(|_| OperationError::Invariant)?;
   1604             RhiPublicationLeaseOwner::from_bytes(bytes)
   1605                 .map(Some)
   1606                 .map_err(|_| OperationError::Invariant)
   1607         }
   1608         _ => Err(OperationError::Invariant),
   1609     }
   1610 }
   1611 
   1612 fn optional_attempt_id(
   1613     row: &sqlx::sqlite::SqliteRow,
   1614 ) -> Result<Option<RhiPublicationAttemptId>, OperationError> {
   1615     let length = row
   1616         .try_get::<Option<i64>, _>("last_attempt_id_bytes")
   1617         .map_err(|_| OperationError::Invariant)?;
   1618     let bytes = row
   1619         .try_get::<Option<Vec<u8>>, _>("last_attempt_id")
   1620         .map_err(|_| OperationError::Invariant)?;
   1621     match (length, bytes) {
   1622         (None, None) => Ok(None),
   1623         (Some(32), Some(bytes)) if bytes.len() == 32 => Ok(Some(attempt_id_from_durable_bytes(
   1624             bytes.try_into().map_err(|_| OperationError::Invariant)?,
   1625         ))),
   1626         _ => Err(OperationError::Invariant),
   1627     }
   1628 }
   1629 
   1630 fn attempt_id_from_durable_bytes(bytes: [u8; 32]) -> RhiPublicationAttemptId {
   1631     // Restricted to bounded database decoding; every caller subsequently
   1632     // compares the value with a freshly domain-derived identity.
   1633     crate::publication_attempt::attempt_id_from_durable_bytes(bytes)
   1634 }
   1635 
   1636 fn i64_value(value: u64) -> Result<i64, OperationError> {
   1637     i64::try_from(value).map_err(|_| OperationError::InvalidInput)
   1638 }
   1639 
   1640 fn optional_i64(
   1641     value: Option<RhiPublicationUnixMilliseconds>,
   1642 ) -> Result<Option<i64>, OperationError> {
   1643     value.map(|value| i64_value(value.get())).transpose()
   1644 }
   1645 
   1646 #[derive(Clone, Copy, PartialEq, Eq)]
   1647 struct OutboxRecord {
   1648     id: RhiPublicationOutboxId,
   1649     event_sha256: [u8; 32],
   1650     authority_sha256: [u8; 32],
   1651     target_set_sha256: [u8; 32],
   1652     target_count: u8,
   1653     required_target_count: u8,
   1654     max_attempts: u16,
   1655     initial_backoff_ms: u64,
   1656     maximum_backoff_ms: u64,
   1657     attempt_deadline_ms: u64,
   1658     state: RhiPublicationOutboxState,
   1659     revision: u64,
   1660     next_attempt: Option<RhiPublicationUnixMilliseconds>,
   1661     lease_owner: Option<RhiPublicationLeaseOwner>,
   1662     lease_expires: Option<RhiPublicationUnixMilliseconds>,
   1663     created_at: RhiPublicationUnixMilliseconds,
   1664     updated_at: RhiPublicationUnixMilliseconds,
   1665 }
   1666 
   1667 #[derive(Clone, Copy)]
   1668 struct AuthorityBinding {
   1669     authority_sha256: [u8; 32],
   1670     target_set_sha256: [u8; 32],
   1671     target_count: u8,
   1672     required_target_count: u8,
   1673     max_attempts: u16,
   1674     initial_backoff_ms: u64,
   1675     maximum_backoff_ms: u64,
   1676     attempt_deadline_ms: u64,
   1677 }
   1678 
   1679 impl AuthorityBinding {
   1680     fn from_authority(authority: &RhiPublicationAuthority) -> Option<Self> {
   1681         if authority.mode() == RhiPublicationMode::Disabled {
   1682             return None;
   1683         }
   1684         let retry = authority.retry_policy()?;
   1685         Some(Self {
   1686             authority_sha256: *authority.authority_sha256(),
   1687             target_set_sha256: *authority.target_set_sha256(),
   1688             target_count: u8::try_from(authority.targets().len()).ok()?,
   1689             required_target_count: u8::try_from(
   1690                 authority
   1691                     .targets()
   1692                     .iter()
   1693                     .filter(|target| target.required())
   1694                     .count(),
   1695             )
   1696             .ok()?,
   1697             max_attempts: retry.maximum_attempts(),
   1698             initial_backoff_ms: retry.initial_backoff_milliseconds(),
   1699             maximum_backoff_ms: retry.maximum_backoff_milliseconds(),
   1700             attempt_deadline_ms: retry.attempt_deadline_milliseconds(),
   1701         })
   1702     }
   1703 }
   1704 
   1705 #[derive(PartialEq, Eq)]
   1706 struct TargetRecord {
   1707     ordinal: u8,
   1708     relay_id: Box<str>,
   1709     required: bool,
   1710     state: RhiPublicationTargetState,
   1711     revision: u64,
   1712     attempt_count: u16,
   1713     next_attempt: Option<RhiPublicationUnixMilliseconds>,
   1714     last_attempt_id: Option<RhiPublicationAttemptId>,
   1715     updated_at: RhiPublicationUnixMilliseconds,
   1716 }
   1717 
   1718 impl Clone for TargetRecord {
   1719     fn clone(&self) -> Self {
   1720         Self {
   1721             ordinal: self.ordinal,
   1722             relay_id: self.relay_id.clone(),
   1723             required: self.required,
   1724             state: self.state,
   1725             revision: self.revision,
   1726             attempt_count: self.attempt_count,
   1727             next_attempt: self.next_attempt,
   1728             last_attempt_id: self.last_attempt_id,
   1729             updated_at: self.updated_at,
   1730         }
   1731     }
   1732 }
   1733 
   1734 #[derive(Clone, Copy)]
   1735 struct RecoveryCandidate {
   1736     outbox: OutboxRecord,
   1737     retry_upper_bound: u64,
   1738 }
   1739 
   1740 #[derive(Clone, Copy)]
   1741 struct RecordInput {
   1742     lease: RhiPublicationLease,
   1743     target_ordinal: u8,
   1744     target_revision: u64,
   1745     attempt_number: u16,
   1746     attempt_id: RhiPublicationAttemptId,
   1747     started_at: RhiPublicationUnixMilliseconds,
   1748     finished_at: RhiPublicationUnixMilliseconds,
   1749     outcome: RhiPublicationAttemptOutcome,
   1750     retry_delay: RhiPublicationRetryDelayMilliseconds,
   1751 }
   1752 
   1753 impl RecordInput {
   1754     fn from_prepared(
   1755         prepared: &RhiPreparedPublicationAttempt,
   1756         finished_at: RhiPublicationUnixMilliseconds,
   1757         outcome: RhiPublicationAttemptOutcome,
   1758         retry_delay: RhiPublicationRetryDelayMilliseconds,
   1759     ) -> Result<Self, RhiPublicationExecutionError> {
   1760         if finished_at < prepared.started_at || outcome == RhiPublicationAttemptOutcome::Submitted {
   1761             return Err(RhiPublicationExecutionError::new(
   1762                 RhiPublicationExecutionErrorKind::InvalidInput,
   1763             ));
   1764         }
   1765         let evidence = RhiPublicationAttemptEvidence::new(
   1766             &prepared.publication,
   1767             u32::from(prepared.target.ordinal),
   1768             prepared.target.attempt_count,
   1769             prepared.started_at,
   1770             finished_at,
   1771             outcome,
   1772         )
   1773         .map_err(|_| {
   1774             RhiPublicationExecutionError::new(RhiPublicationExecutionErrorKind::InvalidInput)
   1775         })?;
   1776         if evidence.id() != prepared.attempt_id {
   1777             return Err(RhiPublicationExecutionError::new(
   1778                 RhiPublicationExecutionErrorKind::Invariant,
   1779             ));
   1780         }
   1781         Ok(Self {
   1782             lease: prepared.lease,
   1783             target_ordinal: prepared.target.ordinal,
   1784             target_revision: prepared.target.revision,
   1785             attempt_number: prepared.target.attempt_count,
   1786             attempt_id: prepared.attempt_id,
   1787             started_at: prepared.started_at,
   1788             finished_at,
   1789             outcome,
   1790             retry_delay,
   1791         })
   1792     }
   1793 }
   1794 
   1795 #[derive(Clone, Copy, PartialEq, Eq)]
   1796 struct AttemptRecord {
   1797     attempt_id: RhiPublicationAttemptId,
   1798     outbox_id: RhiPublicationOutboxId,
   1799     target_ordinal: u8,
   1800     attempt_number: u16,
   1801     event_sha256: [u8; 32],
   1802     lease_owner: RhiPublicationLeaseOwner,
   1803     started_at: RhiPublicationUnixMilliseconds,
   1804     finished_at: RhiPublicationUnixMilliseconds,
   1805     outcome: RhiPublicationAttemptOutcome,
   1806 }
   1807 
   1808 impl AttemptRecord {
   1809     const fn from_input(input: RecordInput) -> Self {
   1810         Self {
   1811             attempt_id: input.attempt_id,
   1812             outbox_id: input.lease.outbox.id,
   1813             target_ordinal: input.target_ordinal,
   1814             attempt_number: input.attempt_number,
   1815             event_sha256: input.lease.outbox.event_sha256,
   1816             lease_owner: input.lease.owner,
   1817             started_at: input.started_at,
   1818             finished_at: input.finished_at,
   1819             outcome: input.outcome,
   1820         }
   1821     }
   1822 }
   1823 
   1824 #[derive(Clone, Copy)]
   1825 struct OutboxDisposition {
   1826     state: RhiPublicationOutboxState,
   1827     next_attempt: Option<RhiPublicationUnixMilliseconds>,
   1828 }
   1829 
   1830 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
   1831 enum OperationError {
   1832     InvalidInput,
   1833     NotReady,
   1834     LeaseLost,
   1835     Invariant,
   1836     Storage,
   1837 }
   1838 
   1839 fn map_transaction_error(
   1840     error: ServiceSqliteTransactionError<OperationError>,
   1841 ) -> RhiPublicationExecutionError {
   1842     if error.kind() == ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown {
   1843         return RhiPublicationExecutionError::new(
   1844             RhiPublicationExecutionErrorKind::CommitOutcomeUnknown,
   1845         );
   1846     }
   1847     RhiPublicationExecutionError::new(match error.operation_error().copied() {
   1848         Some(OperationError::InvalidInput) => RhiPublicationExecutionErrorKind::InvalidInput,
   1849         Some(OperationError::NotReady) => RhiPublicationExecutionErrorKind::NotReady,
   1850         Some(OperationError::LeaseLost) => RhiPublicationExecutionErrorKind::LeaseLost,
   1851         Some(OperationError::Invariant) => RhiPublicationExecutionErrorKind::Invariant,
   1852         Some(OperationError::Storage) | None => RhiPublicationExecutionErrorKind::Storage,
   1853     })
   1854 }
   1855 
   1856 #[cfg(test)]
   1857 mod tests {
   1858     use super::*;
   1859 
   1860     #[test]
   1861     fn retry_bounds_are_saturating_closed_and_errors_are_safe() {
   1862         let outbox = OutboxRecord {
   1863             id: RhiPublicationOutboxId::from_committed_bytes([1; 32]),
   1864             event_sha256: [2; 32],
   1865             authority_sha256: [3; 32],
   1866             target_set_sha256: [4; 32],
   1867             target_count: 1,
   1868             required_target_count: 1,
   1869             max_attempts: 100,
   1870             initial_backoff_ms: 250,
   1871             maximum_backoff_ms: 30_000,
   1872             attempt_deadline_ms: 5_000,
   1873             state: RhiPublicationOutboxState::Pending,
   1874             revision: 1,
   1875             next_attempt: Some(RhiPublicationUnixMilliseconds(0)),
   1876             lease_owner: None,
   1877             lease_expires: None,
   1878             created_at: RhiPublicationUnixMilliseconds(0),
   1879             updated_at: RhiPublicationUnixMilliseconds(0),
   1880         };
   1881         assert_eq!(retry_upper_bound(outbox, 1), 250);
   1882         assert_eq!(retry_upper_bound(outbox, 2), 500);
   1883         assert_eq!(retry_upper_bound(outbox, 100), 30_000);
   1884         assert!(RhiPublicationLeaseOwner::from_bytes([0; 16]).is_err());
   1885         assert!(RhiPublicationRetryDelayMilliseconds::new(3_600_001).is_err());
   1886         for kind in [
   1887             RhiPublicationExecutionErrorKind::InvalidMode,
   1888             RhiPublicationExecutionErrorKind::InvalidInput,
   1889             RhiPublicationExecutionErrorKind::NotReady,
   1890             RhiPublicationExecutionErrorKind::LeaseLost,
   1891             RhiPublicationExecutionErrorKind::Invariant,
   1892             RhiPublicationExecutionErrorKind::ClockUnavailable,
   1893             RhiPublicationExecutionErrorKind::EntropyUnavailable,
   1894             RhiPublicationExecutionErrorKind::Storage,
   1895             RhiPublicationExecutionErrorKind::CommitOutcomeUnknown,
   1896         ] {
   1897             let error = RhiPublicationExecutionError::new(kind);
   1898             assert!(error.code().starts_with("publication_execution_"));
   1899             assert!(Error::source(&error).is_none());
   1900             let rendered = format!("{error} {error:?}");
   1901             assert!(!rendered.contains("relay-primary"));
   1902             assert!(!rendered.contains("SELECT"));
   1903         }
   1904     }
   1905 }