commit 0e6cc7a6871e7fd7ca06291905ebd969a370ea9c
parent 5c8fb3c3b8a93ff806fbd8f9be2653074f48cfd3
Author: triesap <tyson@radroots.org>
Date: Wed, 5 Aug 2026 01:56:22 +0000
refactor(storage): add independent delivery plans
- add artifact-owned request-digested delivery plan state
- preserve typed attempts partial evidence failures and retry schedules
- enforce independent claim fencing exhaustion and checked attempt bounds
- verify reopen serde features package clippy and full storage suites
Diffstat:
4 files changed, 1025 insertions(+), 0 deletions(-)
diff --git a/crates/storage/src/authored_delivery.rs b/crates/storage/src/authored_delivery.rs
@@ -0,0 +1,673 @@
+//! Independent durable delivery-plan, attempt, retry, and evidence models.
+
+use core::num::{NonZeroU32, NonZeroU64};
+use radroots_transport::{
+ DeliveryReceipt, DeliveryRequest, SinkFailure,
+ outcome::Retryability,
+ policy::{SatisfactionClass, SatisfactionPolicy, SatisfactionState, evaluate_satisfaction},
+};
+use sha2::{Digest, Sha256};
+use std::vec::Vec;
+
+use crate::{
+ Error,
+ authored::{
+ AuthoredArtifactId, FailureClass, RetrySchedule, WorkClaim, WorkFailure, WorkPhase,
+ },
+};
+
+pub const DELIVERY_PLAN_ATTEMPTS_MAX: u32 = 1_024;
+
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[cfg_attr(feature = "serde", serde(try_from = "[u8; 16]", into = "[u8; 16]"))]
+#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
+pub struct AuthoredDeliveryPlanId([u8; 16]);
+
+impl AuthoredDeliveryPlanId {
+ pub const fn new(value: [u8; 16]) -> Result<Self, Error> {
+ if all_zero(&value) {
+ Err(Error::InvalidAuthoredDeliveryPlan)
+ } else {
+ Ok(Self(value))
+ }
+ }
+
+ pub const fn as_bytes(&self) -> &[u8; 16] {
+ &self.0
+ }
+}
+
+impl TryFrom<[u8; 16]> for AuthoredDeliveryPlanId {
+ type Error = Error;
+ fn try_from(value: [u8; 16]) -> Result<Self, Self::Error> {
+ Self::new(value)
+ }
+}
+
+impl From<AuthoredDeliveryPlanId> for [u8; 16] {
+ fn from(value: AuthoredDeliveryPlanId) -> Self {
+ value.0
+ }
+}
+
+#[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 AuthoredDeliveryState {
+ Pending,
+ Retryable,
+ Satisfied,
+ Exhausted,
+ FailedTerminal,
+ Cancelled,
+}
+
+impl AuthoredDeliveryState {
+ pub const fn is_terminal(self) -> bool {
+ matches!(
+ self,
+ Self::Satisfied | Self::Exhausted | Self::FailedTerminal | Self::Cancelled
+ )
+ }
+}
+
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub enum DeliveryAttemptOutcome {
+ Receipt(DeliveryReceipt),
+ SinkFailure(SinkFailure),
+}
+
+impl DeliveryAttemptOutcome {
+ fn validate_for(&self, request: &DeliveryRequest) -> Result<(), Error> {
+ match self {
+ Self::Receipt(receipt) => receipt
+ .validate_for_request(request)
+ .map_err(|_| Error::InvalidAuthoredDeliveryPlan),
+ Self::SinkFailure(failure) => failure
+ .validate_for_request(request)
+ .map_err(|_| Error::InvalidAuthoredDeliveryPlan),
+ }
+ }
+}
+
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct AuthoredDeliveryAttempt {
+ attempt: NonZeroU32,
+ recorded_at_unix_ms: u64,
+ outcome: DeliveryAttemptOutcome,
+ satisfaction: SatisfactionState,
+}
+
+impl AuthoredDeliveryAttempt {
+ pub fn reconstruct(
+ attempt: NonZeroU32,
+ recorded_at_unix_ms: u64,
+ outcome: DeliveryAttemptOutcome,
+ satisfaction: SatisfactionState,
+ ) -> Result<Self, Error> {
+ if recorded_at_unix_ms == 0 {
+ return Err(Error::InvalidAuthoredDeliveryPlan);
+ }
+ Ok(Self {
+ attempt,
+ recorded_at_unix_ms,
+ outcome,
+ satisfaction,
+ })
+ }
+
+ pub const fn attempt(&self) -> NonZeroU32 {
+ self.attempt
+ }
+ pub const fn recorded_at_unix_ms(&self) -> u64 {
+ self.recorded_at_unix_ms
+ }
+ pub const fn outcome(&self) -> &DeliveryAttemptOutcome {
+ &self.outcome
+ }
+ pub const fn satisfaction(&self) -> SatisfactionState {
+ self.satisfaction
+ }
+}
+
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[cfg_attr(
+ feature = "serde",
+ serde(
+ try_from = "AuthoredDeliveryPlanWire",
+ into = "AuthoredDeliveryPlanWire"
+ )
+)]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct AuthoredDeliveryPlan {
+ plan_id: AuthoredDeliveryPlanId,
+ artifact_id: AuthoredArtifactId,
+ request_digest: [u8; 32],
+ request: DeliveryRequest,
+ state: AuthoredDeliveryState,
+ attempts: Vec<AuthoredDeliveryAttempt>,
+ attempt_count: u32,
+ retry: Option<RetrySchedule>,
+ claim: Option<WorkClaim>,
+ last_failure: Option<WorkFailure>,
+ created_at_unix_ms: u64,
+ updated_at_unix_ms: u64,
+ revision: NonZeroU64,
+}
+
+impl AuthoredDeliveryPlan {
+ pub fn new(
+ plan_id: AuthoredDeliveryPlanId,
+ artifact_id: AuthoredArtifactId,
+ request: DeliveryRequest,
+ created_at_unix_ms: u64,
+ ) -> Result<Self, Error> {
+ let request_digest = delivery_request_digest(&request);
+ Self::reconstruct(Self {
+ plan_id,
+ artifact_id,
+ request_digest,
+ request,
+ state: AuthoredDeliveryState::Pending,
+ attempts: Vec::new(),
+ attempt_count: 0,
+ retry: None,
+ claim: None,
+ last_failure: None,
+ created_at_unix_ms,
+ updated_at_unix_ms: created_at_unix_ms,
+ revision: NonZeroU64::MIN,
+ })
+ }
+
+ pub fn reconstruct(value: Self) -> Result<Self, Error> {
+ value.validate()?;
+ Ok(value)
+ }
+
+ pub fn validate(&self) -> Result<(), Error> {
+ if self.created_at_unix_ms == 0
+ || self.updated_at_unix_ms < self.created_at_unix_ms
+ || self.request_digest != delivery_request_digest(&self.request)
+ || self.attempt_count > DELIVERY_PLAN_ATTEMPTS_MAX
+ || usize::try_from(self.attempt_count).ok() != Some(self.attempts.len())
+ || (matches!(self.state, AuthoredDeliveryState::Retryable) != self.retry.is_some())
+ || (self.state.is_terminal() && self.claim.is_some())
+ {
+ return Err(Error::InvalidAuthoredDeliveryPlan);
+ }
+ for (index, attempt) in self.attempts.iter().enumerate() {
+ if attempt.attempt.get() != u32::try_from(index + 1).unwrap_or(u32::MAX)
+ || attempt.recorded_at_unix_ms < self.created_at_unix_ms
+ || attempt.outcome.validate_for(&self.request).is_err()
+ || self
+ .evaluate_outcomes(self.attempts[..=index].iter().map(|value| &value.outcome))
+ .ok()
+ != Some(attempt.satisfaction)
+ {
+ return Err(Error::InvalidAuthoredDeliveryPlan);
+ }
+ }
+ if self
+ .attempts
+ .windows(2)
+ .any(|pair| pair[0].recorded_at_unix_ms > pair[1].recorded_at_unix_ms)
+ || self
+ .claim
+ .as_ref()
+ .is_some_and(|claim| claim.validate().is_err())
+ || self.retry.as_ref().is_some_and(|retry| {
+ retry.failure().phase() != WorkPhase::Delivery
+ || retry.attempt().get() != self.attempt_count
+ || self.attempts.last().is_none_or(|attempt| {
+ retry.not_before_unix_ms() <= attempt.recorded_at_unix_ms
+ })
+ })
+ || self.claim.as_ref().is_some_and(|claim| {
+ claim.row_revision().get().checked_add(1) != Some(self.revision.get())
+ || claim.acquired_at_unix_ms() != self.updated_at_unix_ms
+ })
+ {
+ return Err(Error::InvalidAuthoredDeliveryPlan);
+ }
+ match self.state {
+ AuthoredDeliveryState::Pending
+ | AuthoredDeliveryState::Satisfied
+ | AuthoredDeliveryState::Cancelled => {
+ if self.retry.is_some() || self.last_failure.is_some() {
+ return Err(Error::InvalidAuthoredDeliveryPlan);
+ }
+ }
+ AuthoredDeliveryState::Exhausted => {
+ if self.retry.is_some()
+ || self.last_failure.as_ref().is_some_and(|failure| {
+ failure.phase() != WorkPhase::Delivery
+ || failure.class() != FailureClass::Terminal
+ })
+ {
+ return Err(Error::InvalidAuthoredDeliveryPlan);
+ }
+ }
+ AuthoredDeliveryState::Retryable => {
+ if self.retry.as_ref().map(RetrySchedule::failure) != self.last_failure.as_ref() {
+ return Err(Error::InvalidAuthoredDeliveryPlan);
+ }
+ }
+ AuthoredDeliveryState::FailedTerminal => {
+ if !matches!(
+ self.last_failure.as_ref().map(WorkFailure::class),
+ Some(FailureClass::Terminal)
+ ) {
+ return Err(Error::InvalidAuthoredDeliveryPlan);
+ }
+ }
+ }
+ Ok(())
+ }
+
+ pub fn claim(&mut self, claim: WorkClaim, now_unix_ms: u64) -> Result<(), Error> {
+ if self.state.is_terminal()
+ || self.claim.is_some()
+ || claim.row_revision() != self.revision
+ || claim.acquired_at_unix_ms() != now_unix_ms
+ || self
+ .retry
+ .as_ref()
+ .is_some_and(|retry| now_unix_ms < retry.not_before_unix_ms())
+ {
+ return Err(Error::DeliveryPlanClaimConflict);
+ }
+ claim.validate()?;
+ let previous = self.clone();
+ self.claim = Some(claim);
+ if let Err(error) = self.advance(now_unix_ms) {
+ *self = previous;
+ return Err(error);
+ }
+ if let Err(error) = self.validate() {
+ *self = previous;
+ return Err(error);
+ }
+ Ok(())
+ }
+
+ pub fn apply_receipt(
+ &mut self,
+ token: &[u8; 16],
+ generation: NonZeroU64,
+ claim_revision: NonZeroU64,
+ receipt: DeliveryReceipt,
+ retry: Option<RetrySchedule>,
+ recorded_at_unix_ms: u64,
+ ) -> Result<(), Error> {
+ self.require_claim(token, generation, claim_revision, recorded_at_unix_ms)?;
+ receipt
+ .validate_for_request(&self.request)
+ .map_err(|_| Error::InvalidAuthoredDeliveryPlan)?;
+ let satisfaction = self.evaluate_with(DeliveryAttemptOutcome::Receipt(receipt.clone()))?;
+ let (state, last_failure) = match satisfaction {
+ SatisfactionState::Satisfied if retry.is_none() => {
+ (AuthoredDeliveryState::Satisfied, None)
+ }
+ SatisfactionState::Exhausted if retry.is_none() => {
+ (AuthoredDeliveryState::Exhausted, None)
+ }
+ SatisfactionState::Pending => {
+ let schedule = retry.as_ref().ok_or(Error::InvalidRetrySchedule)?;
+ if schedule.failure().phase() != WorkPhase::Delivery {
+ return Err(Error::InvalidRetrySchedule);
+ }
+ (
+ AuthoredDeliveryState::Retryable,
+ Some(schedule.failure().clone()),
+ )
+ }
+ SatisfactionState::Satisfied | SatisfactionState::Exhausted => {
+ return Err(Error::InvalidRetrySchedule);
+ }
+ };
+ self.apply_attempt(
+ DeliveryAttemptOutcome::Receipt(receipt),
+ satisfaction,
+ state,
+ retry,
+ last_failure,
+ recorded_at_unix_ms,
+ )
+ }
+
+ pub fn apply_sink_failure(
+ &mut self,
+ token: &[u8; 16],
+ generation: NonZeroU64,
+ claim_revision: NonZeroU64,
+ failure: SinkFailure,
+ retry: Option<RetrySchedule>,
+ recorded_at_unix_ms: u64,
+ ) -> Result<(), Error> {
+ self.require_claim(token, generation, claim_revision, recorded_at_unix_ms)?;
+ failure
+ .validate_for_request(&self.request)
+ .map_err(|_| Error::InvalidAuthoredDeliveryPlan)?;
+ let outcome = DeliveryAttemptOutcome::SinkFailure(failure.clone());
+ let satisfaction = self.evaluate_with(outcome.clone())?;
+ let typed = WorkFailure::new(
+ failure.code(),
+ WorkPhase::Delivery,
+ match failure.retryability() {
+ Retryability::Retryable => FailureClass::Retryable,
+ Retryability::Terminal | Retryability::NotApplicable => FailureClass::Terminal,
+ },
+ failure.retry_after_unix_ms(),
+ failure.message().map(str::to_owned),
+ )?;
+ let (state, retry, last_failure) = if satisfaction == SatisfactionState::Satisfied {
+ if retry.is_some() {
+ return Err(Error::InvalidRetrySchedule);
+ }
+ (AuthoredDeliveryState::Satisfied, None, None)
+ } else if satisfaction == SatisfactionState::Exhausted {
+ if retry.is_some() {
+ return Err(Error::InvalidRetrySchedule);
+ }
+ (AuthoredDeliveryState::Exhausted, None, Some(typed))
+ } else if failure.retryability() == Retryability::Retryable {
+ let schedule = retry.ok_or(Error::InvalidRetrySchedule)?;
+ if schedule.failure() != &typed {
+ return Err(Error::InvalidRetrySchedule);
+ }
+ (
+ AuthoredDeliveryState::Retryable,
+ Some(schedule),
+ Some(typed),
+ )
+ } else {
+ if retry.is_some() {
+ return Err(Error::InvalidRetrySchedule);
+ }
+ (AuthoredDeliveryState::FailedTerminal, None, Some(typed))
+ };
+ self.apply_attempt(
+ outcome,
+ satisfaction,
+ state,
+ retry,
+ last_failure,
+ recorded_at_unix_ms,
+ )
+ }
+
+ pub fn cancel(&mut self, cancelled_at_unix_ms: u64) -> Result<(), Error> {
+ if self.state.is_terminal() {
+ return Err(Error::InvalidAuthoredDeliveryPlan);
+ }
+ let previous = self.clone();
+ self.state = AuthoredDeliveryState::Cancelled;
+ self.claim = None;
+ self.retry = None;
+ self.last_failure = None;
+ if let Err(error) = self.advance(cancelled_at_unix_ms) {
+ *self = previous;
+ return Err(error);
+ }
+ self.validate()
+ }
+
+ fn apply_attempt(
+ &mut self,
+ outcome: DeliveryAttemptOutcome,
+ satisfaction: SatisfactionState,
+ state: AuthoredDeliveryState,
+ retry: Option<RetrySchedule>,
+ last_failure: Option<WorkFailure>,
+ recorded_at_unix_ms: u64,
+ ) -> Result<(), Error> {
+ let next = self
+ .attempt_count
+ .checked_add(1)
+ .filter(|attempt| *attempt <= DELIVERY_PLAN_ATTEMPTS_MAX)
+ .and_then(NonZeroU32::new)
+ .ok_or(Error::DeliveryAttemptOverflow)?;
+ let previous = self.clone();
+ self.attempt_count = next.get();
+ self.attempts.push(AuthoredDeliveryAttempt::reconstruct(
+ next,
+ recorded_at_unix_ms,
+ outcome,
+ satisfaction,
+ )?);
+ self.state = state;
+ self.retry = retry;
+ self.claim = None;
+ self.last_failure = last_failure;
+ if let Err(error) = self.advance(recorded_at_unix_ms) {
+ *self = previous;
+ return Err(error);
+ }
+ if let Err(error) = self.validate() {
+ *self = previous;
+ return Err(error);
+ }
+ Ok(())
+ }
+
+ fn evaluate_with(&self, next: DeliveryAttemptOutcome) -> Result<SatisfactionState, Error> {
+ self.evaluate_outcomes(
+ self.attempts
+ .iter()
+ .map(|attempt| &attempt.outcome)
+ .chain(core::iter::once(&next)),
+ )
+ }
+
+ fn evaluate_outcomes<'a, I>(&self, outcomes: I) -> Result<SatisfactionState, Error>
+ where
+ I: IntoIterator<Item = &'a DeliveryAttemptOutcome>,
+ {
+ let mut evidence = Vec::new();
+ for outcome in outcomes {
+ match outcome {
+ DeliveryAttemptOutcome::Receipt(receipt) => {
+ evidence.extend(
+ receipt
+ .target_receipts()
+ .iter()
+ .map(|entry| (entry.target().fingerprint(), entry.outcome())),
+ );
+ }
+ DeliveryAttemptOutcome::SinkFailure(failure) => {
+ evidence.extend(
+ failure
+ .partial_evidence()
+ .iter()
+ .map(|entry| (entry.target().fingerprint(), entry.outcome())),
+ );
+ }
+ }
+ }
+ evaluate_satisfaction(
+ self.request.satisfaction(),
+ self.request.target_set(),
+ evidence,
+ )
+ .map_err(|_| Error::InvalidAuthoredDeliveryPlan)
+ }
+
+ fn require_claim(
+ &self,
+ token: &[u8; 16],
+ generation: NonZeroU64,
+ claim_revision: NonZeroU64,
+ now_unix_ms: u64,
+ ) -> Result<(), Error> {
+ if !self.claim.as_ref().is_some_and(|claim| {
+ claim.matches_fence(token, generation, claim_revision, now_unix_ms)
+ }) {
+ return Err(Error::DeliveryPlanClaimConflict);
+ }
+ Ok(())
+ }
+
+ fn advance(&mut self, at_unix_ms: u64) -> Result<(), Error> {
+ if at_unix_ms < self.updated_at_unix_ms {
+ return Err(Error::InvalidAuthoredDeliveryPlan);
+ }
+ self.revision = self
+ .revision
+ .get()
+ .checked_add(1)
+ .and_then(NonZeroU64::new)
+ .ok_or(Error::InvalidAuthoredDeliveryPlan)?;
+ self.updated_at_unix_ms = at_unix_ms;
+ Ok(())
+ }
+
+ pub const fn plan_id(&self) -> AuthoredDeliveryPlanId {
+ self.plan_id
+ }
+ pub const fn artifact_id(&self) -> AuthoredArtifactId {
+ self.artifact_id
+ }
+ pub const fn request_digest(&self) -> &[u8; 32] {
+ &self.request_digest
+ }
+ pub const fn request(&self) -> &DeliveryRequest {
+ &self.request
+ }
+ pub const fn state(&self) -> AuthoredDeliveryState {
+ self.state
+ }
+ pub fn attempts(&self) -> &[AuthoredDeliveryAttempt] {
+ self.attempts.as_slice()
+ }
+ pub const fn attempt_count(&self) -> u32 {
+ self.attempt_count
+ }
+ pub const fn retry(&self) -> Option<&RetrySchedule> {
+ self.retry.as_ref()
+ }
+ pub const fn claim_evidence(&self) -> Option<&WorkClaim> {
+ self.claim.as_ref()
+ }
+ pub const fn last_failure(&self) -> Option<&WorkFailure> {
+ self.last_failure.as_ref()
+ }
+ pub const fn revision(&self) -> NonZeroU64 {
+ self.revision
+ }
+}
+
+#[cfg(feature = "serde")]
+#[derive(serde::Serialize, serde::Deserialize)]
+struct AuthoredDeliveryPlanWire {
+ plan_id: AuthoredDeliveryPlanId,
+ artifact_id: AuthoredArtifactId,
+ request_digest: [u8; 32],
+ request: DeliveryRequest,
+ state: AuthoredDeliveryState,
+ attempts: Vec<AuthoredDeliveryAttempt>,
+ attempt_count: u32,
+ retry: Option<RetrySchedule>,
+ claim: Option<WorkClaim>,
+ last_failure: Option<WorkFailure>,
+ created_at_unix_ms: u64,
+ updated_at_unix_ms: u64,
+ revision: NonZeroU64,
+}
+
+#[cfg(feature = "serde")]
+impl TryFrom<AuthoredDeliveryPlanWire> for AuthoredDeliveryPlan {
+ type Error = Error;
+ fn try_from(value: AuthoredDeliveryPlanWire) -> Result<Self, Self::Error> {
+ Self::reconstruct(Self {
+ plan_id: value.plan_id,
+ artifact_id: value.artifact_id,
+ request_digest: value.request_digest,
+ request: value.request,
+ state: value.state,
+ attempts: value.attempts,
+ attempt_count: value.attempt_count,
+ retry: value.retry,
+ claim: value.claim,
+ last_failure: value.last_failure,
+ created_at_unix_ms: value.created_at_unix_ms,
+ updated_at_unix_ms: value.updated_at_unix_ms,
+ revision: value.revision,
+ })
+ }
+}
+
+#[cfg(feature = "serde")]
+impl From<AuthoredDeliveryPlan> for AuthoredDeliveryPlanWire {
+ fn from(value: AuthoredDeliveryPlan) -> Self {
+ Self {
+ plan_id: value.plan_id,
+ artifact_id: value.artifact_id,
+ request_digest: value.request_digest,
+ request: value.request,
+ state: value.state,
+ attempts: value.attempts,
+ attempt_count: value.attempt_count,
+ retry: value.retry,
+ claim: value.claim,
+ last_failure: value.last_failure,
+ created_at_unix_ms: value.created_at_unix_ms,
+ updated_at_unix_ms: value.updated_at_unix_ms,
+ revision: value.revision,
+ }
+ }
+}
+
+fn delivery_request_digest(request: &DeliveryRequest) -> [u8; 32] {
+ let mut hasher = Sha256::new();
+ hash_field(&mut hasher, b"radroots.authored.delivery.v2");
+ hash_field(&mut hasher, request.request_id().as_str().as_bytes());
+ hash_field(&mut hasher, request.payload().event().raw_json().as_bytes());
+ for target in request.target_set().targets() {
+ hash_field(&mut hasher, target.fingerprint().as_str().as_bytes());
+ }
+ let policy = request.satisfaction();
+ hasher.update([match policy.class() {
+ SatisfactionClass::Accepted => 0,
+ SatisfactionClass::Delivered => 1,
+ }]);
+ hash_policy(&mut hasher, policy);
+ hasher.update(request.deadline_unix_ms().to_be_bytes());
+ hasher.finalize().into()
+}
+
+fn hash_policy(hasher: &mut Sha256, policy: &SatisfactionPolicy) {
+ let targets = policy.targets();
+ if targets.is_any() {
+ hasher.update([0]);
+ } else if targets.is_all() {
+ hasher.update([1]);
+ } else if let Some(threshold) = targets.quorum_threshold() {
+ hasher.update([2]);
+ hasher.update(threshold.to_be_bytes());
+ } else if let Some(required) = targets.required_targets() {
+ hasher.update([3]);
+ for target in required {
+ hash_field(hasher, target.as_str().as_bytes());
+ }
+ }
+}
+
+fn hash_field(hasher: &mut Sha256, value: &[u8]) {
+ hasher.update(u64::try_from(value.len()).unwrap_or(u64::MAX).to_be_bytes());
+ hasher.update(value);
+}
+
+const fn all_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/src/error.rs b/crates/storage/src/error.rs
@@ -9,6 +9,9 @@ pub enum Error {
InvalidAuthoredOperation,
InvalidAuthoredArtifact,
InvalidAuthoredTransition,
+ InvalidAuthoredDeliveryPlan,
+ DeliveryPlanClaimConflict,
+ DeliveryAttemptOverflow,
InvalidWorkClaim,
InvalidWorkFailure,
InvalidRetrySchedule,
@@ -126,6 +129,9 @@ impl fmt::Display for Error {
Self::InvalidAuthoredOperation => "storage authored operation is invalid",
Self::InvalidAuthoredArtifact => "storage authored artifact is invalid",
Self::InvalidAuthoredTransition => "storage authored transition is invalid",
+ Self::InvalidAuthoredDeliveryPlan => "storage authored delivery plan is invalid",
+ Self::DeliveryPlanClaimConflict => "storage delivery plan claim conflicts",
+ Self::DeliveryAttemptOverflow => "storage delivery attempt counter overflowed",
Self::InvalidWorkClaim => "storage authored work claim is invalid",
Self::InvalidWorkFailure => "storage authored work failure is invalid",
Self::InvalidRetrySchedule => "storage authored retry schedule is invalid",
@@ -271,6 +277,9 @@ mod tests {
InvalidAuthoredOperation,
InvalidAuthoredArtifact,
InvalidAuthoredTransition,
+ InvalidAuthoredDeliveryPlan,
+ DeliveryPlanClaimConflict,
+ DeliveryAttemptOverflow,
InvalidWorkClaim,
InvalidWorkFailure,
InvalidRetrySchedule,
diff --git a/crates/storage/src/lib.rs b/crates/storage/src/lib.rs
@@ -3,6 +3,7 @@
pub mod atomic;
pub mod authored;
+pub mod authored_delivery;
pub mod backup;
mod error;
pub mod event;
diff --git a/crates/storage/tests/authored_delivery.rs b/crates/storage/tests/authored_delivery.rs
@@ -0,0 +1,342 @@
+use core::num::{NonZeroU32, NonZeroU64};
+use radroots_event::{SignedEvent, wire::v1::Nip01EventWire};
+use radroots_storage::{
+ Error,
+ authored::{
+ AuthoredArtifactId, FailureClass, RetrySchedule, WorkClaim, WorkFailure, WorkPhase,
+ },
+ authored_delivery::{
+ AuthoredDeliveryAttempt, AuthoredDeliveryPlan, AuthoredDeliveryPlanId,
+ AuthoredDeliveryState, DELIVERY_PLAN_ATTEMPTS_MAX, DeliveryAttemptOutcome,
+ },
+};
+use radroots_transport::{
+ DeliveryReceipt, DeliveryRequest, SinkFailure, Target, TargetSet,
+ outcome::{DeliveryOutcome, Retryability},
+ policy::{SatisfactionClass, SatisfactionPolicy, SatisfactionState, TargetPolicy},
+ sink::{DeliveryPayload, DeliveryTargetReceipt},
+};
+
+fn signed_event() -> SignedEvent {
+ let mut wire = Nip01EventWire {
+ id: "0".repeat(64),
+ pubkey: "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df".to_owned(),
+ created_at: 1_800_000_100,
+ kind: 20_000,
+ tags: Vec::new(),
+ content: "delivery plan".to_owned(),
+ sig: "dd".repeat(64),
+ extra: Default::default(),
+ };
+ wire.id = wire.computed_event_id().expect("event ID").to_hex();
+ let raw = serde_json::to_string(&wire).expect("raw event");
+ SignedEvent::from_wire_verified_id(wire, raw).expect("signed event")
+}
+
+fn target_set() -> TargetSet {
+ TargetSet::new(vec![
+ Target::nostr_relay("wss://one.example").expect("one"),
+ Target::nostr_relay("wss://two.example").expect("two"),
+ ])
+ .expect("targets")
+}
+
+fn request(policy: TargetPolicy) -> DeliveryRequest {
+ DeliveryRequest::new(
+ "authored-delivery",
+ DeliveryPayload::new(signed_event()),
+ target_set(),
+ SatisfactionPolicy::new(SatisfactionClass::Accepted, policy),
+ 100,
+ )
+ .expect("request")
+}
+
+fn plan(value: u8, policy: TargetPolicy) -> AuthoredDeliveryPlan {
+ AuthoredDeliveryPlan::new(
+ AuthoredDeliveryPlanId::new([value; 16]).expect("plan ID"),
+ AuthoredArtifactId::new([9; 16]).expect("artifact ID"),
+ request(policy),
+ 10,
+ )
+ .expect("plan")
+}
+
+fn claim(plan: &AuthoredDeliveryPlan, token: u8, acquired: u64) -> WorkClaim {
+ WorkClaim::new(
+ [token; 16],
+ format!("worker-{token}"),
+ NonZeroU64::new(u64::from(token)).expect("generation"),
+ acquired,
+ acquired + 20,
+ plan.revision(),
+ )
+ .expect("claim")
+}
+
+fn receipt(request: &DeliveryRequest, outcomes: Vec<DeliveryOutcome>) -> DeliveryReceipt {
+ DeliveryReceipt::for_request(
+ request,
+ request
+ .target_set()
+ .targets()
+ .iter()
+ .cloned()
+ .zip(outcomes)
+ .map(|(target, outcome)| DeliveryTargetReceipt::attempted(target, outcome))
+ .collect(),
+ )
+ .expect("receipt")
+}
+
+fn retry(code: &str, attempt: u32, not_before: u64) -> RetrySchedule {
+ let failure = WorkFailure::new(
+ code,
+ WorkPhase::Delivery,
+ FailureClass::Retryable,
+ Some(not_before),
+ None,
+ )
+ .expect("failure");
+ RetrySchedule::new(
+ NonZeroU32::new(attempt).expect("attempt"),
+ not_before,
+ failure,
+ )
+ .expect("retry")
+}
+
+#[test]
+fn independent_plans_claim_and_progress_without_cross_blocking() {
+ let mut first = plan(1, TargetPolicy::any());
+ let mut second = plan(2, TargetPolicy::all());
+ let first_claim = claim(&first, 1, 11);
+ let second_claim = claim(&second, 2, 11);
+ first.claim(first_claim.clone(), 11).expect("first claim");
+ second
+ .claim(second_claim.clone(), 11)
+ .expect("second claim");
+
+ let first_receipt = receipt(
+ first.request(),
+ vec![DeliveryOutcome::accepted(), DeliveryOutcome::unavailable()],
+ );
+ first
+ .apply_receipt(
+ first_claim.token(),
+ first_claim.generation(),
+ first_claim.row_revision(),
+ first_receipt,
+ None,
+ 12,
+ )
+ .expect("first satisfied");
+ assert_eq!(first.state(), AuthoredDeliveryState::Satisfied);
+ assert!(second.claim_evidence().is_some());
+
+ let second_receipt = receipt(
+ second.request(),
+ vec![DeliveryOutcome::accepted(), DeliveryOutcome::unavailable()],
+ );
+ second
+ .apply_receipt(
+ second_claim.token(),
+ second_claim.generation(),
+ second_claim.row_revision(),
+ second_receipt,
+ Some(retry("delivery_pending", 1, 20)),
+ 12,
+ )
+ .expect("second retry");
+ assert_eq!(second.state(), AuthoredDeliveryState::Retryable);
+ assert_eq!(second.retry().expect("retry").not_before_unix_ms(), 20);
+ let before_early_claim = second.clone();
+ let early = claim(&second, 7, 19);
+ assert_eq!(
+ second.claim(early, 19),
+ Err(Error::DeliveryPlanClaimConflict)
+ );
+ assert_eq!(second, before_early_claim);
+}
+
+#[test]
+fn retry_schedule_and_partial_evidence_survive_reconstruction() {
+ let mut delivery = plan(3, TargetPolicy::all());
+ let active = claim(&delivery, 3, 11);
+ delivery.claim(active.clone(), 11).expect("claim");
+ let partial = DeliveryTargetReceipt::attempted(
+ delivery.request().target_set().targets()[0].clone(),
+ DeliveryOutcome::accepted(),
+ );
+ let sink_failure = SinkFailure::for_request(
+ delivery.request(),
+ "relay_batch_unavailable",
+ Retryability::Retryable,
+ Some(20),
+ None,
+ vec![partial],
+ )
+ .expect("sink failure");
+ delivery
+ .apply_sink_failure(
+ active.token(),
+ active.generation(),
+ active.row_revision(),
+ sink_failure,
+ Some(retry("relay_batch_unavailable", 1, 20)),
+ 12,
+ )
+ .expect("retryable failure");
+
+ let json = serde_json::to_string(&delivery).expect("plan json");
+ let reopened: AuthoredDeliveryPlan = serde_json::from_str(&json).expect("reopen plan");
+ assert_eq!(reopened, delivery);
+ assert_eq!(reopened.attempt_count(), 1);
+ assert_eq!(reopened.state(), AuthoredDeliveryState::Retryable);
+ assert_eq!(
+ reopened.attempts()[0].satisfaction(),
+ SatisfactionState::Pending
+ );
+ assert_eq!(
+ reopened.last_failure().expect("failure").code(),
+ "relay_batch_unavailable"
+ );
+}
+
+#[test]
+fn terminal_partial_success_stale_claim_and_invalid_receipt_fail_closed() {
+ let mut any = plan(4, TargetPolicy::any());
+ let active = claim(&any, 4, 11);
+ any.claim(active.clone(), 11).expect("claim");
+ let partial = DeliveryTargetReceipt::attempted(
+ any.request().target_set().targets()[0].clone(),
+ DeliveryOutcome::accepted(),
+ );
+ let failure = SinkFailure::for_request(
+ any.request(),
+ "terminal_batch_failure",
+ Retryability::Terminal,
+ None,
+ None,
+ vec![partial],
+ )
+ .expect("terminal failure");
+ let other_request = DeliveryRequest::new(
+ "other-request",
+ any.request().payload().clone(),
+ any.request().target_set().clone(),
+ any.request().satisfaction().clone(),
+ any.request().deadline_unix_ms(),
+ )
+ .expect("other request");
+ let invalid_receipt = receipt(
+ &other_request,
+ vec![DeliveryOutcome::accepted(), DeliveryOutcome::accepted()],
+ );
+ assert_eq!(
+ any.apply_receipt(
+ active.token(),
+ active.generation(),
+ active.row_revision(),
+ invalid_receipt,
+ None,
+ 12,
+ ),
+ Err(Error::InvalidAuthoredDeliveryPlan)
+ );
+ assert_eq!(any.attempt_count(), 0);
+ assert!(any.claim_evidence().is_some());
+ assert_eq!(
+ any.apply_sink_failure(
+ &[8; 16],
+ active.generation(),
+ active.row_revision(),
+ failure.clone(),
+ None,
+ 12,
+ ),
+ Err(Error::DeliveryPlanClaimConflict)
+ );
+ assert_eq!(any.attempt_count(), 0);
+ any.apply_sink_failure(
+ active.token(),
+ active.generation(),
+ active.row_revision(),
+ failure,
+ None,
+ 12,
+ )
+ .expect("partial satisfaction");
+ assert_eq!(any.state(), AuthoredDeliveryState::Satisfied);
+
+ let mut all = plan(5, TargetPolicy::all());
+ let all_claim = claim(&all, 5, 11);
+ all.claim(all_claim.clone(), 11).expect("claim");
+ let terminal = SinkFailure::for_request(
+ all.request(),
+ "terminal_batch_failure",
+ Retryability::Terminal,
+ None,
+ None,
+ Vec::new(),
+ )
+ .expect("terminal failure");
+ all.apply_sink_failure(
+ all_claim.token(),
+ all_claim.generation(),
+ all_claim.row_revision(),
+ terminal,
+ None,
+ 12,
+ )
+ .expect("terminal apply");
+ assert_eq!(all.state(), AuthoredDeliveryState::FailedTerminal);
+}
+
+#[test]
+fn attempt_limit_is_checked_without_mutating_claimed_state() {
+ let base = plan(6, TargetPolicy::all());
+ let pending_receipt = receipt(
+ base.request(),
+ vec![
+ DeliveryOutcome::unavailable(),
+ DeliveryOutcome::unavailable(),
+ ],
+ );
+ let attempts: Vec<_> = (1..=DELIVERY_PLAN_ATTEMPTS_MAX)
+ .map(|attempt| {
+ AuthoredDeliveryAttempt::reconstruct(
+ NonZeroU32::new(attempt).expect("attempt"),
+ 12,
+ DeliveryAttemptOutcome::Receipt(pending_receipt.clone()),
+ SatisfactionState::Pending,
+ )
+ .expect("attempt record")
+ })
+ .collect();
+ let mut value = serde_json::to_value(&base).expect("plan value");
+ value["state"] = serde_json::json!("retryable");
+ value["attempts"] = serde_json::to_value(attempts).expect("attempts value");
+ value["attempt_count"] = serde_json::json!(DELIVERY_PLAN_ATTEMPTS_MAX);
+ let schedule = retry("delivery_pending", DELIVERY_PLAN_ATTEMPTS_MAX, 20);
+ value["retry"] = serde_json::to_value(&schedule).expect("retry value");
+ value["last_failure"] = serde_json::to_value(schedule.failure()).expect("failure value");
+ value["updated_at_unix_ms"] = serde_json::json!(12);
+ let mut saturated: AuthoredDeliveryPlan =
+ serde_json::from_value(value).expect("saturated plan");
+ let active = claim(&saturated, 6, 20);
+ saturated.claim(active.clone(), 20).expect("claim");
+ let before = saturated.clone();
+ assert_eq!(
+ saturated.apply_receipt(
+ active.token(),
+ active.generation(),
+ active.row_revision(),
+ pending_receipt,
+ Some(retry("delivery_pending", DELIVERY_PLAN_ATTEMPTS_MAX, 30,)),
+ 21,
+ ),
+ Err(Error::DeliveryAttemptOverflow)
+ );
+ assert_eq!(saturated, before);
+}