rhi

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

presence_publication.rs (116530B)


      1 //! Exact signed presence documents and durable target delivery.
      2 
      3 use core::fmt;
      4 use std::error::Error;
      5 
      6 use nostr::{EventBuilder, Kind, PublicKey as NostrPublicKey, Tag, Timestamp};
      7 use radroots_event::{
      8     GenericEventDraft,
      9     envelope::{EventEnvelope, kind::KIND_APPLICATION_HANDLER},
     10     profile::AuthoredProfile,
     11     wire::{EventWireLimits, Nip01EventWire},
     12 };
     13 use radroots_event_codec::authoring::AuthoredEventPlan;
     14 use radroots_nostr::event::{
     15     ApplicationHandlerSpec, Verification, build_application_handler, verify, verify_id,
     16 };
     17 use radroots_service_host::{EntropySource, UnixTimeSeconds};
     18 use radroots_service_sqlite::{
     19     ServiceSqliteTransaction, ServiceSqliteTransactionError, ServiceSqliteTransactionErrorKind,
     20 };
     21 use radroots_transport::BoxFuture;
     22 use sha2::{Digest as _, Sha256};
     23 use sqlx::Row as _;
     24 use zeroize::Zeroizing;
     25 
     26 use crate::{
     27     RHI_PRESENCE_DESIRED_MAX_TARGETS, RhiDecryptedIdentity, RhiJitterBoundMilliseconds,
     28     RhiPresenceDesiredAuthority, RhiPresenceDesiredCommitOutcome, RhiPresenceDesiredMode,
     29     RhiPresenceDesiredState, RhiPresenceDocumentKind, RhiPresenceOutboxRepository,
     30     RhiPresenceTarget, RhiRuntimeAdapterErrorKind, RhiStateHostMode, RhiTimeEntropyAdapters,
     31     presence_desired::presence_authority_matches_state,
     32 };
     33 
     34 /// Exact version of the signed presence and delivery contract.
     35 pub const RHI_PRESENCE_PUBLICATION_CONTRACT_VERSION: u32 = 1;
     36 /// Maximum canonical bytes retained for one signed presence event.
     37 pub const RHI_PRESENCE_SIGNED_EVENT_MAX_BYTES: usize = 32 * 1024;
     38 /// Maximum durable attempts for one presence target.
     39 pub const RHI_PRESENCE_MAX_ATTEMPTS: u16 = 100;
     40 
     41 const PROFILE_NAME: &str = "rhi";
     42 const PROFILE_DISPLAY_NAME: &str = "Radroots RHI";
     43 const PROFILE_ABOUT: &str = "Radroots evidence reconciliation and attestation service";
     44 const APPLICATION_HANDLER_IDENTIFIER: &str = "rhi";
     45 const APPLICATION_HANDLER_KINDS: [u32; 1] = [3441];
     46 const INITIAL_BACKOFF_MILLISECONDS: u64 = 250;
     47 const MAXIMUM_BACKOFF_MILLISECONDS: u64 = 30_000;
     48 const ATTEMPT_DEADLINE_MILLISECONDS: u64 = 15_000;
     49 const MAX_UNIX_MILLISECONDS: u64 = i64::MAX as u64;
     50 const OUTBOX_ID_DOMAIN: &[u8] = b"radroots.rhi.presence_outbox.v1\0";
     51 const ATTEMPT_ID_DOMAIN: &[u8] = b"radroots.rhi.presence_attempt.v1\0";
     52 
     53 const READ_DESIRED_BINDING_SQL: &str = r#"SELECT singleton, generation, enabled,
     54     profile, application_handler,
     55     CASE WHEN typeof(target_set_sha256) = 'blob' AND length(target_set_sha256) = 32
     56         THEN target_set_sha256 ELSE NULL END AS target_set_sha256,
     57     target_count, required_target_count, queue_capacity,
     58     CASE WHEN typeof(desired_sha256) = 'blob' AND length(desired_sha256) = 32
     59         THEN desired_sha256 ELSE NULL END AS desired_sha256
     60 FROM presence_desired_state
     61 LIMIT 2"#;
     62 
     63 const READ_GENERATION_OUTBOX_SQL: &str = r#"SELECT
     64     CASE WHEN typeof(outbox_id) = 'blob' AND length(outbox_id) = 32
     65         THEN outbox_id ELSE NULL END AS outbox_id,
     66     desired_generation,
     67     length(CAST(document_kind AS BLOB)) AS document_kind_bytes,
     68     substr(document_kind, 1, 20) AS document_kind,
     69     CASE WHEN typeof(desired_sha256) = 'blob' AND length(desired_sha256) = 32
     70         THEN desired_sha256 ELSE NULL END AS desired_sha256,
     71     CASE WHEN typeof(target_set_sha256) = 'blob' AND length(target_set_sha256) = 32
     72         THEN target_set_sha256 ELSE NULL END AS target_set_sha256,
     73     CASE WHEN typeof(event_id) = 'blob' AND length(event_id) = 32
     74         THEN event_id ELSE NULL END AS event_id,
     75     CASE WHEN typeof(event_sha256) = 'blob' AND length(event_sha256) = 32
     76         THEN event_sha256 ELSE NULL END AS event_sha256,
     77     length(exact_signed_event_bytes) AS exact_signed_event_bytes_length,
     78     substr(exact_signed_event_bytes, 1, 32769) AS exact_signed_event_bytes,
     79     authored_at_unix_s,
     80     length(CAST(service_public_key AS BLOB)) AS service_public_key_bytes,
     81     substr(service_public_key, 1, 65) AS service_public_key,
     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, 11) AS state,
     85     revision, next_attempt_unix_ms,
     86     CASE WHEN lease_owner IS NULL THEN NULL ELSE length(lease_owner) END AS lease_owner_bytes,
     87     CASE WHEN lease_owner IS NULL THEN NULL ELSE substr(lease_owner, 1, 17) END AS lease_owner,
     88     lease_expires_unix_ms, created_at_unix_ms, updated_at_unix_ms
     89 FROM presence_outbox
     90 WHERE desired_generation = ?
     91 ORDER BY CASE document_kind
     92     WHEN 'service_profile' THEN 0
     93     WHEN 'application_handler' THEN 1
     94     ELSE 2 END
     95 LIMIT 3"#;
     96 
     97 const READ_OUTBOX_SQL: &str = r#"SELECT
     98     CASE WHEN typeof(outbox_id) = 'blob' AND length(outbox_id) = 32
     99         THEN outbox_id ELSE NULL END AS outbox_id,
    100     desired_generation,
    101     length(CAST(document_kind AS BLOB)) AS document_kind_bytes,
    102     substr(document_kind, 1, 20) AS document_kind,
    103     CASE WHEN typeof(desired_sha256) = 'blob' AND length(desired_sha256) = 32
    104         THEN desired_sha256 ELSE NULL END AS desired_sha256,
    105     CASE WHEN typeof(target_set_sha256) = 'blob' AND length(target_set_sha256) = 32
    106         THEN target_set_sha256 ELSE NULL END AS target_set_sha256,
    107     CASE WHEN typeof(event_id) = 'blob' AND length(event_id) = 32
    108         THEN event_id ELSE NULL END AS event_id,
    109     CASE WHEN typeof(event_sha256) = 'blob' AND length(event_sha256) = 32
    110         THEN event_sha256 ELSE NULL END AS event_sha256,
    111     length(exact_signed_event_bytes) AS exact_signed_event_bytes_length,
    112     substr(exact_signed_event_bytes, 1, 32769) AS exact_signed_event_bytes,
    113     authored_at_unix_s,
    114     length(CAST(service_public_key AS BLOB)) AS service_public_key_bytes,
    115     substr(service_public_key, 1, 65) AS service_public_key,
    116     target_count, required_target_count, max_attempts,
    117     initial_backoff_ms, maximum_backoff_ms, attempt_deadline_ms,
    118     length(CAST(state AS BLOB)) AS state_bytes, substr(state, 1, 11) AS state,
    119     revision, next_attempt_unix_ms,
    120     CASE WHEN lease_owner IS NULL THEN NULL ELSE length(lease_owner) END AS lease_owner_bytes,
    121     CASE WHEN lease_owner IS NULL THEN NULL ELSE substr(lease_owner, 1, 17) END AS lease_owner,
    122     lease_expires_unix_ms, created_at_unix_ms, updated_at_unix_ms
    123 FROM presence_outbox
    124 WHERE outbox_id = ?
    125 LIMIT 2"#;
    126 
    127 const READ_CLAIMABLE_OUTBOX_SQL: &str = r#"SELECT
    128     CASE WHEN typeof(outbox.outbox_id) = 'blob' AND length(outbox.outbox_id) = 32
    129         THEN outbox.outbox_id ELSE NULL END AS outbox_id
    130 FROM presence_outbox AS outbox
    131 JOIN presence_desired_state AS desired
    132   ON desired.singleton = 1
    133  AND desired.generation = outbox.desired_generation
    134  AND desired.desired_sha256 = outbox.desired_sha256
    135 WHERE outbox.state = 'pending' AND outbox.next_attempt_unix_ms <= ?
    136 ORDER BY outbox.next_attempt_unix_ms, outbox.created_at_unix_ms,
    137          CASE outbox.document_kind
    138              WHEN 'service_profile' THEN 0
    139              WHEN 'application_handler' THEN 1
    140              ELSE 2 END,
    141          outbox.outbox_id
    142 LIMIT 1"#;
    143 
    144 const READ_EXPIRED_OUTBOX_SQL: &str = r#"SELECT
    145     CASE WHEN typeof(outbox_id) = 'blob' AND length(outbox_id) = 32
    146         THEN outbox_id ELSE NULL END AS outbox_id
    147 FROM presence_outbox
    148 WHERE state = 'leased' AND lease_expires_unix_ms <= ?
    149 ORDER BY lease_expires_unix_ms, created_at_unix_ms, document_kind, outbox_id
    150 LIMIT 1"#;
    151 
    152 const READ_TARGETS_SQL: &str = r#"SELECT target_ordinal,
    153     length(CAST(relay_id AS BLOB)) AS relay_id_bytes,
    154     substr(relay_id, 1, 65) AS relay_id, required,
    155     length(CAST(state AS BLOB)) AS state_bytes, substr(state, 1, 14) AS state,
    156     revision, attempt_count, next_attempt_unix_ms,
    157     CASE WHEN last_attempt_id IS NULL THEN NULL ELSE length(last_attempt_id) END
    158         AS last_attempt_id_bytes,
    159     CASE WHEN last_attempt_id IS NULL THEN NULL ELSE substr(last_attempt_id, 1, 33) END
    160         AS last_attempt_id,
    161     updated_at_unix_ms
    162 FROM presence_targets
    163 WHERE outbox_id = ?
    164 ORDER BY target_ordinal
    165 LIMIT 33"#;
    166 
    167 const INSERT_OUTBOX_SQL: &str = r#"INSERT INTO presence_outbox (
    168     outbox_id, desired_generation, document_kind, desired_sha256,
    169     target_set_sha256, event_id, event_sha256, exact_signed_event_bytes,
    170     authored_at_unix_s, service_public_key, target_count, required_target_count,
    171     max_attempts, initial_backoff_ms, maximum_backoff_ms, attempt_deadline_ms,
    172     state, revision, next_attempt_unix_ms, lease_owner, lease_expires_unix_ms,
    173     created_at_unix_ms, updated_at_unix_ms
    174 ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 100, 250, 30000, 15000,
    175           'pending', 1, ?, NULL, NULL, ?, ?)"#;
    176 
    177 const INSERT_TARGET_SQL: &str = r#"INSERT INTO presence_targets (
    178     outbox_id, target_ordinal, relay_id, required, state, revision,
    179     attempt_count, next_attempt_unix_ms, last_attempt_id, updated_at_unix_ms
    180 ) VALUES (?, ?, ?, ?, 'pending', 1, 0, ?, NULL, ?)"#;
    181 
    182 const SUPERSEDE_OUTBOX_SQL: &str = r#"UPDATE presence_outbox
    183 SET state = 'superseded', revision = revision + 1,
    184     next_attempt_unix_ms = NULL, lease_owner = NULL,
    185     lease_expires_unix_ms = NULL, updated_at_unix_ms = ?
    186 WHERE desired_generation != ? AND state IN ('pending', 'blocked')
    187   AND updated_at_unix_ms <= ?"#;
    188 
    189 const COUNT_STALE_LEASES_SQL: &str = r#"SELECT COUNT(*) AS lease_count
    190 FROM presence_outbox
    191 WHERE desired_generation != ? AND state = 'leased'"#;
    192 
    193 const READ_MAX_AUTHORED_SQL: &str = r#"SELECT MAX(authored_at_unix_s) AS authored_at_unix_s
    194 FROM presence_outbox WHERE document_kind = ?"#;
    195 
    196 const CLAIM_OUTBOX_SQL: &str = r#"UPDATE presence_outbox
    197 SET state = 'leased', revision = revision + 1,
    198     next_attempt_unix_ms = NULL, lease_owner = ?, lease_expires_unix_ms = ?,
    199     updated_at_unix_ms = ?
    200 WHERE outbox_id = ? AND revision = ? AND state = 'pending'
    201   AND next_attempt_unix_ms <= ? AND updated_at_unix_ms <= ?"#;
    202 
    203 const PREPARE_TARGET_SQL: &str = r#"UPDATE presence_targets
    204 SET state = 'submitted', revision = revision + 1,
    205     attempt_count = attempt_count + 1, next_attempt_unix_ms = NULL,
    206     last_attempt_id = ?, updated_at_unix_ms = ?
    207 WHERE outbox_id = ? AND target_ordinal = ? AND revision = ?
    208   AND state IN ('pending', 'failed', 'rate_limited', 'unknown')
    209   AND attempt_count < ? AND next_attempt_unix_ms <= ?
    210   AND updated_at_unix_ms <= ?"#;
    211 
    212 const INSERT_ATTEMPT_SQL: &str = r#"INSERT INTO presence_attempts (
    213     attempt_id, outbox_id, target_ordinal, attempt_number, event_sha256,
    214     lease_owner, started_at_unix_ms, finished_at_unix_ms, outcome, result_code
    215 ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)"#;
    216 
    217 const READ_ATTEMPT_SQL: &str = r#"SELECT
    218     CASE WHEN typeof(attempt_id) = 'blob' AND length(attempt_id) = 32
    219         THEN attempt_id ELSE NULL END AS attempt_id,
    220     CASE WHEN typeof(outbox_id) = 'blob' AND length(outbox_id) = 32
    221         THEN outbox_id ELSE NULL END AS outbox_id,
    222     target_ordinal, attempt_number,
    223     CASE WHEN typeof(event_sha256) = 'blob' AND length(event_sha256) = 32
    224         THEN event_sha256 ELSE NULL END AS event_sha256,
    225     CASE WHEN typeof(lease_owner) = 'blob' AND length(lease_owner) = 16
    226         THEN lease_owner ELSE NULL END AS lease_owner,
    227     started_at_unix_ms, finished_at_unix_ms,
    228     length(CAST(outcome AS BLOB)) AS outcome_bytes, substr(outcome, 1, 14) AS outcome,
    229     length(CAST(result_code AS BLOB)) AS result_code_bytes,
    230     substr(result_code, 1, 65) AS result_code
    231 FROM presence_attempts
    232 WHERE attempt_id = ?
    233 LIMIT 2"#;
    234 
    235 const UPDATE_TARGET_OUTCOME_SQL: &str = r#"UPDATE presence_targets
    236 SET state = ?, revision = revision + 1, next_attempt_unix_ms = ?,
    237     updated_at_unix_ms = ?
    238 WHERE outbox_id = ? AND target_ordinal = ? AND revision = ?
    239   AND state = 'submitted' AND attempt_count = ? AND last_attempt_id = ?"#;
    240 
    241 const UPDATE_OUTBOX_AFTER_ATTEMPT_SQL: &str = r#"UPDATE presence_outbox
    242 SET state = ?, revision = revision + 1, next_attempt_unix_ms = ?,
    243     lease_owner = NULL, lease_expires_unix_ms = NULL, updated_at_unix_ms = ?
    244 WHERE outbox_id = ? AND revision = ? AND state = 'leased'
    245   AND lease_owner = ? AND lease_expires_unix_ms = ?
    246   AND lease_expires_unix_ms > ? AND updated_at_unix_ms <= ?"#;
    247 
    248 const UPDATE_OUTBOX_RECOVERY_SQL: &str = r#"UPDATE presence_outbox
    249 SET state = ?, revision = revision + 1, next_attempt_unix_ms = ?,
    250     lease_owner = NULL, lease_expires_unix_ms = NULL, updated_at_unix_ms = ?
    251 WHERE outbox_id = ? AND revision = ? AND state = 'leased'
    252   AND lease_owner = ? AND lease_expires_unix_ms = ?
    253   AND lease_expires_unix_ms <= ? AND updated_at_unix_ms <= ?"#;
    254 
    255 /// Stable source-free signed-presence and delivery failure class.
    256 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
    257 pub enum RhiPresencePublicationErrorKind {
    258     InvalidMode,
    259     InvalidInput,
    260     DesiredStateMismatch,
    261     IdentityMismatch,
    262     EntropyUnavailable,
    263     RenderingFailed,
    264     VerificationFailed,
    265     NotReady,
    266     LeaseLost,
    267     Invariant,
    268     ClockUnavailable,
    269     Storage,
    270     CommitOutcomeUnknown,
    271 }
    272 
    273 impl RhiPresencePublicationErrorKind {
    274     /// Returns the stable machine-readable failure code.
    275     #[must_use]
    276     pub const fn code(self) -> &'static str {
    277         match self {
    278             Self::InvalidMode => "presence_publication_mode_invalid",
    279             Self::InvalidInput => "presence_publication_input_invalid",
    280             Self::DesiredStateMismatch => "presence_publication_desired_state_mismatch",
    281             Self::IdentityMismatch => "presence_publication_identity_mismatch",
    282             Self::EntropyUnavailable => "presence_publication_entropy_unavailable",
    283             Self::RenderingFailed => "presence_publication_rendering_failed",
    284             Self::VerificationFailed => "presence_publication_verification_failed",
    285             Self::NotReady => "presence_publication_not_ready",
    286             Self::LeaseLost => "presence_publication_lease_lost",
    287             Self::Invariant => "presence_publication_invariant_failed",
    288             Self::ClockUnavailable => "presence_publication_clock_unavailable",
    289             Self::Storage => "presence_publication_storage_failed",
    290             Self::CommitOutcomeUnknown => "presence_publication_commit_outcome_unknown",
    291         }
    292     }
    293 }
    294 
    295 /// Redacted source-free signed-presence and delivery failure.
    296 #[derive(Clone, Copy, PartialEq, Eq)]
    297 pub struct RhiPresencePublicationError {
    298     kind: RhiPresencePublicationErrorKind,
    299 }
    300 
    301 impl RhiPresencePublicationError {
    302     /// Returns the stable failure class.
    303     #[must_use]
    304     pub const fn kind(self) -> RhiPresencePublicationErrorKind {
    305         self.kind
    306     }
    307 
    308     /// Returns the stable machine-readable failure code.
    309     #[must_use]
    310     pub const fn code(self) -> &'static str {
    311         self.kind.code()
    312     }
    313 }
    314 
    315 impl fmt::Display for RhiPresencePublicationError {
    316     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    317         formatter.write_str(match self.kind {
    318             RhiPresencePublicationErrorKind::InvalidMode => {
    319                 "RHI presence publication requires writable state"
    320             }
    321             RhiPresencePublicationErrorKind::InvalidInput => {
    322                 "RHI presence publication input is invalid"
    323             }
    324             RhiPresencePublicationErrorKind::DesiredStateMismatch => {
    325                 "RHI presence publication desired state does not match"
    326             }
    327             RhiPresencePublicationErrorKind::IdentityMismatch => {
    328                 "RHI presence publication identity does not match"
    329             }
    330             RhiPresencePublicationErrorKind::EntropyUnavailable => {
    331                 "RHI presence publication entropy is unavailable"
    332             }
    333             RhiPresencePublicationErrorKind::RenderingFailed => {
    334                 "RHI presence publication rendering failed"
    335             }
    336             RhiPresencePublicationErrorKind::VerificationFailed => {
    337                 "RHI presence publication verification failed"
    338             }
    339             RhiPresencePublicationErrorKind::NotReady => "RHI presence work is not ready",
    340             RhiPresencePublicationErrorKind::LeaseLost => {
    341                 "RHI presence publication lease is no longer authoritative"
    342             }
    343             RhiPresencePublicationErrorKind::Invariant => {
    344                 "RHI presence publication state invariant failed"
    345             }
    346             RhiPresencePublicationErrorKind::ClockUnavailable => {
    347                 "RHI presence publication clock is unavailable"
    348             }
    349             RhiPresencePublicationErrorKind::Storage => {
    350                 "RHI presence publication transaction failed"
    351             }
    352             RhiPresencePublicationErrorKind::CommitOutcomeUnknown => {
    353                 "RHI presence publication commit outcome is unknown"
    354             }
    355         })
    356     }
    357 }
    358 
    359 impl fmt::Debug for RhiPresencePublicationError {
    360     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    361         formatter
    362             .debug_struct("RhiPresencePublicationError")
    363             .field("kind", &self.kind)
    364             .finish()
    365     }
    366 }
    367 
    368 impl Error for RhiPresencePublicationError {}
    369 
    370 /// One independently verified exact signed presence document.
    371 pub struct RhiSignedPresenceDocument {
    372     kind: RhiPresenceDocumentKind,
    373     event_id: [u8; 32],
    374     signed_event_sha256: [u8; 32],
    375     signed_event_bytes: Box<[u8]>,
    376     created_at_unix_seconds: u64,
    377 }
    378 
    379 impl RhiSignedPresenceDocument {
    380     /// Returns the closed document kind.
    381     #[must_use]
    382     pub const fn kind(&self) -> RhiPresenceDocumentKind {
    383         self.kind
    384     }
    385 
    386     /// Returns the independently verified NIP-01 event identifier.
    387     #[must_use]
    388     pub const fn event_id(&self) -> &[u8; 32] {
    389         &self.event_id
    390     }
    391 
    392     /// Returns the digest of the exact retained bytes.
    393     #[must_use]
    394     pub const fn signed_event_sha256(&self) -> &[u8; 32] {
    395         &self.signed_event_sha256
    396     }
    397 
    398     /// Returns the exact canonical signed bytes.
    399     #[must_use]
    400     pub fn signed_event_bytes(&self) -> &[u8] {
    401         &self.signed_event_bytes
    402     }
    403 
    404     /// Returns the injected authored timestamp.
    405     #[must_use]
    406     pub const fn created_at_unix_seconds(&self) -> u64 {
    407         self.created_at_unix_seconds
    408     }
    409 }
    410 
    411 impl fmt::Debug for RhiSignedPresenceDocument {
    412     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    413         formatter
    414             .debug_struct("RhiSignedPresenceDocument")
    415             .field("kind", &self.kind)
    416             .field("signed_event_bytes", &self.signed_event_bytes.len())
    417             .finish_non_exhaustive()
    418     }
    419 }
    420 
    421 /// Sealed exact document set bound to one committed desired-state generation.
    422 pub struct RhiSignedPresenceDocuments {
    423     desired_state: RhiPresenceDesiredState,
    424     service_public_key: Box<str>,
    425     targets: Box<[RhiPresenceTarget]>,
    426     documents: Box<[RhiSignedPresenceDocument]>,
    427 }
    428 
    429 impl RhiSignedPresenceDocuments {
    430     /// Returns the exact committed desired state.
    431     #[must_use]
    432     pub const fn desired_state(&self) -> RhiPresenceDesiredState {
    433         self.desired_state
    434     }
    435 
    436     /// Returns the exact ordered signed-document inventory.
    437     #[must_use]
    438     pub fn documents(&self) -> &[RhiSignedPresenceDocument] {
    439         &self.documents
    440     }
    441 
    442     /// Returns the stable target count without exposing endpoints.
    443     #[must_use]
    444     pub fn target_count(&self) -> usize {
    445         self.targets.len()
    446     }
    447 
    448     fn retained_copy(&self) -> Self {
    449         Self {
    450             desired_state: self.desired_state,
    451             service_public_key: self.service_public_key.clone(),
    452             targets: self.targets.clone(),
    453             documents: self
    454                 .documents
    455                 .iter()
    456                 .map(|document| RhiSignedPresenceDocument {
    457                     kind: document.kind,
    458                     event_id: document.event_id,
    459                     signed_event_sha256: document.signed_event_sha256,
    460                     signed_event_bytes: document.signed_event_bytes.clone(),
    461                     created_at_unix_seconds: document.created_at_unix_seconds,
    462                 })
    463                 .collect(),
    464         }
    465     }
    466 }
    467 
    468 impl fmt::Debug for RhiSignedPresenceDocuments {
    469     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    470         formatter
    471             .debug_struct("RhiSignedPresenceDocuments")
    472             .field("desired_generation", &self.desired_state.generation())
    473             .field("document_count", &self.documents.len())
    474             .field("target_count", &self.targets.len())
    475             .finish_non_exhaustive()
    476     }
    477 }
    478 
    479 /// Opaque identity of one exact durable presence outbox.
    480 #[derive(Clone, Copy, PartialEq, Eq, Hash)]
    481 pub struct RhiPresenceOutboxId([u8; 32]);
    482 
    483 impl RhiPresenceOutboxId {
    484     /// Returns the exact identity bytes.
    485     #[must_use]
    486     pub const fn as_bytes(&self) -> &[u8; 32] {
    487         &self.0
    488     }
    489 }
    490 
    491 impl fmt::Debug for RhiPresenceOutboxId {
    492     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    493         formatter.write_str("RhiPresenceOutboxId([redacted])")
    494     }
    495 }
    496 
    497 /// Opaque identity of one durable presence attempt.
    498 #[derive(Clone, Copy, PartialEq, Eq, Hash)]
    499 pub struct RhiPresenceAttemptId([u8; 32]);
    500 
    501 impl RhiPresenceAttemptId {
    502     /// Returns the exact identity bytes.
    503     #[must_use]
    504     pub const fn as_bytes(&self) -> &[u8; 32] {
    505         &self.0
    506     }
    507 }
    508 
    509 impl fmt::Debug for RhiPresenceAttemptId {
    510     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    511         formatter.write_str("RhiPresenceAttemptId([redacted])")
    512     }
    513 }
    514 
    515 /// Bounded wall-clock instant in whole Unix milliseconds.
    516 #[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord)]
    517 pub struct RhiPresenceUnixMilliseconds(u64);
    518 
    519 impl RhiPresenceUnixMilliseconds {
    520     /// Validates one injected instant against SQLite's signed representation.
    521     pub fn new(value: u64) -> Result<Self, RhiPresencePublicationError> {
    522         if value > MAX_UNIX_MILLISECONDS {
    523             return Err(failure(RhiPresencePublicationErrorKind::InvalidInput));
    524         }
    525         Ok(Self(value))
    526     }
    527 
    528     /// Returns the exact instant.
    529     #[must_use]
    530     pub const fn get(self) -> u64 {
    531         self.0
    532     }
    533 }
    534 
    535 /// Bounded caller-injected retry delay in whole milliseconds.
    536 #[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord)]
    537 pub struct RhiPresenceRetryDelayMilliseconds(u64);
    538 
    539 impl RhiPresenceRetryDelayMilliseconds {
    540     /// Validates one delay against the absolute retry ceiling.
    541     pub fn new(value: u64) -> Result<Self, RhiPresencePublicationError> {
    542         if value > MAXIMUM_BACKOFF_MILLISECONDS {
    543             return Err(failure(RhiPresencePublicationErrorKind::InvalidInput));
    544         }
    545         Ok(Self(value))
    546     }
    547 
    548     /// Returns the exact delay.
    549     #[must_use]
    550     pub const fn get(self) -> u64 {
    551         self.0
    552     }
    553 }
    554 
    555 /// Stable process-local owner token for one compare-and-swap lease.
    556 #[derive(Clone, Copy, PartialEq, Eq, Hash)]
    557 pub struct RhiPresenceLeaseOwner([u8; 16]);
    558 
    559 impl RhiPresenceLeaseOwner {
    560     /// Validates one injected nonzero owner token.
    561     pub fn from_bytes(bytes: [u8; 16]) -> Result<Self, RhiPresencePublicationError> {
    562         if bytes.iter().all(|byte| *byte == 0) {
    563             return Err(failure(RhiPresencePublicationErrorKind::InvalidInput));
    564         }
    565         Ok(Self(bytes))
    566     }
    567 }
    568 
    569 impl fmt::Debug for RhiPresenceLeaseOwner {
    570     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    571         formatter.write_str("RhiPresenceLeaseOwner([redacted])")
    572     }
    573 }
    574 
    575 /// Stable durable presence-outbox lifecycle state.
    576 #[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
    577 pub enum RhiPresenceOutboxState {
    578     Pending,
    579     Leased,
    580     Complete,
    581     Blocked,
    582     Superseded,
    583 }
    584 
    585 impl RhiPresenceOutboxState {
    586     /// Returns the exact machine-contract spelling.
    587     #[must_use]
    588     pub const fn code(self) -> &'static str {
    589         match self {
    590             Self::Pending => "pending",
    591             Self::Leased => "leased",
    592             Self::Complete => "complete",
    593             Self::Blocked => "blocked",
    594             Self::Superseded => "superseded",
    595         }
    596     }
    597 }
    598 
    599 /// Stable durable per-target presence state.
    600 #[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
    601 pub enum RhiPresenceTargetState {
    602     Pending,
    603     Submitted,
    604     Accepted,
    605     Rejected,
    606     RateLimited,
    607     AuthRequired,
    608     Failed,
    609     Unknown,
    610 }
    611 
    612 impl RhiPresenceTargetState {
    613     /// Returns the exact machine-contract spelling.
    614     #[must_use]
    615     pub const fn code(self) -> &'static str {
    616         match self {
    617             Self::Pending => "pending",
    618             Self::Submitted => "submitted",
    619             Self::Accepted => "accepted",
    620             Self::Rejected => "rejected",
    621             Self::RateLimited => "rate_limited",
    622             Self::AuthRequired => "auth_required",
    623             Self::Failed => "failed",
    624             Self::Unknown => "unknown",
    625         }
    626     }
    627 }
    628 
    629 /// Closed result observed from one exact-byte relay submission.
    630 #[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
    631 pub enum RhiPresenceAttemptOutcome {
    632     Submitted,
    633     Accepted,
    634     Rejected,
    635     RateLimited,
    636     AuthRequired,
    637     Failed,
    638     Unknown,
    639 }
    640 
    641 impl RhiPresenceAttemptOutcome {
    642     /// Returns the exact machine-contract spelling.
    643     #[must_use]
    644     pub const fn code(self) -> &'static str {
    645         match self {
    646             Self::Submitted => "submitted",
    647             Self::Accepted => "accepted",
    648             Self::Rejected => "rejected",
    649             Self::RateLimited => "rate_limited",
    650             Self::AuthRequired => "auth_required",
    651             Self::Failed => "failed",
    652             Self::Unknown => "unknown",
    653         }
    654     }
    655 }
    656 
    657 /// Confirmed result of committing one complete signed presence generation.
    658 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
    659 pub struct RhiPresenceCommitOutcome {
    660     desired_state: RhiPresenceDesiredState,
    661     document_count: u8,
    662     changed: bool,
    663 }
    664 
    665 impl RhiPresenceCommitOutcome {
    666     /// Returns the exact committed desired-state identity.
    667     #[must_use]
    668     pub const fn desired_state(self) -> RhiPresenceDesiredState {
    669         self.desired_state
    670     }
    671 
    672     /// Returns the exact committed document count.
    673     #[must_use]
    674     pub const fn document_count(self) -> u8 {
    675         self.document_count
    676     }
    677 
    678     /// Reports whether any durable presence workflow state changed.
    679     #[must_use]
    680     pub const fn changed(self) -> bool {
    681         self.changed
    682     }
    683 }
    684 
    685 /// Non-forgeable compare-and-swap authority for one claimed presence outbox.
    686 #[derive(PartialEq, Eq)]
    687 pub struct RhiPresenceLease {
    688     outbox: PresenceOutboxRecord,
    689     owner: RhiPresenceLeaseOwner,
    690     expires_at: RhiPresenceUnixMilliseconds,
    691 }
    692 
    693 impl RhiPresenceLease {
    694     /// Returns the exact claimed outbox identity.
    695     #[must_use]
    696     pub const fn outbox_id(&self) -> RhiPresenceOutboxId {
    697         self.outbox.id
    698     }
    699 
    700     /// Returns the exact lease expiry.
    701     #[must_use]
    702     pub const fn expires_at(&self) -> RhiPresenceUnixMilliseconds {
    703         self.expires_at
    704     }
    705 }
    706 
    707 impl fmt::Debug for RhiPresenceLease {
    708     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    709         formatter
    710             .debug_struct("RhiPresenceLease")
    711             .field("identity", &"[redacted]")
    712             .field("revision", &self.outbox.revision)
    713             .field("expires_at", &self.expires_at)
    714             .finish()
    715     }
    716 }
    717 
    718 /// Exact signed bytes exposed only after durable Submitted state exists.
    719 #[must_use = "prepared presence must be executed exactly or left for unknown recovery"]
    720 pub struct RhiPreparedPresenceAttempt {
    721     lease: RhiPresenceLease,
    722     target: PresenceTargetRecord,
    723     attempt_id: RhiPresenceAttemptId,
    724     started_at: RhiPresenceUnixMilliseconds,
    725     deadline_at: RhiPresenceUnixMilliseconds,
    726     exact_signed_event_bytes: Box<[u8]>,
    727 }
    728 
    729 impl RhiPreparedPresenceAttempt {
    730     /// Returns the exact attempt identity.
    731     #[must_use]
    732     pub const fn attempt_id(&self) -> RhiPresenceAttemptId {
    733         self.attempt_id
    734     }
    735 
    736     /// Returns the stable relay identity without an endpoint or secret.
    737     #[must_use]
    738     pub fn relay_id(&self) -> &str {
    739         &self.target.relay_id
    740     }
    741 
    742     /// Returns the original committed bytes with no transformation.
    743     #[must_use]
    744     pub const fn exact_signed_event_bytes(&self) -> &[u8] {
    745         &self.exact_signed_event_bytes
    746     }
    747 
    748     /// Returns the absolute attempt deadline.
    749     #[must_use]
    750     pub const fn deadline_at(&self) -> RhiPresenceUnixMilliseconds {
    751         self.deadline_at
    752     }
    753 
    754     /// Returns the one-based durable attempt number.
    755     #[must_use]
    756     pub const fn attempt_number(&self) -> u16 {
    757         self.target.attempt_count
    758     }
    759 
    760     fn retry_upper_bound(&self) -> u64 {
    761         retry_upper_bound(self.target.attempt_count)
    762     }
    763 }
    764 
    765 impl fmt::Debug for RhiPreparedPresenceAttempt {
    766     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    767         formatter
    768             .debug_struct("RhiPreparedPresenceAttempt")
    769             .field("identity", &"[redacted]")
    770             .field("target_ordinal", &self.target.ordinal)
    771             .field("attempt_number", &self.target.attempt_count)
    772             .field("deadline_at", &self.deadline_at)
    773             .finish()
    774     }
    775 }
    776 
    777 /// Closed exact-byte transport boundary for durable presence delivery.
    778 pub trait RhiExactPresenceSink: Send + Sync {
    779     /// Submits the exact retained payload and returns one closed observation.
    780     fn submit_exact<'a>(
    781         &'a self,
    782         attempt: &'a RhiPreparedPresenceAttempt,
    783     ) -> BoxFuture<'a, RhiPresenceAttemptOutcome>;
    784 }
    785 
    786 /// Confirmed durable result for one exact presence attempt.
    787 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
    788 pub struct RhiPresenceAttemptCommit {
    789     outbox_id: RhiPresenceOutboxId,
    790     attempt_id: RhiPresenceAttemptId,
    791     target_ordinal: u8,
    792     attempt_number: u16,
    793     outcome: RhiPresenceAttemptOutcome,
    794     target_state: RhiPresenceTargetState,
    795     outbox_state: RhiPresenceOutboxState,
    796 }
    797 
    798 impl RhiPresenceAttemptCommit {
    799     #[must_use]
    800     pub const fn outbox_id(self) -> RhiPresenceOutboxId {
    801         self.outbox_id
    802     }
    803 
    804     #[must_use]
    805     pub const fn attempt_id(self) -> RhiPresenceAttemptId {
    806         self.attempt_id
    807     }
    808 
    809     #[must_use]
    810     pub const fn target_ordinal(self) -> u8 {
    811         self.target_ordinal
    812     }
    813 
    814     #[must_use]
    815     pub const fn attempt_number(self) -> u16 {
    816         self.attempt_number
    817     }
    818 
    819     #[must_use]
    820     pub const fn outcome(self) -> RhiPresenceAttemptOutcome {
    821         self.outcome
    822     }
    823 
    824     #[must_use]
    825     pub const fn target_state(self) -> RhiPresenceTargetState {
    826         self.target_state
    827     }
    828 
    829     #[must_use]
    830     pub const fn outbox_state(self) -> RhiPresenceOutboxState {
    831         self.outbox_state
    832     }
    833 }
    834 
    835 /// Builds, signs, and independently revalidates one complete desired document set.
    836 ///
    837 /// Exactly 32 entropy bytes are consumed per enabled document. No clock,
    838 /// persistence, relay, task, or network authority is acquired implicitly.
    839 pub fn build_rhi_signed_presence_documents(
    840     desired: RhiPresenceDesiredCommitOutcome,
    841     authority: &RhiPresenceDesiredAuthority,
    842     identity: &RhiDecryptedIdentity,
    843     created_at: UnixTimeSeconds,
    844     entropy: &dyn EntropySource,
    845 ) -> Result<RhiSignedPresenceDocuments, RhiPresencePublicationError> {
    846     let state = desired.state();
    847     if created_at.get() > i64::MAX as u64 {
    848         return Err(failure(RhiPresencePublicationErrorKind::InvalidInput));
    849     }
    850     if !presence_authority_matches_state(state, authority) {
    851         return Err(failure(
    852             RhiPresencePublicationErrorKind::DesiredStateMismatch,
    853         ));
    854     }
    855     if identity.public_identity().as_hex() != authority.service_public_key() {
    856         return Err(failure(RhiPresencePublicationErrorKind::IdentityMismatch));
    857     }
    858     let mut documents = Vec::with_capacity(authority.document_kinds().len());
    859     for kind in authority.document_kinds() {
    860         let plan = presence_plan(*kind, identity.public_identity().as_hex(), created_at.get())?;
    861         let mut auxiliary = Zeroizing::new([0_u8; 32]);
    862         entropy
    863             .fill_bytes(&mut auxiliary[..])
    864             .map_err(|_| failure(RhiPresencePublicationErrorKind::EntropyUnavailable))?;
    865         let event = sign_plan(identity, &plan, &auxiliary)?;
    866         let bytes = serde_json::to_vec(&event)
    867             .map_err(|_| failure(RhiPresencePublicationErrorKind::RenderingFailed))?;
    868         let verified = validate_signed_document(
    869             *kind,
    870             identity.public_identity().as_hex(),
    871             created_at.get(),
    872             &bytes,
    873         )?;
    874         documents.push(RhiSignedPresenceDocument {
    875             kind: *kind,
    876             event_id: *verified.id().as_bytes(),
    877             signed_event_sha256: Sha256::digest(&bytes).into(),
    878             signed_event_bytes: bytes.into_boxed_slice(),
    879             created_at_unix_seconds: created_at.get(),
    880         });
    881     }
    882     if (state.mode() == RhiPresenceDesiredMode::Disabled && !documents.is_empty())
    883         || (state.mode() == RhiPresenceDesiredMode::Enabled
    884             && documents.len() != authority.document_kinds().len())
    885     {
    886         return Err(failure(RhiPresencePublicationErrorKind::Invariant));
    887     }
    888     Ok(RhiSignedPresenceDocuments {
    889         desired_state: state,
    890         service_public_key: authority.service_public_key().into(),
    891         targets: authority.targets().to_vec().into_boxed_slice(),
    892         documents: documents.into_boxed_slice(),
    893     })
    894 }
    895 
    896 /// Independently verifies every exact signed document against its bound intent.
    897 pub fn validate_rhi_signed_presence_documents(
    898     documents: &RhiSignedPresenceDocuments,
    899     authority: &RhiPresenceDesiredAuthority,
    900 ) -> Result<(), RhiPresencePublicationError> {
    901     if !presence_authority_matches_state(documents.desired_state, authority)
    902         || documents.service_public_key.as_ref() != authority.service_public_key()
    903         || documents.targets.as_ref() != authority.targets()
    904         || documents.documents.len() != authority.document_kinds().len()
    905     {
    906         return Err(failure(
    907             RhiPresencePublicationErrorKind::DesiredStateMismatch,
    908         ));
    909     }
    910     for (document, kind) in documents.documents.iter().zip(authority.document_kinds()) {
    911         if document.kind != *kind
    912             || Sha256::digest(document.signed_event_bytes.as_ref()).as_slice()
    913                 != document.signed_event_sha256
    914         {
    915             return Err(failure(RhiPresencePublicationErrorKind::VerificationFailed));
    916         }
    917         let verified = validate_signed_document(
    918             document.kind,
    919             authority.service_public_key(),
    920             document.created_at_unix_seconds,
    921             &document.signed_event_bytes,
    922         )?;
    923         if verified.id().as_bytes() != &document.event_id {
    924             return Err(failure(RhiPresencePublicationErrorKind::VerificationFailed));
    925         }
    926     }
    927     Ok(())
    928 }
    929 
    930 impl RhiPresenceOutboxRepository<'_> {
    931     /// Atomically commits exact signed bytes and the complete immutable target set.
    932     ///
    933     /// A successful return means every exact payload and initial target state is
    934     /// durable before any relay adapter can observe those bytes. The caller
    935     /// retains the sealed exact-byte capability so an unknown commit result can
    936     /// be reconciled by replaying this same value. Exact replay is idempotent.
    937     pub async fn commit_signed_presence(
    938         &self,
    939         documents: &RhiSignedPresenceDocuments,
    940         committed_at: RhiPresenceUnixMilliseconds,
    941     ) -> Result<RhiPresenceCommitOutcome, RhiPresencePublicationError> {
    942         require_writable(self)?;
    943         let documents = documents.retained_copy();
    944         self.host()
    945             .sqlite_host()
    946             .transaction(move |transaction| {
    947                 Box::pin(async move { commit_signed(transaction, documents, committed_at).await })
    948             })
    949             .await
    950             .map_err(map_transaction_error)
    951     }
    952 
    953     /// Claims the oldest currently desired, due presence outbox.
    954     pub async fn claim_next_presence(
    955         &self,
    956         owner: RhiPresenceLeaseOwner,
    957         now: RhiPresenceUnixMilliseconds,
    958     ) -> Result<Option<RhiPresenceLease>, RhiPresencePublicationError> {
    959         require_writable(self)?;
    960         self.host()
    961             .sqlite_host()
    962             .transaction(move |transaction| {
    963                 Box::pin(async move { claim_next(transaction, owner, now).await })
    964             })
    965             .await
    966             .map_err(map_transaction_error)
    967     }
    968 
    969     /// Persists Submitted before exposing the committed exact bytes.
    970     pub async fn prepare_next_presence_target(
    971         &self,
    972         lease: RhiPresenceLease,
    973         started_at: RhiPresenceUnixMilliseconds,
    974     ) -> Result<RhiPreparedPresenceAttempt, RhiPresencePublicationError> {
    975         require_writable(self)?;
    976         self.host()
    977             .sqlite_host()
    978             .transaction(move |transaction| {
    979                 Box::pin(async move { prepare_next(transaction, lease, started_at).await })
    980             })
    981             .await
    982             .map_err(map_transaction_error)
    983     }
    984 
    985     /// Commits one closed observed outcome by exact compare-and-swap.
    986     pub async fn record_presence_outcome(
    987         &self,
    988         prepared: &RhiPreparedPresenceAttempt,
    989         finished_at: RhiPresenceUnixMilliseconds,
    990         outcome: RhiPresenceAttemptOutcome,
    991         retry_delay: RhiPresenceRetryDelayMilliseconds,
    992     ) -> Result<RhiPresenceAttemptCommit, RhiPresencePublicationError> {
    993         require_writable(self)?;
    994         let input =
    995             PresenceRecordInput::from_prepared(prepared, finished_at, outcome, retry_delay)?;
    996         self.host()
    997             .sqlite_host()
    998             .transaction(move |transaction| {
    999                 Box::pin(async move { record_outcome(transaction, input).await })
   1000             })
   1001             .await
   1002             .map_err(map_transaction_error)
   1003     }
   1004 
   1005     /// Recovers at most one expired lease and records Unknown for submitted work.
   1006     pub async fn recover_one_expired_presence(
   1007         &self,
   1008         adapters: &RhiTimeEntropyAdapters,
   1009         now: RhiPresenceUnixMilliseconds,
   1010     ) -> Result<bool, RhiPresencePublicationError> {
   1011         require_writable(self)?;
   1012         let candidate = self
   1013             .host()
   1014             .sqlite_host()
   1015             .transaction(move |transaction| {
   1016                 Box::pin(async move { read_expired_candidate(transaction, now).await })
   1017             })
   1018             .await
   1019             .map_err(map_transaction_error)?;
   1020         let Some(candidate) = candidate else {
   1021             return Ok(false);
   1022         };
   1023         let delay = sample_retry_delay(adapters, candidate.retry_upper_bound)?;
   1024         self.host()
   1025             .sqlite_host()
   1026             .transaction(move |transaction| {
   1027                 Box::pin(async move { recover_expired(transaction, candidate, now, delay).await })
   1028             })
   1029             .await
   1030             .map_err(map_transaction_error)
   1031     }
   1032 
   1033     /// Executes at most one exact-byte presence attempt through an injected sink.
   1034     ///
   1035     /// Cancellation after durable preparation leaves Submitted evidence. Lease
   1036     /// expiry recovery records Unknown before retrying the same retained bytes.
   1037     pub async fn execute_next_presence(
   1038         &self,
   1039         owner: RhiPresenceLeaseOwner,
   1040         adapters: &RhiTimeEntropyAdapters,
   1041         sink: &dyn RhiExactPresenceSink,
   1042     ) -> Result<Option<RhiPresenceAttemptCommit>, RhiPresencePublicationError> {
   1043         let now = presence_now(adapters)?;
   1044         self.recover_one_expired_presence(adapters, now).await?;
   1045         let Some(lease) = self.claim_next_presence(owner, now).await? else {
   1046             return Ok(None);
   1047         };
   1048         let prepared = self.prepare_next_presence_target(lease, now).await?;
   1049         let outcome = match sink.submit_exact(&prepared).await {
   1050             RhiPresenceAttemptOutcome::Submitted => RhiPresenceAttemptOutcome::Unknown,
   1051             outcome => outcome,
   1052         };
   1053         let finished_at = presence_now(adapters)?;
   1054         let retry_delay = if retryable_outcome(outcome) {
   1055             sample_retry_delay(adapters, prepared.retry_upper_bound())?
   1056         } else {
   1057             RhiPresenceRetryDelayMilliseconds(0)
   1058         };
   1059         self.record_presence_outcome(&prepared, finished_at, outcome, retry_delay)
   1060             .await
   1061             .map(Some)
   1062     }
   1063 }
   1064 
   1065 fn require_writable(
   1066     repository: &RhiPresenceOutboxRepository<'_>,
   1067 ) -> Result<(), RhiPresencePublicationError> {
   1068     if repository.host().mode() == RhiStateHostMode::ReadWriteExisting {
   1069         Ok(())
   1070     } else {
   1071         Err(failure(RhiPresencePublicationErrorKind::InvalidMode))
   1072     }
   1073 }
   1074 
   1075 fn presence_now(
   1076     adapters: &RhiTimeEntropyAdapters,
   1077 ) -> Result<RhiPresenceUnixMilliseconds, RhiPresencePublicationError> {
   1078     let value = adapters.now_utc_milliseconds().map_err(|error| {
   1079         failure(match error.kind() {
   1080             RhiRuntimeAdapterErrorKind::WallClockUnavailable => {
   1081                 RhiPresencePublicationErrorKind::ClockUnavailable
   1082             }
   1083             _ => RhiPresencePublicationErrorKind::ClockUnavailable,
   1084         })
   1085     })?;
   1086     RhiPresenceUnixMilliseconds::new(value)
   1087         .map_err(|_| failure(RhiPresencePublicationErrorKind::ClockUnavailable))
   1088 }
   1089 
   1090 fn sample_retry_delay(
   1091     adapters: &RhiTimeEntropyAdapters,
   1092     cap: u64,
   1093 ) -> Result<RhiPresenceRetryDelayMilliseconds, RhiPresencePublicationError> {
   1094     if cap == 0 {
   1095         return Ok(RhiPresenceRetryDelayMilliseconds(0));
   1096     }
   1097     let bound = RhiJitterBoundMilliseconds::new(cap)
   1098         .map_err(|_| failure(RhiPresencePublicationErrorKind::InvalidInput))?;
   1099     adapters
   1100         .sample_full_jitter(bound)
   1101         .map(|delay| RhiPresenceRetryDelayMilliseconds(delay.get()))
   1102         .map_err(|_| failure(RhiPresencePublicationErrorKind::EntropyUnavailable))
   1103 }
   1104 
   1105 async fn commit_signed(
   1106     transaction: &mut ServiceSqliteTransaction<'_>,
   1107     documents: RhiSignedPresenceDocuments,
   1108     committed_at: RhiPresenceUnixMilliseconds,
   1109 ) -> Result<RhiPresenceCommitOutcome, PresenceOperationError> {
   1110     let desired = read_desired_binding(transaction).await?;
   1111     validate_document_set(&documents, desired)?;
   1112     let current = read_generation_outboxes(transaction, desired.generation).await?;
   1113     if !current.is_empty() && generation_matches(&documents, &current)? {
   1114         return Ok(RhiPresenceCommitOutcome {
   1115             desired_state: documents.desired_state,
   1116             document_count: u8::try_from(documents.documents.len())
   1117                 .map_err(|_| PresenceOperationError::Invariant)?,
   1118             changed: false,
   1119         });
   1120     }
   1121     if !current.is_empty() {
   1122         return Err(PresenceOperationError::Invariant);
   1123     }
   1124     if documents.desired_state.mode() == RhiPresenceDesiredMode::Enabled {
   1125         for document in &documents.documents {
   1126             let rows = sqlx::query(READ_MAX_AUTHORED_SQL)
   1127                 .bind(document.kind.code())
   1128                 .fetch_all(&mut *transaction)
   1129                 .await
   1130                 .map_err(|_| PresenceOperationError::Storage)?;
   1131             let [row] = rows.as_slice() else {
   1132                 return Err(PresenceOperationError::Invariant);
   1133             };
   1134             if let Some(previous) = row
   1135                 .try_get::<Option<i64>, _>("authored_at_unix_s")
   1136                 .map_err(|_| PresenceOperationError::Invariant)?
   1137             {
   1138                 let previous =
   1139                     u64::try_from(previous).map_err(|_| PresenceOperationError::Invariant)?;
   1140                 if document.created_at_unix_seconds <= previous {
   1141                     return Err(PresenceOperationError::InvalidInput);
   1142                 }
   1143             }
   1144         }
   1145     }
   1146     let lease_rows = sqlx::query(COUNT_STALE_LEASES_SQL)
   1147         .bind(i64_value(desired.generation)?)
   1148         .fetch_all(&mut *transaction)
   1149         .await
   1150         .map_err(|_| PresenceOperationError::Storage)?;
   1151     let [lease_row] = lease_rows.as_slice() else {
   1152         return Err(PresenceOperationError::Invariant);
   1153     };
   1154     if lease_row
   1155         .try_get::<i64, _>("lease_count")
   1156         .ok()
   1157         .filter(|value| *value >= 0)
   1158         .ok_or(PresenceOperationError::Invariant)?
   1159         != 0
   1160     {
   1161         return Err(PresenceOperationError::NotReady);
   1162     }
   1163     let superseded = sqlx::query(SUPERSEDE_OUTBOX_SQL)
   1164         .bind(i64_value(committed_at.get())?)
   1165         .bind(i64_value(desired.generation)?)
   1166         .bind(i64_value(committed_at.get())?)
   1167         .execute(&mut *transaction)
   1168         .await
   1169         .map_err(|_| PresenceOperationError::Storage)?
   1170         .rows_affected();
   1171 
   1172     for document in &documents.documents {
   1173         let outbox_id = derive_outbox_id(
   1174             desired.generation,
   1175             document.kind,
   1176             desired.desired_sha256,
   1177             document.signed_event_sha256,
   1178         );
   1179         let inserted = sqlx::query(INSERT_OUTBOX_SQL)
   1180             .bind(outbox_id.as_bytes().as_slice())
   1181             .bind(i64_value(desired.generation)?)
   1182             .bind(document.kind.code())
   1183             .bind(desired.desired_sha256.as_slice())
   1184             .bind(desired.target_set_sha256.as_slice())
   1185             .bind(document.event_id.as_slice())
   1186             .bind(document.signed_event_sha256.as_slice())
   1187             .bind(document.signed_event_bytes.as_ref())
   1188             .bind(i64_value(document.created_at_unix_seconds)?)
   1189             .bind(documents.service_public_key.as_ref())
   1190             .bind(i64::from(desired.target_count))
   1191             .bind(i64::from(desired.required_target_count))
   1192             .bind(i64_value(committed_at.get())?)
   1193             .bind(i64_value(committed_at.get())?)
   1194             .bind(i64_value(committed_at.get())?)
   1195             .execute(&mut *transaction)
   1196             .await
   1197             .map_err(|_| PresenceOperationError::Storage)?;
   1198         if inserted.rows_affected() != 1 {
   1199             return Err(PresenceOperationError::Invariant);
   1200         }
   1201         for target in &documents.targets {
   1202             let inserted = sqlx::query(INSERT_TARGET_SQL)
   1203                 .bind(outbox_id.as_bytes().as_slice())
   1204                 .bind(i64::from(target.ordinal()))
   1205                 .bind(target.relay_id())
   1206                 .bind(i64::from(target.required()))
   1207                 .bind(i64_value(committed_at.get())?)
   1208                 .bind(i64_value(committed_at.get())?)
   1209                 .bind(i64_value(committed_at.get())?)
   1210                 .execute(&mut *transaction)
   1211                 .await
   1212                 .map_err(|_| PresenceOperationError::Storage)?;
   1213             if inserted.rows_affected() != 1 {
   1214                 return Err(PresenceOperationError::Invariant);
   1215             }
   1216         }
   1217     }
   1218     let committed = read_generation_outboxes(transaction, desired.generation).await?;
   1219     if !generation_matches(&documents, &committed)? {
   1220         return Err(PresenceOperationError::Invariant);
   1221     }
   1222     Ok(RhiPresenceCommitOutcome {
   1223         desired_state: documents.desired_state,
   1224         document_count: u8::try_from(documents.documents.len())
   1225             .map_err(|_| PresenceOperationError::Invariant)?,
   1226         changed: !documents.documents.is_empty() || superseded != 0,
   1227     })
   1228 }
   1229 
   1230 fn validate_document_set(
   1231     documents: &RhiSignedPresenceDocuments,
   1232     desired: PresenceDesiredBinding,
   1233 ) -> Result<(), PresenceOperationError> {
   1234     let state = documents.desired_state;
   1235     if state.generation() != desired.generation
   1236         || state.mode() != desired.mode
   1237         || state.profile() != desired.profile
   1238         || state.application_handler() != desired.application_handler
   1239         || state.target_set_sha256() != &desired.target_set_sha256
   1240         || state.target_count() != desired.target_count
   1241         || state.required_target_count() != desired.required_target_count
   1242         || state.queue_capacity() != desired.queue_capacity
   1243         || state.desired_sha256() != &desired.desired_sha256
   1244         || documents.targets.len() != usize::from(desired.target_count)
   1245         || documents.documents.len() != desired.document_count()
   1246         || target_set_digest(&documents.targets)? != desired.target_set_sha256
   1247         || !valid_public_key(&documents.service_public_key)
   1248     {
   1249         return Err(PresenceOperationError::DesiredStateMismatch);
   1250     }
   1251     let expected = desired.document_kinds();
   1252     for (document, kind) in documents.documents.iter().zip(expected) {
   1253         if document.kind != kind
   1254             || document.signed_event_bytes.is_empty()
   1255             || document.signed_event_bytes.len() > RHI_PRESENCE_SIGNED_EVENT_MAX_BYTES
   1256             || <[u8; 32]>::from(Sha256::digest(&document.signed_event_bytes))
   1257                 != document.signed_event_sha256
   1258         {
   1259             return Err(PresenceOperationError::VerificationFailed);
   1260         }
   1261         let event = validate_signed_document(
   1262             document.kind,
   1263             &documents.service_public_key,
   1264             document.created_at_unix_seconds,
   1265             &document.signed_event_bytes,
   1266         )
   1267         .map_err(|_| PresenceOperationError::VerificationFailed)?;
   1268         if event.id().as_bytes() != &document.event_id {
   1269             return Err(PresenceOperationError::VerificationFailed);
   1270         }
   1271     }
   1272     Ok(())
   1273 }
   1274 
   1275 fn generation_matches(
   1276     documents: &RhiSignedPresenceDocuments,
   1277     current: &[PresenceOutboxRecord],
   1278 ) -> Result<bool, PresenceOperationError> {
   1279     if documents.documents.len() != current.len() {
   1280         return Ok(false);
   1281     }
   1282     for (document, outbox) in documents.documents.iter().zip(current) {
   1283         if outbox.desired_generation != documents.desired_state.generation()
   1284             || outbox.desired_sha256 != *documents.desired_state.desired_sha256()
   1285             || outbox.target_set_sha256 != *documents.desired_state.target_set_sha256()
   1286             || outbox.target_count != documents.desired_state.target_count()
   1287             || outbox.required_target_count != documents.desired_state.required_target_count()
   1288             || outbox.id
   1289                 != derive_outbox_id(
   1290                     outbox.desired_generation,
   1291                     document.kind,
   1292                     outbox.desired_sha256,
   1293                     document.signed_event_sha256,
   1294                 )
   1295             || document.kind != outbox.document_kind
   1296             || document.event_id != outbox.event_id
   1297             || document.signed_event_sha256 != outbox.event_sha256
   1298             || document.signed_event_bytes.as_ref() != outbox.exact_signed_event_bytes.as_ref()
   1299             || document.created_at_unix_seconds != outbox.authored_at_unix_s
   1300             || documents.service_public_key.as_ref() != outbox.service_public_key.as_ref()
   1301             || documents
   1302                 .targets
   1303                 .iter()
   1304                 .zip(outbox.targets.iter())
   1305                 .any(|(expected, actual)| {
   1306                     expected.ordinal() != actual.ordinal
   1307                         || expected.relay_id() != actual.relay_id.as_ref()
   1308                         || expected.required() != actual.required
   1309                 })
   1310         {
   1311             return Ok(false);
   1312         }
   1313         validate_targets(outbox, &outbox.targets)?;
   1314     }
   1315     Ok(true)
   1316 }
   1317 
   1318 fn derive_outbox_id(
   1319     generation: u64,
   1320     kind: RhiPresenceDocumentKind,
   1321     desired_sha256: [u8; 32],
   1322     event_sha256: [u8; 32],
   1323 ) -> RhiPresenceOutboxId {
   1324     let mut digest = Sha256::new();
   1325     digest.update(OUTBOX_ID_DOMAIN);
   1326     digest.update(generation.to_be_bytes());
   1327     digest.update([document_kind_tag(kind)]);
   1328     digest.update(desired_sha256);
   1329     digest.update(event_sha256);
   1330     RhiPresenceOutboxId(digest.finalize().into())
   1331 }
   1332 
   1333 fn derive_attempt_id(
   1334     outbox: RhiPresenceOutboxId,
   1335     event_sha256: [u8; 32],
   1336     ordinal: u8,
   1337     attempt_number: u16,
   1338 ) -> RhiPresenceAttemptId {
   1339     let mut digest = Sha256::new();
   1340     digest.update(ATTEMPT_ID_DOMAIN);
   1341     digest.update(outbox.0);
   1342     digest.update(event_sha256);
   1343     digest.update([ordinal]);
   1344     digest.update(attempt_number.to_be_bytes());
   1345     RhiPresenceAttemptId(digest.finalize().into())
   1346 }
   1347 
   1348 const fn document_kind_tag(kind: RhiPresenceDocumentKind) -> u8 {
   1349     match kind {
   1350         RhiPresenceDocumentKind::ServiceProfile => 0,
   1351         RhiPresenceDocumentKind::ApplicationHandler => 1,
   1352     }
   1353 }
   1354 
   1355 fn target_set_digest(targets: &[RhiPresenceTarget]) -> Result<[u8; 32], PresenceOperationError> {
   1356     const DOMAIN: &[u8] = b"radroots.rhi.presence_target_set.v1\0";
   1357     let mut digest = Sha256::new();
   1358     digest.update(DOMAIN);
   1359     digest.update(
   1360         u32::try_from(targets.len())
   1361             .map_err(|_| PresenceOperationError::InvalidInput)?
   1362             .to_be_bytes(),
   1363     );
   1364     for target in targets {
   1365         digest.update(u32::from(target.ordinal()).to_be_bytes());
   1366         digest.update(
   1367             u64::try_from(target.relay_id().len())
   1368                 .map_err(|_| PresenceOperationError::InvalidInput)?
   1369                 .to_be_bytes(),
   1370         );
   1371         digest.update(target.relay_id().as_bytes());
   1372         digest.update([u8::from(target.required())]);
   1373     }
   1374     Ok(digest.finalize().into())
   1375 }
   1376 
   1377 async fn read_desired_binding(
   1378     transaction: &mut ServiceSqliteTransaction<'_>,
   1379 ) -> Result<PresenceDesiredBinding, PresenceOperationError> {
   1380     let rows = sqlx::query(READ_DESIRED_BINDING_SQL)
   1381         .fetch_all(&mut *transaction)
   1382         .await
   1383         .map_err(|_| PresenceOperationError::Storage)?;
   1384     let [row] = rows.as_slice() else {
   1385         return Err(PresenceOperationError::DesiredStateMismatch);
   1386     };
   1387     if row.try_get::<i64, _>("singleton").ok() != Some(1) {
   1388         return Err(PresenceOperationError::Invariant);
   1389     }
   1390     let generation = positive_u64(row, "generation")?;
   1391     let enabled = boolean_i64(row, "enabled")?;
   1392     let profile = boolean_i64(row, "profile")?;
   1393     let application_handler = boolean_i64(row, "application_handler")?;
   1394     let target_set_sha256 = blob::<32>(row, "target_set_sha256")?;
   1395     let target_count = bounded_u8(row, "target_count", RHI_PRESENCE_DESIRED_MAX_TARGETS)?;
   1396     let required_target_count =
   1397         bounded_u8(row, "required_target_count", usize::from(target_count))?;
   1398     let queue_capacity = row
   1399         .try_get::<i64, _>("queue_capacity")
   1400         .ok()
   1401         .and_then(|value| u32::try_from(value).ok())
   1402         .filter(|value| *value <= 4_096)
   1403         .ok_or(PresenceOperationError::Invariant)?;
   1404     let desired_sha256 = blob::<32>(row, "desired_sha256")?;
   1405     let mode = if enabled {
   1406         RhiPresenceDesiredMode::Enabled
   1407     } else {
   1408         RhiPresenceDesiredMode::Disabled
   1409     };
   1410     let valid = if enabled {
   1411         (profile || application_handler) && target_count > 0 && queue_capacity > 0
   1412     } else {
   1413         !profile
   1414             && !application_handler
   1415             && target_count == 0
   1416             && required_target_count == 0
   1417             && queue_capacity == 0
   1418     };
   1419     if !valid {
   1420         return Err(PresenceOperationError::Invariant);
   1421     }
   1422     Ok(PresenceDesiredBinding {
   1423         generation,
   1424         mode,
   1425         profile,
   1426         application_handler,
   1427         target_set_sha256,
   1428         target_count,
   1429         required_target_count,
   1430         queue_capacity,
   1431         desired_sha256,
   1432     })
   1433 }
   1434 
   1435 async fn read_generation_outboxes(
   1436     transaction: &mut ServiceSqliteTransaction<'_>,
   1437     generation: u64,
   1438 ) -> Result<Vec<PresenceOutboxRecord>, PresenceOperationError> {
   1439     let rows = sqlx::query(READ_GENERATION_OUTBOX_SQL)
   1440         .bind(i64_value(generation)?)
   1441         .fetch_all(&mut *transaction)
   1442         .await
   1443         .map_err(|_| PresenceOperationError::Storage)?;
   1444     if rows.len() > 2 {
   1445         return Err(PresenceOperationError::Invariant);
   1446     }
   1447     let mut records = Vec::with_capacity(rows.len());
   1448     for row in rows {
   1449         let mut outbox = decode_outbox(row)?;
   1450         outbox.targets = read_targets(transaction, outbox.id)
   1451             .await?
   1452             .into_boxed_slice();
   1453         records.push(outbox);
   1454     }
   1455     Ok(records)
   1456 }
   1457 
   1458 async fn read_outbox(
   1459     transaction: &mut ServiceSqliteTransaction<'_>,
   1460     id: RhiPresenceOutboxId,
   1461 ) -> Result<Option<PresenceOutboxRecord>, PresenceOperationError> {
   1462     let rows = sqlx::query(READ_OUTBOX_SQL)
   1463         .bind(id.as_bytes().as_slice())
   1464         .fetch_all(&mut *transaction)
   1465         .await
   1466         .map_err(|_| PresenceOperationError::Storage)?;
   1467     let Some(row) = exactly_zero_or_one(rows)? else {
   1468         return Ok(None);
   1469     };
   1470     let mut outbox = decode_outbox(row)?;
   1471     outbox.targets = read_targets(transaction, outbox.id)
   1472         .await?
   1473         .into_boxed_slice();
   1474     Ok(Some(outbox))
   1475 }
   1476 
   1477 async fn read_targets(
   1478     transaction: &mut ServiceSqliteTransaction<'_>,
   1479     outbox_id: RhiPresenceOutboxId,
   1480 ) -> Result<Vec<PresenceTargetRecord>, PresenceOperationError> {
   1481     let rows = sqlx::query(READ_TARGETS_SQL)
   1482         .bind(outbox_id.as_bytes().as_slice())
   1483         .fetch_all(&mut *transaction)
   1484         .await
   1485         .map_err(|_| PresenceOperationError::Storage)?;
   1486     if rows.len() > RHI_PRESENCE_DESIRED_MAX_TARGETS {
   1487         return Err(PresenceOperationError::Invariant);
   1488     }
   1489     rows.into_iter().map(decode_target).collect()
   1490 }
   1491 
   1492 fn decode_outbox(
   1493     row: sqlx::sqlite::SqliteRow,
   1494 ) -> Result<PresenceOutboxRecord, PresenceOperationError> {
   1495     let id = RhiPresenceOutboxId(blob::<32>(&row, "outbox_id")?);
   1496     let desired_generation = positive_u64(&row, "desired_generation")?;
   1497     let document_kind = decode_document_kind(&bounded_text(
   1498         &row,
   1499         "document_kind",
   1500         "document_kind_bytes",
   1501         19,
   1502     )?)?;
   1503     let desired_sha256 = blob::<32>(&row, "desired_sha256")?;
   1504     let target_set_sha256 = blob::<32>(&row, "target_set_sha256")?;
   1505     let event_id = blob::<32>(&row, "event_id")?;
   1506     let event_sha256 = blob::<32>(&row, "event_sha256")?;
   1507     let exact_length = positive_usize(&row, "exact_signed_event_bytes_length")?;
   1508     if exact_length > RHI_PRESENCE_SIGNED_EVENT_MAX_BYTES {
   1509         return Err(PresenceOperationError::Invariant);
   1510     }
   1511     let exact_signed_event_bytes = row
   1512         .try_get::<Vec<u8>, _>("exact_signed_event_bytes")
   1513         .map_err(|_| PresenceOperationError::Invariant)?;
   1514     let authored_at_unix_s = nonnegative_u64(&row, "authored_at_unix_s")?;
   1515     let service_public_key =
   1516         bounded_text(&row, "service_public_key", "service_public_key_bytes", 64)?;
   1517     let target_count = bounded_u8(&row, "target_count", RHI_PRESENCE_DESIRED_MAX_TARGETS)?;
   1518     let required_target_count =
   1519         bounded_u8(&row, "required_target_count", usize::from(target_count))?;
   1520     let max_attempts = bounded_u16(&row, "max_attempts", RHI_PRESENCE_MAX_ATTEMPTS)?;
   1521     let initial_backoff_ms = bounded_u64(&row, "initial_backoff_ms", INITIAL_BACKOFF_MILLISECONDS)?;
   1522     let maximum_backoff_ms = bounded_u64(&row, "maximum_backoff_ms", MAXIMUM_BACKOFF_MILLISECONDS)?;
   1523     let attempt_deadline_ms =
   1524         bounded_u64(&row, "attempt_deadline_ms", ATTEMPT_DEADLINE_MILLISECONDS)?;
   1525     let state = decode_outbox_state(&bounded_text(&row, "state", "state_bytes", 10)?)?;
   1526     let revision = positive_u64(&row, "revision")?;
   1527     let next_attempt = optional_millis(&row, "next_attempt_unix_ms")?;
   1528     let lease_owner = optional_owner(&row)?;
   1529     let lease_expires = optional_millis(&row, "lease_expires_unix_ms")?;
   1530     let created_at = RhiPresenceUnixMilliseconds(nonnegative_u64(&row, "created_at_unix_ms")?);
   1531     let updated_at = RhiPresenceUnixMilliseconds(nonnegative_u64(&row, "updated_at_unix_ms")?);
   1532     if exact_signed_event_bytes.len() != exact_length
   1533         || exact_signed_event_bytes.is_empty()
   1534         || <[u8; 32]>::from(Sha256::digest(&exact_signed_event_bytes)) != event_sha256
   1535         || !valid_public_key(&service_public_key)
   1536         || max_attempts != RHI_PRESENCE_MAX_ATTEMPTS
   1537         || initial_backoff_ms != INITIAL_BACKOFF_MILLISECONDS
   1538         || maximum_backoff_ms != MAXIMUM_BACKOFF_MILLISECONDS
   1539         || attempt_deadline_ms != ATTEMPT_DEADLINE_MILLISECONDS
   1540         || initial_backoff_ms > maximum_backoff_ms
   1541         || updated_at < created_at
   1542         || !valid_outbox_shape(state, next_attempt, lease_owner, lease_expires)
   1543     {
   1544         return Err(PresenceOperationError::Invariant);
   1545     }
   1546     Ok(PresenceOutboxRecord {
   1547         id,
   1548         desired_generation,
   1549         document_kind,
   1550         desired_sha256,
   1551         target_set_sha256,
   1552         event_id,
   1553         event_sha256,
   1554         exact_signed_event_bytes: exact_signed_event_bytes.into_boxed_slice(),
   1555         authored_at_unix_s,
   1556         service_public_key: service_public_key.into_boxed_str(),
   1557         target_count,
   1558         required_target_count,
   1559         max_attempts,
   1560         state,
   1561         revision,
   1562         next_attempt,
   1563         lease_owner,
   1564         lease_expires,
   1565         created_at,
   1566         updated_at,
   1567         targets: Box::new([]),
   1568     })
   1569 }
   1570 
   1571 fn decode_target(
   1572     row: sqlx::sqlite::SqliteRow,
   1573 ) -> Result<PresenceTargetRecord, PresenceOperationError> {
   1574     let ordinal = bounded_u8(&row, "target_ordinal", RHI_PRESENCE_DESIRED_MAX_TARGETS - 1)?;
   1575     let relay_id = bounded_text(&row, "relay_id", "relay_id_bytes", 64)?;
   1576     let required = boolean_i64(&row, "required")?;
   1577     let state = decode_target_state(&bounded_text(&row, "state", "state_bytes", 13)?)?;
   1578     let revision = positive_u64(&row, "revision")?;
   1579     let attempt_count = bounded_u16(&row, "attempt_count", RHI_PRESENCE_MAX_ATTEMPTS)?;
   1580     let next_attempt = optional_millis(&row, "next_attempt_unix_ms")?;
   1581     let last_attempt_id = optional_attempt_id(&row)?;
   1582     let updated_at = RhiPresenceUnixMilliseconds(nonnegative_u64(&row, "updated_at_unix_ms")?);
   1583     if !valid_relay_id(&relay_id)
   1584         || !valid_target_shape(state, attempt_count, next_attempt, last_attempt_id)
   1585     {
   1586         return Err(PresenceOperationError::Invariant);
   1587     }
   1588     Ok(PresenceTargetRecord {
   1589         ordinal,
   1590         relay_id: relay_id.into_boxed_str(),
   1591         required,
   1592         state,
   1593         revision,
   1594         attempt_count,
   1595         next_attempt,
   1596         last_attempt_id,
   1597         updated_at,
   1598     })
   1599 }
   1600 
   1601 async fn claim_next(
   1602     transaction: &mut ServiceSqliteTransaction<'_>,
   1603     owner: RhiPresenceLeaseOwner,
   1604     now: RhiPresenceUnixMilliseconds,
   1605 ) -> Result<Option<RhiPresenceLease>, PresenceOperationError> {
   1606     let rows = sqlx::query(READ_CLAIMABLE_OUTBOX_SQL)
   1607         .bind(i64_value(now.get())?)
   1608         .fetch_all(&mut *transaction)
   1609         .await
   1610         .map_err(|_| PresenceOperationError::Storage)?;
   1611     let Some(row) = exactly_zero_or_one(rows)? else {
   1612         return Ok(None);
   1613     };
   1614     let id = RhiPresenceOutboxId(blob::<32>(&row, "outbox_id")?);
   1615     let outbox = read_outbox(transaction, id)
   1616         .await?
   1617         .ok_or(PresenceOperationError::Invariant)?;
   1618     validate_current_desired(transaction, &outbox).await?;
   1619     validate_targets(&outbox, &outbox.targets)?;
   1620     let expires = now
   1621         .get()
   1622         .checked_add(ATTEMPT_DEADLINE_MILLISECONDS)
   1623         .filter(|value| *value <= MAX_UNIX_MILLISECONDS)
   1624         .ok_or(PresenceOperationError::InvalidInput)?;
   1625     let result = sqlx::query(CLAIM_OUTBOX_SQL)
   1626         .bind(owner.0.as_slice())
   1627         .bind(i64_value(expires)?)
   1628         .bind(i64_value(now.get())?)
   1629         .bind(outbox.id.as_bytes().as_slice())
   1630         .bind(i64_value(outbox.revision)?)
   1631         .bind(i64_value(now.get())?)
   1632         .bind(i64_value(now.get())?)
   1633         .execute(&mut *transaction)
   1634         .await
   1635         .map_err(|_| PresenceOperationError::Storage)?;
   1636     if result.rows_affected() != 1 {
   1637         return Err(PresenceOperationError::LeaseLost);
   1638     }
   1639     let leased = read_outbox(transaction, id)
   1640         .await?
   1641         .filter(|record| {
   1642             record.state == RhiPresenceOutboxState::Leased
   1643                 && record.lease_owner == Some(owner)
   1644                 && record.lease_expires == Some(RhiPresenceUnixMilliseconds(expires))
   1645         })
   1646         .ok_or(PresenceOperationError::Invariant)?;
   1647     Ok(Some(RhiPresenceLease {
   1648         outbox: leased,
   1649         owner,
   1650         expires_at: RhiPresenceUnixMilliseconds(expires),
   1651     }))
   1652 }
   1653 
   1654 async fn prepare_next(
   1655     transaction: &mut ServiceSqliteTransaction<'_>,
   1656     lease: RhiPresenceLease,
   1657     started_at: RhiPresenceUnixMilliseconds,
   1658 ) -> Result<RhiPreparedPresenceAttempt, PresenceOperationError> {
   1659     validate_lease(transaction, &lease, started_at, true).await?;
   1660     let candidate = lease
   1661         .outbox
   1662         .targets
   1663         .iter()
   1664         .find(|target| target_is_due(target, lease.outbox.max_attempts, started_at))
   1665         .ok_or(PresenceOperationError::NotReady)?;
   1666     let candidate_ordinal = candidate.ordinal;
   1667     let candidate_revision = candidate.revision;
   1668     let attempt_number = candidate
   1669         .attempt_count
   1670         .checked_add(1)
   1671         .filter(|value| *value <= lease.outbox.max_attempts)
   1672         .ok_or(PresenceOperationError::Invariant)?;
   1673     let attempt_id = derive_attempt_id(
   1674         lease.outbox.id,
   1675         lease.outbox.event_sha256,
   1676         candidate_ordinal,
   1677         attempt_number,
   1678     );
   1679     let result = sqlx::query(PREPARE_TARGET_SQL)
   1680         .bind(attempt_id.as_bytes().as_slice())
   1681         .bind(i64_value(started_at.get())?)
   1682         .bind(lease.outbox.id.as_bytes().as_slice())
   1683         .bind(i64::from(candidate_ordinal))
   1684         .bind(i64_value(candidate_revision)?)
   1685         .bind(i64::from(lease.outbox.max_attempts))
   1686         .bind(i64_value(started_at.get())?)
   1687         .bind(i64_value(started_at.get())?)
   1688         .execute(&mut *transaction)
   1689         .await
   1690         .map_err(|_| PresenceOperationError::Storage)?;
   1691     if result.rows_affected() != 1 {
   1692         return Err(PresenceOperationError::LeaseLost);
   1693     }
   1694     let current = read_outbox(transaction, lease.outbox.id)
   1695         .await?
   1696         .ok_or(PresenceOperationError::Invariant)?;
   1697     if !same_outbox_identity(&current, &lease.outbox) {
   1698         return Err(PresenceOperationError::LeaseLost);
   1699     }
   1700     let target = current
   1701         .targets
   1702         .into_vec()
   1703         .into_iter()
   1704         .find(|target| target.ordinal == candidate_ordinal)
   1705         .filter(|target| {
   1706             target.state == RhiPresenceTargetState::Submitted
   1707                 && target.attempt_count == attempt_number
   1708                 && target.last_attempt_id == Some(attempt_id)
   1709                 && target.updated_at == started_at
   1710         })
   1711         .ok_or(PresenceOperationError::Invariant)?;
   1712     validate_signed_document(
   1713         lease.outbox.document_kind,
   1714         &lease.outbox.service_public_key,
   1715         lease.outbox.authored_at_unix_s,
   1716         &lease.outbox.exact_signed_event_bytes,
   1717     )
   1718     .map_err(|_| PresenceOperationError::Invariant)?;
   1719     let deadline_at = lease.expires_at;
   1720     let exact_signed_event_bytes = lease.outbox.exact_signed_event_bytes.clone();
   1721     Ok(RhiPreparedPresenceAttempt {
   1722         lease,
   1723         target,
   1724         attempt_id,
   1725         started_at,
   1726         deadline_at,
   1727         exact_signed_event_bytes,
   1728     })
   1729 }
   1730 
   1731 async fn record_outcome(
   1732     transaction: &mut ServiceSqliteTransaction<'_>,
   1733     input: PresenceRecordInput,
   1734 ) -> Result<RhiPresenceAttemptCommit, PresenceOperationError> {
   1735     if let Some(existing) = read_attempt(transaction, input.attempt_id).await? {
   1736         return reconcile_recorded(transaction, input, existing).await;
   1737     }
   1738     validate_input_lease(transaction, input, true).await?;
   1739     let outbox = read_outbox(transaction, input.outbox_id)
   1740         .await?
   1741         .ok_or(PresenceOperationError::Invariant)?;
   1742     let target = outbox
   1743         .targets
   1744         .iter()
   1745         .find(|target| target.ordinal == input.target_ordinal)
   1746         .ok_or(PresenceOperationError::Invariant)?;
   1747     if target.state != RhiPresenceTargetState::Submitted
   1748         || target.revision != input.target_revision
   1749         || target.attempt_count != input.attempt_number
   1750         || target.last_attempt_id != Some(input.attempt_id)
   1751         || target.updated_at != input.started_at
   1752     {
   1753         return Err(PresenceOperationError::LeaseLost);
   1754     }
   1755     insert_attempt(transaction, input).await?;
   1756     let (next_state, next_attempt) = target_outcome_schedule(input)?;
   1757     let result = sqlx::query(UPDATE_TARGET_OUTCOME_SQL)
   1758         .bind(next_state.code())
   1759         .bind(optional_i64(next_attempt)?)
   1760         .bind(i64_value(input.finished_at.get())?)
   1761         .bind(input.outbox_id.as_bytes().as_slice())
   1762         .bind(i64::from(input.target_ordinal))
   1763         .bind(i64_value(input.target_revision)?)
   1764         .bind(i64::from(input.attempt_number))
   1765         .bind(input.attempt_id.as_bytes().as_slice())
   1766         .execute(&mut *transaction)
   1767         .await
   1768         .map_err(|_| PresenceOperationError::Storage)?;
   1769     if result.rows_affected() != 1 {
   1770         return Err(PresenceOperationError::LeaseLost);
   1771     }
   1772     let targets = read_targets(transaction, input.outbox_id).await?;
   1773     let disposition = disposition(&outbox, &targets)?;
   1774     update_outbox_after_attempt(transaction, input, disposition).await?;
   1775     Ok(RhiPresenceAttemptCommit {
   1776         outbox_id: input.outbox_id,
   1777         attempt_id: input.attempt_id,
   1778         target_ordinal: input.target_ordinal,
   1779         attempt_number: input.attempt_number,
   1780         outcome: input.outcome,
   1781         target_state: next_state,
   1782         outbox_state: disposition.state,
   1783     })
   1784 }
   1785 
   1786 async fn read_expired_candidate(
   1787     transaction: &mut ServiceSqliteTransaction<'_>,
   1788     now: RhiPresenceUnixMilliseconds,
   1789 ) -> Result<Option<PresenceRecoveryCandidate>, PresenceOperationError> {
   1790     let rows = sqlx::query(READ_EXPIRED_OUTBOX_SQL)
   1791         .bind(i64_value(now.get())?)
   1792         .fetch_all(&mut *transaction)
   1793         .await
   1794         .map_err(|_| PresenceOperationError::Storage)?;
   1795     let Some(row) = exactly_zero_or_one(rows)? else {
   1796         return Ok(None);
   1797     };
   1798     let id = RhiPresenceOutboxId(blob::<32>(&row, "outbox_id")?);
   1799     let outbox = read_outbox(transaction, id)
   1800         .await?
   1801         .ok_or(PresenceOperationError::Invariant)?;
   1802     let desired = read_desired_binding(transaction).await?;
   1803     let current_desired = desired_matches_outbox(desired, &outbox);
   1804     validate_targets(&outbox, &outbox.targets)?;
   1805     let submitted: Vec<_> = outbox
   1806         .targets
   1807         .iter()
   1808         .filter(|target| target.state == RhiPresenceTargetState::Submitted)
   1809         .map(|target| target.attempt_count)
   1810         .collect();
   1811     if submitted.len() > 1 {
   1812         return Err(PresenceOperationError::Invariant);
   1813     }
   1814     let retry_upper_bound = if current_desired {
   1815         submitted
   1816             .first()
   1817             .copied()
   1818             .filter(|attempt| *attempt < outbox.max_attempts)
   1819             .map_or(0, retry_upper_bound)
   1820     } else {
   1821         0
   1822     };
   1823     Ok(Some(PresenceRecoveryCandidate {
   1824         outbox,
   1825         retry_upper_bound,
   1826         current_desired,
   1827     }))
   1828 }
   1829 
   1830 async fn recover_expired(
   1831     transaction: &mut ServiceSqliteTransaction<'_>,
   1832     candidate: PresenceRecoveryCandidate,
   1833     now: RhiPresenceUnixMilliseconds,
   1834     delay: RhiPresenceRetryDelayMilliseconds,
   1835 ) -> Result<bool, PresenceOperationError> {
   1836     let current = read_outbox(transaction, candidate.outbox.id)
   1837         .await?
   1838         .filter(|current| same_outbox(current, &candidate.outbox))
   1839         .ok_or(PresenceOperationError::LeaseLost)?;
   1840     if current.state != RhiPresenceOutboxState::Leased
   1841         || current.lease_expires.is_none_or(|expires| expires > now)
   1842     {
   1843         return Err(PresenceOperationError::LeaseLost);
   1844     }
   1845     let owner = current
   1846         .lease_owner
   1847         .ok_or(PresenceOperationError::Invariant)?;
   1848     let expires = current
   1849         .lease_expires
   1850         .ok_or(PresenceOperationError::Invariant)?;
   1851     for target in current
   1852         .targets
   1853         .iter()
   1854         .filter(|target| target.state == RhiPresenceTargetState::Submitted)
   1855     {
   1856         let attempt_id = target
   1857             .last_attempt_id
   1858             .ok_or(PresenceOperationError::Invariant)?;
   1859         if read_attempt(transaction, attempt_id).await?.is_some() {
   1860             return Err(PresenceOperationError::Invariant);
   1861         }
   1862         let input = PresenceRecordInput {
   1863             outbox_id: current.id,
   1864             outbox_revision: current.revision,
   1865             event_sha256: current.event_sha256,
   1866             owner,
   1867             lease_expires: expires,
   1868             target_ordinal: target.ordinal,
   1869             target_revision: target.revision,
   1870             attempt_number: target.attempt_count,
   1871             attempt_id,
   1872             started_at: target.updated_at,
   1873             finished_at: now,
   1874             outcome: RhiPresenceAttemptOutcome::Unknown,
   1875             retry_delay: delay,
   1876         };
   1877         insert_attempt(transaction, input).await?;
   1878         let (state, next) = target_outcome_schedule(input)?;
   1879         let result = sqlx::query(UPDATE_TARGET_OUTCOME_SQL)
   1880             .bind(state.code())
   1881             .bind(optional_i64(next)?)
   1882             .bind(i64_value(now.get())?)
   1883             .bind(current.id.as_bytes().as_slice())
   1884             .bind(i64::from(target.ordinal))
   1885             .bind(i64_value(target.revision)?)
   1886             .bind(i64::from(target.attempt_count))
   1887             .bind(attempt_id.as_bytes().as_slice())
   1888             .execute(&mut *transaction)
   1889             .await
   1890             .map_err(|_| PresenceOperationError::Storage)?;
   1891         if result.rows_affected() != 1 {
   1892             return Err(PresenceOperationError::LeaseLost);
   1893         }
   1894     }
   1895     let targets = read_targets(transaction, current.id).await?;
   1896     let disposition = if candidate.current_desired {
   1897         disposition(&current, &targets)?
   1898     } else {
   1899         PresenceOutboxDisposition {
   1900             state: RhiPresenceOutboxState::Superseded,
   1901             next_attempt: None,
   1902         }
   1903     };
   1904     let result = sqlx::query(UPDATE_OUTBOX_RECOVERY_SQL)
   1905         .bind(disposition.state.code())
   1906         .bind(optional_i64(disposition.next_attempt)?)
   1907         .bind(i64_value(now.get())?)
   1908         .bind(current.id.as_bytes().as_slice())
   1909         .bind(i64_value(current.revision)?)
   1910         .bind(owner.0.as_slice())
   1911         .bind(i64_value(expires.get())?)
   1912         .bind(i64_value(now.get())?)
   1913         .bind(i64_value(now.get())?)
   1914         .execute(&mut *transaction)
   1915         .await
   1916         .map_err(|_| PresenceOperationError::Storage)?;
   1917     if result.rows_affected() != 1 {
   1918         return Err(PresenceOperationError::LeaseLost);
   1919     }
   1920     Ok(true)
   1921 }
   1922 
   1923 async fn validate_current_desired(
   1924     transaction: &mut ServiceSqliteTransaction<'_>,
   1925     outbox: &PresenceOutboxRecord,
   1926 ) -> Result<(), PresenceOperationError> {
   1927     let desired = read_desired_binding(transaction).await?;
   1928     if !desired_matches_outbox(desired, outbox) {
   1929         return Err(PresenceOperationError::DesiredStateMismatch);
   1930     }
   1931     Ok(())
   1932 }
   1933 
   1934 fn desired_matches_outbox(desired: PresenceDesiredBinding, outbox: &PresenceOutboxRecord) -> bool {
   1935     desired.mode == RhiPresenceDesiredMode::Enabled
   1936         && desired.generation == outbox.desired_generation
   1937         && desired.desired_sha256 == outbox.desired_sha256
   1938         && desired.target_set_sha256 == outbox.target_set_sha256
   1939         && desired.target_count == outbox.target_count
   1940         && desired.required_target_count == outbox.required_target_count
   1941         && desired
   1942             .document_kinds()
   1943             .any(|kind| kind == outbox.document_kind)
   1944 }
   1945 
   1946 async fn validate_lease(
   1947     transaction: &mut ServiceSqliteTransaction<'_>,
   1948     lease: &RhiPresenceLease,
   1949     now: RhiPresenceUnixMilliseconds,
   1950     require_unexpired: bool,
   1951 ) -> Result<(), PresenceOperationError> {
   1952     let current = read_outbox(transaction, lease.outbox.id)
   1953         .await?
   1954         .filter(|current| same_outbox(current, &lease.outbox))
   1955         .ok_or(PresenceOperationError::LeaseLost)?;
   1956     validate_current_desired(transaction, &current).await?;
   1957     if current.state != RhiPresenceOutboxState::Leased
   1958         || current.lease_owner != Some(lease.owner)
   1959         || current.lease_expires != Some(lease.expires_at)
   1960         || (require_unexpired && lease.expires_at <= now)
   1961     {
   1962         return Err(PresenceOperationError::LeaseLost);
   1963     }
   1964     Ok(())
   1965 }
   1966 
   1967 async fn validate_input_lease(
   1968     transaction: &mut ServiceSqliteTransaction<'_>,
   1969     input: PresenceRecordInput,
   1970     require_unexpired: bool,
   1971 ) -> Result<(), PresenceOperationError> {
   1972     let current = read_outbox(transaction, input.outbox_id)
   1973         .await?
   1974         .ok_or(PresenceOperationError::LeaseLost)?;
   1975     validate_current_desired(transaction, &current).await?;
   1976     if current.revision != input.outbox_revision
   1977         || current.event_sha256 != input.event_sha256
   1978         || current.state != RhiPresenceOutboxState::Leased
   1979         || current.lease_owner != Some(input.owner)
   1980         || current.lease_expires != Some(input.lease_expires)
   1981         || (require_unexpired && input.lease_expires <= input.finished_at)
   1982     {
   1983         return Err(PresenceOperationError::LeaseLost);
   1984     }
   1985     Ok(())
   1986 }
   1987 
   1988 async fn insert_attempt(
   1989     transaction: &mut ServiceSqliteTransaction<'_>,
   1990     input: PresenceRecordInput,
   1991 ) -> Result<(), PresenceOperationError> {
   1992     let result = sqlx::query(INSERT_ATTEMPT_SQL)
   1993         .bind(input.attempt_id.as_bytes().as_slice())
   1994         .bind(input.outbox_id.as_bytes().as_slice())
   1995         .bind(i64::from(input.target_ordinal))
   1996         .bind(i64::from(input.attempt_number))
   1997         .bind(input.event_sha256.as_slice())
   1998         .bind(input.owner.0.as_slice())
   1999         .bind(i64_value(input.started_at.get())?)
   2000         .bind(i64_value(input.finished_at.get())?)
   2001         .bind(input.outcome.code())
   2002         .bind(input.outcome.code())
   2003         .execute(&mut *transaction)
   2004         .await
   2005         .map_err(|_| PresenceOperationError::Storage)?;
   2006     if result.rows_affected() != 1 {
   2007         return Err(PresenceOperationError::Storage);
   2008     }
   2009     Ok(())
   2010 }
   2011 
   2012 async fn update_outbox_after_attempt(
   2013     transaction: &mut ServiceSqliteTransaction<'_>,
   2014     input: PresenceRecordInput,
   2015     disposition: PresenceOutboxDisposition,
   2016 ) -> Result<(), PresenceOperationError> {
   2017     let result = sqlx::query(UPDATE_OUTBOX_AFTER_ATTEMPT_SQL)
   2018         .bind(disposition.state.code())
   2019         .bind(optional_i64(disposition.next_attempt)?)
   2020         .bind(i64_value(input.finished_at.get())?)
   2021         .bind(input.outbox_id.as_bytes().as_slice())
   2022         .bind(i64_value(input.outbox_revision)?)
   2023         .bind(input.owner.0.as_slice())
   2024         .bind(i64_value(input.lease_expires.get())?)
   2025         .bind(i64_value(input.finished_at.get())?)
   2026         .bind(i64_value(input.finished_at.get())?)
   2027         .execute(&mut *transaction)
   2028         .await
   2029         .map_err(|_| PresenceOperationError::Storage)?;
   2030     if result.rows_affected() != 1 {
   2031         return Err(PresenceOperationError::LeaseLost);
   2032     }
   2033     Ok(())
   2034 }
   2035 
   2036 async fn read_attempt(
   2037     transaction: &mut ServiceSqliteTransaction<'_>,
   2038     attempt_id: RhiPresenceAttemptId,
   2039 ) -> Result<Option<PresenceAttemptRecord>, PresenceOperationError> {
   2040     let rows = sqlx::query(READ_ATTEMPT_SQL)
   2041         .bind(attempt_id.as_bytes().as_slice())
   2042         .fetch_all(&mut *transaction)
   2043         .await
   2044         .map_err(|_| PresenceOperationError::Storage)?;
   2045     exactly_zero_or_one(rows)?.map(decode_attempt).transpose()
   2046 }
   2047 
   2048 fn decode_attempt(
   2049     row: sqlx::sqlite::SqliteRow,
   2050 ) -> Result<PresenceAttemptRecord, PresenceOperationError> {
   2051     let attempt_id = RhiPresenceAttemptId(blob::<32>(&row, "attempt_id")?);
   2052     let outbox_id = RhiPresenceOutboxId(blob::<32>(&row, "outbox_id")?);
   2053     let target_ordinal = bounded_u8(&row, "target_ordinal", RHI_PRESENCE_DESIRED_MAX_TARGETS - 1)?;
   2054     let attempt_number = bounded_u16(&row, "attempt_number", RHI_PRESENCE_MAX_ATTEMPTS)?;
   2055     let event_sha256 = blob::<32>(&row, "event_sha256")?;
   2056     let owner = RhiPresenceLeaseOwner(blob::<16>(&row, "lease_owner")?);
   2057     let started_at = RhiPresenceUnixMilliseconds(nonnegative_u64(&row, "started_at_unix_ms")?);
   2058     let finished_at = RhiPresenceUnixMilliseconds(nonnegative_u64(&row, "finished_at_unix_ms")?);
   2059     let outcome = decode_attempt_outcome(&bounded_text(&row, "outcome", "outcome_bytes", 13)?)?;
   2060     let result_code = bounded_text(&row, "result_code", "result_code_bytes", 64)?;
   2061     if owner.0.iter().all(|byte| *byte == 0)
   2062         || started_at > finished_at
   2063         || result_code != outcome.code()
   2064     {
   2065         return Err(PresenceOperationError::Invariant);
   2066     }
   2067     Ok(PresenceAttemptRecord {
   2068         attempt_id,
   2069         outbox_id,
   2070         target_ordinal,
   2071         attempt_number,
   2072         event_sha256,
   2073         owner,
   2074         started_at,
   2075         finished_at,
   2076         outcome,
   2077     })
   2078 }
   2079 
   2080 async fn reconcile_recorded(
   2081     transaction: &mut ServiceSqliteTransaction<'_>,
   2082     input: PresenceRecordInput,
   2083     existing: PresenceAttemptRecord,
   2084 ) -> Result<RhiPresenceAttemptCommit, PresenceOperationError> {
   2085     if existing != PresenceAttemptRecord::from_input(input) {
   2086         return Err(PresenceOperationError::Invariant);
   2087     }
   2088     let outbox = read_outbox(transaction, input.outbox_id)
   2089         .await?
   2090         .ok_or(PresenceOperationError::Invariant)?;
   2091     if outbox.revision
   2092         != input
   2093             .outbox_revision
   2094             .checked_add(1)
   2095             .ok_or(PresenceOperationError::Invariant)?
   2096         || outbox.updated_at != input.finished_at
   2097         || outbox.lease_owner.is_some()
   2098         || outbox.lease_expires.is_some()
   2099         || outbox.event_sha256 != input.event_sha256
   2100     {
   2101         return Err(PresenceOperationError::Invariant);
   2102     }
   2103     validate_targets(&outbox, &outbox.targets)?;
   2104     let (expected_state, expected_next) = target_outcome_schedule(input)?;
   2105     let expected_target_revision = input
   2106         .target_revision
   2107         .checked_add(1)
   2108         .ok_or(PresenceOperationError::Invariant)?;
   2109     let target = outbox
   2110         .targets
   2111         .iter()
   2112         .find(|target| target.ordinal == input.target_ordinal)
   2113         .filter(|target| {
   2114             target.revision == expected_target_revision
   2115                 && target.attempt_count == input.attempt_number
   2116                 && target.last_attempt_id == Some(input.attempt_id)
   2117                 && target.state == expected_state
   2118                 && target.next_attempt == expected_next
   2119                 && target.updated_at == input.finished_at
   2120         })
   2121         .ok_or(PresenceOperationError::Invariant)?;
   2122     let expected_disposition = disposition(&outbox, &outbox.targets)?;
   2123     if outbox.state != expected_disposition.state
   2124         || outbox.next_attempt != expected_disposition.next_attempt
   2125     {
   2126         return Err(PresenceOperationError::Invariant);
   2127     }
   2128     Ok(RhiPresenceAttemptCommit {
   2129         outbox_id: outbox.id,
   2130         attempt_id: input.attempt_id,
   2131         target_ordinal: input.target_ordinal,
   2132         attempt_number: input.attempt_number,
   2133         outcome: input.outcome,
   2134         target_state: target.state,
   2135         outbox_state: outbox.state,
   2136     })
   2137 }
   2138 
   2139 fn disposition(
   2140     outbox: &PresenceOutboxRecord,
   2141     targets: &[PresenceTargetRecord],
   2142 ) -> Result<PresenceOutboxDisposition, PresenceOperationError> {
   2143     validate_targets(outbox, targets)?;
   2144     let required = targets.iter().filter(|target| target.required);
   2145     if required
   2146         .clone()
   2147         .all(|target| target.state == RhiPresenceTargetState::Accepted)
   2148     {
   2149         return Ok(PresenceOutboxDisposition {
   2150             state: RhiPresenceOutboxState::Complete,
   2151             next_attempt: None,
   2152         });
   2153     }
   2154     if required
   2155         .clone()
   2156         .any(|target| target_is_blocking(target, outbox.max_attempts))
   2157     {
   2158         return Ok(PresenceOutboxDisposition {
   2159             state: RhiPresenceOutboxState::Blocked,
   2160             next_attempt: None,
   2161         });
   2162     }
   2163     let next_attempt = required
   2164         .filter_map(|target| target.next_attempt)
   2165         .min()
   2166         .ok_or(PresenceOperationError::Invariant)?;
   2167     Ok(PresenceOutboxDisposition {
   2168         state: RhiPresenceOutboxState::Pending,
   2169         next_attempt: Some(next_attempt),
   2170     })
   2171 }
   2172 
   2173 fn validate_targets(
   2174     outbox: &PresenceOutboxRecord,
   2175     targets: &[PresenceTargetRecord],
   2176 ) -> Result<(), PresenceOperationError> {
   2177     if targets.len() != usize::from(outbox.target_count)
   2178         || targets.is_empty()
   2179         || targets.len() > RHI_PRESENCE_DESIRED_MAX_TARGETS
   2180         || targets.iter().filter(|target| target.required).count()
   2181             != usize::from(outbox.required_target_count)
   2182         || targets
   2183             .iter()
   2184             .enumerate()
   2185             .any(|(ordinal, target)| usize::from(target.ordinal) != ordinal)
   2186         || targets.iter().enumerate().any(|(index, target)| {
   2187             targets[index + 1..]
   2188                 .iter()
   2189                 .any(|later| later.relay_id == target.relay_id)
   2190         })
   2191     {
   2192         return Err(PresenceOperationError::Invariant);
   2193     }
   2194     let mut digest = Sha256::new();
   2195     digest.update(b"radroots.rhi.presence_target_set.v1\0");
   2196     digest.update(
   2197         u32::try_from(targets.len())
   2198             .map_err(|_| PresenceOperationError::Invariant)?
   2199             .to_be_bytes(),
   2200     );
   2201     for target in targets {
   2202         digest.update(u32::from(target.ordinal).to_be_bytes());
   2203         digest.update(
   2204             u64::try_from(target.relay_id.len())
   2205                 .map_err(|_| PresenceOperationError::Invariant)?
   2206                 .to_be_bytes(),
   2207         );
   2208         digest.update(target.relay_id.as_bytes());
   2209         digest.update([u8::from(target.required)]);
   2210     }
   2211     if <[u8; 32]>::from(digest.finalize()) != outbox.target_set_sha256 {
   2212         return Err(PresenceOperationError::Invariant);
   2213     }
   2214     Ok(())
   2215 }
   2216 
   2217 fn target_outcome_schedule(
   2218     input: PresenceRecordInput,
   2219 ) -> Result<(RhiPresenceTargetState, Option<RhiPresenceUnixMilliseconds>), PresenceOperationError> {
   2220     let state = target_state_for_outcome(input.outcome)?;
   2221     let next =
   2222         if retryable_outcome(input.outcome) && input.attempt_number < RHI_PRESENCE_MAX_ATTEMPTS {
   2223             if input.retry_delay.get() > retry_upper_bound(input.attempt_number) {
   2224                 return Err(PresenceOperationError::InvalidInput);
   2225             }
   2226             Some(RhiPresenceUnixMilliseconds(
   2227                 input
   2228                     .finished_at
   2229                     .get()
   2230                     .checked_add(input.retry_delay.get())
   2231                     .filter(|value| *value <= MAX_UNIX_MILLISECONDS)
   2232                     .ok_or(PresenceOperationError::InvalidInput)?,
   2233             ))
   2234         } else {
   2235             if input.retry_delay.get() != 0 {
   2236                 return Err(PresenceOperationError::InvalidInput);
   2237             }
   2238             None
   2239         };
   2240     Ok((state, next))
   2241 }
   2242 
   2243 fn target_is_blocking(target: &PresenceTargetRecord, max_attempts: u16) -> bool {
   2244     matches!(
   2245         target.state,
   2246         RhiPresenceTargetState::Rejected | RhiPresenceTargetState::AuthRequired
   2247     ) || (retryable_state(target.state)
   2248         && target.attempt_count >= max_attempts
   2249         && target.next_attempt.is_none())
   2250 }
   2251 
   2252 fn target_is_due(
   2253     target: &PresenceTargetRecord,
   2254     max_attempts: u16,
   2255     now: RhiPresenceUnixMilliseconds,
   2256 ) -> bool {
   2257     retryable_state(target.state)
   2258         && target.attempt_count < max_attempts
   2259         && target.next_attempt.is_some_and(|next| next <= now)
   2260 }
   2261 
   2262 const fn retryable_state(state: RhiPresenceTargetState) -> bool {
   2263     matches!(
   2264         state,
   2265         RhiPresenceTargetState::Pending
   2266             | RhiPresenceTargetState::Failed
   2267             | RhiPresenceTargetState::RateLimited
   2268             | RhiPresenceTargetState::Unknown
   2269     )
   2270 }
   2271 
   2272 const fn retryable_outcome(outcome: RhiPresenceAttemptOutcome) -> bool {
   2273     matches!(
   2274         outcome,
   2275         RhiPresenceAttemptOutcome::RateLimited
   2276             | RhiPresenceAttemptOutcome::Failed
   2277             | RhiPresenceAttemptOutcome::Unknown
   2278     )
   2279 }
   2280 
   2281 const fn target_state_for_outcome(
   2282     outcome: RhiPresenceAttemptOutcome,
   2283 ) -> Result<RhiPresenceTargetState, PresenceOperationError> {
   2284     match outcome {
   2285         RhiPresenceAttemptOutcome::Submitted => Err(PresenceOperationError::InvalidInput),
   2286         RhiPresenceAttemptOutcome::Accepted => Ok(RhiPresenceTargetState::Accepted),
   2287         RhiPresenceAttemptOutcome::Rejected => Ok(RhiPresenceTargetState::Rejected),
   2288         RhiPresenceAttemptOutcome::RateLimited => Ok(RhiPresenceTargetState::RateLimited),
   2289         RhiPresenceAttemptOutcome::AuthRequired => Ok(RhiPresenceTargetState::AuthRequired),
   2290         RhiPresenceAttemptOutcome::Failed => Ok(RhiPresenceTargetState::Failed),
   2291         RhiPresenceAttemptOutcome::Unknown => Ok(RhiPresenceTargetState::Unknown),
   2292     }
   2293 }
   2294 
   2295 fn retry_upper_bound(attempt_number: u16) -> u64 {
   2296     let mut bound = INITIAL_BACKOFF_MILLISECONDS;
   2297     for _ in 1..attempt_number {
   2298         bound = bound.saturating_mul(2).min(MAXIMUM_BACKOFF_MILLISECONDS);
   2299     }
   2300     bound.min(MAXIMUM_BACKOFF_MILLISECONDS)
   2301 }
   2302 
   2303 fn same_outbox(current: &PresenceOutboxRecord, prior: &PresenceOutboxRecord) -> bool {
   2304     same_outbox_identity(current, prior)
   2305         && current.state == prior.state
   2306         && current.revision == prior.revision
   2307         && current.next_attempt == prior.next_attempt
   2308         && current.lease_owner == prior.lease_owner
   2309         && current.lease_expires == prior.lease_expires
   2310         && current.updated_at == prior.updated_at
   2311         && current.targets == prior.targets
   2312 }
   2313 
   2314 fn same_outbox_identity(current: &PresenceOutboxRecord, prior: &PresenceOutboxRecord) -> bool {
   2315     current.id == prior.id
   2316         && current.desired_generation == prior.desired_generation
   2317         && current.document_kind == prior.document_kind
   2318         && current.desired_sha256 == prior.desired_sha256
   2319         && current.target_set_sha256 == prior.target_set_sha256
   2320         && current.event_id == prior.event_id
   2321         && current.event_sha256 == prior.event_sha256
   2322         && current.exact_signed_event_bytes == prior.exact_signed_event_bytes
   2323         && current.authored_at_unix_s == prior.authored_at_unix_s
   2324         && current.service_public_key == prior.service_public_key
   2325         && current.target_count == prior.target_count
   2326         && current.required_target_count == prior.required_target_count
   2327         && current.max_attempts == prior.max_attempts
   2328         && current.created_at == prior.created_at
   2329 }
   2330 
   2331 fn valid_outbox_shape(
   2332     state: RhiPresenceOutboxState,
   2333     next_attempt: Option<RhiPresenceUnixMilliseconds>,
   2334     lease_owner: Option<RhiPresenceLeaseOwner>,
   2335     lease_expires: Option<RhiPresenceUnixMilliseconds>,
   2336 ) -> bool {
   2337     match state {
   2338         RhiPresenceOutboxState::Pending => {
   2339             next_attempt.is_some() && lease_owner.is_none() && lease_expires.is_none()
   2340         }
   2341         RhiPresenceOutboxState::Leased => {
   2342             next_attempt.is_none() && lease_owner.is_some() && lease_expires.is_some()
   2343         }
   2344         RhiPresenceOutboxState::Complete
   2345         | RhiPresenceOutboxState::Blocked
   2346         | RhiPresenceOutboxState::Superseded => {
   2347             next_attempt.is_none() && lease_owner.is_none() && lease_expires.is_none()
   2348         }
   2349     }
   2350 }
   2351 
   2352 fn valid_target_shape(
   2353     state: RhiPresenceTargetState,
   2354     attempt_count: u16,
   2355     next_attempt: Option<RhiPresenceUnixMilliseconds>,
   2356     last_attempt_id: Option<RhiPresenceAttemptId>,
   2357 ) -> bool {
   2358     match state {
   2359         RhiPresenceTargetState::Pending => {
   2360             attempt_count == 0 && next_attempt.is_some() && last_attempt_id.is_none()
   2361         }
   2362         RhiPresenceTargetState::Submitted => {
   2363             attempt_count > 0 && next_attempt.is_none() && last_attempt_id.is_some()
   2364         }
   2365         RhiPresenceTargetState::Accepted
   2366         | RhiPresenceTargetState::Rejected
   2367         | RhiPresenceTargetState::AuthRequired => {
   2368             attempt_count > 0 && next_attempt.is_none() && last_attempt_id.is_some()
   2369         }
   2370         RhiPresenceTargetState::RateLimited
   2371         | RhiPresenceTargetState::Failed
   2372         | RhiPresenceTargetState::Unknown => {
   2373             attempt_count > 0
   2374                 && last_attempt_id.is_some()
   2375                 && if attempt_count < RHI_PRESENCE_MAX_ATTEMPTS {
   2376                     next_attempt.is_some()
   2377                 } else {
   2378                     next_attempt.is_none()
   2379                 }
   2380         }
   2381     }
   2382 }
   2383 
   2384 fn decode_document_kind(value: &str) -> Result<RhiPresenceDocumentKind, PresenceOperationError> {
   2385     match value {
   2386         "service_profile" => Ok(RhiPresenceDocumentKind::ServiceProfile),
   2387         "application_handler" => Ok(RhiPresenceDocumentKind::ApplicationHandler),
   2388         _ => Err(PresenceOperationError::Invariant),
   2389     }
   2390 }
   2391 
   2392 fn decode_outbox_state(value: &str) -> Result<RhiPresenceOutboxState, PresenceOperationError> {
   2393     match value {
   2394         "pending" => Ok(RhiPresenceOutboxState::Pending),
   2395         "leased" => Ok(RhiPresenceOutboxState::Leased),
   2396         "complete" => Ok(RhiPresenceOutboxState::Complete),
   2397         "blocked" => Ok(RhiPresenceOutboxState::Blocked),
   2398         "superseded" => Ok(RhiPresenceOutboxState::Superseded),
   2399         _ => Err(PresenceOperationError::Invariant),
   2400     }
   2401 }
   2402 
   2403 fn decode_target_state(value: &str) -> Result<RhiPresenceTargetState, PresenceOperationError> {
   2404     match value {
   2405         "pending" => Ok(RhiPresenceTargetState::Pending),
   2406         "submitted" => Ok(RhiPresenceTargetState::Submitted),
   2407         "accepted" => Ok(RhiPresenceTargetState::Accepted),
   2408         "rejected" => Ok(RhiPresenceTargetState::Rejected),
   2409         "rate_limited" => Ok(RhiPresenceTargetState::RateLimited),
   2410         "auth_required" => Ok(RhiPresenceTargetState::AuthRequired),
   2411         "failed" => Ok(RhiPresenceTargetState::Failed),
   2412         "unknown" => Ok(RhiPresenceTargetState::Unknown),
   2413         _ => Err(PresenceOperationError::Invariant),
   2414     }
   2415 }
   2416 
   2417 fn decode_attempt_outcome(
   2418     value: &str,
   2419 ) -> Result<RhiPresenceAttemptOutcome, PresenceOperationError> {
   2420     match value {
   2421         "submitted" => Ok(RhiPresenceAttemptOutcome::Submitted),
   2422         "accepted" => Ok(RhiPresenceAttemptOutcome::Accepted),
   2423         "rejected" => Ok(RhiPresenceAttemptOutcome::Rejected),
   2424         "rate_limited" => Ok(RhiPresenceAttemptOutcome::RateLimited),
   2425         "auth_required" => Ok(RhiPresenceAttemptOutcome::AuthRequired),
   2426         "failed" => Ok(RhiPresenceAttemptOutcome::Failed),
   2427         "unknown" => Ok(RhiPresenceAttemptOutcome::Unknown),
   2428         _ => Err(PresenceOperationError::Invariant),
   2429     }
   2430 }
   2431 
   2432 fn exactly_zero_or_one(
   2433     rows: Vec<sqlx::sqlite::SqliteRow>,
   2434 ) -> Result<Option<sqlx::sqlite::SqliteRow>, PresenceOperationError> {
   2435     match rows.as_slice() {
   2436         [] => Ok(None),
   2437         [_] => Ok(rows.into_iter().next()),
   2438         _ => Err(PresenceOperationError::Invariant),
   2439     }
   2440 }
   2441 
   2442 fn blob<const N: usize>(
   2443     row: &sqlx::sqlite::SqliteRow,
   2444     column: &str,
   2445 ) -> Result<[u8; N], PresenceOperationError> {
   2446     row.try_get::<Option<Vec<u8>>, _>(column)
   2447         .map_err(|_| PresenceOperationError::Invariant)?
   2448         .ok_or(PresenceOperationError::Invariant)?
   2449         .try_into()
   2450         .map_err(|_| PresenceOperationError::Invariant)
   2451 }
   2452 
   2453 fn bounded_text(
   2454     row: &sqlx::sqlite::SqliteRow,
   2455     value_column: &str,
   2456     length_column: &str,
   2457     maximum: usize,
   2458 ) -> Result<String, PresenceOperationError> {
   2459     let length = row
   2460         .try_get::<i64, _>(length_column)
   2461         .ok()
   2462         .and_then(|value| usize::try_from(value).ok())
   2463         .filter(|value| *value > 0 && *value <= maximum)
   2464         .ok_or(PresenceOperationError::Invariant)?;
   2465     let value = row
   2466         .try_get::<String, _>(value_column)
   2467         .map_err(|_| PresenceOperationError::Invariant)?;
   2468     if value.len() != length {
   2469         return Err(PresenceOperationError::Invariant);
   2470     }
   2471     Ok(value)
   2472 }
   2473 
   2474 fn bounded_u8(
   2475     row: &sqlx::sqlite::SqliteRow,
   2476     column: &str,
   2477     maximum: usize,
   2478 ) -> Result<u8, PresenceOperationError> {
   2479     row.try_get::<i64, _>(column)
   2480         .ok()
   2481         .and_then(|value| u8::try_from(value).ok())
   2482         .filter(|value| usize::from(*value) <= maximum)
   2483         .ok_or(PresenceOperationError::Invariant)
   2484 }
   2485 
   2486 fn bounded_u16(
   2487     row: &sqlx::sqlite::SqliteRow,
   2488     column: &str,
   2489     maximum: u16,
   2490 ) -> Result<u16, PresenceOperationError> {
   2491     row.try_get::<i64, _>(column)
   2492         .ok()
   2493         .and_then(|value| u16::try_from(value).ok())
   2494         .filter(|value| *value <= maximum)
   2495         .ok_or(PresenceOperationError::Invariant)
   2496 }
   2497 
   2498 fn bounded_u64(
   2499     row: &sqlx::sqlite::SqliteRow,
   2500     column: &str,
   2501     maximum: u64,
   2502 ) -> Result<u64, PresenceOperationError> {
   2503     row.try_get::<i64, _>(column)
   2504         .ok()
   2505         .and_then(|value| u64::try_from(value).ok())
   2506         .filter(|value| *value <= maximum)
   2507         .ok_or(PresenceOperationError::Invariant)
   2508 }
   2509 
   2510 fn positive_u64(
   2511     row: &sqlx::sqlite::SqliteRow,
   2512     column: &str,
   2513 ) -> Result<u64, PresenceOperationError> {
   2514     nonnegative_u64(row, column).and_then(|value| {
   2515         (value != 0)
   2516             .then_some(value)
   2517             .ok_or(PresenceOperationError::Invariant)
   2518     })
   2519 }
   2520 
   2521 fn nonnegative_u64(
   2522     row: &sqlx::sqlite::SqliteRow,
   2523     column: &str,
   2524 ) -> Result<u64, PresenceOperationError> {
   2525     row.try_get::<i64, _>(column)
   2526         .ok()
   2527         .and_then(|value| u64::try_from(value).ok())
   2528         .ok_or(PresenceOperationError::Invariant)
   2529 }
   2530 
   2531 fn positive_usize(
   2532     row: &sqlx::sqlite::SqliteRow,
   2533     column: &str,
   2534 ) -> Result<usize, PresenceOperationError> {
   2535     row.try_get::<i64, _>(column)
   2536         .ok()
   2537         .and_then(|value| usize::try_from(value).ok())
   2538         .filter(|value| *value != 0)
   2539         .ok_or(PresenceOperationError::Invariant)
   2540 }
   2541 
   2542 fn boolean_i64(
   2543     row: &sqlx::sqlite::SqliteRow,
   2544     column: &str,
   2545 ) -> Result<bool, PresenceOperationError> {
   2546     match row.try_get::<i64, _>(column) {
   2547         Ok(0) => Ok(false),
   2548         Ok(1) => Ok(true),
   2549         Ok(_) | Err(_) => Err(PresenceOperationError::Invariant),
   2550     }
   2551 }
   2552 
   2553 fn optional_millis(
   2554     row: &sqlx::sqlite::SqliteRow,
   2555     column: &str,
   2556 ) -> Result<Option<RhiPresenceUnixMilliseconds>, PresenceOperationError> {
   2557     row.try_get::<Option<i64>, _>(column)
   2558         .map_err(|_| PresenceOperationError::Invariant)?
   2559         .map(|value| {
   2560             u64::try_from(value)
   2561                 .map(RhiPresenceUnixMilliseconds)
   2562                 .map_err(|_| PresenceOperationError::Invariant)
   2563         })
   2564         .transpose()
   2565 }
   2566 
   2567 fn optional_owner(
   2568     row: &sqlx::sqlite::SqliteRow,
   2569 ) -> Result<Option<RhiPresenceLeaseOwner>, PresenceOperationError> {
   2570     let length = row
   2571         .try_get::<Option<i64>, _>("lease_owner_bytes")
   2572         .map_err(|_| PresenceOperationError::Invariant)?;
   2573     let value = row
   2574         .try_get::<Option<Vec<u8>>, _>("lease_owner")
   2575         .map_err(|_| PresenceOperationError::Invariant)?;
   2576     match (length, value) {
   2577         (None, None) => Ok(None),
   2578         (Some(16), Some(value)) => {
   2579             let value: [u8; 16] = value
   2580                 .try_into()
   2581                 .map_err(|_| PresenceOperationError::Invariant)?;
   2582             if value.iter().all(|byte| *byte == 0) {
   2583                 return Err(PresenceOperationError::Invariant);
   2584             }
   2585             Ok(Some(RhiPresenceLeaseOwner(value)))
   2586         }
   2587         _ => Err(PresenceOperationError::Invariant),
   2588     }
   2589 }
   2590 
   2591 fn optional_attempt_id(
   2592     row: &sqlx::sqlite::SqliteRow,
   2593 ) -> Result<Option<RhiPresenceAttemptId>, PresenceOperationError> {
   2594     let length = row
   2595         .try_get::<Option<i64>, _>("last_attempt_id_bytes")
   2596         .map_err(|_| PresenceOperationError::Invariant)?;
   2597     let value = row
   2598         .try_get::<Option<Vec<u8>>, _>("last_attempt_id")
   2599         .map_err(|_| PresenceOperationError::Invariant)?;
   2600     match (length, value) {
   2601         (None, None) => Ok(None),
   2602         (Some(32), Some(value)) => Ok(Some(RhiPresenceAttemptId(
   2603             value
   2604                 .try_into()
   2605                 .map_err(|_| PresenceOperationError::Invariant)?,
   2606         ))),
   2607         _ => Err(PresenceOperationError::Invariant),
   2608     }
   2609 }
   2610 
   2611 fn valid_public_key(value: &str) -> bool {
   2612     value.len() == 64
   2613         && value
   2614             .bytes()
   2615             .all(|byte| byte.is_ascii_hexdigit() && !byte.is_ascii_uppercase())
   2616         && NostrPublicKey::from_hex(value).is_ok()
   2617 }
   2618 
   2619 fn valid_relay_id(value: &str) -> bool {
   2620     !value.is_empty()
   2621         && value.len() <= 64
   2622         && value.as_bytes()[0].is_ascii_lowercase()
   2623         && value.bytes().all(|byte| {
   2624             byte.is_ascii_lowercase() || byte.is_ascii_digit() || matches!(byte, b'_' | b'-')
   2625         })
   2626 }
   2627 
   2628 fn i64_value(value: u64) -> Result<i64, PresenceOperationError> {
   2629     i64::try_from(value).map_err(|_| PresenceOperationError::InvalidInput)
   2630 }
   2631 
   2632 fn optional_i64(
   2633     value: Option<RhiPresenceUnixMilliseconds>,
   2634 ) -> Result<Option<i64>, PresenceOperationError> {
   2635     value.map(|value| i64_value(value.get())).transpose()
   2636 }
   2637 
   2638 #[derive(Clone, Copy)]
   2639 struct PresenceDesiredBinding {
   2640     generation: u64,
   2641     mode: RhiPresenceDesiredMode,
   2642     profile: bool,
   2643     application_handler: bool,
   2644     target_set_sha256: [u8; 32],
   2645     target_count: u8,
   2646     required_target_count: u8,
   2647     queue_capacity: u32,
   2648     desired_sha256: [u8; 32],
   2649 }
   2650 
   2651 impl PresenceDesiredBinding {
   2652     const fn document_count(self) -> usize {
   2653         self.profile as usize + self.application_handler as usize
   2654     }
   2655 
   2656     fn document_kinds(self) -> impl Iterator<Item = RhiPresenceDocumentKind> {
   2657         [
   2658             self.profile
   2659                 .then_some(RhiPresenceDocumentKind::ServiceProfile),
   2660             self.application_handler
   2661                 .then_some(RhiPresenceDocumentKind::ApplicationHandler),
   2662         ]
   2663         .into_iter()
   2664         .flatten()
   2665     }
   2666 }
   2667 
   2668 #[derive(PartialEq, Eq)]
   2669 struct PresenceOutboxRecord {
   2670     id: RhiPresenceOutboxId,
   2671     desired_generation: u64,
   2672     document_kind: RhiPresenceDocumentKind,
   2673     desired_sha256: [u8; 32],
   2674     target_set_sha256: [u8; 32],
   2675     event_id: [u8; 32],
   2676     event_sha256: [u8; 32],
   2677     exact_signed_event_bytes: Box<[u8]>,
   2678     authored_at_unix_s: u64,
   2679     service_public_key: Box<str>,
   2680     target_count: u8,
   2681     required_target_count: u8,
   2682     max_attempts: u16,
   2683     state: RhiPresenceOutboxState,
   2684     revision: u64,
   2685     next_attempt: Option<RhiPresenceUnixMilliseconds>,
   2686     lease_owner: Option<RhiPresenceLeaseOwner>,
   2687     lease_expires: Option<RhiPresenceUnixMilliseconds>,
   2688     created_at: RhiPresenceUnixMilliseconds,
   2689     updated_at: RhiPresenceUnixMilliseconds,
   2690     targets: Box<[PresenceTargetRecord]>,
   2691 }
   2692 
   2693 #[derive(Clone, PartialEq, Eq)]
   2694 struct PresenceTargetRecord {
   2695     ordinal: u8,
   2696     relay_id: Box<str>,
   2697     required: bool,
   2698     state: RhiPresenceTargetState,
   2699     revision: u64,
   2700     attempt_count: u16,
   2701     next_attempt: Option<RhiPresenceUnixMilliseconds>,
   2702     last_attempt_id: Option<RhiPresenceAttemptId>,
   2703     updated_at: RhiPresenceUnixMilliseconds,
   2704 }
   2705 
   2706 struct PresenceRecoveryCandidate {
   2707     outbox: PresenceOutboxRecord,
   2708     retry_upper_bound: u64,
   2709     current_desired: bool,
   2710 }
   2711 
   2712 #[derive(Clone, Copy)]
   2713 struct PresenceRecordInput {
   2714     outbox_id: RhiPresenceOutboxId,
   2715     outbox_revision: u64,
   2716     event_sha256: [u8; 32],
   2717     owner: RhiPresenceLeaseOwner,
   2718     lease_expires: RhiPresenceUnixMilliseconds,
   2719     target_ordinal: u8,
   2720     target_revision: u64,
   2721     attempt_number: u16,
   2722     attempt_id: RhiPresenceAttemptId,
   2723     started_at: RhiPresenceUnixMilliseconds,
   2724     finished_at: RhiPresenceUnixMilliseconds,
   2725     outcome: RhiPresenceAttemptOutcome,
   2726     retry_delay: RhiPresenceRetryDelayMilliseconds,
   2727 }
   2728 
   2729 impl PresenceRecordInput {
   2730     fn from_prepared(
   2731         prepared: &RhiPreparedPresenceAttempt,
   2732         finished_at: RhiPresenceUnixMilliseconds,
   2733         outcome: RhiPresenceAttemptOutcome,
   2734         retry_delay: RhiPresenceRetryDelayMilliseconds,
   2735     ) -> Result<Self, RhiPresencePublicationError> {
   2736         if finished_at < prepared.started_at || outcome == RhiPresenceAttemptOutcome::Submitted {
   2737             return Err(failure(RhiPresencePublicationErrorKind::InvalidInput));
   2738         }
   2739         let expected = derive_attempt_id(
   2740             prepared.lease.outbox.id,
   2741             prepared.lease.outbox.event_sha256,
   2742             prepared.target.ordinal,
   2743             prepared.target.attempt_count,
   2744         );
   2745         if expected != prepared.attempt_id {
   2746             return Err(failure(RhiPresencePublicationErrorKind::Invariant));
   2747         }
   2748         Ok(Self {
   2749             outbox_id: prepared.lease.outbox.id,
   2750             outbox_revision: prepared.lease.outbox.revision,
   2751             event_sha256: prepared.lease.outbox.event_sha256,
   2752             owner: prepared.lease.owner,
   2753             lease_expires: prepared.lease.expires_at,
   2754             target_ordinal: prepared.target.ordinal,
   2755             target_revision: prepared.target.revision,
   2756             attempt_number: prepared.target.attempt_count,
   2757             attempt_id: prepared.attempt_id,
   2758             started_at: prepared.started_at,
   2759             finished_at,
   2760             outcome,
   2761             retry_delay,
   2762         })
   2763     }
   2764 }
   2765 
   2766 #[derive(Clone, Copy, PartialEq, Eq)]
   2767 struct PresenceAttemptRecord {
   2768     attempt_id: RhiPresenceAttemptId,
   2769     outbox_id: RhiPresenceOutboxId,
   2770     target_ordinal: u8,
   2771     attempt_number: u16,
   2772     event_sha256: [u8; 32],
   2773     owner: RhiPresenceLeaseOwner,
   2774     started_at: RhiPresenceUnixMilliseconds,
   2775     finished_at: RhiPresenceUnixMilliseconds,
   2776     outcome: RhiPresenceAttemptOutcome,
   2777 }
   2778 
   2779 impl PresenceAttemptRecord {
   2780     const fn from_input(input: PresenceRecordInput) -> Self {
   2781         Self {
   2782             attempt_id: input.attempt_id,
   2783             outbox_id: input.outbox_id,
   2784             target_ordinal: input.target_ordinal,
   2785             attempt_number: input.attempt_number,
   2786             event_sha256: input.event_sha256,
   2787             owner: input.owner,
   2788             started_at: input.started_at,
   2789             finished_at: input.finished_at,
   2790             outcome: input.outcome,
   2791         }
   2792     }
   2793 }
   2794 
   2795 #[derive(Clone, Copy)]
   2796 struct PresenceOutboxDisposition {
   2797     state: RhiPresenceOutboxState,
   2798     next_attempt: Option<RhiPresenceUnixMilliseconds>,
   2799 }
   2800 
   2801 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
   2802 enum PresenceOperationError {
   2803     InvalidInput,
   2804     DesiredStateMismatch,
   2805     VerificationFailed,
   2806     NotReady,
   2807     LeaseLost,
   2808     Invariant,
   2809     Storage,
   2810 }
   2811 
   2812 fn map_transaction_error(
   2813     error: ServiceSqliteTransactionError<PresenceOperationError>,
   2814 ) -> RhiPresencePublicationError {
   2815     if error.kind() == ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown {
   2816         return failure(RhiPresencePublicationErrorKind::CommitOutcomeUnknown);
   2817     }
   2818     failure(match error.operation_error().copied() {
   2819         Some(PresenceOperationError::InvalidInput) => RhiPresencePublicationErrorKind::InvalidInput,
   2820         Some(PresenceOperationError::DesiredStateMismatch) => {
   2821             RhiPresencePublicationErrorKind::DesiredStateMismatch
   2822         }
   2823         Some(PresenceOperationError::VerificationFailed) => {
   2824             RhiPresencePublicationErrorKind::VerificationFailed
   2825         }
   2826         Some(PresenceOperationError::NotReady) => RhiPresencePublicationErrorKind::NotReady,
   2827         Some(PresenceOperationError::LeaseLost) => RhiPresencePublicationErrorKind::LeaseLost,
   2828         Some(PresenceOperationError::Invariant) => RhiPresencePublicationErrorKind::Invariant,
   2829         Some(PresenceOperationError::Storage) | None => RhiPresencePublicationErrorKind::Storage,
   2830     })
   2831 }
   2832 
   2833 fn presence_plan(
   2834     kind: RhiPresenceDocumentKind,
   2835     author: &str,
   2836     created_at: u64,
   2837 ) -> Result<AuthoredEventPlan, RhiPresencePublicationError> {
   2838     match kind {
   2839         RhiPresenceDocumentKind::ServiceProfile => {
   2840             let profile = AuthoredProfile::new(PROFILE_NAME)
   2841                 .map_err(|_| failure(RhiPresencePublicationErrorKind::RenderingFailed))?
   2842                 .with_display_name(PROFILE_DISPLAY_NAME)
   2843                 .with_about(PROFILE_ABOUT)
   2844                 .with_bot(true);
   2845             AuthoredEventPlan::from_profile(&profile, created_at, author)
   2846                 .map_err(|_| failure(RhiPresencePublicationErrorKind::RenderingFailed))
   2847         }
   2848         RhiPresenceDocumentKind::ApplicationHandler => {
   2849             let metadata = nostr::Metadata::new()
   2850                 .name(PROFILE_DISPLAY_NAME)
   2851                 .about(PROFILE_ABOUT);
   2852             let content = serde_json::to_string(&metadata)
   2853                 .map_err(|_| failure(RhiPresencePublicationErrorKind::RenderingFailed))?;
   2854             let spec = ApplicationHandlerSpec::new(APPLICATION_HANDLER_KINDS.to_vec())
   2855                 .with_identifier(APPLICATION_HANDLER_IDENTIFIER)
   2856                 .with_metadata(metadata);
   2857             let builder = build_application_handler(&spec)
   2858                 .map_err(|_| failure(RhiPresencePublicationErrorKind::RenderingFailed))?
   2859                 .custom_created_at(Timestamp::from_secs(created_at));
   2860             let public_key = NostrPublicKey::from_hex(author)
   2861                 .map_err(|_| failure(RhiPresencePublicationErrorKind::IdentityMismatch))?;
   2862             let request = builder
   2863                 .into_external_signing_request(public_key)
   2864                 .map_err(|_| failure(RhiPresencePublicationErrorKind::RenderingFailed))?;
   2865             let draft = GenericEventDraft::new(
   2866                 "radroots.application.handler.v1",
   2867                 KIND_APPLICATION_HANDLER,
   2868                 created_at,
   2869                 vec![
   2870                     vec!["d".to_owned(), APPLICATION_HANDLER_IDENTIFIER.to_owned()],
   2871                     vec!["k".to_owned(), APPLICATION_HANDLER_KINDS[0].to_string()],
   2872                 ],
   2873                 content,
   2874                 author,
   2875             )
   2876             .map_err(|_| failure(RhiPresencePublicationErrorKind::RenderingFailed))?;
   2877             let plan = AuthoredEventPlan::from_generic(draft)
   2878                 .map_err(|_| failure(RhiPresencePublicationErrorKind::RenderingFailed))?;
   2879             if request.expected_event_id().to_bytes() != *plan.expected_event_id().as_bytes() {
   2880                 return Err(failure(RhiPresencePublicationErrorKind::RenderingFailed));
   2881             }
   2882             Ok(plan)
   2883         }
   2884     }
   2885 }
   2886 
   2887 fn sign_plan(
   2888     identity: &RhiDecryptedIdentity,
   2889     plan: &AuthoredEventPlan,
   2890     auxiliary: &[u8; 32],
   2891 ) -> Result<nostr::Event, RhiPresencePublicationError> {
   2892     let kind = u16::try_from(plan.body().kind())
   2893         .map(Kind::Custom)
   2894         .map_err(|_| failure(RhiPresencePublicationErrorKind::RenderingFailed))?;
   2895     let tags = plan
   2896         .body()
   2897         .tags()
   2898         .iter()
   2899         .cloned()
   2900         .map(Tag::parse)
   2901         .collect::<Result<Vec<_>, _>>()
   2902         .map_err(|_| failure(RhiPresencePublicationErrorKind::RenderingFailed))?;
   2903     let author = NostrPublicKey::from_hex(identity.public_identity().as_hex())
   2904         .map_err(|_| failure(RhiPresencePublicationErrorKind::IdentityMismatch))?;
   2905     let unsigned = EventBuilder::new(kind, plan.body().content())
   2906         .tags(tags)
   2907         .custom_created_at(Timestamp::from_secs(plan.created_at()))
   2908         .build(author);
   2909     if unsigned.id.as_ref().map(|id| id.to_bytes()) != Some(*plan.expected_event_id().as_bytes()) {
   2910         return Err(failure(RhiPresencePublicationErrorKind::RenderingFailed));
   2911     }
   2912     identity
   2913         .sign_nostr_event(unsigned, auxiliary)
   2914         .map_err(|_| failure(RhiPresencePublicationErrorKind::RenderingFailed))
   2915 }
   2916 
   2917 fn validate_signed_document(
   2918     kind: RhiPresenceDocumentKind,
   2919     author: &str,
   2920     created_at: u64,
   2921     bytes: &[u8],
   2922 ) -> Result<EventEnvelope, RhiPresencePublicationError> {
   2923     if bytes.is_empty() || bytes.len() > RHI_PRESENCE_SIGNED_EVENT_MAX_BYTES {
   2924         return Err(failure(RhiPresencePublicationErrorKind::VerificationFailed));
   2925     }
   2926     let source = core::str::from_utf8(bytes)
   2927         .map_err(|_| failure(RhiPresencePublicationErrorKind::VerificationFailed))?;
   2928     let wire = Nip01EventWire::parse_json_unverified_with_limits(source, presence_wire_limits())
   2929         .map_err(|_| failure(RhiPresencePublicationErrorKind::VerificationFailed))?;
   2930     let event = wire
   2931         .into_unverified_envelope()
   2932         .map_err(|_| failure(RhiPresencePublicationErrorKind::VerificationFailed))?;
   2933     if verify_id(&event) != Verification::IdVerified
   2934         || verify(&event) != Verification::Verified
   2935         || event.author().to_hex() != author
   2936         || event.created_at_u64() != created_at
   2937     {
   2938         return Err(failure(RhiPresencePublicationErrorKind::VerificationFailed));
   2939     }
   2940     let expected = presence_plan(kind, author, created_at)?;
   2941     if event.id().as_bytes() != expected.expected_event_id().as_bytes()
   2942         || event.kind_u32() != expected.body().kind()
   2943         || event.tags_as_vec() != expected.body().tags()
   2944         || event.content() != expected.body().content()
   2945     {
   2946         return Err(failure(RhiPresencePublicationErrorKind::VerificationFailed));
   2947     }
   2948     Ok(event)
   2949 }
   2950 
   2951 const fn presence_wire_limits() -> EventWireLimits {
   2952     EventWireLimits {
   2953         max_raw_json_bytes: RHI_PRESENCE_SIGNED_EVENT_MAX_BYTES,
   2954         max_content_bytes: 4 * 1024,
   2955         max_tag_count: 4,
   2956         max_total_tag_elements: 8,
   2957         max_tag_element_bytes: 64,
   2958         max_total_tag_bytes: 256,
   2959         max_extra_fields: 0,
   2960         max_total_extra_json_bytes: 0,
   2961     }
   2962 }
   2963 
   2964 const fn failure(kind: RhiPresencePublicationErrorKind) -> RhiPresencePublicationError {
   2965     RhiPresencePublicationError { kind }
   2966 }
   2967 
   2968 #[cfg(test)]
   2969 mod tests {
   2970     use std::{collections::BTreeSet, error::Error as _};
   2971 
   2972     use super::*;
   2973 
   2974     #[test]
   2975     fn public_vocabularies_bounds_and_diagnostics_are_closed() {
   2976         let kinds = [
   2977             RhiPresencePublicationErrorKind::InvalidMode,
   2978             RhiPresencePublicationErrorKind::InvalidInput,
   2979             RhiPresencePublicationErrorKind::DesiredStateMismatch,
   2980             RhiPresencePublicationErrorKind::IdentityMismatch,
   2981             RhiPresencePublicationErrorKind::EntropyUnavailable,
   2982             RhiPresencePublicationErrorKind::RenderingFailed,
   2983             RhiPresencePublicationErrorKind::VerificationFailed,
   2984             RhiPresencePublicationErrorKind::NotReady,
   2985             RhiPresencePublicationErrorKind::LeaseLost,
   2986             RhiPresencePublicationErrorKind::Invariant,
   2987             RhiPresencePublicationErrorKind::ClockUnavailable,
   2988             RhiPresencePublicationErrorKind::Storage,
   2989             RhiPresencePublicationErrorKind::CommitOutcomeUnknown,
   2990         ];
   2991         let codes: BTreeSet<_> = kinds.iter().map(|kind| kind.code()).collect();
   2992         assert_eq!(codes.len(), kinds.len());
   2993         for kind in kinds {
   2994             let error = failure(kind);
   2995             assert_eq!(error.kind(), kind);
   2996             assert_eq!(error.code(), kind.code());
   2997             assert!(error.source().is_none());
   2998             let rendered = format!("{error} {error:?}");
   2999             for forbidden in ["relay", "wss://", "state.sqlite", "010101", "sqlx"] {
   3000                 assert!(!rendered.contains(forbidden));
   3001             }
   3002         }
   3003 
   3004         assert_eq!(
   3005             RhiPresenceUnixMilliseconds::new(i64::MAX as u64)
   3006                 .expect("maximum time")
   3007                 .get(),
   3008             i64::MAX as u64
   3009         );
   3010         assert!(RhiPresenceUnixMilliseconds::new(i64::MAX as u64 + 1).is_err());
   3011         assert_eq!(
   3012             RhiPresenceRetryDelayMilliseconds::new(MAXIMUM_BACKOFF_MILLISECONDS)
   3013                 .expect("maximum delay")
   3014                 .get(),
   3015             MAXIMUM_BACKOFF_MILLISECONDS
   3016         );
   3017         assert!(RhiPresenceRetryDelayMilliseconds::new(MAXIMUM_BACKOFF_MILLISECONDS + 1).is_err());
   3018         assert!(RhiPresenceLeaseOwner::from_bytes([0; 16]).is_err());
   3019         assert_eq!(
   3020             format!("{:?}", RhiPresenceLeaseOwner::from_bytes([1; 16]).unwrap()),
   3021             "RhiPresenceLeaseOwner([redacted])"
   3022         );
   3023         assert_eq!(
   3024             format!("{:?}", RhiPresenceOutboxId([0x5a; 32])),
   3025             "RhiPresenceOutboxId([redacted])"
   3026         );
   3027         assert_eq!(
   3028             format!("{:?}", RhiPresenceAttemptId([0x6a; 32])),
   3029             "RhiPresenceAttemptId([redacted])"
   3030         );
   3031     }
   3032 
   3033     #[test]
   3034     fn state_and_outcome_codes_are_exact() {
   3035         assert_eq!(
   3036             [
   3037                 RhiPresenceOutboxState::Pending,
   3038                 RhiPresenceOutboxState::Leased,
   3039                 RhiPresenceOutboxState::Complete,
   3040                 RhiPresenceOutboxState::Blocked,
   3041                 RhiPresenceOutboxState::Superseded,
   3042             ]
   3043             .map(RhiPresenceOutboxState::code),
   3044             ["pending", "leased", "complete", "blocked", "superseded"]
   3045         );
   3046         assert_eq!(
   3047             [
   3048                 RhiPresenceTargetState::Pending,
   3049                 RhiPresenceTargetState::Submitted,
   3050                 RhiPresenceTargetState::Accepted,
   3051                 RhiPresenceTargetState::Rejected,
   3052                 RhiPresenceTargetState::RateLimited,
   3053                 RhiPresenceTargetState::AuthRequired,
   3054                 RhiPresenceTargetState::Failed,
   3055                 RhiPresenceTargetState::Unknown,
   3056             ]
   3057             .map(RhiPresenceTargetState::code),
   3058             [
   3059                 "pending",
   3060                 "submitted",
   3061                 "accepted",
   3062                 "rejected",
   3063                 "rate_limited",
   3064                 "auth_required",
   3065                 "failed",
   3066                 "unknown",
   3067             ]
   3068         );
   3069         assert_eq!(
   3070             [
   3071                 RhiPresenceAttemptOutcome::Submitted,
   3072                 RhiPresenceAttemptOutcome::Accepted,
   3073                 RhiPresenceAttemptOutcome::Rejected,
   3074                 RhiPresenceAttemptOutcome::RateLimited,
   3075                 RhiPresenceAttemptOutcome::AuthRequired,
   3076                 RhiPresenceAttemptOutcome::Failed,
   3077                 RhiPresenceAttemptOutcome::Unknown,
   3078             ]
   3079             .map(RhiPresenceAttemptOutcome::code),
   3080             [
   3081                 "submitted",
   3082                 "accepted",
   3083                 "rejected",
   3084                 "rate_limited",
   3085                 "auth_required",
   3086                 "failed",
   3087                 "unknown",
   3088             ]
   3089         );
   3090     }
   3091 }