rhi

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

commit 0726c6766fbbc0037ace1490d0f808e41739f354
parent da1a1c901fd1aa8028af33a2321dc7645aac2af0
Author: triesap <tyson@radroots.org>
Date:   Mon, 24 Aug 2026 15:34:23 +0000

refactor(rhi): persist exact presence publication

Diffstat:
MAGENTS.md | 18+++++++++++++++---
MREADME | 43+++++++++++++++++++++++++++++++++++++------
Mcontracts/api_baselines/rhi.txt | 165++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcontracts/services_hardening/presence_desired_state.v1.json | 2+-
Acontracts/services_hardening/presence_publication.v1.json | 115+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcontracts/services_hardening/state_repository_topology.v1.json | 9++++++---
Msrc/lib.rs | 20+++++++++++++++++---
Msrc/presence_desired.rs | 11+++++++++++
Asrc/presence_publication.rs | 3091+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Msrc/state_catalog.rs | 476+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
Msrc/state_repository.rs | 50+++++++++++++++++++++++++++++++++++++++++++++++++-
Mtests/package_boundary.rs | 117++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mtests/services_hardening_presence_desired_state.rs | 12+++++++++---
Atests/services_hardening_presence_publication.rs | 650+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mtests/services_hardening_state_catalog.rs | 58+++++++++++++++++++++++++++++++++++++++++++++++++---------
Mtests/services_hardening_state_host.rs | 7+++++--
Mtests/services_hardening_state_repository_topology.rs | 23++++++++++++++++++++++-
Mtests/services_hardening_state_resilience.rs | 4++--
18 files changed, 4825 insertions(+), 46 deletions(-)

