lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

commit 3fe6d63fef9a64884c9914211d90b8127cdc6cda
parent 7776c100b01f609adf4ab96e1cb77e47fa5f12bc
Author: triesap <tyson@radroots.org>
Date:   Sat,  1 Aug 2026 20:27:20 +0000

storage: define outbox and delivery evidence contracts

- define idempotent durable delivery plans and status records
- enforce exclusive expiring leases and optimistic attempt order
- preserve normalized per-target evidence and satisfaction results
- prove replay partial retry terminal and expiry behavior

Diffstat:
Mcrates/storage/src/error.rs | 32++++++++++++++++++++++++++++++++
Mcrates/storage/src/lib.rs | 1+
Mcrates/storage/src/outbox.rs | 747+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acrates/storage/tests/outbox.rs | 397+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
4 files changed, 1177 insertions(+), 0 deletions(-)

diff --git a/crates/storage/src/error.rs b/crates/storage/src/error.rs @@ -35,6 +35,22 @@ pub enum Error { InvalidJournalTransition, JournalOperationCommitted, CorruptJournalRecord, + InvalidOutboxItemId, + InvalidOutboxRevision, + InvalidOutboxTimestamp, + InvalidOutboxLease, + InvalidOutboxLeaseOwner, + InvalidOutboxClaimLimit, + InvalidDeliveryAttempt, + InvalidDeliveryEvidence, + OutboxItemNotFound, + OutboxItemNotReady, + OutboxItemTerminal, + OutboxPlanConflict, + OutboxLeaseConflict, + OutboxLeaseExpired, + OutboxRevisionConflict, + CorruptOutboxRecord, } impl fmt::Display for Error { @@ -71,6 +87,22 @@ impl fmt::Display for Error { Self::InvalidJournalTransition => "storage journal transition is invalid", Self::JournalOperationCommitted => "storage journal operation is already committed", Self::CorruptJournalRecord => "storage journal record is corrupt", + Self::InvalidOutboxItemId => "storage outbox item id is invalid", + Self::InvalidOutboxRevision => "storage outbox revision is invalid", + Self::InvalidOutboxTimestamp => "storage outbox timestamp is invalid", + Self::InvalidOutboxLease => "storage outbox lease is invalid", + Self::InvalidOutboxLeaseOwner => "storage outbox lease owner is invalid", + Self::InvalidOutboxClaimLimit => "storage outbox claim limit is invalid", + Self::InvalidDeliveryAttempt => "storage delivery attempt is invalid", + Self::InvalidDeliveryEvidence => "storage delivery evidence is invalid", + Self::OutboxItemNotFound => "storage outbox item was not found", + Self::OutboxItemNotReady => "storage outbox item is not ready", + Self::OutboxItemTerminal => "storage outbox item is terminal", + Self::OutboxPlanConflict => "storage outbox plan conflicts with durable state", + Self::OutboxLeaseConflict => "storage outbox lease conflicts with durable state", + Self::OutboxLeaseExpired => "storage outbox lease expired", + Self::OutboxRevisionConflict => "storage outbox revision conflicts with durable state", + Self::CorruptOutboxRecord => "storage outbox record is corrupt", }) } } diff --git a/crates/storage/src/lib.rs b/crates/storage/src/lib.rs @@ -17,3 +17,4 @@ pub mod status; pub use error::Error; pub use event::EventStore; pub use journal::Journal; +pub use outbox::Outbox; diff --git a/crates/storage/src/outbox.rs b/crates/storage/src/outbox.rs @@ -1 +1,748 @@ //! Durable outbox and delivery-evidence contracts. +//! +//! This module stores delivery intent and normalized evidence. It deliberately +//! owns no transport adapter and performs no transport I/O. + +use core::fmt; +use radroots_transport::{ + BoxFuture, DeliveryReceipt, DeliveryRequest, + outcome::{DeliveryOutcome, Retryability}, + target::TargetFingerprint, +}; + +use crate::{Error, journal::OperationInstanceId}; + +/// Maximum records claimed by one bounded outbox query. +pub const OUTBOX_CLAIM_LIMIT_MAX: u16 = 256; +/// Maximum UTF-8 bytes in a lease owner identity. +pub const LEASE_OWNER_MAX_BYTES: usize = 128; + +/// Stable host-generated identity for one durable delivery plan. +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] +pub struct OutboxItemId([u8; 16]); + +impl OutboxItemId { + pub const fn new(bytes: [u8; 16]) -> Result<Self, Error> { + if bytes_are_zero(&bytes) { + return Err(Error::InvalidOutboxItemId); + } + Ok(Self(bytes)) + } + + pub const fn as_bytes(&self) -> &[u8; 16] { + &self.0 + } +} + +/// Digest of the canonical delivery request, computed by its owning workflow. +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] +pub struct DeliveryPlanDigest([u8; 32]); + +impl DeliveryPlanDigest { + pub const fn new(bytes: [u8; 32]) -> Self { + Self(bytes) + } + + pub const fn as_bytes(&self) -> &[u8; 32] { + &self.0 + } +} + +/// Non-zero optimistic outbox record revision. +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)] +pub struct OutboxRevision(u64); + +impl OutboxRevision { + pub const INITIAL: Self = Self(1); + + pub const fn new(value: u64) -> Result<Self, Error> { + if value == 0 { + return Err(Error::InvalidOutboxRevision); + } + Ok(Self(value)) + } + + pub const fn get(self) -> u64 { + self.0 + } + + fn next(self) -> Result<Self, Error> { + self.0 + .checked_add(1) + .map(Self) + .ok_or(Error::CorruptOutboxRecord) + } +} + +/// Non-zero delivery attempt sequence. +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)] +pub struct DeliveryAttempt(u32); + +impl DeliveryAttempt { + pub const FIRST: Self = Self(1); + + pub const fn new(value: u32) -> Result<Self, Error> { + if value == 0 { + return Err(Error::InvalidDeliveryAttempt); + } + Ok(Self(value)) + } + + pub const fn get(self) -> u32 { + self.0 + } + + fn next(self) -> Result<Self, Error> { + self.0 + .checked_add(1) + .map(Self) + .ok_or(Error::CorruptOutboxRecord) + } +} + +/// Opaque, caller-generated lease identity. +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] +pub struct LeaseId([u8; 16]); + +impl LeaseId { + pub const fn new(bytes: [u8; 16]) -> Result<Self, Error> { + if bytes_are_zero(&bytes) { + return Err(Error::InvalidOutboxLease); + } + Ok(Self(bytes)) + } + + pub const fn as_bytes(&self) -> &[u8; 16] { + &self.0 + } +} + +/// Validated worker identity recorded with a lease. +#[derive(Clone, Eq, Hash, Ord, PartialEq, PartialOrd)] +pub struct LeaseOwner(String); + +impl LeaseOwner { + pub fn parse(value: impl Into<String>) -> Result<Self, Error> { + let value = value.into(); + if value.is_empty() + || value.len() > LEASE_OWNER_MAX_BYTES + || value != value.trim() + || value.chars().any(char::is_control) + { + return Err(Error::InvalidOutboxLeaseOwner); + } + Ok(Self(value)) + } + + pub fn as_str(&self) -> &str { + self.0.as_str() + } +} + +impl fmt::Debug for LeaseOwner { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("LeaseOwner") + .field("value", &"[REDACTED]") + .field("bytes", &self.0.len()) + .finish() + } +} + +#[cfg(feature = "serde")] +impl serde::Serialize for LeaseOwner { + fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error> + where + S: serde::Serializer, + { + serializer.serialize_str(self.as_str()) + } +} + +#[cfg(feature = "serde")] +impl<'de> serde::Deserialize<'de> for LeaseOwner { + fn deserialize<D>(deserializer: D) -> Result<Self, D::Error> + where + D: serde::Deserializer<'de>, + { + let value = <String as serde::Deserialize>::deserialize(deserializer)?; + Self::parse(value).map_err(serde::de::Error::custom) + } +} + +/// Exclusive, expiring authority to mutate one outbox item. +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct OutboxLease { + id: LeaseId, + owner: LeaseOwner, + acquired_at_unix_ms: u64, + expires_at_unix_ms: u64, +} + +impl OutboxLease { + pub fn new( + id: LeaseId, + owner: LeaseOwner, + acquired_at_unix_ms: u64, + expires_at_unix_ms: u64, + ) -> Result<Self, Error> { + if acquired_at_unix_ms == 0 || expires_at_unix_ms <= acquired_at_unix_ms { + return Err(Error::InvalidOutboxLease); + } + Ok(Self { + id, + owner, + acquired_at_unix_ms, + expires_at_unix_ms, + }) + } + + pub const fn id(&self) -> LeaseId { + self.id + } + + pub const fn owner(&self) -> &LeaseOwner { + &self.owner + } + + pub const fn acquired_at_unix_ms(&self) -> u64 { + self.acquired_at_unix_ms + } + + pub const fn expires_at_unix_ms(&self) -> u64 { + self.expires_at_unix_ms + } + + pub const fn is_active_at(&self, unix_ms: u64) -> bool { + unix_ms >= self.acquired_at_unix_ms && unix_ms < self.expires_at_unix_ms + } +} + +/// Durable lifecycle of one delivery plan. +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum OutboxStage { + Pending, + Leased, + Retryable, + Satisfied, + Exhausted, +} + +impl OutboxStage { + pub const fn is_terminal(self) -> bool { + matches!(self, Self::Satisfied | Self::Exhausted) + } +} + +/// Latest durable evidence for one requested target. +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct TargetDeliveryEvidence { + target: TargetFingerprint, + attempt: DeliveryAttempt, + attempted: bool, + outcome: DeliveryOutcome, + recorded_at_unix_ms: u64, +} + +impl TargetDeliveryEvidence { + pub const fn target(&self) -> &TargetFingerprint { + &self.target + } + + pub const fn attempt(&self) -> DeliveryAttempt { + self.attempt + } + + pub const fn was_attempted(&self) -> bool { + self.attempted + } + + pub const fn outcome(&self) -> &DeliveryOutcome { + &self.outcome + } + + pub const fn recorded_at_unix_ms(&self) -> u64 { + self.recorded_at_unix_ms + } +} + +/// Result of evaluating the latest complete transport receipt. +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum SatisfactionResult { + Pending, + Satisfied, + Exhausted, +} + +/// Durable outbox item and all current delivery evidence. +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct OutboxRecord { + item_id: OutboxItemId, + operation_instance_id: OperationInstanceId, + plan_digest: DeliveryPlanDigest, + request: DeliveryRequest, + revision: OutboxRevision, + stage: OutboxStage, + lease: Option<OutboxLease>, + last_attempt: Option<DeliveryAttempt>, + evidence: Vec<TargetDeliveryEvidence>, + satisfaction: SatisfactionResult, + retry_not_before_unix_ms: Option<u64>, + created_at_unix_ms: u64, + updated_at_unix_ms: u64, +} + +impl OutboxRecord { + fn from_enqueue(value: EnqueueOutboxItem) -> Self { + Self { + item_id: value.item_id, + operation_instance_id: value.operation_instance_id, + plan_digest: value.plan_digest, + request: value.request, + revision: OutboxRevision::INITIAL, + stage: OutboxStage::Pending, + lease: None, + last_attempt: None, + evidence: Vec::new(), + satisfaction: SatisfactionResult::Pending, + retry_not_before_unix_ms: None, + created_at_unix_ms: value.created_at_unix_ms, + updated_at_unix_ms: value.created_at_unix_ms, + } + } + + pub const fn item_id(&self) -> OutboxItemId { + self.item_id + } + pub const fn operation_instance_id(&self) -> OperationInstanceId { + self.operation_instance_id + } + pub const fn plan_digest(&self) -> DeliveryPlanDigest { + self.plan_digest + } + pub const fn request(&self) -> &DeliveryRequest { + &self.request + } + pub const fn revision(&self) -> OutboxRevision { + self.revision + } + pub const fn stage(&self) -> OutboxStage { + self.stage + } + pub const fn lease(&self) -> Option<&OutboxLease> { + self.lease.as_ref() + } + pub const fn last_attempt(&self) -> Option<DeliveryAttempt> { + self.last_attempt + } + pub fn evidence(&self) -> &[TargetDeliveryEvidence] { + self.evidence.as_slice() + } + /// Returns the latest evidence for one target without discarding history. + pub fn latest_target_evidence( + &self, + target: &TargetFingerprint, + ) -> Option<&TargetDeliveryEvidence> { + self.evidence + .iter() + .rev() + .find(|evidence| evidence.target() == target) + } + pub const fn satisfaction(&self) -> SatisfactionResult { + self.satisfaction + } + pub const fn retry_not_before_unix_ms(&self) -> Option<u64> { + self.retry_not_before_unix_ms + } + pub const fn created_at_unix_ms(&self) -> u64 { + self.created_at_unix_ms + } + pub const fn updated_at_unix_ms(&self) -> u64 { + self.updated_at_unix_ms + } + + /// Claims this item if it is ready and has no active lease. + pub fn claim(&mut self, lease: OutboxLease) -> Result<(), Error> { + if self.stage.is_terminal() { + return Err(Error::OutboxItemTerminal); + } + if matches!(self.retry_not_before_unix_ms, Some(not_before) if lease.acquired_at_unix_ms() < not_before) + { + return Err(Error::OutboxItemNotReady); + } + if self + .lease + .as_ref() + .is_some_and(|current| current.is_active_at(lease.acquired_at_unix_ms())) + { + return Err(Error::OutboxLeaseConflict); + } + self.revision = self.revision.next()?; + self.updated_at_unix_ms = lease.acquired_at_unix_ms(); + self.stage = OutboxStage::Leased; + self.lease = Some(lease); + Ok(()) + } + + /// Applies one request-bound transport receipt under the active lease. + pub fn record_attempt(&mut self, value: DeliveryAttemptEvidence) -> Result<(), Error> { + self.validate_lease(value.lease_id, value.recorded_at_unix_ms)?; + if value.item_id != self.item_id || value.expected_revision != self.revision { + return Err(Error::OutboxRevisionConflict); + } + let expected_attempt = self + .last_attempt + .map_or(Ok(DeliveryAttempt::FIRST), DeliveryAttempt::next)?; + if value.attempt != expected_attempt { + return Err(Error::InvalidDeliveryAttempt); + } + value + .receipt + .validate_for_request(&self.request) + .map_err(|_| Error::InvalidDeliveryEvidence)?; + let satisfied = value + .receipt + .is_satisfied(&self.request) + .map_err(|_| Error::InvalidDeliveryEvidence)?; + let retryable = value + .receipt + .target_receipts() + .iter() + .any(|receipt| matches!(receipt.outcome().retryability(), Retryability::Retryable)); + self.evidence + .extend( + value + .receipt + .target_receipts() + .iter() + .map(|receipt| TargetDeliveryEvidence { + target: receipt.target().fingerprint().clone(), + attempt: value.attempt, + attempted: receipt.was_attempted(), + outcome: receipt.outcome().clone(), + recorded_at_unix_ms: value.recorded_at_unix_ms, + }), + ); + self.last_attempt = Some(value.attempt); + self.satisfaction = if satisfied { + SatisfactionResult::Satisfied + } else if retryable { + SatisfactionResult::Pending + } else { + SatisfactionResult::Exhausted + }; + self.stage = match self.satisfaction { + SatisfactionResult::Pending => OutboxStage::Retryable, + SatisfactionResult::Satisfied => OutboxStage::Satisfied, + SatisfactionResult::Exhausted => OutboxStage::Exhausted, + }; + self.lease = None; + self.retry_not_before_unix_ms = None; + self.updated_at_unix_ms = value.recorded_at_unix_ms; + self.revision = self.revision.next()?; + Ok(()) + } + + /// Releases an active lease and optionally defers the next claim. + pub fn release( + &mut self, + lease_id: LeaseId, + expected_revision: OutboxRevision, + released_at_unix_ms: u64, + retry_not_before_unix_ms: Option<u64>, + ) -> Result<(), Error> { + self.validate_lease(lease_id, released_at_unix_ms)?; + if expected_revision != self.revision { + return Err(Error::OutboxRevisionConflict); + } + if released_at_unix_ms == 0 + || matches!(retry_not_before_unix_ms, Some(value) if value <= released_at_unix_ms) + { + return Err(Error::InvalidOutboxTimestamp); + } + self.lease = None; + self.stage = if self.last_attempt.is_some() { + OutboxStage::Retryable + } else { + OutboxStage::Pending + }; + self.retry_not_before_unix_ms = retry_not_before_unix_ms; + self.updated_at_unix_ms = released_at_unix_ms; + self.revision = self.revision.next()?; + Ok(()) + } + + fn validate_lease(&self, lease_id: LeaseId, at_unix_ms: u64) -> Result<(), Error> { + let lease = self.lease.as_ref().ok_or(Error::OutboxLeaseConflict)?; + if lease.id() != lease_id { + return Err(Error::OutboxLeaseConflict); + } + if !lease.is_active_at(at_unix_ms) { + return Err(Error::OutboxLeaseExpired); + } + Ok(()) + } +} + +/// Validated durable plan enqueue request. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct EnqueueOutboxItem { + item_id: OutboxItemId, + operation_instance_id: OperationInstanceId, + plan_digest: DeliveryPlanDigest, + request: DeliveryRequest, + created_at_unix_ms: u64, +} + +impl EnqueueOutboxItem { + pub fn new( + item_id: OutboxItemId, + operation_instance_id: OperationInstanceId, + plan_digest: DeliveryPlanDigest, + request: DeliveryRequest, + created_at_unix_ms: u64, + ) -> Result<Self, Error> { + if created_at_unix_ms == 0 { + return Err(Error::InvalidOutboxTimestamp); + } + Ok(Self { + item_id, + operation_instance_id, + plan_digest, + request, + created_at_unix_ms, + }) + } + + pub const fn item_id(&self) -> OutboxItemId { + self.item_id + } + pub const fn operation_instance_id(&self) -> OperationInstanceId { + self.operation_instance_id + } + pub const fn plan_digest(&self) -> DeliveryPlanDigest { + self.plan_digest + } + pub const fn request(&self) -> &DeliveryRequest { + &self.request + } + pub const fn created_at_unix_ms(&self) -> u64 { + self.created_at_unix_ms + } + pub fn into_record(self) -> OutboxRecord { + OutboxRecord::from_enqueue(self) + } +} + +/// Idempotent enqueue result. +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum EnqueueDisposition { + Created, + Replay, +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct EnqueueReceipt { + disposition: EnqueueDisposition, + record: OutboxRecord, +} + +impl EnqueueReceipt { + pub const fn new(disposition: EnqueueDisposition, record: OutboxRecord) -> Self { + Self { + disposition, + record, + } + } + pub const fn disposition(&self) -> EnqueueDisposition { + self.disposition + } + pub const fn record(&self) -> &OutboxRecord { + &self.record + } +} + +/// Bounded lease acquisition request. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct ClaimOutboxItems { + owner: LeaseOwner, + lease_id_seed: LeaseId, + now_unix_ms: u64, + lease_expires_at_unix_ms: u64, + limit: u16, +} + +impl ClaimOutboxItems { + pub fn new( + owner: LeaseOwner, + lease_id_seed: LeaseId, + now_unix_ms: u64, + lease_expires_at_unix_ms: u64, + limit: u16, + ) -> Result<Self, Error> { + if now_unix_ms == 0 || lease_expires_at_unix_ms <= now_unix_ms { + return Err(Error::InvalidOutboxLease); + } + if limit == 0 || limit > OUTBOX_CLAIM_LIMIT_MAX { + return Err(Error::InvalidOutboxClaimLimit); + } + Ok(Self { + owner, + lease_id_seed, + now_unix_ms, + lease_expires_at_unix_ms, + limit, + }) + } + pub const fn owner(&self) -> &LeaseOwner { + &self.owner + } + pub const fn lease_id_seed(&self) -> LeaseId { + self.lease_id_seed + } + /// Derives a stable, item-specific token from the caller's unique seed. + pub fn lease_id_for(&self, item_id: OutboxItemId) -> LeaseId { + let mut bytes = *self.lease_id_seed.as_bytes(); + for (byte, item_byte) in bytes.iter_mut().zip(item_id.as_bytes()) { + *byte ^= item_byte; + } + if bytes_are_zero(&bytes) { + bytes[0] = 1; + } + LeaseId(bytes) + } + pub const fn now_unix_ms(&self) -> u64 { + self.now_unix_ms + } + pub const fn lease_expires_at_unix_ms(&self) -> u64 { + self.lease_expires_at_unix_ms + } + pub const fn limit(&self) -> u16 { + self.limit + } +} + +/// Claimed outbox item with exact lease authority. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct ClaimedOutboxItem { + record: OutboxRecord, + lease: OutboxLease, +} + +impl ClaimedOutboxItem { + pub const fn new(record: OutboxRecord, lease: OutboxLease) -> Self { + Self { record, lease } + } + pub const fn record(&self) -> &OutboxRecord { + &self.record + } + pub const fn lease(&self) -> &OutboxLease { + &self.lease + } +} + +/// Request-bound evidence for one complete adapter attempt. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct DeliveryAttemptEvidence { + item_id: OutboxItemId, + lease_id: LeaseId, + expected_revision: OutboxRevision, + attempt: DeliveryAttempt, + receipt: DeliveryReceipt, + recorded_at_unix_ms: u64, +} + +impl DeliveryAttemptEvidence { + pub fn new( + item_id: OutboxItemId, + lease_id: LeaseId, + expected_revision: OutboxRevision, + attempt: DeliveryAttempt, + receipt: DeliveryReceipt, + recorded_at_unix_ms: u64, + ) -> Result<Self, Error> { + if recorded_at_unix_ms == 0 { + return Err(Error::InvalidOutboxTimestamp); + } + Ok(Self { + item_id, + lease_id, + expected_revision, + attempt, + receipt, + recorded_at_unix_ms, + }) + } + pub const fn item_id(&self) -> OutboxItemId { + self.item_id + } +} + +/// Passive outbox state summary. +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub struct OutboxStatus { + pub pending: u64, + pub leased: u64, + pub retryable: u64, + pub satisfied: u64, + pub exhausted: u64, +} + +impl OutboxStatus { + pub fn total(self) -> Option<u64> { + self.pending + .checked_add(self.leased)? + .checked_add(self.retryable)? + .checked_add(self.satisfied)? + .checked_add(self.exhausted) + } +} + +/// Backend-neutral durable delivery-plan SPI. +pub trait Outbox: Send + Sync { + fn enqueue(&self, item: EnqueueOutboxItem) -> BoxFuture<'_, Result<EnqueueReceipt, Error>>; + fn item(&self, item_id: OutboxItemId) -> BoxFuture<'_, Result<Option<OutboxRecord>, Error>>; + fn claim( + &self, + request: ClaimOutboxItems, + ) -> BoxFuture<'_, Result<Vec<ClaimedOutboxItem>, Error>>; + fn record_attempt( + &self, + evidence: DeliveryAttemptEvidence, + ) -> BoxFuture<'_, Result<OutboxRecord, Error>>; + fn release( + &self, + item_id: OutboxItemId, + lease_id: LeaseId, + expected_revision: OutboxRevision, + released_at_unix_ms: u64, + retry_not_before_unix_ms: Option<u64>, + ) -> BoxFuture<'_, Result<OutboxRecord, Error>>; + fn status(&self) -> BoxFuture<'_, Result<OutboxStatus, Error>>; +} + +const fn bytes_are_zero(bytes: &[u8; 16]) -> bool { + let mut index = 0; + while index < bytes.len() { + if bytes[index] != 0 { + return false; + } + index += 1; + } + true +} diff --git a/crates/storage/tests/outbox.rs b/crates/storage/tests/outbox.rs @@ -0,0 +1,397 @@ +use futures_executor::block_on; +use radroots_event::{SignedEvent, wire::Nip01EventWire}; +use radroots_storage::{ + Error, Outbox, + journal::OperationInstanceId, + outbox::{ + ClaimOutboxItems, ClaimedOutboxItem, DeliveryAttempt, DeliveryAttemptEvidence, + DeliveryPlanDigest, EnqueueDisposition, EnqueueOutboxItem, EnqueueReceipt, LeaseId, + LeaseOwner, OutboxItemId, OutboxLease, OutboxRecord, OutboxStage, OutboxStatus, + SatisfactionResult, + }, +}; +use radroots_transport::{ + BoxFuture, DeliveryReceipt, DeliveryRequest, Target, TargetSet, + outcome::DeliveryOutcome, + policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy}, + sink::{DeliveryPayload, DeliveryTargetReceipt}, +}; +use std::{collections::BTreeMap, sync::Mutex}; + +struct MemoryOutbox { + records: Mutex<BTreeMap<OutboxItemId, OutboxRecord>>, +} + +impl MemoryOutbox { + fn new() -> Self { + Self { + records: Mutex::new(BTreeMap::new()), + } + } +} + +impl Outbox for MemoryOutbox { + fn enqueue(&self, item: EnqueueOutboxItem) -> BoxFuture<'_, Result<EnqueueReceipt, Error>> { + Box::pin(async move { + let mut records = self.records.lock().expect("test outbox lock"); + if let Some(existing) = records.get(&item.item_id()) { + let exact = existing.operation_instance_id() == item.operation_instance_id() + && existing.plan_digest() == item.plan_digest() + && existing.request() == item.request(); + return if exact { + Ok(EnqueueReceipt::new( + EnqueueDisposition::Replay, + existing.clone(), + )) + } else { + Err(Error::OutboxPlanConflict) + }; + } + let record = item.into_record(); + records.insert(record.item_id(), record.clone()); + Ok(EnqueueReceipt::new(EnqueueDisposition::Created, record)) + }) + } + + fn item(&self, item_id: OutboxItemId) -> BoxFuture<'_, Result<Option<OutboxRecord>, Error>> { + Box::pin(async move { + Ok(self + .records + .lock() + .expect("test outbox lock") + .get(&item_id) + .cloned()) + }) + } + + fn claim( + &self, + request: ClaimOutboxItems, + ) -> BoxFuture<'_, Result<Vec<ClaimedOutboxItem>, Error>> { + Box::pin(async move { + let mut records = self.records.lock().expect("test outbox lock"); + let mut claimed = Vec::new(); + for record in records.values_mut() { + if claimed.len() >= usize::from(request.limit()) || record.stage().is_terminal() { + continue; + } + if matches!(record.retry_not_before_unix_ms(), Some(value) if value > request.now_unix_ms()) + { + continue; + } + if record + .lease() + .is_some_and(|lease| lease.is_active_at(request.now_unix_ms())) + { + continue; + } + let lease = OutboxLease::new( + request.lease_id_for(record.item_id()), + request.owner().clone(), + request.now_unix_ms(), + request.lease_expires_at_unix_ms(), + )?; + record.claim(lease.clone())?; + claimed.push(ClaimedOutboxItem::new(record.clone(), lease)); + } + Ok(claimed) + }) + } + + fn record_attempt( + &self, + evidence: DeliveryAttemptEvidence, + ) -> BoxFuture<'_, Result<OutboxRecord, Error>> { + Box::pin(async move { + let mut records = self.records.lock().expect("test outbox lock"); + let record = records + .get_mut(&evidence.item_id()) + .ok_or(Error::OutboxItemNotFound)?; + record.record_attempt(evidence)?; + Ok(record.clone()) + }) + } + + fn release( + &self, + item_id: OutboxItemId, + lease_id: LeaseId, + expected_revision: radroots_storage::outbox::OutboxRevision, + released_at_unix_ms: u64, + retry_not_before_unix_ms: Option<u64>, + ) -> BoxFuture<'_, Result<OutboxRecord, Error>> { + Box::pin(async move { + let mut records = self.records.lock().expect("test outbox lock"); + let record = records.get_mut(&item_id).ok_or(Error::OutboxItemNotFound)?; + record.release( + lease_id, + expected_revision, + released_at_unix_ms, + retry_not_before_unix_ms, + )?; + Ok(record.clone()) + }) + } + + fn status(&self) -> BoxFuture<'_, Result<OutboxStatus, Error>> { + Box::pin(async move { + let mut status = OutboxStatus { + pending: 0, + leased: 0, + retryable: 0, + satisfied: 0, + exhausted: 0, + }; + for record in self.records.lock().expect("test outbox lock").values() { + match record.stage() { + OutboxStage::Pending => status.pending += 1, + OutboxStage::Leased => status.leased += 1, + OutboxStage::Retryable => status.retryable += 1, + OutboxStage::Satisfied => status.satisfied += 1, + OutboxStage::Exhausted => status.exhausted += 1, + } + } + Ok(status) + }) + } +} + +fn signed_event() -> SignedEvent { + let mut wire = Nip01EventWire { + id: "0".repeat(64), + pubkey: "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df".to_owned(), + created_at: 1_800_000_100, + kind: 0, + tags: vec![], + content: "{\"display_name\":\"Moss Street Farm\",\"bot\":false}".to_owned(), + sig: "42".repeat(64), + extra: Default::default(), + }; + wire.id = wire + .computed_event_id() + .expect("canonical event id") + .to_hex(); + let raw_json = serde_json::to_string(&wire).expect("event JSON"); + SignedEvent::from_wire_verified_id(wire, raw_json).expect("signed event") +} + +fn targets() -> Vec<Target> { + vec![ + Target::nostr_relay("wss://one.example").expect("first target"), + Target::nostr_relay("wss://two.example").expect("second target"), + ] +} + +fn request() -> DeliveryRequest { + DeliveryRequest::new( + "outbox-test-request", + DeliveryPayload::new(signed_event()), + TargetSet::new(targets()).expect("target set"), + SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::all()), + 10_000, + ) + .expect("delivery request") +} + +fn enqueue(item_byte: u8, digest_byte: u8) -> EnqueueOutboxItem { + EnqueueOutboxItem::new( + OutboxItemId::new([item_byte; 16]).expect("item id"), + OperationInstanceId::new([9; 16]).expect("operation instance"), + DeliveryPlanDigest::new([digest_byte; 32]), + request(), + 10, + ) + .expect("enqueue request") +} + +fn claim(store: &dyn Outbox, now: u64, expiry: u64, seed: u8) -> ClaimedOutboxItem { + let request = ClaimOutboxItems::new( + LeaseOwner::parse("worker-a").expect("owner"), + LeaseId::new([seed; 16]).expect("lease seed"), + now, + expiry, + 1, + ) + .expect("claim request"); + block_on(store.claim(request)) + .expect("claim") + .pop() + .expect("claimed item") +} + +fn receipt(request: &DeliveryRequest, outcomes: [DeliveryOutcome; 2]) -> DeliveryReceipt { + let receipts = request + .target_set() + .targets() + .iter() + .cloned() + .zip(outcomes) + .map(|(target, outcome)| DeliveryTargetReceipt::attempted(target, outcome)) + .collect(); + DeliveryReceipt::for_request(request, receipts).expect("delivery receipt") +} + +#[test] +fn enqueue_is_idempotent_and_rejects_a_conflicting_plan() { + let store = MemoryOutbox::new(); + let first = enqueue(1, 2); + let created = block_on(store.enqueue(first.clone())).expect("created"); + assert_eq!(created.disposition(), EnqueueDisposition::Created); + let replay = block_on(store.enqueue(first)).expect("exact replay"); + assert_eq!(replay.disposition(), EnqueueDisposition::Replay); + + let conflict = enqueue(1, 3); + assert_eq!( + block_on(store.enqueue(conflict)), + Err(Error::OutboxPlanConflict) + ); +} + +#[test] +fn leases_exclude_concurrent_claims_expire_and_defer_retries() { + let store = MemoryOutbox::new(); + block_on(store.enqueue(enqueue(1, 2))).expect("enqueue"); + let first = claim(&store, 100, 200, 3); + + let concurrent = ClaimOutboxItems::new( + LeaseOwner::parse("worker-b").expect("owner"), + LeaseId::new([4; 16]).expect("seed"), + 150, + 250, + 1, + ) + .expect("claim request"); + assert!(block_on(store.claim(concurrent)).expect("claim").is_empty()); + + let stale = DeliveryAttemptEvidence::new( + first.record().item_id(), + first.lease().id(), + first.record().revision(), + DeliveryAttempt::FIRST, + receipt( + first.record().request(), + [DeliveryOutcome::accepted(), DeliveryOutcome::unavailable()], + ), + 200, + ) + .expect("evidence"); + assert_eq!( + block_on(store.record_attempt(stale)), + Err(Error::OutboxLeaseExpired) + ); + + let reclaimed = claim(&store, 200, 300, 5); + let released = block_on(store.release( + reclaimed.record().item_id(), + reclaimed.lease().id(), + reclaimed.record().revision(), + 210, + Some(250), + )) + .expect("release"); + assert_eq!(released.stage(), OutboxStage::Pending); + assert_eq!(released.retry_not_before_unix_ms(), Some(250)); +} + +#[test] +fn partial_retryable_evidence_advances_to_satisfaction() { + let store = MemoryOutbox::new(); + block_on(store.enqueue(enqueue(1, 2))).expect("enqueue"); + let first = claim(&store, 100, 200, 3); + let first_result = block_on( + store.record_attempt( + DeliveryAttemptEvidence::new( + first.record().item_id(), + first.lease().id(), + first.record().revision(), + DeliveryAttempt::FIRST, + receipt( + first.record().request(), + [DeliveryOutcome::accepted(), DeliveryOutcome::unavailable()], + ), + 150, + ) + .expect("attempt evidence"), + ), + ) + .expect("record partial attempt"); + assert_eq!(first_result.stage(), OutboxStage::Retryable); + assert_eq!(first_result.satisfaction(), SatisfactionResult::Pending); + assert_eq!(first_result.evidence().len(), 2); + assert!(first_result.evidence()[1].outcome().is_retryable()); + + let second = claim(&store, 250, 350, 4); + let second_result = block_on( + store.record_attempt( + DeliveryAttemptEvidence::new( + second.record().item_id(), + second.lease().id(), + second.record().revision(), + DeliveryAttempt::new(2).expect("second attempt"), + receipt( + second.record().request(), + [DeliveryOutcome::accepted(), DeliveryOutcome::delivered()], + ), + 300, + ) + .expect("attempt evidence"), + ), + ) + .expect("record successful attempt"); + assert_eq!(second_result.stage(), OutboxStage::Satisfied); + assert_eq!(second_result.satisfaction(), SatisfactionResult::Satisfied); + assert_eq!(second_result.evidence().len(), 4); + let second_target = &second_result.request().target_set().targets()[1]; + assert_eq!( + second_result + .latest_target_evidence(second_target.fingerprint()) + .expect("latest target evidence") + .attempt() + .get(), + 2 + ); + assert_eq!(block_on(store.status()).expect("status").satisfied, 1); +} + +#[test] +fn terminal_outcomes_exhaust_the_plan_and_models_are_bounded() { + let store = MemoryOutbox::new(); + block_on(store.enqueue(enqueue(1, 2))).expect("enqueue"); + let claimed = claim(&store, 100, 200, 3); + let exhausted = block_on( + store.record_attempt( + DeliveryAttemptEvidence::new( + claimed.record().item_id(), + claimed.lease().id(), + claimed.record().revision(), + DeliveryAttempt::FIRST, + receipt( + claimed.record().request(), + [DeliveryOutcome::rejected(), DeliveryOutcome::rejected()], + ), + 150, + ) + .expect("attempt evidence"), + ), + ) + .expect("record terminal attempt"); + assert_eq!(exhausted.stage(), OutboxStage::Exhausted); + assert_eq!(exhausted.satisfaction(), SatisfactionResult::Exhausted); + assert_eq!(block_on(store.status()).expect("status").total(), Some(1)); + + assert_eq!(OutboxItemId::new([0; 16]), Err(Error::InvalidOutboxItemId)); + assert_eq!(DeliveryAttempt::new(0), Err(Error::InvalidDeliveryAttempt)); + assert_eq!( + ClaimOutboxItems::new( + LeaseOwner::parse("worker").expect("owner"), + LeaseId::new([1; 16]).expect("seed"), + 1, + 2, + 0, + ), + Err(Error::InvalidOutboxClaimLimit) + ); + assert!( + format!("{:?}", LeaseOwner::parse("secret-worker").expect("owner")).contains("[REDACTED]") + ); +}