diff --git a/AGENTS.md b/AGENTS.md @@ -217,9 +217,21 @@ Durable desired state contains only the closed mode, document-presence bits, bounded target counts, queue capacity, and governed digests; it never stores relay URLs, rendered or signed events, attempts, schedules, or outcomes. -- Keep rendered and independently verified signed presence bytes, per-document - target state, attempt evidence, retry scheduling, and relay I/O with the next - ordered checkpoint. Never reconstruct exact committed bytes during retry. +- Build service-profile and application-handler presence only through the + governed typed Lib plans, then independently revalidate exact ID, signature, + author, kind, authored time, ordered tags, and content before persistence. + Commit each verified signed byte sequence and the complete immutable target + inventory before an injected presence sink can observe it. Preserve the + caller's sealed exact-byte capability across an unknown commit result so + reconciliation never depends on re-signing. +- Persist presence `submitted` in a short SQLx-owned transaction before relay + I/O, keep the remote await outside every transaction, and atomically append + the closed attempt outcome with target scheduling, outbox disposition, and + lease release. Cancellation or lost acknowledgement becomes durable + `unknown` before retry of the same retained bytes. Recover an expired stale + desired generation to `unknown` and supersede it without retry entropy so it + cannot block the current generation indefinitely. Never reconstruct exact + committed presence bytes or persist raw relay diagnostics. ## 7. Configuration, identity, state, and process boundaries diff --git a/README b/README @@ -359,9 +359,39 @@ must still match the latest durable configuration binding. The table stores no relay URL, filesystem path, rendered event, signature, authored time, delivery attempt, or network result, and its triggers reject deletion and ungoverned updates. Rendering, signing, target delivery state, retries, and relay I/O remain -owned by Step 205. The exact machine contract is +separate from this semantic-intent commit. The exact machine contract is [`presence_desired_state.v1.json`](contracts/services_hardening/presence_desired_state.v1.json). +## Durable exact-byte presence publication + +`build_rhi_signed_presence_documents` consumes one committed desired-state +generation, the independently revalidated complete presence authority, the +matching decrypted service identity, caller-injected authored time, and +caller-injected Schnorr auxiliary entropy. It builds the kind-0 service profile +through the typed `radroots_event` profile plan and cross-checks the kind-31990 +application-handler plan produced by `radroots_nostr` against the typed +`radroots_event` plan. Every signed document is then independently parsed under +the bounded NIP-01 wire limits and revalidated for event ID, signature, author, +kind, time, ordered tags, and exact content before exposure. + +Schema v10 retains each independently verified exact signed byte sequence, +digest, document identity, and complete immutable target inventory. One short +SQLx-owned transaction commits those bytes and initial target state before any +injected `RhiExactPresenceSink` can observe them. A second short transaction +commits `submitted` before remote I/O. The remote await owns no transaction; +the first commit borrows rather than consumes its sealed document set, so an +unknown commit result can be reconciled by replaying the same retained bytes; +the resulting closed outcome, immutable attempt evidence, target schedule, +outbox disposition, and lease release are committed together afterward. +Cancellation or acknowledgement loss after `submitted` becomes durable +`unknown` before the same retained bytes can be retried. Expired work from a +superseded desired generation is recorded as `unknown` and the old outbox is +superseded without sampling retry entropy, preventing stale leases from +blocking the current generation forever. No retry parses, rebuilds, +reserializes, or re-signs an event, and no raw relay diagnostic is persisted. +The exact machine contract and deterministic signed vectors are +[`presence_publication.v1.json`](contracts/services_hardening/presence_publication.v1.json). + ## Atomic reconciliation finalization commit `RhiReconciliationAttemptRepository::commit_finalization` borrows one sealed, @@ -620,17 +650,18 @@ credential remain excluded from the state-backup contract. That earlier wave-two boundary proof did not itself execute a backup; the governed resilience boundary below now owns the actual backup and recovery mechanics. -RHI's state surface is partitioned into eighteen distinct non-forgeable typed +RHI's state surface is partitioned into twenty-one distinct non-forgeable typed repository capabilities bound to one already-opened `RhiStateHost`. Their closed topology covers source results, cursors, completions, admitted signed events, canonical mutations, provenance, dirty generations, reconciliation jobs and attempts, immutable manifests/projections/reports, supersession, -signed attestation events, publication outbox/targets/attempts, and desired -presence. Each capability has one exact backing-table identity and an +signed attestation events, publication outbox/targets/attempts, desired +presence, and presence outbox/targets/attempts. Each capability has one exact backing-table identity and an append-only, compare-and-swap, or immutable write class. It exposes no raw pool, connection, transaction, SQL, path, or cloneable write authority. -Schema migration and verification, CRUD behavior, backup/restore, network I/O, -and task supervision remain owned by their later ordered checkpoints. +Schema migration and verification, remaining repository behavior, +backup/restore, service task ownership, and admin routing remain owned by their +later ordered checkpoints. RHI now composes the shared SQLx-owned resilience boundary without exposing a second database authority. A writable `RhiStateHost` can capture one governed diff --git a/contracts/api_baselines/rhi.txt b/contracts/api_baselines/rhi.txt @@ -129,6 +129,16 @@ impl rhi::RhiIdentityRole pub const fn rhi::RhiIdentityRole::as_str(self) -> &'static str pub enum rhi::RhiMetricsCommandV1 pub rhi::RhiMetricsCommandV1::Snapshot +pub enum rhi::RhiPresenceAttemptOutcome +pub rhi::RhiPresenceAttemptOutcome::Accepted +pub rhi::RhiPresenceAttemptOutcome::AuthRequired +pub rhi::RhiPresenceAttemptOutcome::Failed +pub rhi::RhiPresenceAttemptOutcome::RateLimited +pub rhi::RhiPresenceAttemptOutcome::Rejected +pub rhi::RhiPresenceAttemptOutcome::Submitted +pub rhi::RhiPresenceAttemptOutcome::Unknown +impl rhi::RhiPresenceAttemptOutcome +pub const fn rhi::RhiPresenceAttemptOutcome::code(self) -> &'static str pub enum rhi::RhiPresenceCommandV1 pub rhi::RhiPresenceCommandV1::Desired pub rhi::RhiPresenceCommandV1::Refresh @@ -153,6 +163,41 @@ pub rhi::RhiPresenceDocumentKind::ApplicationHandler pub rhi::RhiPresenceDocumentKind::ServiceProfile impl rhi::RhiPresenceDocumentKind pub const fn rhi::RhiPresenceDocumentKind::code(self) -> &'static str +pub enum rhi::RhiPresenceOutboxState +pub rhi::RhiPresenceOutboxState::Blocked +pub rhi::RhiPresenceOutboxState::Complete +pub rhi::RhiPresenceOutboxState::Leased +pub rhi::RhiPresenceOutboxState::Pending +pub rhi::RhiPresenceOutboxState::Superseded +impl rhi::RhiPresenceOutboxState +pub const fn rhi::RhiPresenceOutboxState::code(self) -> &'static str +pub enum rhi::RhiPresencePublicationErrorKind +pub rhi::RhiPresencePublicationErrorKind::ClockUnavailable +pub rhi::RhiPresencePublicationErrorKind::CommitOutcomeUnknown +pub rhi::RhiPresencePublicationErrorKind::DesiredStateMismatch +pub rhi::RhiPresencePublicationErrorKind::EntropyUnavailable +pub rhi::RhiPresencePublicationErrorKind::IdentityMismatch +pub rhi::RhiPresencePublicationErrorKind::InvalidInput +pub rhi::RhiPresencePublicationErrorKind::InvalidMode +pub rhi::RhiPresencePublicationErrorKind::Invariant +pub rhi::RhiPresencePublicationErrorKind::LeaseLost +pub rhi::RhiPresencePublicationErrorKind::NotReady +pub rhi::RhiPresencePublicationErrorKind::RenderingFailed +pub rhi::RhiPresencePublicationErrorKind::Storage +pub rhi::RhiPresencePublicationErrorKind::VerificationFailed +impl rhi::RhiPresencePublicationErrorKind +pub const fn rhi::RhiPresencePublicationErrorKind::code(self) -> &'static str +pub enum rhi::RhiPresenceTargetState +pub rhi::RhiPresenceTargetState::Accepted +pub rhi::RhiPresenceTargetState::AuthRequired +pub rhi::RhiPresenceTargetState::Failed +pub rhi::RhiPresenceTargetState::Pending +pub rhi::RhiPresenceTargetState::RateLimited +pub rhi::RhiPresenceTargetState::Rejected +pub rhi::RhiPresenceTargetState::Submitted +pub rhi::RhiPresenceTargetState::Unknown +impl rhi::RhiPresenceTargetState +pub const fn rhi::RhiPresenceTargetState::code(self) -> &'static str pub enum rhi::RhiPublicationAttemptEvidenceErrorKind pub rhi::RhiPublicationAttemptEvidenceErrorKind::InvalidAttemptNumber pub rhi::RhiPublicationAttemptEvidenceErrorKind::InvalidTargetOrdinal @@ -434,6 +479,9 @@ pub rhi::RhiStateRepositoryKind::DesiredPresence pub rhi::RhiStateRepositoryKind::DirtyTrade pub rhi::RhiStateRepositoryKind::EvidenceManifest pub rhi::RhiStateRepositoryKind::Mutation +pub rhi::RhiStateRepositoryKind::PresenceAttempt +pub rhi::RhiStateRepositoryKind::PresenceOutbox +pub rhi::RhiStateRepositoryKind::PresenceTarget pub rhi::RhiStateRepositoryKind::Projection pub rhi::RhiStateRepositoryKind::Provenance pub rhi::RhiStateRepositoryKind::PublicationAttempt @@ -728,6 +776,15 @@ impl rhi::RhiNormalizedConfigDigest pub const fn rhi::RhiNormalizedConfigDigest::as_bytes(&self) -> &[u8; 32] impl core::fmt::Debug for rhi::RhiNormalizedConfigDigest pub fn rhi::RhiNormalizedConfigDigest::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiPreparedPresenceAttempt +impl rhi::RhiPreparedPresenceAttempt +pub const fn rhi::RhiPreparedPresenceAttempt::attempt_id(&self) -> rhi::RhiPresenceAttemptId +pub const fn rhi::RhiPreparedPresenceAttempt::attempt_number(&self) -> u16 +pub const fn rhi::RhiPreparedPresenceAttempt::deadline_at(&self) -> rhi::RhiPresenceUnixMilliseconds +pub const fn rhi::RhiPreparedPresenceAttempt::exact_signed_event_bytes(&self) -> &[u8] +pub fn rhi::RhiPreparedPresenceAttempt::relay_id(&self) -> &str +impl core::fmt::Debug for rhi::RhiPreparedPresenceAttempt +pub fn rhi::RhiPreparedPresenceAttempt::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result pub struct rhi::RhiPreparedPublicationAttempt impl rhi::RhiPreparedPublicationAttempt pub const fn rhi::RhiPreparedPublicationAttempt::attempt_id(&self) -> rhi::RhiPublicationAttemptId @@ -737,6 +794,31 @@ pub const fn rhi::RhiPreparedPublicationAttempt::exact_signed_event_bytes(&self) pub fn rhi::RhiPreparedPublicationAttempt::relay_id(&self) -> &str impl core::fmt::Debug for rhi::RhiPreparedPublicationAttempt pub fn rhi::RhiPreparedPublicationAttempt::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiPresenceAttemptCommit +impl rhi::RhiPresenceAttemptCommit +pub const fn rhi::RhiPresenceAttemptCommit::attempt_id(self) -> rhi::RhiPresenceAttemptId +pub const fn rhi::RhiPresenceAttemptCommit::attempt_number(self) -> u16 +pub const fn rhi::RhiPresenceAttemptCommit::outbox_id(self) -> rhi::RhiPresenceOutboxId +pub const fn rhi::RhiPresenceAttemptCommit::outbox_state(self) -> rhi::RhiPresenceOutboxState +pub const fn rhi::RhiPresenceAttemptCommit::outcome(self) -> rhi::RhiPresenceAttemptOutcome +pub const fn rhi::RhiPresenceAttemptCommit::target_ordinal(self) -> u8 +pub const fn rhi::RhiPresenceAttemptCommit::target_state(self) -> rhi::RhiPresenceTargetState +pub struct rhi::RhiPresenceAttemptId(_) +impl rhi::RhiPresenceAttemptId +pub const fn rhi::RhiPresenceAttemptId::as_bytes(&self) -> &[u8; 32] +impl core::fmt::Debug for rhi::RhiPresenceAttemptId +pub fn rhi::RhiPresenceAttemptId::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiPresenceAttemptRepository<'host> +impl rhi::RhiPresenceAttemptRepository<'_> +pub const fn rhi::RhiPresenceAttemptRepository<'_>::descriptor(&self) -> rhi::RhiStateRepositoryDescriptor +pub const fn rhi::RhiPresenceAttemptRepository<'_>::kind(&self) -> rhi::RhiStateRepositoryKind +impl core::fmt::Debug for rhi::RhiPresenceAttemptRepository<'_> +pub fn rhi::RhiPresenceAttemptRepository<'_>::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiPresenceCommitOutcome +impl rhi::RhiPresenceCommitOutcome +pub const fn rhi::RhiPresenceCommitOutcome::changed(self) -> bool +pub const fn rhi::RhiPresenceCommitOutcome::desired_state(self) -> rhi::RhiPresenceDesiredState +pub const fn rhi::RhiPresenceCommitOutcome::document_count(self) -> u8 pub struct rhi::RhiPresenceDesiredAuthority impl rhi::RhiPresenceDesiredAuthority pub const fn rhi::RhiPresenceDesiredAuthority::desired_sha256(&self) -> &[u8; 32] @@ -774,6 +856,48 @@ pub const fn rhi::RhiPresenceDesiredState::target_count(self) -> u8 pub const fn rhi::RhiPresenceDesiredState::target_set_sha256(&self) -> &[u8; 32] impl core::fmt::Debug for rhi::RhiPresenceDesiredState pub fn rhi::RhiPresenceDesiredState::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiPresenceLease +impl rhi::RhiPresenceLease +pub const fn rhi::RhiPresenceLease::expires_at(&self) -> rhi::RhiPresenceUnixMilliseconds +pub const fn rhi::RhiPresenceLease::outbox_id(&self) -> rhi::RhiPresenceOutboxId +impl core::fmt::Debug for rhi::RhiPresenceLease +pub fn rhi::RhiPresenceLease::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiPresenceLeaseOwner(_) +impl rhi::RhiPresenceLeaseOwner +pub fn rhi::RhiPresenceLeaseOwner::from_bytes([u8; 16]) -> core::result::Result<Self, rhi::RhiPresencePublicationError> +impl core::fmt::Debug for rhi::RhiPresenceLeaseOwner +pub fn rhi::RhiPresenceLeaseOwner::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiPresenceOutboxId(_) +impl rhi::RhiPresenceOutboxId +pub const fn rhi::RhiPresenceOutboxId::as_bytes(&self) -> &[u8; 32] +impl core::fmt::Debug for rhi::RhiPresenceOutboxId +pub fn rhi::RhiPresenceOutboxId::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiPresenceOutboxRepository<'host> +impl rhi::RhiPresenceOutboxRepository<'_> +pub async fn rhi::RhiPresenceOutboxRepository<'_>::claim_next_presence(&self, rhi::RhiPresenceLeaseOwner, rhi::RhiPresenceUnixMilliseconds) -> core::result::Result<core::option::Option<rhi::RhiPresenceLease>, rhi::RhiPresencePublicationError> +pub async fn rhi::RhiPresenceOutboxRepository<'_>::commit_signed_presence(&self, &rhi::RhiSignedPresenceDocuments, rhi::RhiPresenceUnixMilliseconds) -> core::result::Result<rhi::RhiPresenceCommitOutcome, rhi::RhiPresencePublicationError> +pub async fn rhi::RhiPresenceOutboxRepository<'_>::execute_next_presence(&self, rhi::RhiPresenceLeaseOwner, &rhi::RhiTimeEntropyAdapters, &dyn rhi::RhiExactPresenceSink) -> core::result::Result<core::option::Option<rhi::RhiPresenceAttemptCommit>, rhi::RhiPresencePublicationError> +pub async fn rhi::RhiPresenceOutboxRepository<'_>::prepare_next_presence_target(&self, rhi::RhiPresenceLease, rhi::RhiPresenceUnixMilliseconds) -> core::result::Result<rhi::RhiPreparedPresenceAttempt, rhi::RhiPresencePublicationError> +pub async fn rhi::RhiPresenceOutboxRepository<'_>::record_presence_outcome(&self, &rhi::RhiPreparedPresenceAttempt, rhi::RhiPresenceUnixMilliseconds, rhi::RhiPresenceAttemptOutcome, rhi::RhiPresenceRetryDelayMilliseconds) -> core::result::Result<rhi::RhiPresenceAttemptCommit, rhi::RhiPresencePublicationError> +pub async fn rhi::RhiPresenceOutboxRepository<'_>::recover_one_expired_presence(&self, &rhi::RhiTimeEntropyAdapters, rhi::RhiPresenceUnixMilliseconds) -> core::result::Result<bool, rhi::RhiPresencePublicationError> +impl rhi::RhiPresenceOutboxRepository<'_> +pub const fn rhi::RhiPresenceOutboxRepository<'_>::descriptor(&self) -> rhi::RhiStateRepositoryDescriptor +pub const fn rhi::RhiPresenceOutboxRepository<'_>::kind(&self) -> rhi::RhiStateRepositoryKind +impl core::fmt::Debug for rhi::RhiPresenceOutboxRepository<'_> +pub fn rhi::RhiPresenceOutboxRepository<'_>::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiPresencePublicationError +impl rhi::RhiPresencePublicationError +pub const fn rhi::RhiPresencePublicationError::code(self) -> &'static str +pub const fn rhi::RhiPresencePublicationError::kind(self) -> rhi::RhiPresencePublicationErrorKind +impl core::error::Error for rhi::RhiPresencePublicationError +impl core::fmt::Debug for rhi::RhiPresencePublicationError +pub fn rhi::RhiPresencePublicationError::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +impl core::fmt::Display for rhi::RhiPresencePublicationError +pub fn rhi::RhiPresencePublicationError::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiPresenceRetryDelayMilliseconds(_) +impl rhi::RhiPresenceRetryDelayMilliseconds +pub const fn rhi::RhiPresenceRetryDelayMilliseconds::get(self) -> u64 +pub fn rhi::RhiPresenceRetryDelayMilliseconds::new(u64) -> core::result::Result<Self, rhi::RhiPresencePublicationError> pub struct rhi::RhiPresenceTarget impl rhi::RhiPresenceTarget pub const fn rhi::RhiPresenceTarget::ordinal(&self) -> u8 @@ -781,6 +905,16 @@ pub fn rhi::RhiPresenceTarget::relay_id(&self) -> &str pub const fn rhi::RhiPresenceTarget::required(&self) -> bool impl core::fmt::Debug for rhi::RhiPresenceTarget pub fn rhi::RhiPresenceTarget::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiPresenceTargetRepository<'host> +impl rhi::RhiPresenceTargetRepository<'_> +pub const fn rhi::RhiPresenceTargetRepository<'_>::descriptor(&self) -> rhi::RhiStateRepositoryDescriptor +pub const fn rhi::RhiPresenceTargetRepository<'_>::kind(&self) -> rhi::RhiStateRepositoryKind +impl core::fmt::Debug for rhi::RhiPresenceTargetRepository<'_> +pub fn rhi::RhiPresenceTargetRepository<'_>::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiPresenceUnixMilliseconds(_) +impl rhi::RhiPresenceUnixMilliseconds +pub const fn rhi::RhiPresenceUnixMilliseconds::get(self) -> u64 +pub fn rhi::RhiPresenceUnixMilliseconds::new(u64) -> core::result::Result<Self, rhi::RhiPresencePublicationError> pub struct rhi::RhiProjectionRepository<'host> impl rhi::RhiProjectionRepository<'_> pub const fn rhi::RhiProjectionRepository<'_>::descriptor(&self) -> rhi::RhiStateRepositoryDescriptor @@ -1343,6 +1477,22 @@ pub const fn rhi::RhiSignedEvidenceAttestation::signed_event_sha256(&self) -> &[ pub fn rhi::RhiSignedEvidenceAttestation::statement_digest(&self) -> [u8; 32] impl core::fmt::Debug for rhi::RhiSignedEvidenceAttestation pub fn rhi::RhiSignedEvidenceAttestation::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiSignedPresenceDocument +impl rhi::RhiSignedPresenceDocument +pub const fn rhi::RhiSignedPresenceDocument::created_at_unix_seconds(&self) -> u64 +pub const fn rhi::RhiSignedPresenceDocument::event_id(&self) -> &[u8; 32] +pub const fn rhi::RhiSignedPresenceDocument::kind(&self) -> rhi::RhiPresenceDocumentKind +pub fn rhi::RhiSignedPresenceDocument::signed_event_bytes(&self) -> &[u8] +pub const fn rhi::RhiSignedPresenceDocument::signed_event_sha256(&self) -> &[u8; 32] +impl core::fmt::Debug for rhi::RhiSignedPresenceDocument +pub fn rhi::RhiSignedPresenceDocument::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct rhi::RhiSignedPresenceDocuments +impl rhi::RhiSignedPresenceDocuments +pub const fn rhi::RhiSignedPresenceDocuments::desired_state(&self) -> rhi::RhiPresenceDesiredState +pub fn rhi::RhiSignedPresenceDocuments::documents(&self) -> &[rhi::RhiSignedPresenceDocument] +pub fn rhi::RhiSignedPresenceDocuments::target_count(&self) -> usize +impl core::fmt::Debug for rhi::RhiSignedPresenceDocuments +pub fn rhi::RhiSignedPresenceDocuments::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result pub struct rhi::RhiSourceCompletionRepository<'host> impl rhi::RhiSourceCompletionRepository<'_> pub const fn rhi::RhiSourceCompletionRepository<'_>::descriptor(&self) -> rhi::RhiStateRepositoryDescriptor @@ -1436,6 +1586,9 @@ pub const fn rhi::RhiStateRepositories<'host>::desired_presence(&self) -> rhi::R pub const fn rhi::RhiStateRepositories<'host>::dirty_trades(&self) -> rhi::RhiDirtyTradeRepository<'host> pub const fn rhi::RhiStateRepositories<'host>::evidence_manifests(&self) -> rhi::RhiEvidenceManifestRepository<'host> pub const fn rhi::RhiStateRepositories<'host>::mutations(&self) -> rhi::RhiMutationRepository<'host> +pub const fn rhi::RhiStateRepositories<'host>::presence_attempts(&self) -> rhi::RhiPresenceAttemptRepository<'host> +pub const fn rhi::RhiStateRepositories<'host>::presence_outbox(&self) -> rhi::RhiPresenceOutboxRepository<'host> +pub const fn rhi::RhiStateRepositories<'host>::presence_targets(&self) -> rhi::RhiPresenceTargetRepository<'host> pub const fn rhi::RhiStateRepositories<'host>::projections(&self) -> rhi::RhiProjectionRepository<'host> pub const fn rhi::RhiStateRepositories<'host>::provenance(&self) -> rhi::RhiProvenanceRepository<'host> pub const fn rhi::RhiStateRepositories<'host>::publication_attempts(&self) -> rhi::RhiPublicationAttemptRepository<'host> @@ -1639,6 +1792,9 @@ pub const rhi::RHI_ENCRYPTED_IDENTITY_ENVELOPE_MAX_BYTES: usize pub const rhi::RHI_MIGRATION_CATALOG_SHA256: [u8; 32] pub const rhi::RHI_PRESENCE_DESIRED_CONTRACT_VERSION: u32 pub const rhi::RHI_PRESENCE_DESIRED_MAX_TARGETS: usize +pub const rhi::RHI_PRESENCE_MAX_ATTEMPTS: u16 +pub const rhi::RHI_PRESENCE_PUBLICATION_CONTRACT_VERSION: u32 +pub const rhi::RHI_PRESENCE_SIGNED_EVENT_MAX_BYTES: usize pub const rhi::RHI_PROVIDER_CONTRACT_VERSION: u32 pub const rhi::RHI_PUBLICATION_ATTEMPT_EVIDENCE_CONTRACT_VERSION: u32 pub const rhi::RHI_PUBLICATION_ATTEMPT_NUMBER_MAXIMUM: u16 @@ -1671,6 +1827,9 @@ pub const rhi::RHI_STATE_REPOSITORY_CONTRACT_VERSION: u32 pub const rhi::RHI_STATE_REPOSITORY_COUNT: usize pub const rhi::RHI_STATE_SCHEMA_CATALOG_SHA256: [u8; 32] pub const rhi::RHI_STATE_SCHEMA_VERSION: u32 +pub const rhi::RHI_STATE_SCHEMA_VERSION_10_MIGRATION_SHA256: [u8; 32] +pub const rhi::RHI_STATE_SCHEMA_VERSION_10_OBJECT_COUNT: u32 +pub const rhi::RHI_STATE_SCHEMA_VERSION_10_SHA256: [u8; 32] pub const rhi::RHI_STATE_SCHEMA_VERSION_1_OBJECT_COUNT: u32 pub const rhi::RHI_STATE_SCHEMA_VERSION_1_SHA256: [u8; 32] pub const rhi::RHI_STATE_SCHEMA_VERSION_2_MIGRATION_SHA256: [u8; 32] @@ -1714,6 +1873,8 @@ pub trait rhi::RhiCredentialAccess: core::marker::Send + core::marker::Sync pub fn rhi::RhiCredentialAccess::resolve_existing(&self, &rhi::RhiRuntimeContext, &rhi::RhiIdentityEnvelopeBinding) -> core::result::Result<rhi::RhiWrappingCredential, rhi::RhiCredentialResolutionError> impl rhi::RhiCredentialAccess for rhi::CanonicalRhiCredentialAccess pub fn rhi::CanonicalRhiCredentialAccess::resolve_existing(&self, &rhi::RhiRuntimeContext, &rhi::RhiIdentityEnvelopeBinding) -> core::result::Result<rhi::RhiWrappingCredential, rhi::RhiCredentialResolutionError> +pub trait rhi::RhiExactPresenceSink: core::marker::Send + core::marker::Sync +pub fn rhi::RhiExactPresenceSink::submit_exact<'a>(&'a self, &'a rhi::RhiPreparedPresenceAttempt) -> radroots_transport::source::BoxFuture<'a, rhi::RhiPresenceAttemptOutcome> pub trait rhi::RhiExactPublicationSink: core::marker::Send + core::marker::Sync pub fn rhi::RhiExactPublicationSink::submit_exact<'a>(&'a self, &'a rhi::RhiPreparedPublicationAttempt) -> radroots_transport::source::BoxFuture<'a, rhi::RhiPublicationAttemptOutcome> pub trait rhi::RhiIdentityAccess: core::marker::Send + core::marker::Sync @@ -1724,6 +1885,7 @@ pub fn rhi::admit_rhi_trade_mutation_event(rhi::RhiTradeMutationAdmissionLimits, pub async fn rhi::apply_rhi_configuration(&rhi::RhiRuntimeContext, &rhi::RhiConfigDocumentV1, &rhi::RhiConfigDocumentV1, radroots_service_sqlite::migration::MigrationAppliedAtUnixSeconds, &radroots_service_sqlite::migration::MigrationBuildIdentity) -> core::result::Result<rhi::RhiConfigApplyOutcome, rhi::RhiConfigApplyError> pub fn rhi::attest_projection_claim(&radroots_trade::trade_contract_v1::RadrootsTradeProjectionV1, &radroots_event::id::MutationId, &rhi::TradeAgreementAttestationPolicy) -> core::result::Result<rhi::TradeAgreementAttestationReportV1, rhi::TradeAgreementAttestationError> pub fn rhi::build_rhi_signed_evidence_attestation(rhi::RhiReconciliationFinalizationFence, &rhi::RhiDecryptedIdentity, radroots_service_host::time::UnixTimeSeconds, &dyn radroots_service_host::entropy::EntropySource, core::option::Option<rhi::RhiEvidenceAttestationSupersession>) -> core::result::Result<rhi::RhiSignedEvidenceAttestation, rhi::RhiReconciliationAttestationError> +pub fn rhi::build_rhi_signed_presence_documents(rhi::RhiPresenceDesiredCommitOutcome, &rhi::RhiPresenceDesiredAuthority, &rhi::RhiDecryptedIdentity, radroots_service_host::time::UnixTimeSeconds, &dyn radroots_service_host::entropy::EntropySource) -> core::result::Result<rhi::RhiSignedPresenceDocuments, rhi::RhiPresencePublicationError> pub fn rhi::evaluate_rhi_reconciliation_claim(rhi::RhiReconciliationProjection, radroots_event::id::MutationId) -> rhi::RhiReconciliationEvaluation pub async fn rhi::finalize_rhi_state_restore(rhi::RhiStagedStateRestore) -> core::result::Result<(), rhi::RhiStateMaintenanceError> pub async fn rhi::ingest_rhi_trade_source(&rhi::RhiStateRepositories<'_>, &rhi::RhiTransportAdapters, &rhi::RhiConfigDocumentV1, &str, radroots_event::id::TradeId, rhi::RhiTradeSourceAttempt) -> core::result::Result<rhi::RhiTradeSourceIngestOutcome, rhi::RhiTradeSourceIngestError> @@ -1741,10 +1903,11 @@ pub fn rhi::resolve_rhi_runtime_context(&radroots_runtime_paths::roots::Radroots pub fn rhi::resolve_rhi_wrapping_credential(&rhi::RhiRuntimeContext, &rhi::RhiIdentityEnvelopeBinding) -> core::result::Result<rhi::RhiWrappingCredential, rhi::RhiCredentialResolutionError> pub fn rhi::rhi_migration_catalog() -> core::result::Result<radroots_service_sqlite::migration::MigrationCatalog, rhi::RhiStateCatalogError> pub fn rhi::rhi_schema_catalog() -> core::result::Result<radroots_service_sqlite::integrity::catalog::SchemaCatalog, rhi::RhiStateCatalogError> -pub const fn rhi::rhi_state_repository_descriptors() -> &'static [rhi::RhiStateRepositoryDescriptor; 18] +pub const fn rhi::rhi_state_repository_descriptors() -> &'static [rhi::RhiStateRepositoryDescriptor; 21] pub async fn rhi::stage_rhi_state_restore(&rhi::RhiRuntimeContext, &rhi::RhiStateMetadata, rhi::RhiVerifiedStateBackup) -> core::result::Result<rhi::RhiStagedStateRestore, rhi::RhiStateMaintenanceError> pub fn rhi::trade_mutation_subscription_kinds() -> alloc::vec::Vec<u32> pub fn rhi::validate_rhi_presence_desired_authority(&rhi::RhiConfigDocumentV1, &rhi::RhiPresenceDesiredAuthority) -> core::result::Result<(), rhi::RhiPresenceDesiredError> +pub fn rhi::validate_rhi_signed_presence_documents(&rhi::RhiSignedPresenceDocuments, &rhi::RhiPresenceDesiredAuthority) -> core::result::Result<(), rhi::RhiPresencePublicationError> pub fn rhi::validate_rhi_state_catalogs(&radroots_service_sqlite::migration::MigrationCatalog, &radroots_service_sqlite::integrity::catalog::SchemaCatalog) -> core::result::Result<(), rhi::RhiStateCatalogError> pub fn rhi::verify_rhi_state_backup(&[u8], radroots_service_sqlite::backup::manifest::BackupManifestSha256, &std::path::Path, &rhi::RhiStateMetadata, core::num::nonzero::NonZeroU64) -> core::result::Result<rhi::RhiVerifiedStateBackup, rhi::RhiStateMaintenanceError> pub type rhi::RhiReconciliationCoverage = radroots_trade::evidence::RadrootsTradeEvidenceCoverageV1 diff --git a/contracts/services_hardening/presence_desired_state.v1.json b/contracts/services_hardening/presence_desired_state.v1.json @@ -68,7 +68,7 @@ "before_relay_io": true, "commit_outcome_unknown": "typed_and_not_retried_as_uncommitted_without_reread" }, - "step_205_deferrals": [ + "separate_publication_authority": [ "rendered_typed_events", "exact_signed_bytes", "per_document_target_state", diff --git a/contracts/services_hardening/presence_publication.v1.json b/contracts/services_hardening/presence_publication.v1.json @@ -0,0 +1,115 @@ +{ + "schema": "radroots.rhi.presence-publication", + "schema_version": 1, + "contract_version": 1, + "step": 205, + "service": "rhi", + "construction": { + "desired_state": "exact_committed_presence_desired_generation", + "profile_plan": "typed_radroots_event_authored_profile", + "application_handler_plan": "typed_radroots_nostr_application_handler_cross_checked_with_radroots_event_plan", + "identity": "decrypted_identity_bound_to_current_service_public_key", + "authored_time": "caller_injected_whole_unix_seconds", + "entropy": "caller_injected_for_signing_only", + "validation_input": "sealed_documents_and_public_desired_authority_only", + "independent_validation": [ + "nip01_bounded_parse", + "event_id", + "signature", + "author", + "kind", + "authored_time", + "ordered_tags", + "exact_content" + ] + }, + "documents": [ + { + "document_kind": "service_profile", + "nostr_kind": 0, + "tags": [], + "content": { + "name": "rhi", + "display_name": "Radroots RHI", + "about": "Radroots evidence reconciliation and attestation service", + "bot": true + } + }, + { + "document_kind": "application_handler", + "nostr_kind": 31990, + "tags": [["d", "rhi"], ["k", "3441"]], + "content": { + "name": "Radroots RHI", + "about": "Radroots evidence reconciliation and attestation service" + } + } + ], + "durable_workflow": { + "schema_version": 10, + "outbox_table": "presence_outbox", + "target_table": "presence_targets", + "attempt_table": "presence_attempts", + "write_classes": { + "presence_outbox": "compare_and_swap", + "presence_targets": "compare_and_swap", + "presence_attempts": "append_only" + }, + "exact_signed_bytes_committed_before_io": true, + "target_submitted_committed_before_io": true, + "remote_io_inside_sql_transaction": false, + "outbox_states": ["pending", "leased", "complete", "blocked", "superseded"], + "target_states": ["pending", "submitted", "accepted", "rejected", "rate_limited", "auth_required", "failed", "unknown"], + "attempt_outcomes": ["accepted", "rejected", "rate_limited", "auth_required", "failed", "unknown"], + "commit_outcome_unknown": "caller_retains_sealed_exact_bytes_and_rereads_exact_identity_before_retry", + "cancelled_or_expired_submitted": "durably_record_unknown_then_retry_same_exact_bytes", + "desired_generation_change": "reject_unexpired_stale_lease_then_record_expired_submitted_as_unknown_and_supersede_without_retry_entropy" + }, + "resource_bounds": { + "maximum_signed_event_bytes": 32768, + "maximum_targets_per_document": 32, + "maximum_attempts_per_target": 100, + "maximum_authored_unix_seconds": 9223372036854775807, + "initial_backoff_milliseconds": 250, + "maximum_backoff_milliseconds": 30000, + "attempt_deadline_milliseconds": 15000, + "maximum_unix_milliseconds": 9223372036854775807 + }, + "reference_vector": { + "secret_key_fixture": "0101010101010101010101010101010101010101010101010101010101010101", + "service_public_key": "1b84c5567b126440995d3ed5aaba0565d71e1834604819ff9c17f5e9d5dd078f", + "authored_at_unix_s": 1725000100, + "profile": { + "event_id": "90dea82b6b86799ae2936a256bd08158fa37a0f4f1e34c4061c652bc3492aad7", + "exact_signed_event_sha256": "f627a73ae309da87889cf5a636ddd06d9ea6bb879da972572928fe4940d885d2", + "exact_signed_event_json": "{\"id\":\"90dea82b6b86799ae2936a256bd08158fa37a0f4f1e34c4061c652bc3492aad7\",\"pubkey\":\"1b84c5567b126440995d3ed5aaba0565d71e1834604819ff9c17f5e9d5dd078f\",\"created_at\":1725000100,\"kind\":0,\"tags\":[],\"content\":\"{\\\"name\\\":\\\"rhi\\\",\\\"display_name\\\":\\\"Radroots RHI\\\",\\\"about\\\":\\\"Radroots evidence reconciliation and attestation service\\\",\\\"bot\\\":true}\",\"sig\":\"abee71a0dd9fdfb04821964f701d9d27e5198c0b7f85c9ff2b13d88b5c844ee6a296dd60d37827327e9267a79a1aa401029049778523add3f991ddefbf447d92\"}" + }, + "application_handler": { + "event_id": "a8a2897ca4aee8a57888b5df1251d2b8cd486c7d9e9d02d25b50c8a79e9a5278", + "exact_signed_event_sha256": "6d32dd2e70368406696a8ec447bbe9d6083a75f9144b3ca1371758d2373247d4", + "exact_signed_event_json": "{\"id\":\"a8a2897ca4aee8a57888b5df1251d2b8cd486c7d9e9d02d25b50c8a79e9a5278\",\"pubkey\":\"1b84c5567b126440995d3ed5aaba0565d71e1834604819ff9c17f5e9d5dd078f\",\"created_at\":1725000100,\"kind\":31990,\"tags\":[[\"d\",\"rhi\"],[\"k\",\"3441\"]],\"content\":\"{\\\"name\\\":\\\"Radroots RHI\\\",\\\"about\\\":\\\"Radroots evidence reconciliation and attestation service\\\"}\",\"sig\":\"eda18aa5589b9fafbdb70c72093e1e25125af5d1ec017c984f01cea785e005b6f652519d3fac8f6a30a8705def1d0f8f87fcee08188004c835807fa29eb764d0\"}" + } + }, + "effects": { + "sqlite": "typed_repository_transactions_only", + "network": "injected_exact_presence_sink_only", + "clock": "injected", + "entropy": "injected", + "filesystem": false, + "task_spawn": false, + "ambient_process_state": false + }, + "deferred": [ + "service_runtime_task_ownership", + "unix_admin_routes", + "status_and_metrics_projection", + "process_supervision", + "rcld_promotion", + "parent_pin_alignment", + "nix", + "oci", + "signing", + "publication", + "deployment" + ] +} diff --git a/contracts/services_hardening/state_repository_topology.v1.json b/contracts/services_hardening/state_repository_topology.v1.json @@ -2,7 +2,7 @@ "schema": "radroots.rhi.state-repository-topology", "schema_version": 1, "contract_version": 1, - "repository_count": 18, + "repository_count": 21, "construction": "sealed_to_open_rhi_state_host", "raw_sqlite_authority_exposed": false, "repositories": [ @@ -23,11 +23,14 @@ { "kind": "publication_outbox", "backing_table": "publication_outbox", "write_class": "compare_and_swap" }, { "kind": "publication_target", "backing_table": "publication_targets", "write_class": "compare_and_swap" }, { "kind": "publication_attempt", "backing_table": "publication_attempts", "write_class": "append_only" }, - { "kind": "desired_presence", "backing_table": "presence_desired_state", "write_class": "compare_and_swap" } + { "kind": "desired_presence", "backing_table": "presence_desired_state", "write_class": "compare_and_swap" }, + { "kind": "presence_outbox", "backing_table": "presence_outbox", "write_class": "compare_and_swap" }, + { "kind": "presence_target", "backing_table": "presence_targets", "write_class": "compare_and_swap" }, + { "kind": "presence_attempt", "backing_table": "presence_attempts", "write_class": "append_only" } ], "deferred_behavior": [ "later_schema_migrations", - "non_job_repository_crud", + "remaining_repository_crud", "backup_restore_and_recovery", "network_io", "task_supervision" diff --git a/src/lib.rs b/src/lib.rs @@ -9,6 +9,7 @@ mod features; mod identity_credential; mod identity_envelope; mod presence_desired; +mod presence_publication; mod publication; mod publication_attempt; mod publication_execution; @@ -76,6 +77,16 @@ pub use presence_desired::{ RhiPresenceDesiredErrorKind, RhiPresenceDesiredMode, RhiPresenceDesiredState, RhiPresenceDocumentKind, RhiPresenceTarget, validate_rhi_presence_desired_authority, }; +pub use presence_publication::{ + RHI_PRESENCE_MAX_ATTEMPTS, RHI_PRESENCE_PUBLICATION_CONTRACT_VERSION, + RHI_PRESENCE_SIGNED_EVENT_MAX_BYTES, RhiExactPresenceSink, RhiPreparedPresenceAttempt, + RhiPresenceAttemptCommit, RhiPresenceAttemptId, RhiPresenceAttemptOutcome, + RhiPresenceCommitOutcome, RhiPresenceLease, RhiPresenceLeaseOwner, RhiPresenceOutboxId, + RhiPresenceOutboxState, RhiPresencePublicationError, RhiPresencePublicationErrorKind, + RhiPresenceRetryDelayMilliseconds, RhiPresenceTargetState, RhiPresenceUnixMilliseconds, + RhiSignedPresenceDocument, RhiSignedPresenceDocuments, build_rhi_signed_presence_documents, + validate_rhi_signed_presence_documents, +}; pub use publication::{ RHI_PUBLICATION_CONTRACT_VERSION, RHI_PUBLICATION_MAX_ATTEMPTS, RHI_PUBLICATION_MAX_TARGETS, RhiPublicationAuthority, RhiPublicationError, RhiPublicationErrorKind, RhiPublicationMode, @@ -196,8 +207,10 @@ pub use state_catalog::{ RHI_STATE_SCHEMA_VERSION_7_SHA256, RHI_STATE_SCHEMA_VERSION_8_MIGRATION_SHA256, RHI_STATE_SCHEMA_VERSION_8_OBJECT_COUNT, RHI_STATE_SCHEMA_VERSION_8_SHA256, RHI_STATE_SCHEMA_VERSION_9_MIGRATION_SHA256, RHI_STATE_SCHEMA_VERSION_9_OBJECT_COUNT, - RHI_STATE_SCHEMA_VERSION_9_SHA256, RhiStateCatalogError, RhiStateCatalogErrorKind, - rhi_migration_catalog, rhi_schema_catalog, validate_rhi_state_catalogs, + RHI_STATE_SCHEMA_VERSION_9_SHA256, RHI_STATE_SCHEMA_VERSION_10_MIGRATION_SHA256, + RHI_STATE_SCHEMA_VERSION_10_OBJECT_COUNT, RHI_STATE_SCHEMA_VERSION_10_SHA256, + RhiStateCatalogError, RhiStateCatalogErrorKind, rhi_migration_catalog, rhi_schema_catalog, + validate_rhi_state_catalogs, }; pub use state_config::{ RHI_CONFIG_BINDING_MAX_GENERATIONS, RhiConfigApplyError, RhiConfigApplyErrorKind, @@ -222,7 +235,8 @@ pub use state_metadata::{ pub use state_repository::{ RHI_STATE_REPOSITORY_CONTRACT_VERSION, RHI_STATE_REPOSITORY_COUNT, RhiDesiredPresenceRepository, RhiDirtyTradeRepository, RhiEvidenceManifestRepository, - RhiMutationRepository, RhiProjectionRepository, RhiProvenanceRepository, + RhiMutationRepository, RhiPresenceAttemptRepository, RhiPresenceOutboxRepository, + RhiPresenceTargetRepository, RhiProjectionRepository, RhiProvenanceRepository, RhiPublicationAttemptRepository, RhiPublicationOutboxRepository, RhiPublicationTargetRepository, RhiReconciliationAttemptRepository, RhiReconciliationJobRepository, RhiReportRepository, RhiSignedAttestationEventRepository, diff --git a/src/presence_desired.rs b/src/presence_desired.rs @@ -196,6 +196,10 @@ impl RhiPresenceDesiredAuthority { pub const fn desired_sha256(&self) -> &[u8; 32] { &self.desired_sha256 } + + pub(crate) fn service_public_key(&self) -> &str { + &self.service_public_key + } } impl fmt::Debug for RhiPresenceDesiredAuthority { @@ -844,6 +848,13 @@ fn matches_authority( && state.desired_sha256 == authority.desired_sha256 } +pub(crate) fn presence_authority_matches_state( + state: RhiPresenceDesiredState, + authority: &RhiPresenceDesiredAuthority, +) -> bool { + matches_authority(state, authority) +} + fn has_document(authority: &RhiPresenceDesiredAuthority, kind: RhiPresenceDocumentKind) -> bool { authority.document_kinds.contains(&kind) } diff --git a/src/presence_publication.rs b/src/presence_publication.rs @@ -0,0 +1,3091 @@ +//! Exact signed presence documents and durable target delivery. + +use core::fmt; +use std::error::Error; + +use nostr::{EventBuilder, Kind, PublicKey as NostrPublicKey, Tag, Timestamp}; +use radroots_event::{ + GenericEventDraft, + envelope::{EventEnvelope, kind::KIND_APPLICATION_HANDLER}, + profile::AuthoredProfile, + wire::{EventWireLimits, Nip01EventWire}, +}; +use radroots_event_codec::authoring::AuthoredEventPlan; +use radroots_nostr::event::{ + ApplicationHandlerSpec, Verification, build_application_handler, verify, verify_id, +}; +use radroots_service_host::{EntropySource, UnixTimeSeconds}; +use radroots_service_sqlite::{ + ServiceSqliteTransaction, ServiceSqliteTransactionError, ServiceSqliteTransactionErrorKind, +}; +use radroots_transport::BoxFuture; +use sha2::{Digest as _, Sha256}; +use sqlx::Row as _; +use zeroize::Zeroizing; + +use crate::{ + RHI_PRESENCE_DESIRED_MAX_TARGETS, RhiDecryptedIdentity, RhiJitterBoundMilliseconds, + RhiPresenceDesiredAuthority, RhiPresenceDesiredCommitOutcome, RhiPresenceDesiredMode, + RhiPresenceDesiredState, RhiPresenceDocumentKind, RhiPresenceOutboxRepository, + RhiPresenceTarget, RhiRuntimeAdapterErrorKind, RhiStateHostMode, RhiTimeEntropyAdapters, + presence_desired::presence_authority_matches_state, +}; + +/// Exact version of the signed presence and delivery contract. +pub const RHI_PRESENCE_PUBLICATION_CONTRACT_VERSION: u32 = 1; +/// Maximum canonical bytes retained for one signed presence event. +pub const RHI_PRESENCE_SIGNED_EVENT_MAX_BYTES: usize = 32 * 1024; +/// Maximum durable attempts for one presence target. +pub const RHI_PRESENCE_MAX_ATTEMPTS: u16 = 100; + +const PROFILE_NAME: &str = "rhi"; +const PROFILE_DISPLAY_NAME: &str = "Radroots RHI"; +const PROFILE_ABOUT: &str = "Radroots evidence reconciliation and attestation service"; +const APPLICATION_HANDLER_IDENTIFIER: &str = "rhi"; +const APPLICATION_HANDLER_KINDS: [u32; 1] = [3441]; +const INITIAL_BACKOFF_MILLISECONDS: u64 = 250; +const MAXIMUM_BACKOFF_MILLISECONDS: u64 = 30_000; +const ATTEMPT_DEADLINE_MILLISECONDS: u64 = 15_000; +const MAX_UNIX_MILLISECONDS: u64 = i64::MAX as u64; +const OUTBOX_ID_DOMAIN: &[u8] = b"radroots.rhi.presence_outbox.v1\0"; +const ATTEMPT_ID_DOMAIN: &[u8] = b"radroots.rhi.presence_attempt.v1\0"; + +const READ_DESIRED_BINDING_SQL: &str = r#"SELECT singleton, generation, enabled, + profile, application_handler, + CASE WHEN typeof(target_set_sha256) = 'blob' AND length(target_set_sha256) = 32 + THEN target_set_sha256 ELSE NULL END AS target_set_sha256, + target_count, required_target_count, queue_capacity, + CASE WHEN typeof(desired_sha256) = 'blob' AND length(desired_sha256) = 32 + THEN desired_sha256 ELSE NULL END AS desired_sha256 +FROM presence_desired_state +LIMIT 2"#; + +const READ_GENERATION_OUTBOX_SQL: &str = r#"SELECT + CASE WHEN typeof(outbox_id) = 'blob' AND length(outbox_id) = 32 + THEN outbox_id ELSE NULL END AS outbox_id, + desired_generation, + length(CAST(document_kind AS BLOB)) AS document_kind_bytes, + substr(document_kind, 1, 20) AS document_kind, + CASE WHEN typeof(desired_sha256) = 'blob' AND length(desired_sha256) = 32 + THEN desired_sha256 ELSE NULL END AS desired_sha256, + CASE WHEN typeof(target_set_sha256) = 'blob' AND length(target_set_sha256) = 32 + THEN target_set_sha256 ELSE NULL END AS target_set_sha256, + CASE WHEN typeof(event_id) = 'blob' AND length(event_id) = 32 + THEN event_id ELSE NULL END AS event_id, + CASE WHEN typeof(event_sha256) = 'blob' AND length(event_sha256) = 32 + THEN event_sha256 ELSE NULL END AS event_sha256, + length(exact_signed_event_bytes) AS exact_signed_event_bytes_length, + substr(exact_signed_event_bytes, 1, 32769) AS exact_signed_event_bytes, + authored_at_unix_s, + length(CAST(service_public_key AS BLOB)) AS service_public_key_bytes, + substr(service_public_key, 1, 65) AS service_public_key, + target_count, required_target_count, max_attempts, + initial_backoff_ms, maximum_backoff_ms, attempt_deadline_ms, + length(CAST(state AS BLOB)) AS state_bytes, substr(state, 1, 11) AS state, + revision, next_attempt_unix_ms, + CASE WHEN lease_owner IS NULL THEN NULL ELSE length(lease_owner) END AS lease_owner_bytes, + CASE WHEN lease_owner IS NULL THEN NULL ELSE substr(lease_owner, 1, 17) END AS lease_owner, + lease_expires_unix_ms, created_at_unix_ms, updated_at_unix_ms +FROM presence_outbox +WHERE desired_generation = ? +ORDER BY CASE document_kind + WHEN 'service_profile' THEN 0 + WHEN 'application_handler' THEN 1 + ELSE 2 END +LIMIT 3"#; + +const READ_OUTBOX_SQL: &str = r#"SELECT + CASE WHEN typeof(outbox_id) = 'blob' AND length(outbox_id) = 32 + THEN outbox_id ELSE NULL END AS outbox_id, + desired_generation, + length(CAST(document_kind AS BLOB)) AS document_kind_bytes, + substr(document_kind, 1, 20) AS document_kind, + CASE WHEN typeof(desired_sha256) = 'blob' AND length(desired_sha256) = 32 + THEN desired_sha256 ELSE NULL END AS desired_sha256, + CASE WHEN typeof(target_set_sha256) = 'blob' AND length(target_set_sha256) = 32 + THEN target_set_sha256 ELSE NULL END AS target_set_sha256, + CASE WHEN typeof(event_id) = 'blob' AND length(event_id) = 32 + THEN event_id ELSE NULL END AS event_id, + CASE WHEN typeof(event_sha256) = 'blob' AND length(event_sha256) = 32 + THEN event_sha256 ELSE NULL END AS event_sha256, + length(exact_signed_event_bytes) AS exact_signed_event_bytes_length, + substr(exact_signed_event_bytes, 1, 32769) AS exact_signed_event_bytes, + authored_at_unix_s, + length(CAST(service_public_key AS BLOB)) AS service_public_key_bytes, + substr(service_public_key, 1, 65) AS service_public_key, + target_count, required_target_count, max_attempts, + initial_backoff_ms, maximum_backoff_ms, attempt_deadline_ms, + length(CAST(state AS BLOB)) AS state_bytes, substr(state, 1, 11) AS state, + revision, next_attempt_unix_ms, + CASE WHEN lease_owner IS NULL THEN NULL ELSE length(lease_owner) END AS lease_owner_bytes, + CASE WHEN lease_owner IS NULL THEN NULL ELSE substr(lease_owner, 1, 17) END AS lease_owner, + lease_expires_unix_ms, created_at_unix_ms, updated_at_unix_ms +FROM presence_outbox +WHERE outbox_id = ? +LIMIT 2"#; + +const READ_CLAIMABLE_OUTBOX_SQL: &str = r#"SELECT + CASE WHEN typeof(outbox.outbox_id) = 'blob' AND length(outbox.outbox_id) = 32 + THEN outbox.outbox_id ELSE NULL END AS outbox_id +FROM presence_outbox AS outbox +JOIN presence_desired_state AS desired + ON desired.singleton = 1 + AND desired.generation = outbox.desired_generation + AND desired.desired_sha256 = outbox.desired_sha256 +WHERE outbox.state = 'pending' AND outbox.next_attempt_unix_ms <= ? +ORDER BY outbox.next_attempt_unix_ms, outbox.created_at_unix_ms, + CASE outbox.document_kind + WHEN 'service_profile' THEN 0 + WHEN 'application_handler' THEN 1 + ELSE 2 END, + outbox.outbox_id +LIMIT 1"#; + +const READ_EXPIRED_OUTBOX_SQL: &str = r#"SELECT + CASE WHEN typeof(outbox_id) = 'blob' AND length(outbox_id) = 32 + THEN outbox_id ELSE NULL END AS outbox_id +FROM presence_outbox +WHERE state = 'leased' AND lease_expires_unix_ms <= ? +ORDER BY lease_expires_unix_ms, created_at_unix_ms, document_kind, outbox_id +LIMIT 1"#; + +const READ_TARGETS_SQL: &str = r#"SELECT target_ordinal, + length(CAST(relay_id AS BLOB)) AS relay_id_bytes, + substr(relay_id, 1, 65) AS relay_id, required, + length(CAST(state AS BLOB)) AS state_bytes, substr(state, 1, 14) AS state, + revision, attempt_count, next_attempt_unix_ms, + CASE WHEN last_attempt_id IS NULL THEN NULL ELSE length(last_attempt_id) END + AS last_attempt_id_bytes, + CASE WHEN last_attempt_id IS NULL THEN NULL ELSE substr(last_attempt_id, 1, 33) END + AS last_attempt_id, + updated_at_unix_ms +FROM presence_targets +WHERE outbox_id = ? +ORDER BY target_ordinal +LIMIT 33"#; + +const INSERT_OUTBOX_SQL: &str = r#"INSERT INTO presence_outbox ( + outbox_id, desired_generation, document_kind, desired_sha256, + target_set_sha256, event_id, event_sha256, exact_signed_event_bytes, + authored_at_unix_s, service_public_key, target_count, required_target_count, + max_attempts, initial_backoff_ms, maximum_backoff_ms, attempt_deadline_ms, + state, revision, next_attempt_unix_ms, lease_owner, lease_expires_unix_ms, + created_at_unix_ms, updated_at_unix_ms +) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 100, 250, 30000, 15000, + 'pending', 1, ?, NULL, NULL, ?, ?)"#; + +const INSERT_TARGET_SQL: &str = r#"INSERT INTO presence_targets ( + outbox_id, target_ordinal, relay_id, required, state, revision, + attempt_count, next_attempt_unix_ms, last_attempt_id, updated_at_unix_ms +) VALUES (?, ?, ?, ?, 'pending', 1, 0, ?, NULL, ?)"#; + +const SUPERSEDE_OUTBOX_SQL: &str = r#"UPDATE presence_outbox +SET state = 'superseded', revision = revision + 1, + next_attempt_unix_ms = NULL, lease_owner = NULL, + lease_expires_unix_ms = NULL, updated_at_unix_ms = ? +WHERE desired_generation != ? AND state IN ('pending', 'blocked') + AND updated_at_unix_ms <= ?"#; + +const COUNT_STALE_LEASES_SQL: &str = r#"SELECT COUNT(*) AS lease_count +FROM presence_outbox +WHERE desired_generation != ? AND state = 'leased'"#; + +const READ_MAX_AUTHORED_SQL: &str = r#"SELECT MAX(authored_at_unix_s) AS authored_at_unix_s +FROM presence_outbox WHERE document_kind = ?"#; + +const CLAIM_OUTBOX_SQL: &str = r#"UPDATE presence_outbox +SET state = 'leased', revision = revision + 1, + next_attempt_unix_ms = NULL, lease_owner = ?, lease_expires_unix_ms = ?, + updated_at_unix_ms = ? +WHERE outbox_id = ? AND revision = ? AND state = 'pending' + AND next_attempt_unix_ms <= ? AND updated_at_unix_ms <= ?"#; + +const PREPARE_TARGET_SQL: &str = r#"UPDATE presence_targets +SET state = 'submitted', revision = revision + 1, + attempt_count = attempt_count + 1, next_attempt_unix_ms = NULL, + last_attempt_id = ?, updated_at_unix_ms = ? +WHERE outbox_id = ? AND target_ordinal = ? AND revision = ? + AND state IN ('pending', 'failed', 'rate_limited', 'unknown') + AND attempt_count < ? AND next_attempt_unix_ms <= ? + AND updated_at_unix_ms <= ?"#; + +const INSERT_ATTEMPT_SQL: &str = r#"INSERT INTO presence_attempts ( + attempt_id, outbox_id, target_ordinal, attempt_number, event_sha256, + lease_owner, started_at_unix_ms, finished_at_unix_ms, outcome, result_code +) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)"#; + +const READ_ATTEMPT_SQL: &str = r#"SELECT + CASE WHEN typeof(attempt_id) = 'blob' AND length(attempt_id) = 32 + THEN attempt_id ELSE NULL END AS attempt_id, + CASE WHEN typeof(outbox_id) = 'blob' AND length(outbox_id) = 32 + THEN outbox_id ELSE NULL END AS outbox_id, + target_ordinal, attempt_number, + CASE WHEN typeof(event_sha256) = 'blob' AND length(event_sha256) = 32 + THEN event_sha256 ELSE NULL END AS event_sha256, + CASE WHEN typeof(lease_owner) = 'blob' AND length(lease_owner) = 16 + THEN lease_owner ELSE NULL END AS lease_owner, + started_at_unix_ms, finished_at_unix_ms, + length(CAST(outcome AS BLOB)) AS outcome_bytes, substr(outcome, 1, 14) AS outcome, + length(CAST(result_code AS BLOB)) AS result_code_bytes, + substr(result_code, 1, 65) AS result_code +FROM presence_attempts +WHERE attempt_id = ? +LIMIT 2"#; + +const UPDATE_TARGET_OUTCOME_SQL: &str = r#"UPDATE presence_targets +SET state = ?, revision = revision + 1, next_attempt_unix_ms = ?, + updated_at_unix_ms = ? +WHERE outbox_id = ? AND target_ordinal = ? AND revision = ? + AND state = 'submitted' AND attempt_count = ? AND last_attempt_id = ?"#; + +const UPDATE_OUTBOX_AFTER_ATTEMPT_SQL: &str = r#"UPDATE presence_outbox +SET state = ?, revision = revision + 1, next_attempt_unix_ms = ?, + lease_owner = NULL, lease_expires_unix_ms = NULL, updated_at_unix_ms = ? +WHERE outbox_id = ? AND revision = ? AND state = 'leased' + AND lease_owner = ? AND lease_expires_unix_ms = ? + AND lease_expires_unix_ms > ? AND updated_at_unix_ms <= ?"#; + +const UPDATE_OUTBOX_RECOVERY_SQL: &str = r#"UPDATE presence_outbox +SET state = ?, revision = revision + 1, next_attempt_unix_ms = ?, + lease_owner = NULL, lease_expires_unix_ms = NULL, updated_at_unix_ms = ? +WHERE outbox_id = ? AND revision = ? AND state = 'leased' + AND lease_owner = ? AND lease_expires_unix_ms = ? + AND lease_expires_unix_ms <= ? AND updated_at_unix_ms <= ?"#; + +/// Stable source-free signed-presence and delivery failure class. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum RhiPresencePublicationErrorKind { + InvalidMode, + InvalidInput, + DesiredStateMismatch, + IdentityMismatch, + EntropyUnavailable, + RenderingFailed, + VerificationFailed, + NotReady, + LeaseLost, + Invariant, + ClockUnavailable, + Storage, + CommitOutcomeUnknown, +} + +impl RhiPresencePublicationErrorKind { + /// Returns the stable machine-readable failure code. + #[must_use] + pub const fn code(self) -> &'static str { + match self { + Self::InvalidMode => "presence_publication_mode_invalid", + Self::InvalidInput => "presence_publication_input_invalid", + Self::DesiredStateMismatch => "presence_publication_desired_state_mismatch", + Self::IdentityMismatch => "presence_publication_identity_mismatch", + Self::EntropyUnavailable => "presence_publication_entropy_unavailable", + Self::RenderingFailed => "presence_publication_rendering_failed", + Self::VerificationFailed => "presence_publication_verification_failed", + Self::NotReady => "presence_publication_not_ready", + Self::LeaseLost => "presence_publication_lease_lost", + Self::Invariant => "presence_publication_invariant_failed", + Self::ClockUnavailable => "presence_publication_clock_unavailable", + Self::Storage => "presence_publication_storage_failed", + Self::CommitOutcomeUnknown => "presence_publication_commit_outcome_unknown", + } + } +} + +/// Redacted source-free signed-presence and delivery failure. +#[derive(Clone, Copy, PartialEq, Eq)] +pub struct RhiPresencePublicationError { + kind: RhiPresencePublicationErrorKind, +} + +impl RhiPresencePublicationError { + /// Returns the stable failure class. + #[must_use] + pub const fn kind(self) -> RhiPresencePublicationErrorKind { + self.kind + } + + /// Returns the stable machine-readable failure code. + #[must_use] + pub const fn code(self) -> &'static str { + self.kind.code() + } +} + +impl fmt::Display for RhiPresencePublicationError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str(match self.kind { + RhiPresencePublicationErrorKind::InvalidMode => { + "RHI presence publication requires writable state" + } + RhiPresencePublicationErrorKind::InvalidInput => { + "RHI presence publication input is invalid" + } + RhiPresencePublicationErrorKind::DesiredStateMismatch => { + "RHI presence publication desired state does not match" + } + RhiPresencePublicationErrorKind::IdentityMismatch => { + "RHI presence publication identity does not match" + } + RhiPresencePublicationErrorKind::EntropyUnavailable => { + "RHI presence publication entropy is unavailable" + } + RhiPresencePublicationErrorKind::RenderingFailed => { + "RHI presence publication rendering failed" + } + RhiPresencePublicationErrorKind::VerificationFailed => { + "RHI presence publication verification failed" + } + RhiPresencePublicationErrorKind::NotReady => "RHI presence work is not ready", + RhiPresencePublicationErrorKind::LeaseLost => { + "RHI presence publication lease is no longer authoritative" + } + RhiPresencePublicationErrorKind::Invariant => { + "RHI presence publication state invariant failed" + } + RhiPresencePublicationErrorKind::ClockUnavailable => { + "RHI presence publication clock is unavailable" + } + RhiPresencePublicationErrorKind::Storage => { + "RHI presence publication transaction failed" + } + RhiPresencePublicationErrorKind::CommitOutcomeUnknown => { + "RHI presence publication commit outcome is unknown" + } + }) + } +} + +impl fmt::Debug for RhiPresencePublicationError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("RhiPresencePublicationError") + .field("kind", &self.kind) + .finish() + } +} + +impl Error for RhiPresencePublicationError {} + +/// One independently verified exact signed presence document. +pub struct RhiSignedPresenceDocument { + kind: RhiPresenceDocumentKind, + event_id: [u8; 32], + signed_event_sha256: [u8; 32], + signed_event_bytes: Box<[u8]>, + created_at_unix_seconds: u64, +} + +impl RhiSignedPresenceDocument { + /// Returns the closed document kind. + #[must_use] + pub const fn kind(&self) -> RhiPresenceDocumentKind { + self.kind + } + + /// Returns the independently verified NIP-01 event identifier. + #[must_use] + pub const fn event_id(&self) -> &[u8; 32] { + &self.event_id + } + + /// Returns the digest of the exact retained bytes. + #[must_use] + pub const fn signed_event_sha256(&self) -> &[u8; 32] { + &self.signed_event_sha256 + } + + /// Returns the exact canonical signed bytes. + #[must_use] + pub fn signed_event_bytes(&self) -> &[u8] { + &self.signed_event_bytes + } + + /// Returns the injected authored timestamp. + #[must_use] + pub const fn created_at_unix_seconds(&self) -> u64 { + self.created_at_unix_seconds + } +} + +impl fmt::Debug for RhiSignedPresenceDocument { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("RhiSignedPresenceDocument") + .field("kind", &self.kind) + .field("signed_event_bytes", &self.signed_event_bytes.len()) + .finish_non_exhaustive() + } +} + +/// Sealed exact document set bound to one committed desired-state generation. +pub struct RhiSignedPresenceDocuments { + desired_state: RhiPresenceDesiredState, + service_public_key: Box<str>, + targets: Box<[RhiPresenceTarget]>, + documents: Box<[RhiSignedPresenceDocument]>, +} + +impl RhiSignedPresenceDocuments { + /// Returns the exact committed desired state. + #[must_use] + pub const fn desired_state(&self) -> RhiPresenceDesiredState { + self.desired_state + } + + /// Returns the exact ordered signed-document inventory. + #[must_use] + pub fn documents(&self) -> &[RhiSignedPresenceDocument] { + &self.documents + } + + /// Returns the stable target count without exposing endpoints. + #[must_use] + pub fn target_count(&self) -> usize { + self.targets.len() + } + + fn retained_copy(&self) -> Self { + Self { + desired_state: self.desired_state, + service_public_key: self.service_public_key.clone(), + targets: self.targets.clone(), + documents: self + .documents + .iter() + .map(|document| RhiSignedPresenceDocument { + kind: document.kind, + event_id: document.event_id, + signed_event_sha256: document.signed_event_sha256, + signed_event_bytes: document.signed_event_bytes.clone(), + created_at_unix_seconds: document.created_at_unix_seconds, + }) + .collect(), + } + } +} + +impl fmt::Debug for RhiSignedPresenceDocuments { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("RhiSignedPresenceDocuments") + .field("desired_generation", &self.desired_state.generation()) + .field("document_count", &self.documents.len()) + .field("target_count", &self.targets.len()) + .finish_non_exhaustive() + } +} + +/// Opaque identity of one exact durable presence outbox. +#[derive(Clone, Copy, PartialEq, Eq, Hash)] +pub struct RhiPresenceOutboxId([u8; 32]); + +impl RhiPresenceOutboxId { + /// Returns the exact identity bytes. + #[must_use] + pub const fn as_bytes(&self) -> &[u8; 32] { + &self.0 + } +} + +impl fmt::Debug for RhiPresenceOutboxId { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("RhiPresenceOutboxId([redacted])") + } +} + +/// Opaque identity of one durable presence attempt. +#[derive(Clone, Copy, PartialEq, Eq, Hash)] +pub struct RhiPresenceAttemptId([u8; 32]); + +impl RhiPresenceAttemptId { + /// Returns the exact identity bytes. + #[must_use] + pub const fn as_bytes(&self) -> &[u8; 32] { + &self.0 + } +} + +impl fmt::Debug for RhiPresenceAttemptId { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("RhiPresenceAttemptId([redacted])") + } +} + +/// Bounded wall-clock instant in whole Unix milliseconds. +#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord)] +pub struct RhiPresenceUnixMilliseconds(u64); + +impl RhiPresenceUnixMilliseconds { + /// Validates one injected instant against SQLite's signed representation. + pub fn new(value: u64) -> Result<Self, RhiPresencePublicationError> { + if value > MAX_UNIX_MILLISECONDS { + return Err(failure(RhiPresencePublicationErrorKind::InvalidInput)); + } + Ok(Self(value)) + } + + /// Returns the exact instant. + #[must_use] + pub const fn get(self) -> u64 { + self.0 + } +} + +/// Bounded caller-injected retry delay in whole milliseconds. +#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord)] +pub struct RhiPresenceRetryDelayMilliseconds(u64); + +impl RhiPresenceRetryDelayMilliseconds { + /// Validates one delay against the absolute retry ceiling. + pub fn new(value: u64) -> Result<Self, RhiPresencePublicationError> { + if value > MAXIMUM_BACKOFF_MILLISECONDS { + return Err(failure(RhiPresencePublicationErrorKind::InvalidInput)); + } + Ok(Self(value)) + } + + /// Returns the exact delay. + #[must_use] + pub const fn get(self) -> u64 { + self.0 + } +} + +/// Stable process-local owner token for one compare-and-swap lease. +#[derive(Clone, Copy, PartialEq, Eq, Hash)] +pub struct RhiPresenceLeaseOwner([u8; 16]); + +impl RhiPresenceLeaseOwner { + /// Validates one injected nonzero owner token. + pub fn from_bytes(bytes: [u8; 16]) -> Result<Self, RhiPresencePublicationError> { + if bytes.iter().all(|byte| *byte == 0) { + return Err(failure(RhiPresencePublicationErrorKind::InvalidInput)); + } + Ok(Self(bytes)) + } +} + +impl fmt::Debug for RhiPresenceLeaseOwner { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("RhiPresenceLeaseOwner([redacted])") + } +} + +/// Stable durable presence-outbox lifecycle state. +#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)] +pub enum RhiPresenceOutboxState { + Pending, + Leased, + Complete, + Blocked, + Superseded, +} + +impl RhiPresenceOutboxState { + /// Returns the exact machine-contract spelling. + #[must_use] + pub const fn code(self) -> &'static str { + match self { + Self::Pending => "pending", + Self::Leased => "leased", + Self::Complete => "complete", + Self::Blocked => "blocked", + Self::Superseded => "superseded", + } + } +} + +/// Stable durable per-target presence state. +#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)] +pub enum RhiPresenceTargetState { + Pending, + Submitted, + Accepted, + Rejected, + RateLimited, + AuthRequired, + Failed, + Unknown, +} + +impl RhiPresenceTargetState { + /// Returns the exact machine-contract spelling. + #[must_use] + pub const fn code(self) -> &'static str { + match self { + Self::Pending => "pending", + Self::Submitted => "submitted", + Self::Accepted => "accepted", + Self::Rejected => "rejected", + Self::RateLimited => "rate_limited", + Self::AuthRequired => "auth_required", + Self::Failed => "failed", + Self::Unknown => "unknown", + } + } +} + +/// Closed result observed from one exact-byte relay submission. +#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)] +pub enum RhiPresenceAttemptOutcome { + Submitted, + Accepted, + Rejected, + RateLimited, + AuthRequired, + Failed, + Unknown, +} + +impl RhiPresenceAttemptOutcome { + /// Returns the exact machine-contract spelling. + #[must_use] + pub const fn code(self) -> &'static str { + match self { + Self::Submitted => "submitted", + Self::Accepted => "accepted", + Self::Rejected => "rejected", + Self::RateLimited => "rate_limited", + Self::AuthRequired => "auth_required", + Self::Failed => "failed", + Self::Unknown => "unknown", + } + } +} + +/// Confirmed result of committing one complete signed presence generation. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct RhiPresenceCommitOutcome { + desired_state: RhiPresenceDesiredState, + document_count: u8, + changed: bool, +} + +impl RhiPresenceCommitOutcome { + /// Returns the exact committed desired-state identity. + #[must_use] + pub const fn desired_state(self) -> RhiPresenceDesiredState { + self.desired_state + } + + /// Returns the exact committed document count. + #[must_use] + pub const fn document_count(self) -> u8 { + self.document_count + } + + /// Reports whether any durable presence workflow state changed. + #[must_use] + pub const fn changed(self) -> bool { + self.changed + } +} + +/// Non-forgeable compare-and-swap authority for one claimed presence outbox. +#[derive(PartialEq, Eq)] +pub struct RhiPresenceLease { + outbox: PresenceOutboxRecord, + owner: RhiPresenceLeaseOwner, + expires_at: RhiPresenceUnixMilliseconds, +} + +impl RhiPresenceLease { + /// Returns the exact claimed outbox identity. + #[must_use] + pub const fn outbox_id(&self) -> RhiPresenceOutboxId { + self.outbox.id + } + + /// Returns the exact lease expiry. + #[must_use] + pub const fn expires_at(&self) -> RhiPresenceUnixMilliseconds { + self.expires_at + } +} + +impl fmt::Debug for RhiPresenceLease { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("RhiPresenceLease") + .field("identity", &"[redacted]") + .field("revision", &self.outbox.revision) + .field("expires_at", &self.expires_at) + .finish() + } +} + +/// Exact signed bytes exposed only after durable Submitted state exists. +#[must_use = "prepared presence must be executed exactly or left for unknown recovery"] +pub struct RhiPreparedPresenceAttempt { + lease: RhiPresenceLease, + target: PresenceTargetRecord, + attempt_id: RhiPresenceAttemptId, + started_at: RhiPresenceUnixMilliseconds, + deadline_at: RhiPresenceUnixMilliseconds, + exact_signed_event_bytes: Box<[u8]>, +} + +impl RhiPreparedPresenceAttempt { + /// Returns the exact attempt identity. + #[must_use] + pub const fn attempt_id(&self) -> RhiPresenceAttemptId { + self.attempt_id + } + + /// Returns the stable relay identity without an endpoint or secret. + #[must_use] + pub fn relay_id(&self) -> &str { + &self.target.relay_id + } + + /// Returns the original committed bytes with no transformation. + #[must_use] + pub const fn exact_signed_event_bytes(&self) -> &[u8] { + &self.exact_signed_event_bytes + } + + /// Returns the absolute attempt deadline. + #[must_use] + pub const fn deadline_at(&self) -> RhiPresenceUnixMilliseconds { + self.deadline_at + } + + /// Returns the one-based durable attempt number. + #[must_use] + pub const fn attempt_number(&self) -> u16 { + self.target.attempt_count + } + + fn retry_upper_bound(&self) -> u64 { + retry_upper_bound(self.target.attempt_count) + } +} + +impl fmt::Debug for RhiPreparedPresenceAttempt { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("RhiPreparedPresenceAttempt") + .field("identity", &"[redacted]") + .field("target_ordinal", &self.target.ordinal) + .field("attempt_number", &self.target.attempt_count) + .field("deadline_at", &self.deadline_at) + .finish() + } +} + +/// Closed exact-byte transport boundary for durable presence delivery. +pub trait RhiExactPresenceSink: Send + Sync { + /// Submits the exact retained payload and returns one closed observation. + fn submit_exact<'a>( + &'a self, + attempt: &'a RhiPreparedPresenceAttempt, + ) -> BoxFuture<'a, RhiPresenceAttemptOutcome>; +} + +/// Confirmed durable result for one exact presence attempt. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct RhiPresenceAttemptCommit { + outbox_id: RhiPresenceOutboxId, + attempt_id: RhiPresenceAttemptId, + target_ordinal: u8, + attempt_number: u16, + outcome: RhiPresenceAttemptOutcome, + target_state: RhiPresenceTargetState, + outbox_state: RhiPresenceOutboxState, +} + +impl RhiPresenceAttemptCommit { + #[must_use] + pub const fn outbox_id(self) -> RhiPresenceOutboxId { + self.outbox_id + } + + #[must_use] + pub const fn attempt_id(self) -> RhiPresenceAttemptId { + self.attempt_id + } + + #[must_use] + pub const fn target_ordinal(self) -> u8 { + self.target_ordinal + } + + #[must_use] + pub const fn attempt_number(self) -> u16 { + self.attempt_number + } + + #[must_use] + pub const fn outcome(self) -> RhiPresenceAttemptOutcome { + self.outcome + } + + #[must_use] + pub const fn target_state(self) -> RhiPresenceTargetState { + self.target_state + } + + #[must_use] + pub const fn outbox_state(self) -> RhiPresenceOutboxState { + self.outbox_state + } +} + +/// Builds, signs, and independently revalidates one complete desired document set. +/// +/// Exactly 32 entropy bytes are consumed per enabled document. No clock, +/// persistence, relay, task, or network authority is acquired implicitly. +pub fn build_rhi_signed_presence_documents( + desired: RhiPresenceDesiredCommitOutcome, + authority: &RhiPresenceDesiredAuthority, + identity: &RhiDecryptedIdentity, + created_at: UnixTimeSeconds, + entropy: &dyn EntropySource, +) -> Result<RhiSignedPresenceDocuments, RhiPresencePublicationError> { + let state = desired.state(); + if created_at.get() > i64::MAX as u64 { + return Err(failure(RhiPresencePublicationErrorKind::InvalidInput)); + } + if !presence_authority_matches_state(state, authority) { + return Err(failure( + RhiPresencePublicationErrorKind::DesiredStateMismatch, + )); + } + if identity.public_identity().as_hex() != authority.service_public_key() { + return Err(failure(RhiPresencePublicationErrorKind::IdentityMismatch)); + } + let mut documents = Vec::with_capacity(authority.document_kinds().len()); + for kind in authority.document_kinds() { + let plan = presence_plan(*kind, identity.public_identity().as_hex(), created_at.get())?; + let mut auxiliary = Zeroizing::new([0_u8; 32]); + entropy + .fill_bytes(&mut auxiliary[..]) + .map_err(|_| failure(RhiPresencePublicationErrorKind::EntropyUnavailable))?; + let event = sign_plan(identity, &plan, &auxiliary)?; + let bytes = serde_json::to_vec(&event) + .map_err(|_| failure(RhiPresencePublicationErrorKind::RenderingFailed))?; + let verified = validate_signed_document( + *kind, + identity.public_identity().as_hex(), + created_at.get(), + &bytes, + )?; + documents.push(RhiSignedPresenceDocument { + kind: *kind, + event_id: *verified.id().as_bytes(), + signed_event_sha256: Sha256::digest(&bytes).into(), + signed_event_bytes: bytes.into_boxed_slice(), + created_at_unix_seconds: created_at.get(), + }); + } + if (state.mode() == RhiPresenceDesiredMode::Disabled && !documents.is_empty()) + || (state.mode() == RhiPresenceDesiredMode::Enabled + && documents.len() != authority.document_kinds().len()) + { + return Err(failure(RhiPresencePublicationErrorKind::Invariant)); + } + Ok(RhiSignedPresenceDocuments { + desired_state: state, + service_public_key: authority.service_public_key().into(), + targets: authority.targets().to_vec().into_boxed_slice(), + documents: documents.into_boxed_slice(), + }) +} + +/// Independently verifies every exact signed document against its bound intent. +pub fn validate_rhi_signed_presence_documents( + documents: &RhiSignedPresenceDocuments, + authority: &RhiPresenceDesiredAuthority, +) -> Result<(), RhiPresencePublicationError> { + if !presence_authority_matches_state(documents.desired_state, authority) + || documents.service_public_key.as_ref() != authority.service_public_key() + || documents.targets.as_ref() != authority.targets() + || documents.documents.len() != authority.document_kinds().len() + { + return Err(failure( + RhiPresencePublicationErrorKind::DesiredStateMismatch, + )); + } + for (document, kind) in documents.documents.iter().zip(authority.document_kinds()) { + if document.kind != *kind + || Sha256::digest(document.signed_event_bytes.as_ref()).as_slice() + != document.signed_event_sha256 + { + return Err(failure(RhiPresencePublicationErrorKind::VerificationFailed)); + } + let verified = validate_signed_document( + document.kind, + authority.service_public_key(), + document.created_at_unix_seconds, + &document.signed_event_bytes, + )?; + if verified.id().as_bytes() != &document.event_id { + return Err(failure(RhiPresencePublicationErrorKind::VerificationFailed)); + } + } + Ok(()) +} + +impl RhiPresenceOutboxRepository<'_> { + /// Atomically commits exact signed bytes and the complete immutable target set. + /// + /// A successful return means every exact payload and initial target state is + /// durable before any relay adapter can observe those bytes. The caller + /// retains the sealed exact-byte capability so an unknown commit result can + /// be reconciled by replaying this same value. Exact replay is idempotent. + pub async fn commit_signed_presence( + &self, + documents: &RhiSignedPresenceDocuments, + committed_at: RhiPresenceUnixMilliseconds, + ) -> Result<RhiPresenceCommitOutcome, RhiPresencePublicationError> { + require_writable(self)?; + let documents = documents.retained_copy(); + self.host() + .sqlite_host() + .transaction(move |transaction| { + Box::pin(async move { commit_signed(transaction, documents, committed_at).await }) + }) + .await + .map_err(map_transaction_error) + } + + /// Claims the oldest currently desired, due presence outbox. + pub async fn claim_next_presence( + &self, + owner: RhiPresenceLeaseOwner, + now: RhiPresenceUnixMilliseconds, + ) -> Result<Option<RhiPresenceLease>, RhiPresencePublicationError> { + require_writable(self)?; + self.host() + .sqlite_host() + .transaction(move |transaction| { + Box::pin(async move { claim_next(transaction, owner, now).await }) + }) + .await + .map_err(map_transaction_error) + } + + /// Persists Submitted before exposing the committed exact bytes. + pub async fn prepare_next_presence_target( + &self, + lease: RhiPresenceLease, + started_at: RhiPresenceUnixMilliseconds, + ) -> Result<RhiPreparedPresenceAttempt, RhiPresencePublicationError> { + require_writable(self)?; + self.host() + .sqlite_host() + .transaction(move |transaction| { + Box::pin(async move { prepare_next(transaction, lease, started_at).await }) + }) + .await + .map_err(map_transaction_error) + } + + /// Commits one closed observed outcome by exact compare-and-swap. + pub async fn record_presence_outcome( + &self, + prepared: &RhiPreparedPresenceAttempt, + finished_at: RhiPresenceUnixMilliseconds, + outcome: RhiPresenceAttemptOutcome, + retry_delay: RhiPresenceRetryDelayMilliseconds, + ) -> Result<RhiPresenceAttemptCommit, RhiPresencePublicationError> { + require_writable(self)?; + let input = + PresenceRecordInput::from_prepared(prepared, finished_at, outcome, retry_delay)?; + self.host() + .sqlite_host() + .transaction(move |transaction| { + Box::pin(async move { record_outcome(transaction, input).await }) + }) + .await + .map_err(map_transaction_error) + } + + /// Recovers at most one expired lease and records Unknown for submitted work. + pub async fn recover_one_expired_presence( + &self, + adapters: &RhiTimeEntropyAdapters, + now: RhiPresenceUnixMilliseconds, + ) -> Result<bool, RhiPresencePublicationError> { + require_writable(self)?; + let candidate = self + .host() + .sqlite_host() + .transaction(move |transaction| { + Box::pin(async move { read_expired_candidate(transaction, now).await }) + }) + .await + .map_err(map_transaction_error)?; + let Some(candidate) = candidate else { + return Ok(false); + }; + let delay = sample_retry_delay(adapters, candidate.retry_upper_bound)?; + self.host() + .sqlite_host() + .transaction(move |transaction| { + Box::pin(async move { recover_expired(transaction, candidate, now, delay).await }) + }) + .await + .map_err(map_transaction_error) + } + + /// Executes at most one exact-byte presence attempt through an injected sink. + /// + /// Cancellation after durable preparation leaves Submitted evidence. Lease + /// expiry recovery records Unknown before retrying the same retained bytes. + pub async fn execute_next_presence( + &self, + owner: RhiPresenceLeaseOwner, + adapters: &RhiTimeEntropyAdapters, + sink: &dyn RhiExactPresenceSink, + ) -> Result<Option<RhiPresenceAttemptCommit>, RhiPresencePublicationError> { + let now = presence_now(adapters)?; + self.recover_one_expired_presence(adapters, now).await?; + let Some(lease) = self.claim_next_presence(owner, now).await? else { + return Ok(None); + }; + let prepared = self.prepare_next_presence_target(lease, now).await?; + let outcome = match sink.submit_exact(&prepared).await { + RhiPresenceAttemptOutcome::Submitted => RhiPresenceAttemptOutcome::Unknown, + outcome => outcome, + }; + let finished_at = presence_now(adapters)?; + let retry_delay = if retryable_outcome(outcome) { + sample_retry_delay(adapters, prepared.retry_upper_bound())? + } else { + RhiPresenceRetryDelayMilliseconds(0) + }; + self.record_presence_outcome(&prepared, finished_at, outcome, retry_delay) + .await + .map(Some) + } +} + +fn require_writable( + repository: &RhiPresenceOutboxRepository<'_>, +) -> Result<(), RhiPresencePublicationError> { + if repository.host().mode() == RhiStateHostMode::ReadWriteExisting { + Ok(()) + } else { + Err(failure(RhiPresencePublicationErrorKind::InvalidMode)) + } +} + +fn presence_now( + adapters: &RhiTimeEntropyAdapters, +) -> Result<RhiPresenceUnixMilliseconds, RhiPresencePublicationError> { + let value = adapters.now_utc_milliseconds().map_err(|error| { + failure(match error.kind() { + RhiRuntimeAdapterErrorKind::WallClockUnavailable => { + RhiPresencePublicationErrorKind::ClockUnavailable + } + _ => RhiPresencePublicationErrorKind::ClockUnavailable, + }) + })?; + RhiPresenceUnixMilliseconds::new(value) + .map_err(|_| failure(RhiPresencePublicationErrorKind::ClockUnavailable)) +} + +fn sample_retry_delay( + adapters: &RhiTimeEntropyAdapters, + cap: u64, +) -> Result<RhiPresenceRetryDelayMilliseconds, RhiPresencePublicationError> { + if cap == 0 { + return Ok(RhiPresenceRetryDelayMilliseconds(0)); + } + let bound = RhiJitterBoundMilliseconds::new(cap) + .map_err(|_| failure(RhiPresencePublicationErrorKind::InvalidInput))?; + adapters + .sample_full_jitter(bound) + .map(|delay| RhiPresenceRetryDelayMilliseconds(delay.get())) + .map_err(|_| failure(RhiPresencePublicationErrorKind::EntropyUnavailable)) +} + +async fn commit_signed( + transaction: &mut ServiceSqliteTransaction<'_>, + documents: RhiSignedPresenceDocuments, + committed_at: RhiPresenceUnixMilliseconds, +) -> Result<RhiPresenceCommitOutcome, PresenceOperationError> { + let desired = read_desired_binding(transaction).await?; + validate_document_set(&documents, desired)?; + let current = read_generation_outboxes(transaction, desired.generation).await?; + if !current.is_empty() && generation_matches(&documents, &current)? { + return Ok(RhiPresenceCommitOutcome { + desired_state: documents.desired_state, + document_count: u8::try_from(documents.documents.len()) + .map_err(|_| PresenceOperationError::Invariant)?, + changed: false, + }); + } + if !current.is_empty() { + return Err(PresenceOperationError::Invariant); + } + if documents.desired_state.mode() == RhiPresenceDesiredMode::Enabled { + for document in &documents.documents { + let rows = sqlx::query(READ_MAX_AUTHORED_SQL) + .bind(document.kind.code()) + .fetch_all(&mut *transaction) + .await + .map_err(|_| PresenceOperationError::Storage)?; + let [row] = rows.as_slice() else { + return Err(PresenceOperationError::Invariant); + }; + if let Some(previous) = row + .try_get::<Option<i64>, _>("authored_at_unix_s") + .map_err(|_| PresenceOperationError::Invariant)? + { + let previous = + u64::try_from(previous).map_err(|_| PresenceOperationError::Invariant)?; + if document.created_at_unix_seconds <= previous { + return Err(PresenceOperationError::InvalidInput); + } + } + } + } + let lease_rows = sqlx::query(COUNT_STALE_LEASES_SQL) + .bind(i64_value(desired.generation)?) + .fetch_all(&mut *transaction) + .await + .map_err(|_| PresenceOperationError::Storage)?; + let [lease_row] = lease_rows.as_slice() else { + return Err(PresenceOperationError::Invariant); + }; + if lease_row + .try_get::<i64, _>("lease_count") + .ok() + .filter(|value| *value >= 0) + .ok_or(PresenceOperationError::Invariant)? + != 0 + { + return Err(PresenceOperationError::NotReady); + } + let superseded = sqlx::query(SUPERSEDE_OUTBOX_SQL) + .bind(i64_value(committed_at.get())?) + .bind(i64_value(desired.generation)?) + .bind(i64_value(committed_at.get())?) + .execute(&mut *transaction) + .await + .map_err(|_| PresenceOperationError::Storage)? + .rows_affected(); + + for document in &documents.documents { + let outbox_id = derive_outbox_id( + desired.generation, + document.kind, + desired.desired_sha256, + document.signed_event_sha256, + ); + let inserted = sqlx::query(INSERT_OUTBOX_SQL) + .bind(outbox_id.as_bytes().as_slice()) + .bind(i64_value(desired.generation)?) + .bind(document.kind.code()) + .bind(desired.desired_sha256.as_slice()) + .bind(desired.target_set_sha256.as_slice()) + .bind(document.event_id.as_slice()) + .bind(document.signed_event_sha256.as_slice()) + .bind(document.signed_event_bytes.as_ref()) + .bind(i64_value(document.created_at_unix_seconds)?) + .bind(documents.service_public_key.as_ref()) + .bind(i64::from(desired.target_count)) + .bind(i64::from(desired.required_target_count)) + .bind(i64_value(committed_at.get())?) + .bind(i64_value(committed_at.get())?) + .bind(i64_value(committed_at.get())?) + .execute(&mut *transaction) + .await + .map_err(|_| PresenceOperationError::Storage)?; + if inserted.rows_affected() != 1 { + return Err(PresenceOperationError::Invariant); + } + for target in &documents.targets { + let inserted = sqlx::query(INSERT_TARGET_SQL) + .bind(outbox_id.as_bytes().as_slice()) + .bind(i64::from(target.ordinal())) + .bind(target.relay_id()) + .bind(i64::from(target.required())) + .bind(i64_value(committed_at.get())?) + .bind(i64_value(committed_at.get())?) + .bind(i64_value(committed_at.get())?) + .execute(&mut *transaction) + .await + .map_err(|_| PresenceOperationError::Storage)?; + if inserted.rows_affected() != 1 { + return Err(PresenceOperationError::Invariant); + } + } + } + let committed = read_generation_outboxes(transaction, desired.generation).await?; + if !generation_matches(&documents, &committed)? { + return Err(PresenceOperationError::Invariant); + } + Ok(RhiPresenceCommitOutcome { + desired_state: documents.desired_state, + document_count: u8::try_from(documents.documents.len()) + .map_err(|_| PresenceOperationError::Invariant)?, + changed: !documents.documents.is_empty() || superseded != 0, + }) +} + +fn validate_document_set( + documents: &RhiSignedPresenceDocuments, + desired: PresenceDesiredBinding, +) -> Result<(), PresenceOperationError> { + let state = documents.desired_state; + if state.generation() != desired.generation + || state.mode() != desired.mode + || state.profile() != desired.profile + || state.application_handler() != desired.application_handler + || state.target_set_sha256() != &desired.target_set_sha256 + || state.target_count() != desired.target_count + || state.required_target_count() != desired.required_target_count + || state.queue_capacity() != desired.queue_capacity + || state.desired_sha256() != &desired.desired_sha256 + || documents.targets.len() != usize::from(desired.target_count) + || documents.documents.len() != desired.document_count() + || target_set_digest(&documents.targets)? != desired.target_set_sha256 + || !valid_public_key(&documents.service_public_key) + { + return Err(PresenceOperationError::DesiredStateMismatch); + } + let expected = desired.document_kinds(); + for (document, kind) in documents.documents.iter().zip(expected) { + if document.kind != kind + || document.signed_event_bytes.is_empty() + || document.signed_event_bytes.len() > RHI_PRESENCE_SIGNED_EVENT_MAX_BYTES + || <[u8; 32]>::from(Sha256::digest(&document.signed_event_bytes)) + != document.signed_event_sha256 + { + return Err(PresenceOperationError::VerificationFailed); + } + let event = validate_signed_document( + document.kind, + &documents.service_public_key, + document.created_at_unix_seconds, + &document.signed_event_bytes, + ) + .map_err(|_| PresenceOperationError::VerificationFailed)?; + if event.id().as_bytes() != &document.event_id { + return Err(PresenceOperationError::VerificationFailed); + } + } + Ok(()) +} + +fn generation_matches( + documents: &RhiSignedPresenceDocuments, + current: &[PresenceOutboxRecord], +) -> Result<bool, PresenceOperationError> { + if documents.documents.len() != current.len() { + return Ok(false); + } + for (document, outbox) in documents.documents.iter().zip(current) { + if outbox.desired_generation != documents.desired_state.generation() + || outbox.desired_sha256 != *documents.desired_state.desired_sha256() + || outbox.target_set_sha256 != *documents.desired_state.target_set_sha256() + || outbox.target_count != documents.desired_state.target_count() + || outbox.required_target_count != documents.desired_state.required_target_count() + || outbox.id + != derive_outbox_id( + outbox.desired_generation, + document.kind, + outbox.desired_sha256, + document.signed_event_sha256, + ) + || document.kind != outbox.document_kind + || document.event_id != outbox.event_id + || document.signed_event_sha256 != outbox.event_sha256 + || document.signed_event_bytes.as_ref() != outbox.exact_signed_event_bytes.as_ref() + || document.created_at_unix_seconds != outbox.authored_at_unix_s + || documents.service_public_key.as_ref() != outbox.service_public_key.as_ref() + || documents + .targets + .iter() + .zip(outbox.targets.iter()) + .any(|(expected, actual)| { + expected.ordinal() != actual.ordinal + || expected.relay_id() != actual.relay_id.as_ref() + || expected.required() != actual.required + }) + { + return Ok(false); + } + validate_targets(outbox, &outbox.targets)?; + } + Ok(true) +} + +fn derive_outbox_id( + generation: u64, + kind: RhiPresenceDocumentKind, + desired_sha256: [u8; 32], + event_sha256: [u8; 32], +) -> RhiPresenceOutboxId { + let mut digest = Sha256::new(); + digest.update(OUTBOX_ID_DOMAIN); + digest.update(generation.to_be_bytes()); + digest.update([document_kind_tag(kind)]); + digest.update(desired_sha256); + digest.update(event_sha256); + RhiPresenceOutboxId(digest.finalize().into()) +} + +fn derive_attempt_id( + outbox: RhiPresenceOutboxId, + event_sha256: [u8; 32], + ordinal: u8, + attempt_number: u16, +) -> RhiPresenceAttemptId { + let mut digest = Sha256::new(); + digest.update(ATTEMPT_ID_DOMAIN); + digest.update(outbox.0); + digest.update(event_sha256); + digest.update([ordinal]); + digest.update(attempt_number.to_be_bytes()); + RhiPresenceAttemptId(digest.finalize().into()) +} + +const fn document_kind_tag(kind: RhiPresenceDocumentKind) -> u8 { + match kind { + RhiPresenceDocumentKind::ServiceProfile => 0, + RhiPresenceDocumentKind::ApplicationHandler => 1, + } +} + +fn target_set_digest(targets: &[RhiPresenceTarget]) -> Result<[u8; 32], PresenceOperationError> { + const DOMAIN: &[u8] = b"radroots.rhi.presence_target_set.v1\0"; + let mut digest = Sha256::new(); + digest.update(DOMAIN); + digest.update( + u32::try_from(targets.len()) + .map_err(|_| PresenceOperationError::InvalidInput)? + .to_be_bytes(), + ); + for target in targets { + digest.update(u32::from(target.ordinal()).to_be_bytes()); + digest.update( + u64::try_from(target.relay_id().len()) + .map_err(|_| PresenceOperationError::InvalidInput)? + .to_be_bytes(), + ); + digest.update(target.relay_id().as_bytes()); + digest.update([u8::from(target.required())]); + } + Ok(digest.finalize().into()) +} + +async fn read_desired_binding( + transaction: &mut ServiceSqliteTransaction<'_>, +) -> Result<PresenceDesiredBinding, PresenceOperationError> { + let rows = sqlx::query(READ_DESIRED_BINDING_SQL) + .fetch_all(&mut *transaction) + .await + .map_err(|_| PresenceOperationError::Storage)?; + let [row] = rows.as_slice() else { + return Err(PresenceOperationError::DesiredStateMismatch); + }; + if row.try_get::<i64, _>("singleton").ok() != Some(1) { + return Err(PresenceOperationError::Invariant); + } + let generation = positive_u64(row, "generation")?; + let enabled = boolean_i64(row, "enabled")?; + let profile = boolean_i64(row, "profile")?; + let application_handler = boolean_i64(row, "application_handler")?; + let target_set_sha256 = blob::<32>(row, "target_set_sha256")?; + let target_count = bounded_u8(row, "target_count", RHI_PRESENCE_DESIRED_MAX_TARGETS)?; + let required_target_count = + bounded_u8(row, "required_target_count", usize::from(target_count))?; + let queue_capacity = row + .try_get::<i64, _>("queue_capacity") + .ok() + .and_then(|value| u32::try_from(value).ok()) + .filter(|value| *value <= 4_096) + .ok_or(PresenceOperationError::Invariant)?; + let desired_sha256 = blob::<32>(row, "desired_sha256")?; + let mode = if enabled { + RhiPresenceDesiredMode::Enabled + } else { + RhiPresenceDesiredMode::Disabled + }; + let valid = if enabled { + (profile || application_handler) && target_count > 0 && queue_capacity > 0 + } else { + !profile + && !application_handler + && target_count == 0 + && required_target_count == 0 + && queue_capacity == 0 + }; + if !valid { + return Err(PresenceOperationError::Invariant); + } + Ok(PresenceDesiredBinding { + generation, + mode, + profile, + application_handler, + target_set_sha256, + target_count, + required_target_count, + queue_capacity, + desired_sha256, + }) +} + +async fn read_generation_outboxes( + transaction: &mut ServiceSqliteTransaction<'_>, + generation: u64, +) -> Result<Vec<PresenceOutboxRecord>, PresenceOperationError> { + let rows = sqlx::query(READ_GENERATION_OUTBOX_SQL) + .bind(i64_value(generation)?) + .fetch_all(&mut *transaction) + .await + .map_err(|_| PresenceOperationError::Storage)?; + if rows.len() > 2 { + return Err(PresenceOperationError::Invariant); + } + let mut records = Vec::with_capacity(rows.len()); + for row in rows { + let mut outbox = decode_outbox(row)?; + outbox.targets = read_targets(transaction, outbox.id) + .await? + .into_boxed_slice(); + records.push(outbox); + } + Ok(records) +} + +async fn read_outbox( + transaction: &mut ServiceSqliteTransaction<'_>, + id: RhiPresenceOutboxId, +) -> Result<Option<PresenceOutboxRecord>, PresenceOperationError> { + let rows = sqlx::query(READ_OUTBOX_SQL) + .bind(id.as_bytes().as_slice()) + .fetch_all(&mut *transaction) + .await + .map_err(|_| PresenceOperationError::Storage)?; + let Some(row) = exactly_zero_or_one(rows)? else { + return Ok(None); + }; + let mut outbox = decode_outbox(row)?; + outbox.targets = read_targets(transaction, outbox.id) + .await? + .into_boxed_slice(); + Ok(Some(outbox)) +} + +async fn read_targets( + transaction: &mut ServiceSqliteTransaction<'_>, + outbox_id: RhiPresenceOutboxId, +) -> Result<Vec<PresenceTargetRecord>, PresenceOperationError> { + let rows = sqlx::query(READ_TARGETS_SQL) + .bind(outbox_id.as_bytes().as_slice()) + .fetch_all(&mut *transaction) + .await + .map_err(|_| PresenceOperationError::Storage)?; + if rows.len() > RHI_PRESENCE_DESIRED_MAX_TARGETS { + return Err(PresenceOperationError::Invariant); + } + rows.into_iter().map(decode_target).collect() +} + +fn decode_outbox( + row: sqlx::sqlite::SqliteRow, +) -> Result<PresenceOutboxRecord, PresenceOperationError> { + let id = RhiPresenceOutboxId(blob::<32>(&row, "outbox_id")?); + let desired_generation = positive_u64(&row, "desired_generation")?; + let document_kind = decode_document_kind(&bounded_text( + &row, + "document_kind", + "document_kind_bytes", + 19, + )?)?; + let desired_sha256 = blob::<32>(&row, "desired_sha256")?; + let target_set_sha256 = blob::<32>(&row, "target_set_sha256")?; + let event_id = blob::<32>(&row, "event_id")?; + let event_sha256 = blob::<32>(&row, "event_sha256")?; + let exact_length = positive_usize(&row, "exact_signed_event_bytes_length")?; + if exact_length > RHI_PRESENCE_SIGNED_EVENT_MAX_BYTES { + return Err(PresenceOperationError::Invariant); + } + let exact_signed_event_bytes = row + .try_get::<Vec<u8>, _>("exact_signed_event_bytes") + .map_err(|_| PresenceOperationError::Invariant)?; + let authored_at_unix_s = nonnegative_u64(&row, "authored_at_unix_s")?; + let service_public_key = + bounded_text(&row, "service_public_key", "service_public_key_bytes", 64)?; + let target_count = bounded_u8(&row, "target_count", RHI_PRESENCE_DESIRED_MAX_TARGETS)?; + let required_target_count = + bounded_u8(&row, "required_target_count", usize::from(target_count))?; + let max_attempts = bounded_u16(&row, "max_attempts", RHI_PRESENCE_MAX_ATTEMPTS)?; + let initial_backoff_ms = bounded_u64(&row, "initial_backoff_ms", INITIAL_BACKOFF_MILLISECONDS)?; + let maximum_backoff_ms = bounded_u64(&row, "maximum_backoff_ms", MAXIMUM_BACKOFF_MILLISECONDS)?; + let attempt_deadline_ms = + bounded_u64(&row, "attempt_deadline_ms", ATTEMPT_DEADLINE_MILLISECONDS)?; + let state = decode_outbox_state(&bounded_text(&row, "state", "state_bytes", 10)?)?; + let revision = positive_u64(&row, "revision")?; + let next_attempt = optional_millis(&row, "next_attempt_unix_ms")?; + let lease_owner = optional_owner(&row)?; + let lease_expires = optional_millis(&row, "lease_expires_unix_ms")?; + let created_at = RhiPresenceUnixMilliseconds(nonnegative_u64(&row, "created_at_unix_ms")?); + let updated_at = RhiPresenceUnixMilliseconds(nonnegative_u64(&row, "updated_at_unix_ms")?); + if exact_signed_event_bytes.len() != exact_length + || exact_signed_event_bytes.is_empty() + || <[u8; 32]>::from(Sha256::digest(&exact_signed_event_bytes)) != event_sha256 + || !valid_public_key(&service_public_key) + || max_attempts != RHI_PRESENCE_MAX_ATTEMPTS + || initial_backoff_ms != INITIAL_BACKOFF_MILLISECONDS + || maximum_backoff_ms != MAXIMUM_BACKOFF_MILLISECONDS + || attempt_deadline_ms != ATTEMPT_DEADLINE_MILLISECONDS + || initial_backoff_ms > maximum_backoff_ms + || updated_at < created_at + || !valid_outbox_shape(state, next_attempt, lease_owner, lease_expires) + { + return Err(PresenceOperationError::Invariant); + } + Ok(PresenceOutboxRecord { + id, + desired_generation, + document_kind, + desired_sha256, + target_set_sha256, + event_id, + event_sha256, + exact_signed_event_bytes: exact_signed_event_bytes.into_boxed_slice(), + authored_at_unix_s, + service_public_key: service_public_key.into_boxed_str(), + target_count, + required_target_count, + max_attempts, + state, + revision, + next_attempt, + lease_owner, + lease_expires, + created_at, + updated_at, + targets: Box::new([]), + }) +} + +fn decode_target( + row: sqlx::sqlite::SqliteRow, +) -> Result<PresenceTargetRecord, PresenceOperationError> { + let ordinal = bounded_u8(&row, "target_ordinal", RHI_PRESENCE_DESIRED_MAX_TARGETS - 1)?; + let relay_id = bounded_text(&row, "relay_id", "relay_id_bytes", 64)?; + let required = boolean_i64(&row, "required")?; + let state = decode_target_state(&bounded_text(&row, "state", "state_bytes", 13)?)?; + let revision = positive_u64(&row, "revision")?; + let attempt_count = bounded_u16(&row, "attempt_count", RHI_PRESENCE_MAX_ATTEMPTS)?; + let next_attempt = optional_millis(&row, "next_attempt_unix_ms")?; + let last_attempt_id = optional_attempt_id(&row)?; + let updated_at = RhiPresenceUnixMilliseconds(nonnegative_u64(&row, "updated_at_unix_ms")?); + if !valid_relay_id(&relay_id) + || !valid_target_shape(state, attempt_count, next_attempt, last_attempt_id) + { + return Err(PresenceOperationError::Invariant); + } + Ok(PresenceTargetRecord { + ordinal, + relay_id: relay_id.into_boxed_str(), + required, + state, + revision, + attempt_count, + next_attempt, + last_attempt_id, + updated_at, + }) +} + +async fn claim_next( + transaction: &mut ServiceSqliteTransaction<'_>, + owner: RhiPresenceLeaseOwner, + now: RhiPresenceUnixMilliseconds, +) -> Result<Option<RhiPresenceLease>, PresenceOperationError> { + let rows = sqlx::query(READ_CLAIMABLE_OUTBOX_SQL) + .bind(i64_value(now.get())?) + .fetch_all(&mut *transaction) + .await + .map_err(|_| PresenceOperationError::Storage)?; + let Some(row) = exactly_zero_or_one(rows)? else { + return Ok(None); + }; + let id = RhiPresenceOutboxId(blob::<32>(&row, "outbox_id")?); + let outbox = read_outbox(transaction, id) + .await? + .ok_or(PresenceOperationError::Invariant)?; + validate_current_desired(transaction, &outbox).await?; + validate_targets(&outbox, &outbox.targets)?; + let expires = now + .get() + .checked_add(ATTEMPT_DEADLINE_MILLISECONDS) + .filter(|value| *value <= MAX_UNIX_MILLISECONDS) + .ok_or(PresenceOperationError::InvalidInput)?; + let result = sqlx::query(CLAIM_OUTBOX_SQL) + .bind(owner.0.as_slice()) + .bind(i64_value(expires)?) + .bind(i64_value(now.get())?) + .bind(outbox.id.as_bytes().as_slice()) + .bind(i64_value(outbox.revision)?) + .bind(i64_value(now.get())?) + .bind(i64_value(now.get())?) + .execute(&mut *transaction) + .await + .map_err(|_| PresenceOperationError::Storage)?; + if result.rows_affected() != 1 { + return Err(PresenceOperationError::LeaseLost); + } + let leased = read_outbox(transaction, id) + .await? + .filter(|record| { + record.state == RhiPresenceOutboxState::Leased + && record.lease_owner == Some(owner) + && record.lease_expires == Some(RhiPresenceUnixMilliseconds(expires)) + }) + .ok_or(PresenceOperationError::Invariant)?; + Ok(Some(RhiPresenceLease { + outbox: leased, + owner, + expires_at: RhiPresenceUnixMilliseconds(expires), + })) +} + +async fn prepare_next( + transaction: &mut ServiceSqliteTransaction<'_>, + lease: RhiPresenceLease, + started_at: RhiPresenceUnixMilliseconds, +) -> Result<RhiPreparedPresenceAttempt, PresenceOperationError> { + validate_lease(transaction, &lease, started_at, true).await?; + let candidate = lease + .outbox + .targets + .iter() + .find(|target| target_is_due(target, lease.outbox.max_attempts, started_at)) + .ok_or(PresenceOperationError::NotReady)?; + let candidate_ordinal = candidate.ordinal; + let candidate_revision = candidate.revision; + let attempt_number = candidate + .attempt_count + .checked_add(1) + .filter(|value| *value <= lease.outbox.max_attempts) + .ok_or(PresenceOperationError::Invariant)?; + let attempt_id = derive_attempt_id( + lease.outbox.id, + lease.outbox.event_sha256, + candidate_ordinal, + attempt_number, + ); + let result = sqlx::query(PREPARE_TARGET_SQL) + .bind(attempt_id.as_bytes().as_slice()) + .bind(i64_value(started_at.get())?) + .bind(lease.outbox.id.as_bytes().as_slice()) + .bind(i64::from(candidate_ordinal)) + .bind(i64_value(candidate_revision)?) + .bind(i64::from(lease.outbox.max_attempts)) + .bind(i64_value(started_at.get())?) + .bind(i64_value(started_at.get())?) + .execute(&mut *transaction) + .await + .map_err(|_| PresenceOperationError::Storage)?; + if result.rows_affected() != 1 { + return Err(PresenceOperationError::LeaseLost); + } + let current = read_outbox(transaction, lease.outbox.id) + .await? + .ok_or(PresenceOperationError::Invariant)?; + if !same_outbox_identity(&current, &lease.outbox) { + return Err(PresenceOperationError::LeaseLost); + } + let target = current + .targets + .into_vec() + .into_iter() + .find(|target| target.ordinal == candidate_ordinal) + .filter(|target| { + target.state == RhiPresenceTargetState::Submitted + && target.attempt_count == attempt_number + && target.last_attempt_id == Some(attempt_id) + && target.updated_at == started_at + }) + .ok_or(PresenceOperationError::Invariant)?; + validate_signed_document( + lease.outbox.document_kind, + &lease.outbox.service_public_key, + lease.outbox.authored_at_unix_s, + &lease.outbox.exact_signed_event_bytes, + ) + .map_err(|_| PresenceOperationError::Invariant)?; + let deadline_at = lease.expires_at; + let exact_signed_event_bytes = lease.outbox.exact_signed_event_bytes.clone(); + Ok(RhiPreparedPresenceAttempt { + lease, + target, + attempt_id, + started_at, + deadline_at, + exact_signed_event_bytes, + }) +} + +async fn record_outcome( + transaction: &mut ServiceSqliteTransaction<'_>, + input: PresenceRecordInput, +) -> Result<RhiPresenceAttemptCommit, PresenceOperationError> { + if let Some(existing) = read_attempt(transaction, input.attempt_id).await? { + return reconcile_recorded(transaction, input, existing).await; + } + validate_input_lease(transaction, input, true).await?; + let outbox = read_outbox(transaction, input.outbox_id) + .await? + .ok_or(PresenceOperationError::Invariant)?; + let target = outbox + .targets + .iter() + .find(|target| target.ordinal == input.target_ordinal) + .ok_or(PresenceOperationError::Invariant)?; + if target.state != RhiPresenceTargetState::Submitted + || target.revision != input.target_revision + || target.attempt_count != input.attempt_number + || target.last_attempt_id != Some(input.attempt_id) + || target.updated_at != input.started_at + { + return Err(PresenceOperationError::LeaseLost); + } + insert_attempt(transaction, input).await?; + let (next_state, next_attempt) = target_outcome_schedule(input)?; + let result = sqlx::query(UPDATE_TARGET_OUTCOME_SQL) + .bind(next_state.code()) + .bind(optional_i64(next_attempt)?) + .bind(i64_value(input.finished_at.get())?) + .bind(input.outbox_id.as_bytes().as_slice()) + .bind(i64::from(input.target_ordinal)) + .bind(i64_value(input.target_revision)?) + .bind(i64::from(input.attempt_number)) + .bind(input.attempt_id.as_bytes().as_slice()) + .execute(&mut *transaction) + .await + .map_err(|_| PresenceOperationError::Storage)?; + if result.rows_affected() != 1 { + return Err(PresenceOperationError::LeaseLost); + } + let targets = read_targets(transaction, input.outbox_id).await?; + let disposition = disposition(&outbox, &targets)?; + update_outbox_after_attempt(transaction, input, disposition).await?; + Ok(RhiPresenceAttemptCommit { + outbox_id: input.outbox_id, + attempt_id: input.attempt_id, + target_ordinal: input.target_ordinal, + attempt_number: input.attempt_number, + outcome: input.outcome, + target_state: next_state, + outbox_state: disposition.state, + }) +} + +async fn read_expired_candidate( + transaction: &mut ServiceSqliteTransaction<'_>, + now: RhiPresenceUnixMilliseconds, +) -> Result<Option<PresenceRecoveryCandidate>, PresenceOperationError> { + let rows = sqlx::query(READ_EXPIRED_OUTBOX_SQL) + .bind(i64_value(now.get())?) + .fetch_all(&mut *transaction) + .await + .map_err(|_| PresenceOperationError::Storage)?; + let Some(row) = exactly_zero_or_one(rows)? else { + return Ok(None); + }; + let id = RhiPresenceOutboxId(blob::<32>(&row, "outbox_id")?); + let outbox = read_outbox(transaction, id) + .await? + .ok_or(PresenceOperationError::Invariant)?; + let desired = read_desired_binding(transaction).await?; + let current_desired = desired_matches_outbox(desired, &outbox); + validate_targets(&outbox, &outbox.targets)?; + let submitted: Vec<_> = outbox + .targets + .iter() + .filter(|target| target.state == RhiPresenceTargetState::Submitted) + .map(|target| target.attempt_count) + .collect(); + if submitted.len() > 1 { + return Err(PresenceOperationError::Invariant); + } + let retry_upper_bound = if current_desired { + submitted + .first() + .copied() + .filter(|attempt| *attempt < outbox.max_attempts) + .map_or(0, retry_upper_bound) + } else { + 0 + }; + Ok(Some(PresenceRecoveryCandidate { + outbox, + retry_upper_bound, + current_desired, + })) +} + +async fn recover_expired( + transaction: &mut ServiceSqliteTransaction<'_>, + candidate: PresenceRecoveryCandidate, + now: RhiPresenceUnixMilliseconds, + delay: RhiPresenceRetryDelayMilliseconds, +) -> Result<bool, PresenceOperationError> { + let current = read_outbox(transaction, candidate.outbox.id) + .await? + .filter(|current| same_outbox(current, &candidate.outbox)) + .ok_or(PresenceOperationError::LeaseLost)?; + if current.state != RhiPresenceOutboxState::Leased + || current.lease_expires.is_none_or(|expires| expires > now) + { + return Err(PresenceOperationError::LeaseLost); + } + let owner = current + .lease_owner + .ok_or(PresenceOperationError::Invariant)?; + let expires = current + .lease_expires + .ok_or(PresenceOperationError::Invariant)?; + for target in current + .targets + .iter() + .filter(|target| target.state == RhiPresenceTargetState::Submitted) + { + let attempt_id = target + .last_attempt_id + .ok_or(PresenceOperationError::Invariant)?; + if read_attempt(transaction, attempt_id).await?.is_some() { + return Err(PresenceOperationError::Invariant); + } + let input = PresenceRecordInput { + outbox_id: current.id, + outbox_revision: current.revision, + event_sha256: current.event_sha256, + owner, + lease_expires: expires, + target_ordinal: target.ordinal, + target_revision: target.revision, + attempt_number: target.attempt_count, + attempt_id, + started_at: target.updated_at, + finished_at: now, + outcome: RhiPresenceAttemptOutcome::Unknown, + retry_delay: delay, + }; + insert_attempt(transaction, input).await?; + let (state, next) = target_outcome_schedule(input)?; + let result = sqlx::query(UPDATE_TARGET_OUTCOME_SQL) + .bind(state.code()) + .bind(optional_i64(next)?) + .bind(i64_value(now.get())?) + .bind(current.id.as_bytes().as_slice()) + .bind(i64::from(target.ordinal)) + .bind(i64_value(target.revision)?) + .bind(i64::from(target.attempt_count)) + .bind(attempt_id.as_bytes().as_slice()) + .execute(&mut *transaction) + .await + .map_err(|_| PresenceOperationError::Storage)?; + if result.rows_affected() != 1 { + return Err(PresenceOperationError::LeaseLost); + } + } + let targets = read_targets(transaction, current.id).await?; + let disposition = if candidate.current_desired { + disposition(&current, &targets)? + } else { + PresenceOutboxDisposition { + state: RhiPresenceOutboxState::Superseded, + next_attempt: None, + } + }; + let result = sqlx::query(UPDATE_OUTBOX_RECOVERY_SQL) + .bind(disposition.state.code()) + .bind(optional_i64(disposition.next_attempt)?) + .bind(i64_value(now.get())?) + .bind(current.id.as_bytes().as_slice()) + .bind(i64_value(current.revision)?) + .bind(owner.0.as_slice()) + .bind(i64_value(expires.get())?) + .bind(i64_value(now.get())?) + .bind(i64_value(now.get())?) + .execute(&mut *transaction) + .await + .map_err(|_| PresenceOperationError::Storage)?; + if result.rows_affected() != 1 { + return Err(PresenceOperationError::LeaseLost); + } + Ok(true) +} + +async fn validate_current_desired( + transaction: &mut ServiceSqliteTransaction<'_>, + outbox: &PresenceOutboxRecord, +) -> Result<(), PresenceOperationError> { + let desired = read_desired_binding(transaction).await?; + if !desired_matches_outbox(desired, outbox) { + return Err(PresenceOperationError::DesiredStateMismatch); + } + Ok(()) +} + +fn desired_matches_outbox(desired: PresenceDesiredBinding, outbox: &PresenceOutboxRecord) -> bool { + desired.mode == RhiPresenceDesiredMode::Enabled + && desired.generation == outbox.desired_generation + && desired.desired_sha256 == outbox.desired_sha256 + && desired.target_set_sha256 == outbox.target_set_sha256 + && desired.target_count == outbox.target_count + && desired.required_target_count == outbox.required_target_count + && desired + .document_kinds() + .any(|kind| kind == outbox.document_kind) +} + +async fn validate_lease( + transaction: &mut ServiceSqliteTransaction<'_>, + lease: &RhiPresenceLease, + now: RhiPresenceUnixMilliseconds, + require_unexpired: bool, +) -> Result<(), PresenceOperationError> { + let current = read_outbox(transaction, lease.outbox.id) + .await? + .filter(|current| same_outbox(current, &lease.outbox)) + .ok_or(PresenceOperationError::LeaseLost)?; + validate_current_desired(transaction, &current).await?; + if current.state != RhiPresenceOutboxState::Leased + || current.lease_owner != Some(lease.owner) + || current.lease_expires != Some(lease.expires_at) + || (require_unexpired && lease.expires_at <= now) + { + return Err(PresenceOperationError::LeaseLost); + } + Ok(()) +} + +async fn validate_input_lease( + transaction: &mut ServiceSqliteTransaction<'_>, + input: PresenceRecordInput, + require_unexpired: bool, +) -> Result<(), PresenceOperationError> { + let current = read_outbox(transaction, input.outbox_id) + .await? + .ok_or(PresenceOperationError::LeaseLost)?; + validate_current_desired(transaction, &current).await?; + if current.revision != input.outbox_revision + || current.event_sha256 != input.event_sha256 + || current.state != RhiPresenceOutboxState::Leased + || current.lease_owner != Some(input.owner) + || current.lease_expires != Some(input.lease_expires) + || (require_unexpired && input.lease_expires <= input.finished_at) + { + return Err(PresenceOperationError::LeaseLost); + } + Ok(()) +} + +async fn insert_attempt( + transaction: &mut ServiceSqliteTransaction<'_>, + input: PresenceRecordInput, +) -> Result<(), PresenceOperationError> { + let result = sqlx::query(INSERT_ATTEMPT_SQL) + .bind(input.attempt_id.as_bytes().as_slice()) + .bind(input.outbox_id.as_bytes().as_slice()) + .bind(i64::from(input.target_ordinal)) + .bind(i64::from(input.attempt_number)) + .bind(input.event_sha256.as_slice()) + .bind(input.owner.0.as_slice()) + .bind(i64_value(input.started_at.get())?) + .bind(i64_value(input.finished_at.get())?) + .bind(input.outcome.code()) + .bind(input.outcome.code()) + .execute(&mut *transaction) + .await + .map_err(|_| PresenceOperationError::Storage)?; + if result.rows_affected() != 1 { + return Err(PresenceOperationError::Storage); + } + Ok(()) +} + +async fn update_outbox_after_attempt( + transaction: &mut ServiceSqliteTransaction<'_>, + input: PresenceRecordInput, + disposition: PresenceOutboxDisposition, +) -> Result<(), PresenceOperationError> { + let result = sqlx::query(UPDATE_OUTBOX_AFTER_ATTEMPT_SQL) + .bind(disposition.state.code()) + .bind(optional_i64(disposition.next_attempt)?) + .bind(i64_value(input.finished_at.get())?) + .bind(input.outbox_id.as_bytes().as_slice()) + .bind(i64_value(input.outbox_revision)?) + .bind(input.owner.0.as_slice()) + .bind(i64_value(input.lease_expires.get())?) + .bind(i64_value(input.finished_at.get())?) + .bind(i64_value(input.finished_at.get())?) + .execute(&mut *transaction) + .await + .map_err(|_| PresenceOperationError::Storage)?; + if result.rows_affected() != 1 { + return Err(PresenceOperationError::LeaseLost); + } + Ok(()) +} + +async fn read_attempt( + transaction: &mut ServiceSqliteTransaction<'_>, + attempt_id: RhiPresenceAttemptId, +) -> Result<Option<PresenceAttemptRecord>, PresenceOperationError> { + let rows = sqlx::query(READ_ATTEMPT_SQL) + .bind(attempt_id.as_bytes().as_slice()) + .fetch_all(&mut *transaction) + .await + .map_err(|_| PresenceOperationError::Storage)?; + exactly_zero_or_one(rows)?.map(decode_attempt).transpose() +} + +fn decode_attempt( + row: sqlx::sqlite::SqliteRow, +) -> Result<PresenceAttemptRecord, PresenceOperationError> { + let attempt_id = RhiPresenceAttemptId(blob::<32>(&row, "attempt_id")?); + let outbox_id = RhiPresenceOutboxId(blob::<32>(&row, "outbox_id")?); + let target_ordinal = bounded_u8(&row, "target_ordinal", RHI_PRESENCE_DESIRED_MAX_TARGETS - 1)?; + let attempt_number = bounded_u16(&row, "attempt_number", RHI_PRESENCE_MAX_ATTEMPTS)?; + let event_sha256 = blob::<32>(&row, "event_sha256")?; + let owner = RhiPresenceLeaseOwner(blob::<16>(&row, "lease_owner")?); + let started_at = RhiPresenceUnixMilliseconds(nonnegative_u64(&row, "started_at_unix_ms")?); + let finished_at = RhiPresenceUnixMilliseconds(nonnegative_u64(&row, "finished_at_unix_ms")?); + let outcome = decode_attempt_outcome(&bounded_text(&row, "outcome", "outcome_bytes", 13)?)?; + let result_code = bounded_text(&row, "result_code", "result_code_bytes", 64)?; + if owner.0.iter().all(|byte| *byte == 0) + || started_at > finished_at + || result_code != outcome.code() + { + return Err(PresenceOperationError::Invariant); + } + Ok(PresenceAttemptRecord { + attempt_id, + outbox_id, + target_ordinal, + attempt_number, + event_sha256, + owner, + started_at, + finished_at, + outcome, + }) +} + +async fn reconcile_recorded( + transaction: &mut ServiceSqliteTransaction<'_>, + input: PresenceRecordInput, + existing: PresenceAttemptRecord, +) -> Result<RhiPresenceAttemptCommit, PresenceOperationError> { + if existing != PresenceAttemptRecord::from_input(input) { + return Err(PresenceOperationError::Invariant); + } + let outbox = read_outbox(transaction, input.outbox_id) + .await? + .ok_or(PresenceOperationError::Invariant)?; + if outbox.revision + != input + .outbox_revision + .checked_add(1) + .ok_or(PresenceOperationError::Invariant)? + || outbox.updated_at != input.finished_at + || outbox.lease_owner.is_some() + || outbox.lease_expires.is_some() + || outbox.event_sha256 != input.event_sha256 + { + return Err(PresenceOperationError::Invariant); + } + validate_targets(&outbox, &outbox.targets)?; + let (expected_state, expected_next) = target_outcome_schedule(input)?; + let expected_target_revision = input + .target_revision + .checked_add(1) + .ok_or(PresenceOperationError::Invariant)?; + let target = outbox + .targets + .iter() + .find(|target| target.ordinal == input.target_ordinal) + .filter(|target| { + target.revision == expected_target_revision + && target.attempt_count == input.attempt_number + && target.last_attempt_id == Some(input.attempt_id) + && target.state == expected_state + && target.next_attempt == expected_next + && target.updated_at == input.finished_at + }) + .ok_or(PresenceOperationError::Invariant)?; + let expected_disposition = disposition(&outbox, &outbox.targets)?; + if outbox.state != expected_disposition.state + || outbox.next_attempt != expected_disposition.next_attempt + { + return Err(PresenceOperationError::Invariant); + } + Ok(RhiPresenceAttemptCommit { + outbox_id: outbox.id, + attempt_id: input.attempt_id, + target_ordinal: input.target_ordinal, + attempt_number: input.attempt_number, + outcome: input.outcome, + target_state: target.state, + outbox_state: outbox.state, + }) +} + +fn disposition( + outbox: &PresenceOutboxRecord, + targets: &[PresenceTargetRecord], +) -> Result<PresenceOutboxDisposition, PresenceOperationError> { + validate_targets(outbox, targets)?; + let required = targets.iter().filter(|target| target.required); + if required + .clone() + .all(|target| target.state == RhiPresenceTargetState::Accepted) + { + return Ok(PresenceOutboxDisposition { + state: RhiPresenceOutboxState::Complete, + next_attempt: None, + }); + } + if required + .clone() + .any(|target| target_is_blocking(target, outbox.max_attempts)) + { + return Ok(PresenceOutboxDisposition { + state: RhiPresenceOutboxState::Blocked, + next_attempt: None, + }); + } + let next_attempt = required + .filter_map(|target| target.next_attempt) + .min() + .ok_or(PresenceOperationError::Invariant)?; + Ok(PresenceOutboxDisposition { + state: RhiPresenceOutboxState::Pending, + next_attempt: Some(next_attempt), + }) +} + +fn validate_targets( + outbox: &PresenceOutboxRecord, + targets: &[PresenceTargetRecord], +) -> Result<(), PresenceOperationError> { + if targets.len() != usize::from(outbox.target_count) + || targets.is_empty() + || targets.len() > RHI_PRESENCE_DESIRED_MAX_TARGETS + || targets.iter().filter(|target| target.required).count() + != usize::from(outbox.required_target_count) + || targets + .iter() + .enumerate() + .any(|(ordinal, target)| usize::from(target.ordinal) != ordinal) + || targets.iter().enumerate().any(|(index, target)| { + targets[index + 1..] + .iter() + .any(|later| later.relay_id == target.relay_id) + }) + { + return Err(PresenceOperationError::Invariant); + } + let mut digest = Sha256::new(); + digest.update(b"radroots.rhi.presence_target_set.v1\0"); + digest.update( + u32::try_from(targets.len()) + .map_err(|_| PresenceOperationError::Invariant)? + .to_be_bytes(), + ); + for target in targets { + digest.update(u32::from(target.ordinal).to_be_bytes()); + digest.update( + u64::try_from(target.relay_id.len()) + .map_err(|_| PresenceOperationError::Invariant)? + .to_be_bytes(), + ); + digest.update(target.relay_id.as_bytes()); + digest.update([u8::from(target.required)]); + } + if <[u8; 32]>::from(digest.finalize()) != outbox.target_set_sha256 { + return Err(PresenceOperationError::Invariant); + } + Ok(()) +} + +fn target_outcome_schedule( + input: PresenceRecordInput, +) -> Result<(RhiPresenceTargetState, Option<RhiPresenceUnixMilliseconds>), PresenceOperationError> { + let state = target_state_for_outcome(input.outcome)?; + let next = + if retryable_outcome(input.outcome) && input.attempt_number < RHI_PRESENCE_MAX_ATTEMPTS { + if input.retry_delay.get() > retry_upper_bound(input.attempt_number) { + return Err(PresenceOperationError::InvalidInput); + } + Some(RhiPresenceUnixMilliseconds( + input + .finished_at + .get() + .checked_add(input.retry_delay.get()) + .filter(|value| *value <= MAX_UNIX_MILLISECONDS) + .ok_or(PresenceOperationError::InvalidInput)?, + )) + } else { + if input.retry_delay.get() != 0 { + return Err(PresenceOperationError::InvalidInput); + } + None + }; + Ok((state, next)) +} + +fn target_is_blocking(target: &PresenceTargetRecord, max_attempts: u16) -> bool { + matches!( + target.state, + RhiPresenceTargetState::Rejected | RhiPresenceTargetState::AuthRequired + ) || (retryable_state(target.state) + && target.attempt_count >= max_attempts + && target.next_attempt.is_none()) +} + +fn target_is_due( + target: &PresenceTargetRecord, + max_attempts: u16, + now: RhiPresenceUnixMilliseconds, +) -> bool { + retryable_state(target.state) + && target.attempt_count < max_attempts + && target.next_attempt.is_some_and(|next| next <= now) +} + +const fn retryable_state(state: RhiPresenceTargetState) -> bool { + matches!( + state, + RhiPresenceTargetState::Pending + | RhiPresenceTargetState::Failed + | RhiPresenceTargetState::RateLimited + | RhiPresenceTargetState::Unknown + ) +} + +const fn retryable_outcome(outcome: RhiPresenceAttemptOutcome) -> bool { + matches!( + outcome, + RhiPresenceAttemptOutcome::RateLimited + | RhiPresenceAttemptOutcome::Failed + | RhiPresenceAttemptOutcome::Unknown + ) +} + +const fn target_state_for_outcome( + outcome: RhiPresenceAttemptOutcome, +) -> Result<RhiPresenceTargetState, PresenceOperationError> { + match outcome { + RhiPresenceAttemptOutcome::Submitted => Err(PresenceOperationError::InvalidInput), + RhiPresenceAttemptOutcome::Accepted => Ok(RhiPresenceTargetState::Accepted), + RhiPresenceAttemptOutcome::Rejected => Ok(RhiPresenceTargetState::Rejected), + RhiPresenceAttemptOutcome::RateLimited => Ok(RhiPresenceTargetState::RateLimited), + RhiPresenceAttemptOutcome::AuthRequired => Ok(RhiPresenceTargetState::AuthRequired), + RhiPresenceAttemptOutcome::Failed => Ok(RhiPresenceTargetState::Failed), + RhiPresenceAttemptOutcome::Unknown => Ok(RhiPresenceTargetState::Unknown), + } +} + +fn retry_upper_bound(attempt_number: u16) -> u64 { + let mut bound = INITIAL_BACKOFF_MILLISECONDS; + for _ in 1..attempt_number { + bound = bound.saturating_mul(2).min(MAXIMUM_BACKOFF_MILLISECONDS); + } + bound.min(MAXIMUM_BACKOFF_MILLISECONDS) +} + +fn same_outbox(current: &PresenceOutboxRecord, prior: &PresenceOutboxRecord) -> bool { + same_outbox_identity(current, prior) + && current.state == prior.state + && current.revision == prior.revision + && current.next_attempt == prior.next_attempt + && current.lease_owner == prior.lease_owner + && current.lease_expires == prior.lease_expires + && current.updated_at == prior.updated_at + && current.targets == prior.targets +} + +fn same_outbox_identity(current: &PresenceOutboxRecord, prior: &PresenceOutboxRecord) -> bool { + current.id == prior.id + && current.desired_generation == prior.desired_generation + && current.document_kind == prior.document_kind + && current.desired_sha256 == prior.desired_sha256 + && current.target_set_sha256 == prior.target_set_sha256 + && current.event_id == prior.event_id + && current.event_sha256 == prior.event_sha256 + && current.exact_signed_event_bytes == prior.exact_signed_event_bytes + && current.authored_at_unix_s == prior.authored_at_unix_s + && current.service_public_key == prior.service_public_key + && current.target_count == prior.target_count + && current.required_target_count == prior.required_target_count + && current.max_attempts == prior.max_attempts + && current.created_at == prior.created_at +} + +fn valid_outbox_shape( + state: RhiPresenceOutboxState, + next_attempt: Option<RhiPresenceUnixMilliseconds>, + lease_owner: Option<RhiPresenceLeaseOwner>, + lease_expires: Option<RhiPresenceUnixMilliseconds>, +) -> bool { + match state { + RhiPresenceOutboxState::Pending => { + next_attempt.is_some() && lease_owner.is_none() && lease_expires.is_none() + } + RhiPresenceOutboxState::Leased => { + next_attempt.is_none() && lease_owner.is_some() && lease_expires.is_some() + } + RhiPresenceOutboxState::Complete + | RhiPresenceOutboxState::Blocked + | RhiPresenceOutboxState::Superseded => { + next_attempt.is_none() && lease_owner.is_none() && lease_expires.is_none() + } + } +} + +fn valid_target_shape( + state: RhiPresenceTargetState, + attempt_count: u16, + next_attempt: Option<RhiPresenceUnixMilliseconds>, + last_attempt_id: Option<RhiPresenceAttemptId>, +) -> bool { + match state { + RhiPresenceTargetState::Pending => { + attempt_count == 0 && next_attempt.is_some() && last_attempt_id.is_none() + } + RhiPresenceTargetState::Submitted => { + attempt_count > 0 && next_attempt.is_none() && last_attempt_id.is_some() + } + RhiPresenceTargetState::Accepted + | RhiPresenceTargetState::Rejected + | RhiPresenceTargetState::AuthRequired => { + attempt_count > 0 && next_attempt.is_none() && last_attempt_id.is_some() + } + RhiPresenceTargetState::RateLimited + | RhiPresenceTargetState::Failed + | RhiPresenceTargetState::Unknown => { + attempt_count > 0 + && last_attempt_id.is_some() + && if attempt_count < RHI_PRESENCE_MAX_ATTEMPTS { + next_attempt.is_some() + } else { + next_attempt.is_none() + } + } + } +} + +fn decode_document_kind(value: &str) -> Result<RhiPresenceDocumentKind, PresenceOperationError> { + match value { + "service_profile" => Ok(RhiPresenceDocumentKind::ServiceProfile), + "application_handler" => Ok(RhiPresenceDocumentKind::ApplicationHandler), + _ => Err(PresenceOperationError::Invariant), + } +} + +fn decode_outbox_state(value: &str) -> Result<RhiPresenceOutboxState, PresenceOperationError> { + match value { + "pending" => Ok(RhiPresenceOutboxState::Pending), + "leased" => Ok(RhiPresenceOutboxState::Leased), + "complete" => Ok(RhiPresenceOutboxState::Complete), + "blocked" => Ok(RhiPresenceOutboxState::Blocked), + "superseded" => Ok(RhiPresenceOutboxState::Superseded), + _ => Err(PresenceOperationError::Invariant), + } +} + +fn decode_target_state(value: &str) -> Result<RhiPresenceTargetState, PresenceOperationError> { + match value { + "pending" => Ok(RhiPresenceTargetState::Pending), + "submitted" => Ok(RhiPresenceTargetState::Submitted), + "accepted" => Ok(RhiPresenceTargetState::Accepted), + "rejected" => Ok(RhiPresenceTargetState::Rejected), + "rate_limited" => Ok(RhiPresenceTargetState::RateLimited), + "auth_required" => Ok(RhiPresenceTargetState::AuthRequired), + "failed" => Ok(RhiPresenceTargetState::Failed), + "unknown" => Ok(RhiPresenceTargetState::Unknown), + _ => Err(PresenceOperationError::Invariant), + } +} + +fn decode_attempt_outcome( + value: &str, +) -> Result<RhiPresenceAttemptOutcome, PresenceOperationError> { + match value { + "submitted" => Ok(RhiPresenceAttemptOutcome::Submitted), + "accepted" => Ok(RhiPresenceAttemptOutcome::Accepted), + "rejected" => Ok(RhiPresenceAttemptOutcome::Rejected), + "rate_limited" => Ok(RhiPresenceAttemptOutcome::RateLimited), + "auth_required" => Ok(RhiPresenceAttemptOutcome::AuthRequired), + "failed" => Ok(RhiPresenceAttemptOutcome::Failed), + "unknown" => Ok(RhiPresenceAttemptOutcome::Unknown), + _ => Err(PresenceOperationError::Invariant), + } +} + +fn exactly_zero_or_one( + rows: Vec<sqlx::sqlite::SqliteRow>, +) -> Result<Option<sqlx::sqlite::SqliteRow>, PresenceOperationError> { + match rows.as_slice() { + [] => Ok(None), + [_] => Ok(rows.into_iter().next()), + _ => Err(PresenceOperationError::Invariant), + } +} + +fn blob<const N: usize>( + row: &sqlx::sqlite::SqliteRow, + column: &str, +) -> Result<[u8; N], PresenceOperationError> { + row.try_get::<Option<Vec<u8>>, _>(column) + .map_err(|_| PresenceOperationError::Invariant)? + .ok_or(PresenceOperationError::Invariant)? + .try_into() + .map_err(|_| PresenceOperationError::Invariant) +} + +fn bounded_text( + row: &sqlx::sqlite::SqliteRow, + value_column: &str, + length_column: &str, + maximum: usize, +) -> Result<String, PresenceOperationError> { + let length = row + .try_get::<i64, _>(length_column) + .ok() + .and_then(|value| usize::try_from(value).ok()) + .filter(|value| *value > 0 && *value <= maximum) + .ok_or(PresenceOperationError::Invariant)?; + let value = row + .try_get::<String, _>(value_column) + .map_err(|_| PresenceOperationError::Invariant)?; + if value.len() != length { + return Err(PresenceOperationError::Invariant); + } + Ok(value) +} + +fn bounded_u8( + row: &sqlx::sqlite::SqliteRow, + column: &str, + maximum: usize, +) -> Result<u8, PresenceOperationError> { + row.try_get::<i64, _>(column) + .ok() + .and_then(|value| u8::try_from(value).ok()) + .filter(|value| usize::from(*value) <= maximum) + .ok_or(PresenceOperationError::Invariant) +} + +fn bounded_u16( + row: &sqlx::sqlite::SqliteRow, + column: &str, + maximum: u16, +) -> Result<u16, PresenceOperationError> { + row.try_get::<i64, _>(column) + .ok() + .and_then(|value| u16::try_from(value).ok()) + .filter(|value| *value <= maximum) + .ok_or(PresenceOperationError::Invariant) +} + +fn bounded_u64( + row: &sqlx::sqlite::SqliteRow, + column: &str, + maximum: u64, +) -> Result<u64, PresenceOperationError> { + row.try_get::<i64, _>(column) + .ok() + .and_then(|value| u64::try_from(value).ok()) + .filter(|value| *value <= maximum) + .ok_or(PresenceOperationError::Invariant) +} + +fn positive_u64( + row: &sqlx::sqlite::SqliteRow, + column: &str, +) -> Result<u64, PresenceOperationError> { + nonnegative_u64(row, column).and_then(|value| { + (value != 0) + .then_some(value) + .ok_or(PresenceOperationError::Invariant) + }) +} + +fn nonnegative_u64( + row: &sqlx::sqlite::SqliteRow, + column: &str, +) -> Result<u64, PresenceOperationError> { + row.try_get::<i64, _>(column) + .ok() + .and_then(|value| u64::try_from(value).ok()) + .ok_or(PresenceOperationError::Invariant) +} + +fn positive_usize( + row: &sqlx::sqlite::SqliteRow, + column: &str, +) -> Result<usize, PresenceOperationError> { + row.try_get::<i64, _>(column) + .ok() + .and_then(|value| usize::try_from(value).ok()) + .filter(|value| *value != 0) + .ok_or(PresenceOperationError::Invariant) +} + +fn boolean_i64( + row: &sqlx::sqlite::SqliteRow, + column: &str, +) -> Result<bool, PresenceOperationError> { + match row.try_get::<i64, _>(column) { + Ok(0) => Ok(false), + Ok(1) => Ok(true), + Ok(_) | Err(_) => Err(PresenceOperationError::Invariant), + } +} + +fn optional_millis( + row: &sqlx::sqlite::SqliteRow, + column: &str, +) -> Result<Option<RhiPresenceUnixMilliseconds>, PresenceOperationError> { + row.try_get::<Option<i64>, _>(column) + .map_err(|_| PresenceOperationError::Invariant)? + .map(|value| { + u64::try_from(value) + .map(RhiPresenceUnixMilliseconds) + .map_err(|_| PresenceOperationError::Invariant) + }) + .transpose() +} + +fn optional_owner( + row: &sqlx::sqlite::SqliteRow, +) -> Result<Option<RhiPresenceLeaseOwner>, PresenceOperationError> { + let length = row + .try_get::<Option<i64>, _>("lease_owner_bytes") + .map_err(|_| PresenceOperationError::Invariant)?; + let value = row + .try_get::<Option<Vec<u8>>, _>("lease_owner") + .map_err(|_| PresenceOperationError::Invariant)?; + match (length, value) { + (None, None) => Ok(None), + (Some(16), Some(value)) => { + let value: [u8; 16] = value + .try_into() + .map_err(|_| PresenceOperationError::Invariant)?; + if value.iter().all(|byte| *byte == 0) { + return Err(PresenceOperationError::Invariant); + } + Ok(Some(RhiPresenceLeaseOwner(value))) + } + _ => Err(PresenceOperationError::Invariant), + } +} + +fn optional_attempt_id( + row: &sqlx::sqlite::SqliteRow, +) -> Result<Option<RhiPresenceAttemptId>, PresenceOperationError> { + let length = row + .try_get::<Option<i64>, _>("last_attempt_id_bytes") + .map_err(|_| PresenceOperationError::Invariant)?; + let value = row + .try_get::<Option<Vec<u8>>, _>("last_attempt_id") + .map_err(|_| PresenceOperationError::Invariant)?; + match (length, value) { + (None, None) => Ok(None), + (Some(32), Some(value)) => Ok(Some(RhiPresenceAttemptId( + value + .try_into() + .map_err(|_| PresenceOperationError::Invariant)?, + ))), + _ => Err(PresenceOperationError::Invariant), + } +} + +fn valid_public_key(value: &str) -> bool { + value.len() == 64 + && value + .bytes() + .all(|byte| byte.is_ascii_hexdigit() && !byte.is_ascii_uppercase()) + && NostrPublicKey::from_hex(value).is_ok() +} + +fn valid_relay_id(value: &str) -> bool { + !value.is_empty() + && value.len() <= 64 + && value.as_bytes()[0].is_ascii_lowercase() + && value.bytes().all(|byte| { + byte.is_ascii_lowercase() || byte.is_ascii_digit() || matches!(byte, b'_' | b'-') + }) +} + +fn i64_value(value: u64) -> Result<i64, PresenceOperationError> { + i64::try_from(value).map_err(|_| PresenceOperationError::InvalidInput) +} + +fn optional_i64( + value: Option<RhiPresenceUnixMilliseconds>, +) -> Result<Option<i64>, PresenceOperationError> { + value.map(|value| i64_value(value.get())).transpose() +} + +#[derive(Clone, Copy)] +struct PresenceDesiredBinding { + generation: u64, + mode: RhiPresenceDesiredMode, + profile: bool, + application_handler: bool, + target_set_sha256: [u8; 32], + target_count: u8, + required_target_count: u8, + queue_capacity: u32, + desired_sha256: [u8; 32], +} + +impl PresenceDesiredBinding { + const fn document_count(self) -> usize { + self.profile as usize + self.application_handler as usize + } + + fn document_kinds(self) -> impl Iterator<Item = RhiPresenceDocumentKind> { + [ + self.profile + .then_some(RhiPresenceDocumentKind::ServiceProfile), + self.application_handler + .then_some(RhiPresenceDocumentKind::ApplicationHandler), + ] + .into_iter() + .flatten() + } +} + +#[derive(PartialEq, Eq)] +struct PresenceOutboxRecord { + id: RhiPresenceOutboxId, + desired_generation: u64, + document_kind: RhiPresenceDocumentKind, + desired_sha256: [u8; 32], + target_set_sha256: [u8; 32], + event_id: [u8; 32], + event_sha256: [u8; 32], + exact_signed_event_bytes: Box<[u8]>, + authored_at_unix_s: u64, + service_public_key: Box<str>, + target_count: u8, + required_target_count: u8, + max_attempts: u16, + state: RhiPresenceOutboxState, + revision: u64, + next_attempt: Option<RhiPresenceUnixMilliseconds>, + lease_owner: Option<RhiPresenceLeaseOwner>, + lease_expires: Option<RhiPresenceUnixMilliseconds>, + created_at: RhiPresenceUnixMilliseconds, + updated_at: RhiPresenceUnixMilliseconds, + targets: Box<[PresenceTargetRecord]>, +} + +#[derive(Clone, PartialEq, Eq)] +struct PresenceTargetRecord { + ordinal: u8, + relay_id: Box<str>, + required: bool, + state: RhiPresenceTargetState, + revision: u64, + attempt_count: u16, + next_attempt: Option<RhiPresenceUnixMilliseconds>, + last_attempt_id: Option<RhiPresenceAttemptId>, + updated_at: RhiPresenceUnixMilliseconds, +} + +struct PresenceRecoveryCandidate { + outbox: PresenceOutboxRecord, + retry_upper_bound: u64, + current_desired: bool, +} + +#[derive(Clone, Copy)] +struct PresenceRecordInput { + outbox_id: RhiPresenceOutboxId, + outbox_revision: u64, + event_sha256: [u8; 32], + owner: RhiPresenceLeaseOwner, + lease_expires: RhiPresenceUnixMilliseconds, + target_ordinal: u8, + target_revision: u64, + attempt_number: u16, + attempt_id: RhiPresenceAttemptId, + started_at: RhiPresenceUnixMilliseconds, + finished_at: RhiPresenceUnixMilliseconds, + outcome: RhiPresenceAttemptOutcome, + retry_delay: RhiPresenceRetryDelayMilliseconds, +} + +impl PresenceRecordInput { + fn from_prepared( + prepared: &RhiPreparedPresenceAttempt, + finished_at: RhiPresenceUnixMilliseconds, + outcome: RhiPresenceAttemptOutcome, + retry_delay: RhiPresenceRetryDelayMilliseconds, + ) -> Result<Self, RhiPresencePublicationError> { + if finished_at < prepared.started_at || outcome == RhiPresenceAttemptOutcome::Submitted { + return Err(failure(RhiPresencePublicationErrorKind::InvalidInput)); + } + let expected = derive_attempt_id( + prepared.lease.outbox.id, + prepared.lease.outbox.event_sha256, + prepared.target.ordinal, + prepared.target.attempt_count, + ); + if expected != prepared.attempt_id { + return Err(failure(RhiPresencePublicationErrorKind::Invariant)); + } + Ok(Self { + outbox_id: prepared.lease.outbox.id, + outbox_revision: prepared.lease.outbox.revision, + event_sha256: prepared.lease.outbox.event_sha256, + owner: prepared.lease.owner, + lease_expires: prepared.lease.expires_at, + target_ordinal: prepared.target.ordinal, + target_revision: prepared.target.revision, + attempt_number: prepared.target.attempt_count, + attempt_id: prepared.attempt_id, + started_at: prepared.started_at, + finished_at, + outcome, + retry_delay, + }) + } +} + +#[derive(Clone, Copy, PartialEq, Eq)] +struct PresenceAttemptRecord { + attempt_id: RhiPresenceAttemptId, + outbox_id: RhiPresenceOutboxId, + target_ordinal: u8, + attempt_number: u16, + event_sha256: [u8; 32], + owner: RhiPresenceLeaseOwner, + started_at: RhiPresenceUnixMilliseconds, + finished_at: RhiPresenceUnixMilliseconds, + outcome: RhiPresenceAttemptOutcome, +} + +impl PresenceAttemptRecord { + const fn from_input(input: PresenceRecordInput) -> Self { + Self { + attempt_id: input.attempt_id, + outbox_id: input.outbox_id, + target_ordinal: input.target_ordinal, + attempt_number: input.attempt_number, + event_sha256: input.event_sha256, + owner: input.owner, + started_at: input.started_at, + finished_at: input.finished_at, + outcome: input.outcome, + } + } +} + +#[derive(Clone, Copy)] +struct PresenceOutboxDisposition { + state: RhiPresenceOutboxState, + next_attempt: Option<RhiPresenceUnixMilliseconds>, +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum PresenceOperationError { + InvalidInput, + DesiredStateMismatch, + VerificationFailed, + NotReady, + LeaseLost, + Invariant, + Storage, +} + +fn map_transaction_error( + error: ServiceSqliteTransactionError<PresenceOperationError>, +) -> RhiPresencePublicationError { + if error.kind() == ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown { + return failure(RhiPresencePublicationErrorKind::CommitOutcomeUnknown); + } + failure(match error.operation_error().copied() { + Some(PresenceOperationError::InvalidInput) => RhiPresencePublicationErrorKind::InvalidInput, + Some(PresenceOperationError::DesiredStateMismatch) => { + RhiPresencePublicationErrorKind::DesiredStateMismatch + } + Some(PresenceOperationError::VerificationFailed) => { + RhiPresencePublicationErrorKind::VerificationFailed + } + Some(PresenceOperationError::NotReady) => RhiPresencePublicationErrorKind::NotReady, + Some(PresenceOperationError::LeaseLost) => RhiPresencePublicationErrorKind::LeaseLost, + Some(PresenceOperationError::Invariant) => RhiPresencePublicationErrorKind::Invariant, + Some(PresenceOperationError::Storage) | None => RhiPresencePublicationErrorKind::Storage, + }) +} + +fn presence_plan( + kind: RhiPresenceDocumentKind, + author: &str, + created_at: u64, +) -> Result<AuthoredEventPlan, RhiPresencePublicationError> { + match kind { + RhiPresenceDocumentKind::ServiceProfile => { + let profile = AuthoredProfile::new(PROFILE_NAME) + .map_err(|_| failure(RhiPresencePublicationErrorKind::RenderingFailed))? + .with_display_name(PROFILE_DISPLAY_NAME) + .with_about(PROFILE_ABOUT) + .with_bot(true); + AuthoredEventPlan::from_profile(&profile, created_at, author) + .map_err(|_| failure(RhiPresencePublicationErrorKind::RenderingFailed)) + } + RhiPresenceDocumentKind::ApplicationHandler => { + let metadata = nostr::Metadata::new() + .name(PROFILE_DISPLAY_NAME) + .about(PROFILE_ABOUT); + let content = serde_json::to_string(&metadata) + .map_err(|_| failure(RhiPresencePublicationErrorKind::RenderingFailed))?; + let spec = ApplicationHandlerSpec::new(APPLICATION_HANDLER_KINDS.to_vec()) + .with_identifier(APPLICATION_HANDLER_IDENTIFIER) + .with_metadata(metadata); + let builder = build_application_handler(&spec) + .map_err(|_| failure(RhiPresencePublicationErrorKind::RenderingFailed))? + .custom_created_at(Timestamp::from_secs(created_at)); + let public_key = NostrPublicKey::from_hex(author) + .map_err(|_| failure(RhiPresencePublicationErrorKind::IdentityMismatch))?; + let request = builder + .into_external_signing_request(public_key) + .map_err(|_| failure(RhiPresencePublicationErrorKind::RenderingFailed))?; + let draft = GenericEventDraft::new( + "radroots.application.handler.v1", + KIND_APPLICATION_HANDLER, + created_at, + vec![ + vec!["d".to_owned(), APPLICATION_HANDLER_IDENTIFIER.to_owned()], + vec!["k".to_owned(), APPLICATION_HANDLER_KINDS[0].to_string()], + ], + content, + author, + ) + .map_err(|_| failure(RhiPresencePublicationErrorKind::RenderingFailed))?; + let plan = AuthoredEventPlan::from_generic(draft) + .map_err(|_| failure(RhiPresencePublicationErrorKind::RenderingFailed))?; + if request.expected_event_id().to_bytes() != *plan.expected_event_id().as_bytes() { + return Err(failure(RhiPresencePublicationErrorKind::RenderingFailed)); + } + Ok(plan) + } + } +} + +fn sign_plan( + identity: &RhiDecryptedIdentity, + plan: &AuthoredEventPlan, + auxiliary: &[u8; 32], +) -> Result<nostr::Event, RhiPresencePublicationError> { + let kind = u16::try_from(plan.body().kind()) + .map(Kind::Custom) + .map_err(|_| failure(RhiPresencePublicationErrorKind::RenderingFailed))?; + let tags = plan + .body() + .tags() + .iter() + .cloned() + .map(Tag::parse) + .collect::<Result<Vec<_>, _>>() + .map_err(|_| failure(RhiPresencePublicationErrorKind::RenderingFailed))?; + let author = NostrPublicKey::from_hex(identity.public_identity().as_hex()) + .map_err(|_| failure(RhiPresencePublicationErrorKind::IdentityMismatch))?; + let unsigned = EventBuilder::new(kind, plan.body().content()) + .tags(tags) + .custom_created_at(Timestamp::from_secs(plan.created_at())) + .build(author); + if unsigned.id.as_ref().map(|id| id.to_bytes()) != Some(*plan.expected_event_id().as_bytes()) { + return Err(failure(RhiPresencePublicationErrorKind::RenderingFailed)); + } + identity + .sign_nostr_event(unsigned, auxiliary) + .map_err(|_| failure(RhiPresencePublicationErrorKind::RenderingFailed)) +} + +fn validate_signed_document( + kind: RhiPresenceDocumentKind, + author: &str, + created_at: u64, + bytes: &[u8], +) -> Result<EventEnvelope, RhiPresencePublicationError> { + if bytes.is_empty() || bytes.len() > RHI_PRESENCE_SIGNED_EVENT_MAX_BYTES { + return Err(failure(RhiPresencePublicationErrorKind::VerificationFailed)); + } + let source = core::str::from_utf8(bytes) + .map_err(|_| failure(RhiPresencePublicationErrorKind::VerificationFailed))?; + let wire = Nip01EventWire::parse_json_unverified_with_limits(source, presence_wire_limits()) + .map_err(|_| failure(RhiPresencePublicationErrorKind::VerificationFailed))?; + let event = wire + .into_unverified_envelope() + .map_err(|_| failure(RhiPresencePublicationErrorKind::VerificationFailed))?; + if verify_id(&event) != Verification::IdVerified + || verify(&event) != Verification::Verified + || event.author().to_hex() != author + || event.created_at_u64() != created_at + { + return Err(failure(RhiPresencePublicationErrorKind::VerificationFailed)); + } + let expected = presence_plan(kind, author, created_at)?; + if event.id().as_bytes() != expected.expected_event_id().as_bytes() + || event.kind_u32() != expected.body().kind() + || event.tags_as_vec() != expected.body().tags() + || event.content() != expected.body().content() + { + return Err(failure(RhiPresencePublicationErrorKind::VerificationFailed)); + } + Ok(event) +} + +const fn presence_wire_limits() -> EventWireLimits { + EventWireLimits { + max_raw_json_bytes: RHI_PRESENCE_SIGNED_EVENT_MAX_BYTES, + max_content_bytes: 4 * 1024, + max_tag_count: 4, + max_total_tag_elements: 8, + max_tag_element_bytes: 64, + max_total_tag_bytes: 256, + max_extra_fields: 0, + max_total_extra_json_bytes: 0, + } +} + +const fn failure(kind: RhiPresencePublicationErrorKind) -> RhiPresencePublicationError { + RhiPresencePublicationError { kind } +} + +#[cfg(test)] +mod tests { + use std::{collections::BTreeSet, error::Error as _}; + + use super::*; + + #[test] + fn public_vocabularies_bounds_and_diagnostics_are_closed() { + let kinds = [ + RhiPresencePublicationErrorKind::InvalidMode, + RhiPresencePublicationErrorKind::InvalidInput, + RhiPresencePublicationErrorKind::DesiredStateMismatch, + RhiPresencePublicationErrorKind::IdentityMismatch, + RhiPresencePublicationErrorKind::EntropyUnavailable, + RhiPresencePublicationErrorKind::RenderingFailed, + RhiPresencePublicationErrorKind::VerificationFailed, + RhiPresencePublicationErrorKind::NotReady, + RhiPresencePublicationErrorKind::LeaseLost, + RhiPresencePublicationErrorKind::Invariant, + RhiPresencePublicationErrorKind::ClockUnavailable, + RhiPresencePublicationErrorKind::Storage, + RhiPresencePublicationErrorKind::CommitOutcomeUnknown, + ]; + let codes: BTreeSet<_> = kinds.iter().map(|kind| kind.code()).collect(); + assert_eq!(codes.len(), kinds.len()); + for kind in kinds { + let error = failure(kind); + assert_eq!(error.kind(), kind); + assert_eq!(error.code(), kind.code()); + assert!(error.source().is_none()); + let rendered = format!("{error} {error:?}"); + for forbidden in ["relay", "wss://", "state.sqlite", "010101", "sqlx"] { + assert!(!rendered.contains(forbidden)); + } + } + + assert_eq!( + RhiPresenceUnixMilliseconds::new(i64::MAX as u64) + .expect("maximum time") + .get(), + i64::MAX as u64 + ); + assert!(RhiPresenceUnixMilliseconds::new(i64::MAX as u64 + 1).is_err()); + assert_eq!( + RhiPresenceRetryDelayMilliseconds::new(MAXIMUM_BACKOFF_MILLISECONDS) + .expect("maximum delay") + .get(), + MAXIMUM_BACKOFF_MILLISECONDS + ); + assert!(RhiPresenceRetryDelayMilliseconds::new(MAXIMUM_BACKOFF_MILLISECONDS + 1).is_err()); + assert!(RhiPresenceLeaseOwner::from_bytes([0; 16]).is_err()); + assert_eq!( + format!("{:?}", RhiPresenceLeaseOwner::from_bytes([1; 16]).unwrap()), + "RhiPresenceLeaseOwner([redacted])" + ); + assert_eq!( + format!("{:?}", RhiPresenceOutboxId([0x5a; 32])), + "RhiPresenceOutboxId([redacted])" + ); + assert_eq!( + format!("{:?}", RhiPresenceAttemptId([0x6a; 32])), + "RhiPresenceAttemptId([redacted])" + ); + } + + #[test] + fn state_and_outcome_codes_are_exact() { + assert_eq!( + [ + RhiPresenceOutboxState::Pending, + RhiPresenceOutboxState::Leased, + RhiPresenceOutboxState::Complete, + RhiPresenceOutboxState::Blocked, + RhiPresenceOutboxState::Superseded, + ] + .map(RhiPresenceOutboxState::code), + ["pending", "leased", "complete", "blocked", "superseded"] + ); + assert_eq!( + [ + RhiPresenceTargetState::Pending, + RhiPresenceTargetState::Submitted, + RhiPresenceTargetState::Accepted, + RhiPresenceTargetState::Rejected, + RhiPresenceTargetState::RateLimited, + RhiPresenceTargetState::AuthRequired, + RhiPresenceTargetState::Failed, + RhiPresenceTargetState::Unknown, + ] + .map(RhiPresenceTargetState::code), + [ + "pending", + "submitted", + "accepted", + "rejected", + "rate_limited", + "auth_required", + "failed", + "unknown", + ] + ); + assert_eq!( + [ + RhiPresenceAttemptOutcome::Submitted, + RhiPresenceAttemptOutcome::Accepted, + RhiPresenceAttemptOutcome::Rejected, + RhiPresenceAttemptOutcome::RateLimited, + RhiPresenceAttemptOutcome::AuthRequired, + RhiPresenceAttemptOutcome::Failed, + RhiPresenceAttemptOutcome::Unknown, + ] + .map(RhiPresenceAttemptOutcome::code), + [ + "submitted", + "accepted", + "rejected", + "rate_limited", + "auth_required", + "failed", + "unknown", + ] + ); + } +} diff --git a/src/state_catalog.rs b/src/state_catalog.rs @@ -12,7 +12,7 @@ use radroots_service_sqlite::{ pub const RHI_STATE_BASE_SCHEMA_VERSION: u32 = 1; /// The newest governed RHI state schema understood by this binary. -pub const RHI_STATE_SCHEMA_VERSION: u32 = 9; +pub const RHI_STATE_SCHEMA_VERSION: u32 = 10; /// The shared metadata and migration-ledger objects present at schema v1. pub const RHI_STATE_SCHEMA_VERSION_1_OBJECT_COUNT: u32 = 6; @@ -41,10 +41,13 @@ pub const RHI_STATE_SCHEMA_VERSION_8_OBJECT_COUNT: u32 = 65; /// The shared objects plus deterministic durable presence desired state. pub const RHI_STATE_SCHEMA_VERSION_9_OBJECT_COUNT: u32 = 69; +/// The shared objects plus durable exact-byte presence delivery state. +pub const RHI_STATE_SCHEMA_VERSION_10_OBJECT_COUNT: u32 = 80; + /// SHA-256 identity of the ordered migration catalog rooted at schema v1. pub const RHI_MIGRATION_CATALOG_SHA256: [u8; 32] = [ - 0x7f, 0x8c, 0x03, 0xb4, 0x81, 0x84, 0x08, 0xb5, 0x86, 0x74, 0x74, 0x12, 0xc2, 0x65, 0xac, 0x72, - 0xcd, 0x4d, 0xd8, 0x8a, 0x29, 0x0d, 0x8b, 0xbe, 0xe1, 0xd5, 0xf6, 0x8a, 0xdb, 0xf7, 0xaf, 0xd0, + 0x25, 0xe5, 0xba, 0x77, 0x3e, 0xf3, 0xdb, 0x01, 0x33, 0xa8, 0x07, 0x7a, 0x08, 0x3e, 0x88, 0xb4, + 0x0f, 0xc6, 0xda, 0xdf, 0x9f, 0xb6, 0xf0, 0xcf, 0xde, 0x0e, 0x4b, 0xa4, 0xd3, 0xe0, 0x81, 0xd9, ]; /// SHA-256 identity of the exact schema-v1 object snapshot. @@ -149,10 +152,22 @@ pub const RHI_STATE_SCHEMA_VERSION_9_SHA256: [u8; 32] = [ 0xa8, 0x48, 0x6d, 0x2d, 0xc0, 0x8b, 0x6a, 0x78, 0x08, 0xbb, 0x20, 0x07, 0xc9, 0x43, 0x35, 0xec, ]; +/// SHA-256 identity of the schema-v10 presence-publication migration. +pub const RHI_STATE_SCHEMA_VERSION_10_MIGRATION_SHA256: [u8; 32] = [ + 0x54, 0x1a, 0xd1, 0x3b, 0x2c, 0xb0, 0x8d, 0x59, 0x85, 0x72, 0x05, 0xe6, 0xff, 0xae, 0xc1, 0x9b, + 0x74, 0xde, 0x08, 0x53, 0x24, 0x0f, 0xdf, 0x11, 0xfa, 0x13, 0x40, 0xe4, 0x06, 0xc7, 0x11, 0x4e, +]; + +/// SHA-256 identity of the exact schema-v10 object snapshot. +pub const RHI_STATE_SCHEMA_VERSION_10_SHA256: [u8; 32] = [ + 0xd2, 0xae, 0xd5, 0x1d, 0x0a, 0x6a, 0x2c, 0x01, 0xed, 0xa1, 0x84, 0x46, 0x08, 0x47, 0x2b, 0x2d, + 0xcd, 0x50, 0x2a, 0xba, 0xa8, 0xa4, 0xca, 0x30, 0xa4, 0x82, 0x35, 0x35, 0xb8, 0xbd, 0x0e, 0x45, +]; + /// SHA-256 identity of the schema catalog bound to the migration catalog. pub const RHI_STATE_SCHEMA_CATALOG_SHA256: [u8; 32] = [ - 0x5b, 0xea, 0x3e, 0x3e, 0xc6, 0xbe, 0x3f, 0x12, 0x49, 0xad, 0x19, 0xed, 0x7b, 0x2d, 0x2e, 0x87, - 0x2c, 0x34, 0x02, 0x55, 0x79, 0xa5, 0x99, 0xf2, 0x54, 0xde, 0x8e, 0x80, 0xcb, 0xa9, 0xb3, 0xc7, + 0x4f, 0x4d, 0x5f, 0x55, 0x46, 0xa7, 0xc9, 0x4e, 0xf3, 0xcd, 0xab, 0xe2, 0x3e, 0xad, 0xd0, 0xc9, + 0x7a, 0x64, 0x98, 0x09, 0x84, 0xee, 0x38, 0xca, 0xcb, 0x82, 0x8f, 0x11, 0x25, 0x91, 0x34, 0x88, ]; macro_rules! rhi_config_bindings_table_sql { @@ -846,6 +861,286 @@ const CREATE_PRESENCE_DESIRED_STATE_MIGRATION_SQL: &str = concat!( ";", ); +macro_rules! presence_outbox_table_sql { + () => { + r#"CREATE TABLE presence_outbox ( + outbox_id BLOB NOT NULL PRIMARY KEY CHECK (length(outbox_id) = 32), + desired_generation INTEGER NOT NULL + CHECK (desired_generation BETWEEN 1 AND 9223372036854775807), + document_kind TEXT NOT NULL + CHECK (document_kind IN ('service_profile', 'application_handler')), + desired_sha256 BLOB NOT NULL CHECK (length(desired_sha256) = 32), + target_set_sha256 BLOB NOT NULL CHECK (length(target_set_sha256) = 32), + event_id BLOB NOT NULL CHECK (length(event_id) = 32), + event_sha256 BLOB NOT NULL CHECK (length(event_sha256) = 32), + exact_signed_event_bytes BLOB NOT NULL + CHECK (length(exact_signed_event_bytes) BETWEEN 1 AND 32768), + authored_at_unix_s INTEGER NOT NULL + CHECK (authored_at_unix_s BETWEEN 0 AND 9223372036854775807), + service_public_key TEXT NOT NULL + CHECK (length(CAST(service_public_key AS BLOB)) = 64) + CHECK (service_public_key NOT GLOB '*[^0-9a-f]*'), + target_count INTEGER NOT NULL CHECK (target_count BETWEEN 1 AND 32), + required_target_count INTEGER NOT NULL + CHECK (required_target_count BETWEEN 0 AND target_count), + max_attempts INTEGER NOT NULL CHECK (max_attempts = 100), + initial_backoff_ms INTEGER NOT NULL CHECK (initial_backoff_ms = 250), + maximum_backoff_ms INTEGER NOT NULL CHECK (maximum_backoff_ms = 30000), + attempt_deadline_ms INTEGER NOT NULL CHECK (attempt_deadline_ms = 15000), + state TEXT NOT NULL + CHECK (state IN ('pending', 'leased', 'complete', 'blocked', 'superseded')), + revision INTEGER NOT NULL CHECK (revision >= 1), + next_attempt_unix_ms INTEGER + CHECK (next_attempt_unix_ms IS NULL + OR next_attempt_unix_ms BETWEEN 0 AND 9223372036854775807), + lease_owner BLOB CHECK (lease_owner IS NULL OR length(lease_owner) = 16), + lease_expires_unix_ms INTEGER + CHECK (lease_expires_unix_ms IS NULL + OR lease_expires_unix_ms BETWEEN 1 AND 9223372036854775807), + created_at_unix_ms INTEGER NOT NULL + CHECK (created_at_unix_ms BETWEEN 0 AND 9223372036854775807), + updated_at_unix_ms INTEGER NOT NULL + CHECK (updated_at_unix_ms BETWEEN created_at_unix_ms AND 9223372036854775807), + UNIQUE (desired_generation, document_kind), + CHECK ( + (state = 'pending' AND next_attempt_unix_ms IS NOT NULL + AND lease_owner IS NULL AND lease_expires_unix_ms IS NULL) + OR + (state = 'leased' AND next_attempt_unix_ms IS NULL + AND lease_owner IS NOT NULL AND lease_expires_unix_ms IS NOT NULL) + OR + (state IN ('complete', 'blocked', 'superseded') + AND next_attempt_unix_ms IS NULL + AND lease_owner IS NULL AND lease_expires_unix_ms IS NULL) + ) +) STRICT"# + }; +} + +macro_rules! presence_outbox_schedule_sql { + () => { + r#"CREATE INDEX presence_outbox_by_schedule +ON presence_outbox (state, next_attempt_unix_ms, created_at_unix_ms, document_kind, outbox_id)"# + }; +} + +macro_rules! presence_outbox_guard_update_sql { + () => { + r#"CREATE TRIGGER presence_outbox_guard_update +BEFORE UPDATE ON presence_outbox +WHEN NEW.outbox_id != OLD.outbox_id + OR NEW.desired_generation != OLD.desired_generation + OR NEW.document_kind != OLD.document_kind + OR NEW.desired_sha256 != OLD.desired_sha256 + OR NEW.target_set_sha256 != OLD.target_set_sha256 + OR NEW.event_id != OLD.event_id + OR NEW.event_sha256 != OLD.event_sha256 + OR NEW.exact_signed_event_bytes != OLD.exact_signed_event_bytes + OR NEW.authored_at_unix_s != OLD.authored_at_unix_s + OR NEW.service_public_key != OLD.service_public_key + OR NEW.target_count != OLD.target_count + OR NEW.required_target_count != OLD.required_target_count + OR NEW.max_attempts != OLD.max_attempts + OR NEW.initial_backoff_ms != OLD.initial_backoff_ms + OR NEW.maximum_backoff_ms != OLD.maximum_backoff_ms + OR NEW.attempt_deadline_ms != OLD.attempt_deadline_ms + OR NEW.created_at_unix_ms != OLD.created_at_unix_ms + OR NEW.revision != OLD.revision + 1 + OR NEW.updated_at_unix_ms < OLD.updated_at_unix_ms + OR NOT ( + (OLD.state = 'pending' AND NEW.state IN ('leased', 'superseded')) + OR (OLD.state = 'leased' + AND NEW.state IN ('pending', 'complete', 'blocked', 'superseded')) + OR (OLD.state = 'blocked' AND NEW.state IN ('pending', 'superseded')) + ) +BEGIN + SELECT RAISE(ABORT, 'presence outbox transition is invalid'); +END"# + }; +} + +macro_rules! presence_targets_table_sql { + () => { + r#"CREATE TABLE presence_targets ( + outbox_id BLOB NOT NULL CHECK (length(outbox_id) = 32), + target_ordinal INTEGER NOT NULL CHECK (target_ordinal BETWEEN 0 AND 31), + relay_id TEXT NOT NULL + CHECK (length(CAST(relay_id AS BLOB)) BETWEEN 1 AND 64) + CHECK (relay_id NOT GLOB '*[^a-z0-9_-]*') + CHECK (substr(relay_id, 1, 1) GLOB '[a-z]'), + required INTEGER NOT NULL CHECK (required IN (0, 1)), + state TEXT NOT NULL CHECK (state IN ( + 'pending', 'submitted', 'accepted', 'rejected', + 'rate_limited', 'auth_required', 'failed', 'unknown' + )), + revision INTEGER NOT NULL CHECK (revision >= 1), + attempt_count INTEGER NOT NULL CHECK (attempt_count BETWEEN 0 AND 100), + next_attempt_unix_ms INTEGER + CHECK (next_attempt_unix_ms IS NULL + OR next_attempt_unix_ms BETWEEN 0 AND 9223372036854775807), + last_attempt_id BLOB + CHECK (last_attempt_id IS NULL OR length(last_attempt_id) = 32), + updated_at_unix_ms INTEGER NOT NULL + CHECK (updated_at_unix_ms BETWEEN 0 AND 9223372036854775807), + PRIMARY KEY (outbox_id, target_ordinal), + UNIQUE (outbox_id, relay_id), + FOREIGN KEY (outbox_id) REFERENCES presence_outbox (outbox_id), + CHECK ( + (state = 'pending' AND attempt_count = 0 + AND next_attempt_unix_ms IS NOT NULL AND last_attempt_id IS NULL) + OR + (state = 'submitted' AND attempt_count BETWEEN 1 AND 100 + AND next_attempt_unix_ms IS NULL AND last_attempt_id IS NOT NULL) + OR + (state IN ('accepted', 'rejected', 'auth_required') + AND attempt_count BETWEEN 1 AND 100 + AND next_attempt_unix_ms IS NULL AND last_attempt_id IS NOT NULL) + OR + (state IN ('rate_limited', 'failed', 'unknown') + AND attempt_count BETWEEN 1 AND 99 + AND next_attempt_unix_ms IS NOT NULL AND last_attempt_id IS NOT NULL) + OR + (state IN ('rate_limited', 'failed', 'unknown') + AND attempt_count = 100 + AND next_attempt_unix_ms IS NULL AND last_attempt_id IS NOT NULL) + ) +) STRICT"# + }; +} + +macro_rules! presence_targets_schedule_sql { + () => { + r#"CREATE INDEX presence_targets_by_schedule +ON presence_targets (state, next_attempt_unix_ms, outbox_id, target_ordinal)"# + }; +} + +macro_rules! presence_targets_guard_update_sql { + () => { + r#"CREATE TRIGGER presence_targets_guard_update +BEFORE UPDATE ON presence_targets +WHEN NEW.outbox_id != OLD.outbox_id + OR NEW.target_ordinal != OLD.target_ordinal + OR NEW.relay_id != OLD.relay_id + OR NEW.required != OLD.required + OR NEW.revision != OLD.revision + 1 + OR NEW.updated_at_unix_ms < OLD.updated_at_unix_ms + OR NOT ( + (OLD.state IN ('pending', 'rate_limited', 'failed', 'unknown') + AND NEW.state = 'submitted' + AND NEW.attempt_count = OLD.attempt_count + 1 + AND NEW.last_attempt_id IS NOT NULL + AND NEW.last_attempt_id IS NOT OLD.last_attempt_id) + OR + (OLD.state = 'submitted' + AND NEW.state IN ( + 'accepted', 'rejected', 'rate_limited', + 'auth_required', 'failed', 'unknown' + ) + AND NEW.attempt_count = OLD.attempt_count + AND NEW.last_attempt_id = OLD.last_attempt_id) + ) +BEGIN + SELECT RAISE(ABORT, 'presence target transition is invalid'); +END"# + }; +} + +macro_rules! presence_attempts_table_sql { + () => { + r#"CREATE TABLE presence_attempts ( + attempt_id BLOB NOT NULL PRIMARY KEY CHECK (length(attempt_id) = 32), + outbox_id BLOB NOT NULL CHECK (length(outbox_id) = 32), + target_ordinal INTEGER NOT NULL CHECK (target_ordinal BETWEEN 0 AND 31), + attempt_number INTEGER NOT NULL CHECK (attempt_number BETWEEN 1 AND 100), + event_sha256 BLOB NOT NULL CHECK (length(event_sha256) = 32), + lease_owner BLOB NOT NULL CHECK (length(lease_owner) = 16), + started_at_unix_ms INTEGER NOT NULL + CHECK (started_at_unix_ms BETWEEN 0 AND 9223372036854775807), + finished_at_unix_ms INTEGER NOT NULL + CHECK (finished_at_unix_ms BETWEEN started_at_unix_ms AND 9223372036854775807), + outcome TEXT NOT NULL CHECK (outcome IN ( + 'accepted', 'rejected', 'rate_limited', + 'auth_required', 'failed', 'unknown' + )), + result_code TEXT NOT NULL + CHECK (length(CAST(result_code AS BLOB)) BETWEEN 1 AND 64), + UNIQUE (outbox_id, target_ordinal, attempt_number), + FOREIGN KEY (outbox_id, target_ordinal) + REFERENCES presence_targets (outbox_id, target_ordinal) +) STRICT"# + }; +} + +const CREATE_PRESENCE_OUTBOX_TABLE_SQL: &str = presence_outbox_table_sql!(); +const CREATE_PRESENCE_OUTBOX_SCHEDULE_SQL: &str = presence_outbox_schedule_sql!(); +const CREATE_PRESENCE_OUTBOX_GUARD_UPDATE_SQL: &str = presence_outbox_guard_update_sql!(); +const CREATE_PRESENCE_OUTBOX_NO_DELETE_SQL: &str = immutable_no_delete_sql!( + "presence_outbox_no_delete", + "presence_outbox", + "presence outbox rows are retained" +); +const CREATE_PRESENCE_TARGETS_TABLE_SQL: &str = presence_targets_table_sql!(); +const CREATE_PRESENCE_TARGETS_SCHEDULE_SQL: &str = presence_targets_schedule_sql!(); +const CREATE_PRESENCE_TARGETS_GUARD_UPDATE_SQL: &str = presence_targets_guard_update_sql!(); +const CREATE_PRESENCE_TARGETS_NO_DELETE_SQL: &str = immutable_no_delete_sql!( + "presence_targets_no_delete", + "presence_targets", + "presence targets are retained" +); +const CREATE_PRESENCE_ATTEMPTS_TABLE_SQL: &str = presence_attempts_table_sql!(); +const CREATE_PRESENCE_ATTEMPTS_NO_UPDATE_SQL: &str = immutable_no_update_sql!( + "presence_attempts_no_update", + "presence_attempts", + "presence attempts are immutable" +); +const CREATE_PRESENCE_ATTEMPTS_NO_DELETE_SQL: &str = immutable_no_delete_sql!( + "presence_attempts_no_delete", + "presence_attempts", + "presence attempts are retained" +); + +const CREATE_PRESENCE_PUBLICATION_MIGRATION_SQL: &str = concat!( + presence_outbox_table_sql!(), + ";\n", + presence_outbox_schedule_sql!(), + ";\n", + presence_outbox_guard_update_sql!(), + ";\n", + immutable_no_delete_sql!( + "presence_outbox_no_delete", + "presence_outbox", + "presence outbox rows are retained" + ), + ";\n", + presence_targets_table_sql!(), + ";\n", + presence_targets_schedule_sql!(), + ";\n", + presence_targets_guard_update_sql!(), + ";\n", + immutable_no_delete_sql!( + "presence_targets_no_delete", + "presence_targets", + "presence targets are retained" + ), + ";\n", + presence_attempts_table_sql!(), + ";\n", + immutable_no_update_sql!( + "presence_attempts_no_update", + "presence_attempts", + "presence attempts are immutable" + ), + ";\n", + immutable_no_delete_sql!( + "presence_attempts_no_delete", + "presence_attempts", + "presence attempts are retained" + ), + ";" +); + macro_rules! evidence_reconciliations_table_sql { () => { r#"CREATE TABLE evidence_reconciliations ( @@ -1651,6 +1946,51 @@ const PRESENCE_DESIRED_STATE_NO_DELETE_SHA256: [u8; 32] = [ 0xca, 0x44, 0xab, 0x2e, 0x4f, 0x77, 0x65, 0x29, 0x75, 0x02, 0x4b, 0xbf, 0x1f, 0x04, 0xdc, 0x3f, ]; +const PRESENCE_OUTBOX_TABLE_SHA256: [u8; 32] = [ + 0x7e, 0x59, 0xa4, 0xb4, 0x54, 0x09, 0x9e, 0x47, 0xf8, 0xad, 0x0d, 0xb8, 0x7f, 0xd2, 0xa0, 0xc5, + 0xf6, 0x69, 0x8d, 0x70, 0xda, 0x96, 0x8d, 0xb0, 0x57, 0x6a, 0xa8, 0x54, 0x28, 0xd1, 0x40, 0xdd, +]; +const PRESENCE_OUTBOX_SCHEDULE_SHA256: [u8; 32] = [ + 0xed, 0x0e, 0x52, 0xba, 0xc1, 0xa8, 0xc8, 0xe3, 0x25, 0xf0, 0x28, 0xe4, 0x2c, 0x8d, 0x88, 0x8f, + 0x0e, 0x8b, 0x42, 0x63, 0x7c, 0x95, 0x25, 0xf1, 0x02, 0x0c, 0x18, 0x41, 0x3b, 0xf0, 0xf0, 0x41, +]; +const PRESENCE_OUTBOX_GUARD_UPDATE_SHA256: [u8; 32] = [ + 0x1d, 0x80, 0x93, 0xa8, 0xe7, 0x70, 0x82, 0xc4, 0x2f, 0xb3, 0xe4, 0x79, 0xc6, 0x22, 0x6f, 0x39, + 0x92, 0x88, 0x80, 0xd4, 0x92, 0xa9, 0x6c, 0x04, 0x29, 0x6b, 0xf2, 0xcc, 0x56, 0x69, 0x6e, 0xce, +]; +const PRESENCE_OUTBOX_NO_DELETE_SHA256: [u8; 32] = [ + 0xc4, 0xe9, 0xb5, 0x5a, 0x02, 0x0b, 0x67, 0x1d, 0x7d, 0x6c, 0x81, 0xf6, 0x99, 0x90, 0x65, 0x52, + 0xd9, 0x5e, 0x97, 0xfd, 0xba, 0xac, 0xad, 0x7b, 0x01, 0xb5, 0x02, 0x3f, 0x6b, 0x5f, 0x79, 0xeb, +]; +const PRESENCE_TARGETS_TABLE_SHA256: [u8; 32] = [ + 0x1e, 0xfa, 0xc9, 0xe2, 0x8d, 0x55, 0xcf, 0x90, 0xa5, 0x5d, 0xc4, 0x13, 0x1a, 0x8e, 0xcb, 0xc3, + 0x86, 0x2e, 0xff, 0xa9, 0x3d, 0xb3, 0x63, 0xff, 0xed, 0xd4, 0xa1, 0xd6, 0xcd, 0xa4, 0xb7, 0x8c, +]; +const PRESENCE_TARGETS_SCHEDULE_SHA256: [u8; 32] = [ + 0x6d, 0xbc, 0xf6, 0x25, 0x20, 0x1e, 0x11, 0x2e, 0x1a, 0x54, 0xa7, 0xc8, 0x87, 0xaf, 0xbb, 0xb8, + 0x21, 0x0a, 0xdc, 0xf6, 0xfa, 0xc1, 0xc5, 0x99, 0xe1, 0xaa, 0x6c, 0xa5, 0xc4, 0x56, 0x74, 0x6a, +]; +const PRESENCE_TARGETS_GUARD_UPDATE_SHA256: [u8; 32] = [ + 0xa0, 0x8f, 0x28, 0x9e, 0x60, 0x32, 0x27, 0xfa, 0x74, 0x34, 0xea, 0x8d, 0x4a, 0xb5, 0x2a, 0x07, + 0x60, 0xdc, 0xd2, 0x8f, 0xe1, 0x35, 0x33, 0x51, 0xa3, 0x5c, 0x3e, 0x6c, 0x73, 0x5d, 0x12, 0x44, +]; +const PRESENCE_TARGETS_NO_DELETE_SHA256: [u8; 32] = [ + 0x8f, 0x0a, 0x46, 0x61, 0x55, 0x96, 0xa8, 0x33, 0x6f, 0x5f, 0x56, 0xa0, 0x3e, 0xa9, 0xd3, 0x5c, + 0x7d, 0xc3, 0x69, 0xb3, 0xa3, 0x1d, 0x91, 0x0c, 0x23, 0x05, 0xeb, 0xeb, 0xcb, 0x50, 0x71, 0x25, +]; +const PRESENCE_ATTEMPTS_TABLE_SHA256: [u8; 32] = [ + 0xed, 0x26, 0x46, 0x45, 0x6f, 0x96, 0x45, 0xde, 0x57, 0xcb, 0xef, 0x79, 0x62, 0xd6, 0xfe, 0x43, + 0x9d, 0xae, 0x10, 0xfd, 0x5e, 0x4c, 0x25, 0x90, 0x91, 0x62, 0x58, 0x8f, 0x58, 0x0d, 0x94, 0x9c, +]; +const PRESENCE_ATTEMPTS_NO_UPDATE_SHA256: [u8; 32] = [ + 0x8b, 0xe2, 0xd1, 0xa5, 0xe6, 0x5d, 0xb9, 0x96, 0xe3, 0x27, 0xb8, 0xee, 0x27, 0x29, 0xa3, 0xe9, + 0x1d, 0xa0, 0x68, 0x0e, 0x16, 0x05, 0x90, 0x77, 0x3a, 0x78, 0x1e, 0x80, 0x00, 0x35, 0xac, 0x44, +]; +const PRESENCE_ATTEMPTS_NO_DELETE_SHA256: [u8; 32] = [ + 0x8c, 0x59, 0xe5, 0x11, 0x0e, 0xe5, 0xa7, 0x45, 0x3d, 0xde, 0xca, 0x36, 0xc2, 0x87, 0xc3, 0x6e, + 0xbd, 0xf7, 0x74, 0xbb, 0x8d, 0x84, 0xc9, 0xab, 0xd9, 0x76, 0xe2, 0xae, 0x9c, 0xd2, 0x4a, 0x71, +]; + /// Stable classes for invalid embedded RHI catalog definitions. #[derive(Clone, Copy, Debug, PartialEq, Eq)] pub enum RhiStateCatalogErrorKind { @@ -1777,6 +2117,13 @@ fn build_rhi_migration_catalog() -> Result<MigrationCatalog, RhiStateCatalogErro MigrationChecksum::from_bytes(RHI_STATE_SCHEMA_VERSION_9_MIGRATION_SHA256), ) .map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::MigrationCatalog))?; + let presence_publication = MigrationDescriptor::sql( + 10, + "create_presence_publication_workflow", + CREATE_PRESENCE_PUBLICATION_MIGRATION_SQL, + MigrationChecksum::from_bytes(RHI_STATE_SCHEMA_VERSION_10_MIGRATION_SHA256), + ) + .map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::MigrationCatalog))?; let catalog = MigrationCatalog::new([ configuration, trade_evidence, @@ -1786,6 +2133,7 @@ fn build_rhi_migration_catalog() -> Result<MigrationCatalog, RhiStateCatalogErro report_publication, reconciliation_job_shape_guards, presence_desired_state, + presence_publication, ]) .map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::MigrationCatalog))?; Ok(catalog) @@ -1795,7 +2143,7 @@ fn build_rhi_migration_catalog() -> Result<MigrationCatalog, RhiStateCatalogErro pub fn rhi_migration_catalog() -> Result<MigrationCatalog, RhiStateCatalogError> { let catalog = build_rhi_migration_catalog()?; if catalog.current_version() != RHI_STATE_SCHEMA_VERSION - || catalog.descriptors().len() != 8 + || catalog.descriptors().len() != 9 || catalog.digest().as_bytes() != &RHI_MIGRATION_CATALOG_SHA256 { return Err(RhiStateCatalogError::new( @@ -1865,11 +2213,17 @@ fn build_rhi_schema_catalog( ) .map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::SchemaCatalog))?; let version_nine = SchemaVersionCatalog::new( - RHI_STATE_SCHEMA_VERSION, + 9, rhi_schema_version_nine_objects()?, SchemaDigest::from_bytes(RHI_STATE_SCHEMA_VERSION_9_SHA256), ) .map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::SchemaCatalog))?; + let version_ten = SchemaVersionCatalog::new( + 10, + rhi_schema_version_ten_objects()?, + SchemaDigest::from_bytes(RHI_STATE_SCHEMA_VERSION_10_SHA256), + ) + .map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::SchemaCatalog))?; let catalog = SchemaCatalog::new( migrations, [ @@ -1882,6 +2236,7 @@ fn build_rhi_schema_catalog( version_seven, version_eight, version_nine, + version_ten, ], ) .map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::SchemaCatalog))?; @@ -1895,10 +2250,10 @@ pub fn validate_rhi_state_catalogs( ) -> Result<(), RhiStateCatalogError> { let versions = schema.versions(); let valid = migrations.current_version() == RHI_STATE_SCHEMA_VERSION - && migrations.descriptors().len() == 8 + && migrations.descriptors().len() == 9 && migrations.digest().as_bytes() == &RHI_MIGRATION_CATALOG_SHA256 && schema.migration_catalog_digest() == migrations.digest() - && versions.len() == 9 + && versions.len() == 10 && versions[0].version() == RHI_STATE_BASE_SCHEMA_VERSION && versions[0].object_count() == RHI_STATE_SCHEMA_VERSION_1_OBJECT_COUNT && versions[0].digest().as_bytes() == &RHI_STATE_SCHEMA_VERSION_1_SHA256 @@ -1923,9 +2278,12 @@ pub fn validate_rhi_state_catalogs( && versions[7].version() == 8 && versions[7].object_count() == RHI_STATE_SCHEMA_VERSION_8_OBJECT_COUNT && versions[7].digest().as_bytes() == &RHI_STATE_SCHEMA_VERSION_8_SHA256 - && versions[8].version() == RHI_STATE_SCHEMA_VERSION + && versions[8].version() == 9 && versions[8].object_count() == RHI_STATE_SCHEMA_VERSION_9_OBJECT_COUNT && versions[8].digest().as_bytes() == &RHI_STATE_SCHEMA_VERSION_9_SHA256 + && versions[9].version() == RHI_STATE_SCHEMA_VERSION + && versions[9].object_count() == RHI_STATE_SCHEMA_VERSION_10_OBJECT_COUNT + && versions[9].digest().as_bytes() == &RHI_STATE_SCHEMA_VERSION_10_SHA256 && schema.digest().as_bytes() == &RHI_STATE_SCHEMA_CATALOG_SHA256; if valid { Ok(()) @@ -2015,6 +2373,104 @@ fn rhi_schema_version_nine_objects() -> Result<Vec<SchemaObject>, RhiStateCatalo Ok(objects) } +fn rhi_schema_version_ten_objects() -> Result<Vec<SchemaObject>, RhiStateCatalogError> { + let mut objects = rhi_schema_version_nine_objects()?; + objects.extend(rhi_presence_publication_objects()?); + Ok(objects) +} + +fn rhi_presence_publication_objects() -> Result<[SchemaObject; 11], RhiStateCatalogError> { + let object = |kind, name, table_name, sql, digest| { + SchemaObject::new( + kind, + name, + table_name, + sql, + SchemaDigest::from_bytes(digest), + ) + .map_err(|_| RhiStateCatalogError::new(RhiStateCatalogErrorKind::SchemaCatalog)) + }; + Ok([ + object( + SchemaObjectKind::Table, + "presence_outbox", + "presence_outbox", + CREATE_PRESENCE_OUTBOX_TABLE_SQL, + PRESENCE_OUTBOX_TABLE_SHA256, + )?, + object( + SchemaObjectKind::Index, + "presence_outbox_by_schedule", + "presence_outbox", + CREATE_PRESENCE_OUTBOX_SCHEDULE_SQL, + PRESENCE_OUTBOX_SCHEDULE_SHA256, + )?, + object( + SchemaObjectKind::Trigger, + "presence_outbox_guard_update", + "presence_outbox", + CREATE_PRESENCE_OUTBOX_GUARD_UPDATE_SQL, + PRESENCE_OUTBOX_GUARD_UPDATE_SHA256, + )?, + object( + SchemaObjectKind::Trigger, + "presence_outbox_no_delete", + "presence_outbox", + CREATE_PRESENCE_OUTBOX_NO_DELETE_SQL, + PRESENCE_OUTBOX_NO_DELETE_SHA256, + )?, + object( + SchemaObjectKind::Table, + "presence_targets", + "presence_targets", + CREATE_PRESENCE_TARGETS_TABLE_SQL, + PRESENCE_TARGETS_TABLE_SHA256, + )?, + object( + SchemaObjectKind::Index, + "presence_targets_by_schedule", + "presence_targets", + CREATE_PRESENCE_TARGETS_SCHEDULE_SQL, + PRESENCE_TARGETS_SCHEDULE_SHA256, + )?, + object( + SchemaObjectKind::Trigger, + "presence_targets_guard_update", + "presence_targets", + CREATE_PRESENCE_TARGETS_GUARD_UPDATE_SQL, + PRESENCE_TARGETS_GUARD_UPDATE_SHA256, + )?, + object( + SchemaObjectKind::Trigger, + "presence_targets_no_delete", + "presence_targets", + CREATE_PRESENCE_TARGETS_NO_DELETE_SQL, + PRESENCE_TARGETS_NO_DELETE_SHA256, + )?, + object( + SchemaObjectKind::Table, + "presence_attempts", + "presence_attempts", + CREATE_PRESENCE_ATTEMPTS_TABLE_SQL, + PRESENCE_ATTEMPTS_TABLE_SHA256, + )?, + object( + SchemaObjectKind::Trigger, + "presence_attempts_no_update", + "presence_attempts", + CREATE_PRESENCE_ATTEMPTS_NO_UPDATE_SQL, + PRESENCE_ATTEMPTS_NO_UPDATE_SHA256, + )?, + object( + SchemaObjectKind::Trigger, + "presence_attempts_no_delete", + "presence_attempts", + CREATE_PRESENCE_ATTEMPTS_NO_DELETE_SQL, + PRESENCE_ATTEMPTS_NO_DELETE_SHA256, + )?, + ]) +} + fn rhi_presence_desired_state_objects() -> Result<[SchemaObject; 4], RhiStateCatalogError> { let object = |kind, name, sql, digest| { SchemaObject::new( diff --git a/src/state_repository.rs b/src/state_repository.rs @@ -8,7 +8,7 @@ use crate::RhiStateHost; pub const RHI_STATE_REPOSITORY_CONTRACT_VERSION: u32 = 1; /// Number of distinct typed repository capabilities in the v1 topology. -pub const RHI_STATE_REPOSITORY_COUNT: usize = 18; +pub const RHI_STATE_REPOSITORY_COUNT: usize = 21; /// Closed mutation class for one governed repository. #[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)] @@ -52,6 +52,9 @@ pub enum RhiStateRepositoryKind { PublicationTarget, PublicationAttempt, DesiredPresence, + PresenceOutbox, + PresenceTarget, + PresenceAttempt, } impl RhiStateRepositoryKind { @@ -245,6 +248,24 @@ const DESCRIPTORS: [RhiStateRepositoryDescriptor; RHI_STATE_REPOSITORY_COUNT] = "presence_desired_state", Write::CompareAndSwap, ), + RhiStateRepositoryDescriptor::new( + Kind::PresenceOutbox, + "presence_outbox", + "presence_outbox", + Write::CompareAndSwap, + ), + RhiStateRepositoryDescriptor::new( + Kind::PresenceTarget, + "presence_target", + "presence_targets", + Write::CompareAndSwap, + ), + RhiStateRepositoryDescriptor::new( + Kind::PresenceAttempt, + "presence_attempt", + "presence_attempts", + Write::AppendOnly, + ), ]; const fn descriptor(kind: RhiStateRepositoryKind) -> RhiStateRepositoryDescriptor { @@ -397,6 +418,24 @@ impl<'host> RhiStateRepositories<'host> { pub const fn desired_presence(&self) -> RhiDesiredPresenceRepository<'host> { RhiDesiredPresenceRepository { host: self.host } } + + /// Returns typed exact-byte presence-outbox access. + #[must_use] + pub const fn presence_outbox(&self) -> RhiPresenceOutboxRepository<'host> { + RhiPresenceOutboxRepository { host: self.host } + } + + /// Returns typed durable presence-target access. + #[must_use] + pub const fn presence_targets(&self) -> RhiPresenceTargetRepository<'host> { + RhiPresenceTargetRepository { host: self.host } + } + + /// Returns typed append-only presence-attempt access. + #[must_use] + pub const fn presence_attempts(&self) -> RhiPresenceAttemptRepository<'host> { + RhiPresenceAttemptRepository { host: self.host } + } } impl fmt::Debug for RhiStateRepositories<'_> { @@ -460,6 +499,9 @@ repository_handle!(RhiPublicationOutboxRepository, PublicationOutbox); repository_handle!(RhiPublicationTargetRepository, PublicationTarget); repository_handle!(RhiPublicationAttemptRepository, PublicationAttempt); repository_handle!(RhiDesiredPresenceRepository, DesiredPresence); +repository_handle!(RhiPresenceOutboxRepository, PresenceOutbox); +repository_handle!(RhiPresenceTargetRepository, PresenceTarget); +repository_handle!(RhiPresenceAttemptRepository, PresenceAttempt); impl<'host> RhiReconciliationJobRepository<'host> { pub(crate) const fn host(&self) -> &'host RhiStateHost { @@ -485,6 +527,12 @@ impl<'host> RhiDesiredPresenceRepository<'host> { } } +impl<'host> RhiPresenceOutboxRepository<'host> { + pub(crate) const fn host(&self) -> &'host RhiStateHost { + self.host + } +} + #[cfg(test)] mod tests { use super::*; diff --git a/tests/package_boundary.rs b/tests/package_boundary.rs @@ -14,6 +14,9 @@ const RUNTIME_FOUNDATION: &str = include_str!("../src/runtime_foundation.rs"); const PRESENCE_DESIRED: &str = include_str!("../src/presence_desired.rs"); const PRESENCE_DESIRED_CONTRACT: &str = include_str!("../contracts/services_hardening/presence_desired_state.v1.json"); +const PRESENCE_PUBLICATION: &str = include_str!("../src/presence_publication.rs"); +const PRESENCE_PUBLICATION_CONTRACT: &str = + include_str!("../contracts/services_hardening/presence_publication.v1.json"); const PUBLICATION: &str = include_str!("../src/publication.rs"); const PUBLICATION_ATTEMPT: &str = include_str!("../src/publication_attempt.rs"); const PUBLICATION_ATTEMPT_CONTRACT: &str = @@ -73,6 +76,7 @@ const SOURCES: &[&str] = &[ include_str!("../src/identity_credential.rs"), include_str!("../src/identity_envelope.rs"), include_str!("../src/presence_desired.rs"), + include_str!("../src/presence_publication.rs"), include_str!("../src/publication.rs"), include_str!("../src/publication_attempt.rs"), include_str!("../src/publication_execution.rs"), @@ -153,6 +157,7 @@ fn state_catalog_module_is_private_and_root_api_is_curated() { "identity_credential", "identity_envelope", "presence_desired", + "presence_publication", "publication", "publication_attempt", "publication_execution", @@ -208,6 +213,20 @@ fn state_catalog_module_is_private_and_root_api_is_curated() { "RhiPresenceDesiredState", "validate_rhi_presence_desired_authority", "RHI_PRESENCE_DESIRED_CONTRACT_VERSION", + "RhiExactPresenceSink", + "RhiPreparedPresenceAttempt", + "RhiPresenceAttemptCommit", + "RhiPresenceAttemptOutcome", + "RhiPresenceLease", + "RhiPresenceLeaseOwner", + "RhiPresenceOutboxState", + "RhiPresencePublicationErrorKind", + "RhiPresenceTargetState", + "RhiSignedPresenceDocument", + "RhiSignedPresenceDocuments", + "build_rhi_signed_presence_documents", + "validate_rhi_signed_presence_documents", + "RHI_PRESENCE_PUBLICATION_CONTRACT_VERSION", "RhiPublicationErrorKind", "RhiPublicationMode", "RhiPublicationRetryPolicy", @@ -320,6 +339,8 @@ fn state_catalog_module_is_private_and_root_api_is_curated() { assert!(PUBLIC_API.contains("pub struct rhi::RhiPresenceDesiredAuthority")); assert!(PUBLIC_API.contains("pub struct rhi::RhiPresenceDesiredState")); assert!(PUBLIC_API.contains("pub enum rhi::RhiPresenceDesiredErrorKind")); + assert!(PUBLIC_API.contains("pub struct rhi::RhiSignedPresenceDocument")); + assert!(PUBLIC_API.contains("pub enum rhi::RhiPresencePublicationErrorKind")); assert!(!PUBLIC_API.contains("rhi::adapters::")); assert!(!PUBLIC_API.contains("rhi::features::")); assert!(!PUBLIC_API.contains("rhi::runtime_adapters::")); @@ -327,6 +348,7 @@ fn state_catalog_module_is_private_and_root_api_is_curated() { assert!(!PUBLIC_API.contains("rhi::publication_attempt::")); assert!(!PUBLIC_API.contains("rhi::publication_execution::")); assert!(!PUBLIC_API.contains("rhi::publication_submission::")); + assert!(!PUBLIC_API.contains("rhi::presence_publication::")); } #[test] @@ -371,6 +393,91 @@ fn presence_desired_state_is_config_bound_durable_and_effect_free() { } #[test] +fn presence_publication_is_typed_verified_durable_and_exact_byte_only() { + let contract: serde_json::Value = + serde_json::from_str(PRESENCE_PUBLICATION_CONTRACT).expect("presence-publication contract"); + assert_eq!(contract["schema"], "radroots.rhi.presence-publication"); + assert_eq!(contract["schema_version"], 1); + assert_eq!(contract["contract_version"], 1); + assert_eq!(contract["step"], 205); + assert_eq!(contract["durable_workflow"]["schema_version"], 10); + assert_eq!( + contract["durable_workflow"]["exact_signed_bytes_committed_before_io"], + true + ); + assert_eq!( + contract["durable_workflow"]["target_submitted_committed_before_io"], + true + ); + assert_eq!( + contract["durable_workflow"]["remote_io_inside_sql_transaction"], + false + ); + assert_eq!( + contract["durable_workflow"]["commit_outcome_unknown"], + "caller_retains_sealed_exact_bytes_and_rereads_exact_identity_before_retry" + ); + assert_eq!( + contract["construction"]["validation_input"], + "sealed_documents_and_public_desired_authority_only" + ); + assert_eq!( + contract["resource_bounds"]["maximum_signed_event_bytes"], + 32_768 + ); + assert_eq!( + contract["resource_bounds"]["maximum_attempts_per_target"], + 100 + ); + assert_eq!( + contract["resource_bounds"]["maximum_authored_unix_seconds"], + i64::MAX + ); + for required in [ + "AuthoredProfile::new(PROFILE_NAME)", + "ApplicationHandlerSpec::new(APPLICATION_HANDLER_KINDS.to_vec())", + "validate_signed_document(", + "pub fn validate_rhi_signed_presence_documents(\n documents: &RhiSignedPresenceDocuments,\n authority: &RhiPresenceDesiredAuthority,\n)", + "verify_id(&event)", + "verify(&event)", + "pub trait RhiExactPresenceSink: Send + Sync", + "pub async fn commit_signed_presence(\n &self,\n documents: &RhiSignedPresenceDocuments,", + "pub async fn claim_next_presence(", + "pub async fn prepare_next_presence_target(", + "pub async fn record_presence_outcome(", + "pub async fn recover_one_expired_presence(", + "pub async fn execute_next_presence(", + "state = 'submitted'", + "INSERT INTO presence_attempts", + "sink.submit_exact(&prepared).await", + "ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown", + ] { + assert!( + PRESENCE_PUBLICATION.contains(required), + "presence publication is missing {required}" + ); + } + for forbidden in [ + "SystemTime", + "thread_rng", + "OsRng", + "tokio::spawn", + "std::net", + "std::fs", + "SqliteConnection", + "SqlitePool", + "EventSink", + ] { + assert!( + !PRESENCE_PUBLICATION.contains(forbidden), + "presence publication gained forbidden authority {forbidden}" + ); + } + assert!(!ROOT.contains("pub mod presence_publication")); + assert!(!PUBLIC_API.contains("rhi::presence_publication::")); +} + +#[test] fn publication_authority_is_config_derived_sealed_and_effect_free() { let contract: serde_json::Value = serde_json::from_str(PUBLICATION_CONTRACT).expect("publication contract"); @@ -854,7 +961,7 @@ fn public_errors_are_crate_owned_redacted_and_source_free() { .lines() .filter(|line| line.starts_with("pub struct rhi::") && line.ends_with("Error")) .count(); - assert_eq!(public_error_count, 30); + assert_eq!(public_error_count, 31); } #[test] @@ -1290,6 +1397,11 @@ fn readme_freezes_the_root_only_boundary_and_exact_baseline() { "## Explicit publication authority and durable schema", "## Deterministic durable presence intent", "[`presence_desired_state.v1.json`](contracts/services_hardening/presence_desired_state.v1.json)", + "## Durable exact-byte presence publication", + "[`presence_publication.v1.json`](contracts/services_hardening/presence_publication.v1.json)", + "Expired work from a", + "superseded desired generation is recorded as `unknown`", + "unknown commit result can be reconciled by replaying the same retained bytes", "[`publication_outbox.v1.json`](contracts/services_hardening/publication_outbox.v1.json)", "## Bounded publication attempt evidence", "[`publication_attempt_evidence.v1.json`](contracts/services_hardening/publication_attempt_evidence.v1.json)", @@ -1359,6 +1471,9 @@ fn readme_freezes_the_root_only_boundary_and_exact_baseline() { "Persist each canonical", "independently signed Nostr event", "Commit a signed finalization only through the sealed attempt repository", + "Build service-profile and application-handler presence only through the", + "Preserve the\n caller's sealed exact-byte capability", + "Recover an expired stale", ] { assert!(AGENTS.contains(required), "AGENTS is missing {required}"); } diff --git a/tests/services_hardening_presence_desired_state.rs b/tests/services_hardening_presence_desired_state.rs @@ -263,7 +263,7 @@ async fn desired_state_is_durable_exact_replay_and_semantic_change_only() { } #[test] -fn machine_contract_freezes_step_204_and_defers_step_205_effects() { +fn machine_contract_freezes_desired_state_without_publication_effects() { let contract: serde_json::Value = serde_json::from_str(CONTRACT).expect("contract"); assert_eq!(contract["schema"], "radroots.rhi.presence-desired-state"); assert_eq!(contract["contract_version"], 1); @@ -280,7 +280,13 @@ fn machine_contract_freezes_step_204_and_defers_step_205_effects() { assert_eq!(contract["effects"]["entropy"], false); assert_eq!(contract["effects"]["network"], false); assert_eq!(contract["effects"]["relay_io"], false); - assert_eq!(contract["step_205_deferrals"].as_array().unwrap().len(), 7); + assert_eq!( + contract["separate_publication_authority"] + .as_array() + .unwrap() + .len(), + 7 + ); for required in [ "DESIRED_STATE_DOMAIN", "TARGET_SET_DOMAIN", @@ -307,7 +313,7 @@ fn machine_contract_freezes_step_204_and_defers_step_205_effects() { ] { assert!( !SOURCE.contains(forbidden), - "unexpected Step205 effect {forbidden}" + "unexpected desired-state effect {forbidden}" ); } for kind in [ diff --git a/tests/services_hardening_presence_publication.rs b/tests/services_hardening_presence_publication.rs @@ -0,0 +1,650 @@ +#![forbid(unsafe_code)] +#![cfg(any(target_os = "linux", target_os = "macos"))] + +use std::{ + fs, + os::unix::fs::PermissionsExt, + path::{Path, PathBuf}, +}; + +use nostr::{Keys, SecretKey}; +use radroots_service_host::{ + EntropyError, EntropySource, SystemMonotonicClock, UnixTimeSeconds, WallClock, WallClockError, +}; +use radroots_service_sqlite::{MigrationAppliedAtUnixSeconds, MigrationBuildIdentity}; +use radroots_storage::event::SourceGeneration; +use radroots_transport::BoxFuture; +use rhi::{ + RadrootsHostEnvironment, RadrootsPathResolver, RadrootsPlatform, RhiConfigDocumentV1, + RhiConfigProfile, RhiDecryptedIdentity, RhiEncryptedIdentityProvisioningMaterial, + RhiExactPresenceSink, RhiIdentityEnvelopeBinding, RhiPreparedPresenceAttempt, + RhiPresenceAttemptOutcome, RhiPresenceDesiredAuthority, RhiPresenceLeaseOwner, + RhiPresenceOutboxState, RhiPresenceRetryDelayMilliseconds, RhiPresenceTargetState, + RhiPresenceUnixMilliseconds, RhiStateMetadata, RhiTimeEntropyAdapters, apply_rhi_configuration, + build_rhi_signed_presence_documents, initialize_rhi_state, + open_rhi_state_read_write_from_config, parse_rhi_cli_v1_from, parse_rhi_config_v1, + provision_rhi_encrypted_identity, resolve_rhi_runtime_context, resolve_rhi_wrapping_credential, + validate_rhi_signed_presence_documents, +}; +use sqlx::{ConnectOptions, Connection, Row, SqliteConnection, sqlite::SqliteConnectOptions}; + +const EXAMPLE: &str = include_str!("../contracts/services_hardening/config.v1.example.toml"); +const CONTRACT: &str = include_str!("../contracts/services_hardening/presence_publication.v1.json"); + +struct FixedEntropy(u8); + +impl EntropySource for FixedEntropy { + fn fill_bytes(&self, destination: &mut [u8]) -> Result<(), EntropyError> { + destination.fill(self.0); + Ok(()) + } +} + +struct UnavailableEntropy; + +impl EntropySource for UnavailableEntropy { + fn fill_bytes(&self, _destination: &mut [u8]) -> Result<(), EntropyError> { + Err(EntropyError::Unavailable) + } +} + +#[derive(Clone, Copy)] +struct FixedWall(u64); + +impl WallClock for FixedWall { + fn now_utc(&self) -> Result<UnixTimeSeconds, WallClockError> { + Ok(UnixTimeSeconds::new(self.0)) + } +} + +struct InspectingSink { + database: PathBuf, + expected: Box<[u8]>, +} + +impl RhiExactPresenceSink for InspectingSink { + fn submit_exact<'a>( + &'a self, + attempt: &'a RhiPreparedPresenceAttempt, + ) -> BoxFuture<'a, RhiPresenceAttemptOutcome> { + Box::pin(async move { + assert_eq!(attempt.exact_signed_event_bytes(), self.expected.as_ref()); + let mut connection = offline_connection(&self.database).await; + let row = sqlx::query( + "SELECT targets.state, outbox.exact_signed_event_bytes \ + FROM presence_targets AS targets \ + JOIN presence_outbox AS outbox ON outbox.outbox_id = targets.outbox_id \ + WHERE targets.last_attempt_id = ? LIMIT 2", + ) + .bind(attempt.attempt_id().as_bytes().as_slice()) + .fetch_one(&mut connection) + .await + .expect("durable pre-I/O target"); + assert_eq!(row.try_get::<String, _>("state").unwrap(), "submitted"); + assert_eq!( + row.try_get::<Vec<u8>, _>("exact_signed_event_bytes") + .unwrap(), + self.expected.as_ref() + ); + connection.close().await.expect("sink inspection close"); + RhiPresenceAttemptOutcome::Accepted + }) + } +} + +fn runtime(root: &Path) -> rhi::RhiRuntimeContext { + let invocation = parse_rhi_cli_v1_from([ + "rhi", + "--profile", + "repo-local", + "--instance", + "primary", + "--repo-local-root", + root.to_str().expect("UTF-8 root"), + "run", + ]) + .expect("invocation"); + resolve_rhi_runtime_context( + &RadrootsPathResolver::new(RadrootsPlatform::Linux, RadrootsHostEnvironment::default()), + &invocation, + ) + .expect("runtime") +} + +fn evidence(at: u64) -> (MigrationAppliedAtUnixSeconds, MigrationBuildIdentity) { + ( + MigrationAppliedAtUnixSeconds::new(at).expect("migration time"), + MigrationBuildIdentity::new( + env!("CARGO_PKG_VERSION"), + "1111111111111111111111111111111111111111", + "21b11e7a5120ea949f7ad0838c746873fc73aac2", + "rustc-test", + "test-target", + "service-host", + 1, + rhi::RHI_STATE_SCHEMA_VERSION, + 1, + 1, + 1, + ) + .expect("build identity"), + ) +} + +async fn offline_connection(database: &Path) -> SqliteConnection { + let options = SqliteConnectOptions::new() + .filename(database) + .create_if_missing(false) + .foreign_keys(false) + .disable_statement_logging(); + SqliteConnection::connect_with(&options) + .await + .expect("offline connection") +} + +fn secret() -> [u8; 32] { + [1; 32] +} + +fn configuration(runtime: &rhi::RhiRuntimeContext, public_key: &str) -> RhiConfigDocumentV1 { + configuration_from(runtime, public_key, EXAMPLE) +} + +fn configuration_from( + runtime: &rhi::RhiRuntimeContext, + public_key: &str, + source: &str, +) -> RhiConfigDocumentV1 { + let source = source + .replace( + "/var/lib/radroots/services/rhi/default/secrets/service.identity.ncrypt", + runtime + .identity_path() + .to_str() + .expect("UTF-8 identity path"), + ) + .replace(&"2".repeat(64), public_key); + parse_rhi_config_v1(source.as_bytes(), RhiConfigProfile::RepoLocal).expect("configuration") +} + +#[tokio::test] +async fn expired_stale_generation_is_unknown_and_superseded_before_new_bytes_commit() { + let root = tempfile::tempdir().expect("root"); + let runtime = runtime(root.path()); + fs::create_dir_all(runtime.context().paths().state()).expect("state directory"); + fs::set_permissions( + runtime.context().paths().state(), + fs::Permissions::from_mode(0o700), + ) + .expect("state mode"); + let public_key = Keys::new(SecretKey::from_slice(&secret()).expect("secret")) + .public_key() + .to_hex(); + let original = configuration(&runtime, &public_key); + let metadata = RhiStateMetadata::new( + &runtime, + &original, + SourceGeneration::new([0x5b; 32]).expect("source generation"), + 1_725_000_000_000, + ) + .expect("metadata"); + let (first_at, first_build) = evidence(1_725_000_000); + initialize_rhi_state(&runtime, &metadata, first_at, &first_build) + .await + .expect("initialize"); + let identity = provision(&runtime, &original, &metadata); + let original_authority = + RhiPresenceDesiredAuthority::from_config(&original).expect("original authority"); + let host = open_rhi_state_read_write_from_config(&runtime, &original, first_at, &first_build) + .await + .expect("writer"); + let original_desired = host + .repositories() + .desired_presence() + .commit(&original_authority) + .await + .expect("original desired"); + let original_documents = build_rhi_signed_presence_documents( + original_desired, + &original_authority, + &identity, + UnixTimeSeconds::new(1_725_000_100), + &FixedEntropy(0x92), + ) + .expect("original documents"); + host.repositories() + .presence_outbox() + .commit_signed_presence(&original_documents, millis(1_725_000_100_000)) + .await + .expect("original outbox"); + let lease = host + .repositories() + .presence_outbox() + .claim_next_presence( + RhiPresenceLeaseOwner::from_bytes([0x41; 16]).expect("owner"), + millis(1_725_000_101_000), + ) + .await + .expect("claim") + .expect("work"); + let submitted = host + .repositories() + .presence_outbox() + .prepare_next_presence_target(lease, millis(1_725_000_101_000)) + .await + .expect("submitted"); + drop(submitted); + host.close().await.expect("close original host"); + + let changed_source = EXAMPLE.replace("profile = true", "profile = false"); + let changed = configuration_from(&runtime, &public_key, &changed_source); + let (second_at, second_build) = evidence(1_725_000_001); + apply_rhi_configuration(&runtime, &original, &changed, second_at, &second_build) + .await + .expect("apply changed configuration"); + let host = open_rhi_state_read_write_from_config(&runtime, &changed, second_at, &second_build) + .await + .expect("changed writer"); + let changed_authority = + RhiPresenceDesiredAuthority::from_config(&changed).expect("changed authority"); + let changed_desired = host + .repositories() + .desired_presence() + .commit(&changed_authority) + .await + .expect("changed desired"); + assert_eq!(changed_desired.state().generation(), 2); + let changed_documents = build_rhi_signed_presence_documents( + changed_desired, + &changed_authority, + &identity, + UnixTimeSeconds::new(1_725_000_200), + &FixedEntropy(0x93), + ) + .expect("changed documents"); + let blocked = host + .repositories() + .presence_outbox() + .commit_signed_presence(&changed_documents, millis(1_725_000_200_000)) + .await + .expect_err("active stale lease blocks replacement"); + assert_eq!( + blocked.kind(), + rhi::RhiPresencePublicationErrorKind::NotReady + ); + + let recovery_adapters = RhiTimeEntropyAdapters::new( + FixedWall(1_725_000_120), + SystemMonotonicClock::new(), + UnavailableEntropy, + ); + assert!( + host.repositories() + .presence_outbox() + .recover_one_expired_presence(&recovery_adapters, millis(1_725_000_120_000)) + .await + .expect("stale recovery does not sample retry entropy") + ); + assert!( + host.repositories() + .presence_outbox() + .commit_signed_presence(&changed_documents, millis(1_725_000_200_000)) + .await + .expect("new generation commit") + .changed() + ); + host.close().await.expect("close changed host"); + + let mut disabled_table = changed_source + .parse::<toml::Table>() + .expect("changed configuration TOML"); + let disabled_presence = disabled_table + .get_mut("presence") + .and_then(toml::Value::as_table_mut) + .expect("presence table"); + disabled_presence.insert("enabled".to_owned(), toml::Value::Boolean(false)); + disabled_presence.insert("profile".to_owned(), toml::Value::Boolean(false)); + disabled_presence.insert( + "application_handler".to_owned(), + toml::Value::Boolean(false), + ); + disabled_presence.remove("target_relay_ids"); + let disabled_source = toml::to_string(&disabled_table).expect("disabled configuration TOML"); + let disabled = configuration_from(&runtime, &public_key, &disabled_source); + let (third_at, third_build) = evidence(1_725_000_002); + apply_rhi_configuration(&runtime, &changed, &disabled, third_at, &third_build) + .await + .expect("disable presence"); + let host = open_rhi_state_read_write_from_config(&runtime, &disabled, third_at, &third_build) + .await + .expect("disabled writer"); + let disabled_authority = + RhiPresenceDesiredAuthority::from_config(&disabled).expect("disabled authority"); + let disabled_desired = host + .repositories() + .desired_presence() + .commit(&disabled_authority) + .await + .expect("disabled desired"); + assert_eq!(disabled_desired.state().generation(), 3); + let disabled_documents = build_rhi_signed_presence_documents( + disabled_desired, + &disabled_authority, + &identity, + UnixTimeSeconds::new(1_725_000_300), + &UnavailableEntropy, + ) + .expect("disabled document inventory consumes no entropy"); + assert!(disabled_documents.documents().is_empty()); + assert!( + host.repositories() + .presence_outbox() + .commit_signed_presence(&disabled_documents, millis(1_725_000_300_000)) + .await + .expect("disabled generation commit") + .changed() + ); + assert!( + !host + .repositories() + .presence_outbox() + .commit_signed_presence(&disabled_documents, millis(1_725_000_300_000)) + .await + .expect("disabled exact replay") + .changed() + ); + host.close().await.expect("close disabled host"); + + let mut connection = offline_connection(runtime.artifacts().state_database()).await; + let stale: (String, String) = sqlx::query_as( + "SELECT outbox.state, attempts.outcome FROM presence_outbox AS outbox \ + JOIN presence_attempts AS attempts ON attempts.outbox_id = outbox.outbox_id \ + WHERE outbox.desired_generation = 1 AND outbox.document_kind = 'service_profile'", + ) + .fetch_one(&mut connection) + .await + .expect("stale durable result"); + assert_eq!(stale, ("superseded".to_owned(), "unknown".to_owned())); + let disabled_prior_state: String = sqlx::query_scalar( + "SELECT state FROM presence_outbox WHERE desired_generation = 2 LIMIT 2", + ) + .fetch_one(&mut connection) + .await + .expect("disabled prior generation"); + assert_eq!(disabled_prior_state, "superseded"); + let disabled_count: i64 = + sqlx::query_scalar("SELECT COUNT(*) FROM presence_outbox WHERE desired_generation = 3") + .fetch_one(&mut connection) + .await + .expect("disabled outbox count"); + assert_eq!(disabled_count, 0); + connection.close().await.expect("offline close"); +} + +fn provision( + runtime: &rhi::RhiRuntimeContext, + configuration: &RhiConfigDocumentV1, + metadata: &RhiStateMetadata, +) -> RhiDecryptedIdentity { + fs::create_dir_all(runtime.context().paths().secrets()).expect("secrets directory"); + fs::set_permissions( + runtime.context().paths().secrets(), + fs::Permissions::from_mode(0o700), + ) + .expect("secrets mode"); + let credential_path = runtime + .context() + .paths() + .secrets() + .join("service_wrapping_key"); + fs::write(&credential_path, [0x81; 32]).expect("credential"); + fs::set_permissions(&credential_path, fs::Permissions::from_mode(0o600)) + .expect("credential mode"); + let binding = RhiIdentityEnvelopeBinding::from_configuration(configuration, metadata) + .expect("identity binding"); + let credential = + resolve_rhi_wrapping_credential(runtime, &binding).expect("wrapping credential"); + provision_rhi_encrypted_identity( + &binding, + &credential, + RhiEncryptedIdentityProvisioningMaterial::new(secret(), [0x42; 32], [0x43; 24], [0x44; 24]) + .expect("provisioning material"), + ) + .expect("provision identity") +} + +fn millis(value: u64) -> RhiPresenceUnixMilliseconds { + RhiPresenceUnixMilliseconds::new(value).expect("milliseconds") +} + +#[tokio::test] +async fn exact_bytes_are_durable_before_io_and_unknown_recovery_retries_unchanged() { + let root = tempfile::tempdir().expect("root"); + let runtime = runtime(root.path()); + fs::create_dir_all(runtime.context().paths().state()).expect("state directory"); + fs::set_permissions( + runtime.context().paths().state(), + fs::Permissions::from_mode(0o700), + ) + .expect("state mode"); + let public_key = Keys::new(SecretKey::from_slice(&secret()).expect("secret")) + .public_key() + .to_hex(); + let configuration = configuration(&runtime, &public_key); + let metadata = RhiStateMetadata::new( + &runtime, + &configuration, + SourceGeneration::new([0x5a; 32]).expect("source generation"), + 1_725_000_000_000, + ) + .expect("metadata"); + let (applied_at, build) = evidence(1_725_000_000); + initialize_rhi_state(&runtime, &metadata, applied_at, &build) + .await + .expect("initialize"); + let identity = provision(&runtime, &configuration, &metadata); + let authority = RhiPresenceDesiredAuthority::from_config(&configuration).expect("authority"); + let host = open_rhi_state_read_write_from_config(&runtime, &configuration, applied_at, &build) + .await + .expect("writer"); + let desired = host + .repositories() + .desired_presence() + .commit(&authority) + .await + .expect("desired state"); + let authored = UnixTimeSeconds::new(1_725_000_100); + let invalid_time = build_rhi_signed_presence_documents( + desired, + &authority, + &identity, + UnixTimeSeconds::new(i64::MAX as u64 + 1), + &FixedEntropy(0x91), + ) + .expect_err("unrepresentable authored time"); + assert_eq!( + invalid_time.kind(), + rhi::RhiPresencePublicationErrorKind::InvalidInput + ); + let documents = build_rhi_signed_presence_documents( + desired, + &authority, + &identity, + authored, + &FixedEntropy(0x91), + ) + .expect("signed presence"); + validate_rhi_signed_presence_documents(&documents, &authority) + .expect("independently verified signed presence"); + let exact: Vec<Box<[u8]>> = documents + .documents() + .iter() + .map(|document| document.signed_event_bytes().into()) + .collect(); + let contract: serde_json::Value = serde_json::from_str(CONTRACT).expect("machine contract"); + assert_eq!(contract["schema"], "radroots.rhi.presence-publication"); + assert_eq!(contract["schema_version"], 1); + assert_eq!(contract["contract_version"], 1); + assert_eq!(contract["step"], 205); + assert_eq!( + contract["reference_vector"]["profile"]["exact_signed_event_json"], + String::from_utf8_lossy(&exact[0]).as_ref() + ); + assert_eq!( + contract["reference_vector"]["application_handler"]["exact_signed_event_json"], + String::from_utf8_lossy(&exact[1]).as_ref() + ); + assert_eq!( + contract["reference_vector"]["profile"]["event_id"], + lower_hex(documents.documents()[0].event_id()) + ); + assert_eq!( + contract["reference_vector"]["application_handler"]["event_id"], + lower_hex(documents.documents()[1].event_id()) + ); + assert_eq!( + contract["reference_vector"]["profile"]["exact_signed_event_sha256"], + lower_hex(documents.documents()[0].signed_event_sha256()) + ); + assert_eq!( + contract["reference_vector"]["application_handler"]["exact_signed_event_sha256"], + lower_hex(documents.documents()[1].signed_event_sha256()) + ); + let committed_at = millis(1_725_000_100_000); + let committed = host + .repositories() + .presence_outbox() + .commit_signed_presence(&documents, committed_at) + .await + .expect("durable exact presence"); + assert!(committed.changed()); + assert_eq!(committed.document_count(), 2); + + assert!( + !host + .repositories() + .presence_outbox() + .commit_signed_presence(&documents, committed_at) + .await + .expect("exact replay") + .changed() + ); + + let adapters = RhiTimeEntropyAdapters::new( + FixedWall(1_725_000_100), + SystemMonotonicClock::new(), + FixedEntropy(0xff), + ); + let sink = InspectingSink { + database: runtime.artifacts().state_database().to_path_buf(), + expected: exact[0].clone(), + }; + let first = host + .repositories() + .presence_outbox() + .execute_next_presence( + RhiPresenceLeaseOwner::from_bytes([0x11; 16]).expect("owner"), + &adapters, + &sink, + ) + .await + .expect("execute first") + .expect("first work"); + assert_eq!(first.outbox_state(), RhiPresenceOutboxState::Complete); + assert_eq!(first.target_state(), RhiPresenceTargetState::Accepted); + + let repository = host.repositories().presence_outbox(); + let started = millis(1_725_000_101_000); + let second_lease = repository + .claim_next_presence( + RhiPresenceLeaseOwner::from_bytes([0x22; 16]).expect("owner"), + started, + ) + .await + .expect("claim second") + .expect("second work"); + let abandoned = repository + .prepare_next_presence_target(second_lease, started) + .await + .expect("prepare second"); + assert_eq!(abandoned.exact_signed_event_bytes(), exact[1].as_ref()); + drop(abandoned); + host.close().await.expect("close after cancellation"); + + let host = open_rhi_state_read_write_from_config(&runtime, &configuration, applied_at, &build) + .await + .expect("reopen writer"); + let recovery_now = millis(1_725_000_120_000); + assert!( + host.repositories() + .presence_outbox() + .recover_one_expired_presence(&adapters, recovery_now) + .await + .expect("unknown recovery") + ); + let retry_at = millis(1_725_000_150_000); + let lease = host + .repositories() + .presence_outbox() + .claim_next_presence( + RhiPresenceLeaseOwner::from_bytes([0x33; 16]).expect("owner"), + retry_at, + ) + .await + .expect("retry claim") + .expect("retry work"); + let retried = host + .repositories() + .presence_outbox() + .prepare_next_presence_target(lease, retry_at) + .await + .expect("retry prepare"); + assert_eq!(retried.exact_signed_event_bytes(), exact[1].as_ref()); + let accepted = host + .repositories() + .presence_outbox() + .record_presence_outcome( + &retried, + millis(1_725_000_150_001), + RhiPresenceAttemptOutcome::Accepted, + RhiPresenceRetryDelayMilliseconds::new(0).expect("zero delay"), + ) + .await + .expect("retry accepted"); + assert_eq!(accepted.attempt_number(), 2); + assert_eq!(accepted.outbox_state(), RhiPresenceOutboxState::Complete); + host.close().await.expect("final close"); + + let mut connection = offline_connection(runtime.artifacts().state_database()).await; + let counts = sqlx::query( + "SELECT (SELECT COUNT(*) FROM presence_outbox) AS outboxes, \ + (SELECT COUNT(*) FROM presence_targets) AS targets, \ + (SELECT COUNT(*) FROM presence_attempts) AS attempts", + ) + .fetch_one(&mut connection) + .await + .expect("workflow counts"); + assert_eq!(counts.try_get::<i64, _>("outboxes").unwrap(), 2); + assert_eq!(counts.try_get::<i64, _>("targets").unwrap(), 4); + assert_eq!(counts.try_get::<i64, _>("attempts").unwrap(), 3); + let bytes: Vec<Vec<u8>> = + sqlx::query("SELECT exact_signed_event_bytes FROM presence_outbox ORDER BY document_kind") + .fetch_all(&mut connection) + .await + .expect("exact bytes") + .into_iter() + .map(|row| row.try_get("exact_signed_event_bytes").unwrap()) + .collect(); + assert_eq!(bytes, vec![exact[1].to_vec(), exact[0].to_vec()]); + connection.close().await.expect("offline close"); +} + +fn lower_hex(bytes: &[u8]) -> String { + const DIGITS: &[u8; 16] = b"0123456789abcdef"; + let mut rendered = String::with_capacity(bytes.len() * 2); + for byte in bytes { + rendered.push(char::from(DIGITS[usize::from(byte >> 4)])); + rendered.push(char::from(DIGITS[usize::from(byte & 0x0f)])); + } + rendered +} diff --git a/tests/services_hardening_state_catalog.rs b/tests/services_hardening_state_catalog.rs @@ -21,8 +21,10 @@ use rhi::{ RHI_STATE_SCHEMA_VERSION_7_SHA256, RHI_STATE_SCHEMA_VERSION_8_MIGRATION_SHA256, RHI_STATE_SCHEMA_VERSION_8_OBJECT_COUNT, RHI_STATE_SCHEMA_VERSION_8_SHA256, RHI_STATE_SCHEMA_VERSION_9_MIGRATION_SHA256, RHI_STATE_SCHEMA_VERSION_9_OBJECT_COUNT, - RHI_STATE_SCHEMA_VERSION_9_SHA256, RhiStateCatalogErrorKind, rhi_migration_catalog, - rhi_schema_catalog, validate_rhi_state_catalogs, + RHI_STATE_SCHEMA_VERSION_9_SHA256, RHI_STATE_SCHEMA_VERSION_10_MIGRATION_SHA256, + RHI_STATE_SCHEMA_VERSION_10_OBJECT_COUNT, RHI_STATE_SCHEMA_VERSION_10_SHA256, + RhiStateCatalogErrorKind, rhi_migration_catalog, rhi_schema_catalog, + validate_rhi_state_catalogs, }; const CATALOG_SOURCE: &str = include_str!("../src/state_catalog.rs"); @@ -30,14 +32,14 @@ const LIB_SOURCE: &str = include_str!("../src/lib.rs"); const MANIFEST: &str = include_str!("../Cargo.toml"); #[test] -fn schema_v1_through_v9_catalogs_have_exact_literal_identities() { +fn schema_v1_through_v10_catalogs_have_exact_literal_identities() { let migrations = rhi_migration_catalog().expect("RHI migration catalog"); let schema = rhi_schema_catalog().expect("RHI schema catalog"); assert_eq!(RHI_STATE_BASE_SCHEMA_VERSION, 1); - assert_eq!(RHI_STATE_SCHEMA_VERSION, 9); - assert_eq!(migrations.descriptors().len(), 8); - assert_eq!(migrations.current_version(), 9); + assert_eq!(RHI_STATE_SCHEMA_VERSION, 10); + assert_eq!(migrations.descriptors().len(), 9); + assert_eq!(migrations.current_version(), 10); assert_eq!(migrations.descriptors()[0].target_version(), 2); assert_eq!( migrations.descriptors()[0].name().as_str(), @@ -110,12 +112,21 @@ fn schema_v1_through_v9_catalogs_have_exact_literal_identities() { migrations.descriptors()[7].checksum().as_bytes(), &RHI_STATE_SCHEMA_VERSION_9_MIGRATION_SHA256 ); + assert_eq!(migrations.descriptors()[8].target_version(), 10); + assert_eq!( + migrations.descriptors()[8].name().as_str(), + "create_presence_publication_workflow" + ); + assert_eq!( + migrations.descriptors()[8].checksum().as_bytes(), + &RHI_STATE_SCHEMA_VERSION_10_MIGRATION_SHA256 + ); assert_eq!( migrations.digest().as_bytes(), &RHI_MIGRATION_CATALOG_SHA256 ); - assert_eq!(schema.versions().len(), 9); + assert_eq!(schema.versions().len(), 10); let version = schema.versions()[0]; assert_eq!(version.version(), 1); assert_eq!( @@ -215,13 +226,24 @@ fn schema_v1_through_v9_catalogs_have_exact_literal_identities() { version.digest().as_bytes(), &RHI_STATE_SCHEMA_VERSION_9_SHA256 ); + let version = schema.versions()[9]; + assert_eq!(version.version(), 10); + assert_eq!( + version.object_count(), + RHI_STATE_SCHEMA_VERSION_10_OBJECT_COUNT + ); + assert_eq!(version.object_count(), 80); + assert_eq!( + version.digest().as_bytes(), + &RHI_STATE_SCHEMA_VERSION_10_SHA256 + ); assert_eq!(schema.digest().as_bytes(), &RHI_STATE_SCHEMA_CATALOG_SHA256); assert_eq!(schema.migration_catalog_digest(), migrations.digest()); validate_rhi_state_catalogs(&migrations, &schema).expect("exact catalogs"); assert_eq!( lower_hex(&RHI_MIGRATION_CATALOG_SHA256), - "7f8c03b4818408b586747412c265ac72cd4dd88a290d8bbee1d5f68adbf7afd0" + "25e5ba773ef3db0133a8077a083e88b40fc6dadf9fb6f0cfde0e4ba4d3e081d9" ); assert_eq!( lower_hex(&RHI_STATE_SCHEMA_VERSION_1_SHA256), @@ -292,8 +314,16 @@ fn schema_v1_through_v9_catalogs_have_exact_literal_identities() { "5551e8790544a7c78c8376c5ccf83dd2a8486d2dc08b6a7808bb2007c94335ec" ); assert_eq!( + lower_hex(&RHI_STATE_SCHEMA_VERSION_10_MIGRATION_SHA256), + "541ad13b2cb08d59857205e6ffaec19b74de0853240fdf11fa1340e406c7114e" + ); + assert_eq!( + lower_hex(&RHI_STATE_SCHEMA_VERSION_10_SHA256), + "d2aed51d0a6a2c01eda1844608472b2dcd502abaa8a4ca30a4823535b8bd0e45" + ); + assert_eq!( lower_hex(&RHI_STATE_SCHEMA_CATALOG_SHA256), - "5bea3e3ec6be3f1249ad19ed7b2d2e872c34025579a599f254de8e80cba9b3c7" + "4f4d5f5546a7c94ef3cdabe23eadd0c97a64980984ee38cacb828f1125913488" ); } @@ -361,6 +391,10 @@ fn independent_validator_rejects_migration_or_schema_drift() { SchemaVersionCatalog::computed_digest(9, [version_two_object()]).expect("v9 digest"); let version_nine = SchemaVersionCatalog::new(9, [version_two_object()], snapshot_digest) .expect("version nine"); + let snapshot_digest = + SchemaVersionCatalog::computed_digest(10, [version_two_object()]).expect("v10 digest"); + let version_ten = SchemaVersionCatalog::new(10, [version_two_object()], snapshot_digest) + .expect("version ten"); let schema = SchemaCatalog::new( &exact_migrations, [ @@ -373,6 +407,7 @@ fn independent_validator_rejects_migration_or_schema_drift() { version_seven, version_eight, version_nine, + version_ten, ], ) .expect("drift schema catalog"); @@ -439,6 +474,10 @@ fn catalog_errors_are_stable_source_free_and_redacted() { SchemaVersionCatalog::computed_digest(9, [secret_object()]).expect("v9 digest"); let version_nine = SchemaVersionCatalog::new(9, [secret_object()], version_nine_digest).expect("version nine"); + let version_ten_digest = + SchemaVersionCatalog::computed_digest(10, [secret_object()]).expect("v10 digest"); + let version_ten = + SchemaVersionCatalog::new(10, [secret_object()], version_ten_digest).expect("version ten"); let schema = SchemaCatalog::new( &migrations, [ @@ -451,6 +490,7 @@ fn catalog_errors_are_stable_source_free_and_redacted() { version_seven, version_eight, version_nine, + version_ten, ], ) .expect("schema catalog"); diff --git a/tests/services_hardening_state_host.rs b/tests/services_hardening_state_host.rs @@ -106,6 +106,9 @@ async fn downgrade_fixture_to_schema_v7(runtime: &rhi::RhiRuntimeContext) { .await .expect("migration delete guard SQL"); for statement in [ + "DROP TABLE presence_attempts", + "DROP TABLE presence_targets", + "DROP TABLE presence_outbox", "DROP TRIGGER presence_desired_state_guard_insert", "DROP TRIGGER presence_desired_state_guard_update", "DROP TRIGGER presence_desired_state_no_delete", @@ -116,7 +119,7 @@ async fn downgrade_fixture_to_schema_v7(runtime: &rhi::RhiRuntimeContext) { "DROP TRIGGER schema_migrations_no_update", "DROP TRIGGER schema_migrations_no_delete", "UPDATE radroots_service_metadata SET state_schema_version = 7 WHERE singleton = 1", - "DELETE FROM schema_migrations WHERE version IN (8, 9)", + "DELETE FROM schema_migrations WHERE version IN (8, 9, 10)", ] { sqlx::query(statement) .execute(&mut connection) @@ -483,7 +486,7 @@ async fn schema_v8_scans_historical_nullable_job_state_and_installs_permanent_gu .fetch_one(&mut connection) .await .expect("migrated schema state"); - assert_eq!(migrated, (9, 1, 2, 0)); + assert_eq!(migrated, (10, 1, 2, 0)); let invalid_insert = sqlx::query( r#"INSERT INTO reconciliation_jobs ( diff --git a/tests/services_hardening_state_repository_topology.rs b/tests/services_hardening_state_repository_topology.rs @@ -28,7 +28,7 @@ fn machine_contract_and_typed_descriptor_inventory_are_exact() { contract["deferred_behavior"], json!([ "later_schema_migrations", - "non_job_repository_crud", + "remaining_repository_crud", "backup_restore_and_recovery", "network_io", "task_supervision" @@ -205,6 +205,24 @@ fn repository_kinds_are_closed_ordered_and_cross_bound() { "presence_desired_state", Write::CompareAndSwap, ), + ( + Kind::PresenceOutbox, + "presence_outbox", + "presence_outbox", + Write::CompareAndSwap, + ), + ( + Kind::PresenceTarget, + "presence_target", + "presence_targets", + Write::CompareAndSwap, + ), + ( + Kind::PresenceAttempt, + "presence_attempt", + "presence_attempts", + Write::AppendOnly, + ), ]; for (descriptor, (kind, code, table, write_class)) in rhi_state_repository_descriptors().iter().zip(expected) @@ -240,6 +258,9 @@ fn capabilities_are_private_sealed_and_defer_unowned_behavior() { "pub const fn publication_targets(&self) -> RhiPublicationTargetRepository<'host>", "pub const fn publication_attempts(&self) -> RhiPublicationAttemptRepository<'host>", "pub const fn desired_presence(&self) -> RhiDesiredPresenceRepository<'host>", + "pub const fn presence_outbox(&self) -> RhiPresenceOutboxRepository<'host>", + "pub const fn presence_targets(&self) -> RhiPresenceTargetRepository<'host>", + "pub const fn presence_attempts(&self) -> RhiPresenceAttemptRepository<'host>", ] { assert!(REPOSITORY_SOURCE.contains(required), "missing {required}"); } diff --git a/tests/services_hardening_state_resilience.rs b/tests/services_hardening_state_resilience.rs @@ -347,8 +347,8 @@ async fn exact_open_rejects_unexpected_migration_history_without_repair() { service_version, service_commit, lib_revision, rust_version, target, feature_profile, config_contract_version, state_contract_version, admin_contract_version, status_contract_version, provider_contract_version - ) VALUES (10, 'unexpected_schema', ?, 1725000000, '0.1.0', ?, ?, - 'rustc-test', 'test-target', 'service-host', 1, 7, 1, 1, 1)", + ) VALUES (11, 'unexpected_schema', ?, 1725000000, '0.1.0', ?, ?, + 'rustc-test', 'test-target', 'service-host', 1, 10, 1, 1, 1)", ) .bind([0x44_u8; 32].as_slice()) .bind("1111111111111111111111111111111111111111")