commit a956577b4dbd161b4a1b9f66febf97d07eabee69
parent 0e6cc7a6871e7fd7ca06291905ebd969a370ea9c
Author: triesap <tyson@radroots.org>
Date: Wed, 5 Aug 2026 02:15:43 +0000
refactor(storage): atomically persist authored work
- add deterministic authored workflow commands and receipts
- persist preparation, claims, results, failures, and cancellation atomically
- keep memory backend replays conflict-safe without partial mutations
- bind delivery intent to exact signed event bytes before dispatch
Diffstat:
6 files changed, 1522 insertions(+), 32 deletions(-)
diff --git a/crates/storage/src/authored_atomic.rs b/crates/storage/src/authored_atomic.rs
@@ -0,0 +1,617 @@
+//! Atomic authored-operation commands with deterministic phase identities.
+
+use core::num::NonZeroU64;
+use radroots_event::SignedEvent;
+use radroots_transport::BoxFuture;
+use sha2::{Digest, Sha256};
+use std::{collections::BTreeSet, vec::Vec};
+
+use crate::{
+ Error,
+ atomic::{AtomicCommitDigest, AtomicCommitDisposition, AtomicCommitId},
+ authored::{
+ AdmissionState, AuthoredArtifact, AuthoredArtifactId, AuthoredOperation, RetrySchedule,
+ WorkFailure, WorkPhase,
+ },
+ authored_delivery::{AuthoredDeliveryPlan, AuthoredDeliveryPlanId, DeliveryAttemptOutcome},
+ journal::OperationInstanceId,
+};
+
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct WorkFence {
+ token: [u8; 16],
+ generation: NonZeroU64,
+ row_revision: NonZeroU64,
+}
+
+impl WorkFence {
+ pub const fn new(
+ token: [u8; 16],
+ generation: NonZeroU64,
+ row_revision: NonZeroU64,
+ ) -> Result<Self, Error> {
+ if bytes_are_zero(&token) {
+ return Err(Error::InvalidWorkClaim);
+ }
+ Ok(Self {
+ token,
+ generation,
+ row_revision,
+ })
+ }
+ pub const fn token(&self) -> &[u8; 16] {
+ &self.token
+ }
+ pub const fn generation(&self) -> NonZeroU64 {
+ self.generation
+ }
+ pub const fn row_revision(&self) -> NonZeroU64 {
+ self.row_revision
+ }
+}
+
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct PrepareAuthoredOperation {
+ operation: AuthoredOperation,
+ artifacts: Vec<AuthoredArtifact>,
+ delivery_plans: Vec<AuthoredDeliveryPlan>,
+ input_digest: AtomicCommitDigest,
+ requested_at_unix_ms: u64,
+}
+
+impl PrepareAuthoredOperation {
+ pub fn new(
+ operation: AuthoredOperation,
+ artifacts: Vec<AuthoredArtifact>,
+ delivery_plans: Vec<AuthoredDeliveryPlan>,
+ input_digest: AtomicCommitDigest,
+ requested_at_unix_ms: u64,
+ ) -> Result<Self, Error> {
+ if requested_at_unix_ms == 0
+ || artifacts.len() != operation.artifact_ids().len()
+ || artifacts.iter().enumerate().any(|(ordinal, artifact)| {
+ artifact.operation_id() != operation.operation_id()
+ || artifact.artifact_id() != operation.artifact_ids()[ordinal]
+ || usize::from(artifact.ordinal()) != ordinal
+ || artifact.validate().is_err()
+ })
+ {
+ return Err(Error::AtomicWorkflowMismatch);
+ }
+ let artifact_ids: BTreeSet<_> = artifacts
+ .iter()
+ .map(AuthoredArtifact::artifact_id)
+ .collect();
+ let plan_ids: BTreeSet<_> = delivery_plans
+ .iter()
+ .map(AuthoredDeliveryPlan::plan_id)
+ .collect();
+ if plan_ids.len() != delivery_plans.len()
+ || delivery_plans
+ .iter()
+ .any(|plan| !artifact_ids.contains(&plan.artifact_id()) || plan.validate().is_err())
+ {
+ return Err(Error::AtomicWorkflowMismatch);
+ }
+ Ok(Self {
+ operation,
+ artifacts,
+ delivery_plans,
+ input_digest,
+ requested_at_unix_ms,
+ })
+ }
+ pub const fn operation(&self) -> &AuthoredOperation {
+ &self.operation
+ }
+ pub fn artifacts(&self) -> &[AuthoredArtifact] {
+ self.artifacts.as_slice()
+ }
+ pub fn delivery_plans(&self) -> &[AuthoredDeliveryPlan] {
+ self.delivery_plans.as_slice()
+ }
+ pub const fn input_digest(&self) -> AtomicCommitDigest {
+ self.input_digest
+ }
+ pub const fn requested_at_unix_ms(&self) -> u64 {
+ self.requested_at_unix_ms
+ }
+}
+
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct ApplySignedArtifact {
+ artifact_id: AuthoredArtifactId,
+ fence: WorkFence,
+ event: SignedEvent,
+ applied_at_unix_ms: u64,
+}
+
+impl ApplySignedArtifact {
+ pub fn new(
+ artifact_id: AuthoredArtifactId,
+ fence: WorkFence,
+ event: SignedEvent,
+ applied_at_unix_ms: u64,
+ ) -> Result<Self, Error> {
+ if applied_at_unix_ms == 0 {
+ return Err(Error::AtomicWorkflowMismatch);
+ }
+ Ok(Self {
+ artifact_id,
+ fence,
+ event,
+ applied_at_unix_ms,
+ })
+ }
+ pub const fn artifact_id(&self) -> AuthoredArtifactId {
+ self.artifact_id
+ }
+ pub const fn fence(&self) -> &WorkFence {
+ &self.fence
+ }
+ pub const fn event(&self) -> &SignedEvent {
+ &self.event
+ }
+ pub const fn applied_at_unix_ms(&self) -> u64 {
+ self.applied_at_unix_ms
+ }
+}
+
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct ApplyAdmissionResult {
+ artifact_id: AuthoredArtifactId,
+ fence: WorkFence,
+ state: AdmissionState,
+ failure: Option<WorkFailure>,
+ retry: Option<RetrySchedule>,
+ applied_at_unix_ms: u64,
+}
+
+impl ApplyAdmissionResult {
+ pub fn new(
+ artifact_id: AuthoredArtifactId,
+ fence: WorkFence,
+ state: AdmissionState,
+ failure: Option<WorkFailure>,
+ retry: Option<RetrySchedule>,
+ applied_at_unix_ms: u64,
+ ) -> Result<Self, Error> {
+ if applied_at_unix_ms == 0 {
+ return Err(Error::AtomicWorkflowMismatch);
+ }
+ Ok(Self {
+ artifact_id,
+ fence,
+ state,
+ failure,
+ retry,
+ applied_at_unix_ms,
+ })
+ }
+ pub const fn artifact_id(&self) -> AuthoredArtifactId {
+ self.artifact_id
+ }
+ pub const fn fence(&self) -> &WorkFence {
+ &self.fence
+ }
+ pub const fn state(&self) -> AdmissionState {
+ self.state
+ }
+ pub const fn failure(&self) -> Option<&WorkFailure> {
+ self.failure.as_ref()
+ }
+ pub const fn retry(&self) -> Option<&RetrySchedule> {
+ self.retry.as_ref()
+ }
+ pub const fn applied_at_unix_ms(&self) -> u64 {
+ self.applied_at_unix_ms
+ }
+}
+
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct ApplyDeliveryAttempt {
+ plan_id: AuthoredDeliveryPlanId,
+ fence: WorkFence,
+ outcome: DeliveryAttemptOutcome,
+ retry: Option<RetrySchedule>,
+ applied_at_unix_ms: u64,
+}
+
+impl ApplyDeliveryAttempt {
+ pub fn new(
+ plan_id: AuthoredDeliveryPlanId,
+ fence: WorkFence,
+ outcome: DeliveryAttemptOutcome,
+ retry: Option<RetrySchedule>,
+ applied_at_unix_ms: u64,
+ ) -> Result<Self, Error> {
+ if applied_at_unix_ms == 0 {
+ return Err(Error::AtomicWorkflowMismatch);
+ }
+ Ok(Self {
+ plan_id,
+ fence,
+ outcome,
+ retry,
+ applied_at_unix_ms,
+ })
+ }
+ pub const fn plan_id(&self) -> AuthoredDeliveryPlanId {
+ self.plan_id
+ }
+ pub const fn fence(&self) -> &WorkFence {
+ &self.fence
+ }
+ pub const fn outcome(&self) -> &DeliveryAttemptOutcome {
+ &self.outcome
+ }
+ pub const fn retry(&self) -> Option<&RetrySchedule> {
+ self.retry.as_ref()
+ }
+ pub const fn applied_at_unix_ms(&self) -> u64 {
+ self.applied_at_unix_ms
+ }
+}
+
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub enum AuthoredWorkTarget {
+ Artifact(AuthoredArtifactId),
+ DeliveryPlan(AuthoredDeliveryPlanId),
+}
+
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub enum ClaimAuthoredTarget {
+ ArtifactSigning(AuthoredArtifactId),
+ ArtifactAdmission(AuthoredArtifactId),
+ DeliveryPlan(AuthoredDeliveryPlanId),
+}
+
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct ClaimAuthoredWork {
+ target: ClaimAuthoredTarget,
+ claim: crate::authored::WorkClaim,
+}
+
+impl ClaimAuthoredWork {
+ pub const fn new(target: ClaimAuthoredTarget, claim: crate::authored::WorkClaim) -> Self {
+ Self { target, claim }
+ }
+ pub const fn target(&self) -> &ClaimAuthoredTarget {
+ &self.target
+ }
+ pub const fn claim(&self) -> &crate::authored::WorkClaim {
+ &self.claim
+ }
+}
+
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct ApplyWorkFailure {
+ target: AuthoredWorkTarget,
+ fence: WorkFence,
+ failure: WorkFailure,
+ retry: Option<RetrySchedule>,
+ applied_at_unix_ms: u64,
+}
+
+impl ApplyWorkFailure {
+ pub fn new(
+ target: AuthoredWorkTarget,
+ fence: WorkFence,
+ failure: WorkFailure,
+ retry: Option<RetrySchedule>,
+ applied_at_unix_ms: u64,
+ ) -> Result<Self, Error> {
+ if applied_at_unix_ms == 0 {
+ return Err(Error::AtomicWorkflowMismatch);
+ }
+ Ok(Self {
+ target,
+ fence,
+ failure,
+ retry,
+ applied_at_unix_ms,
+ })
+ }
+ pub const fn target(&self) -> &AuthoredWorkTarget {
+ &self.target
+ }
+ pub const fn fence(&self) -> &WorkFence {
+ &self.fence
+ }
+ pub const fn failure(&self) -> &WorkFailure {
+ &self.failure
+ }
+ pub const fn retry(&self) -> Option<&RetrySchedule> {
+ self.retry.as_ref()
+ }
+ pub const fn applied_at_unix_ms(&self) -> u64 {
+ self.applied_at_unix_ms
+ }
+}
+
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub enum CancelAuthoredTarget {
+ ArtifactSigning(AuthoredArtifactId),
+ ArtifactAdmission(AuthoredArtifactId),
+ DeliveryPlan(AuthoredDeliveryPlanId),
+}
+
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct CancelAuthoredWork {
+ target: CancelAuthoredTarget,
+ expected_revision: NonZeroU64,
+ cancelled_at_unix_ms: u64,
+}
+
+impl CancelAuthoredWork {
+ pub const fn new(
+ target: CancelAuthoredTarget,
+ expected_revision: NonZeroU64,
+ cancelled_at_unix_ms: u64,
+ ) -> Result<Self, Error> {
+ if cancelled_at_unix_ms == 0 {
+ return Err(Error::AtomicWorkflowMismatch);
+ }
+ Ok(Self {
+ target,
+ expected_revision,
+ cancelled_at_unix_ms,
+ })
+ }
+ pub const fn target(&self) -> &CancelAuthoredTarget {
+ &self.target
+ }
+ pub const fn expected_revision(&self) -> NonZeroU64 {
+ self.expected_revision
+ }
+ pub const fn cancelled_at_unix_ms(&self) -> u64 {
+ self.cancelled_at_unix_ms
+ }
+}
+
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub enum AuthoredAtomicCommand {
+ Prepare(PrepareAuthoredOperation),
+ Claim(ClaimAuthoredWork),
+ ApplySigned(ApplySignedArtifact),
+ ApplyAdmission(ApplyAdmissionResult),
+ ApplyDelivery(ApplyDeliveryAttempt),
+ ApplyFailure(ApplyWorkFailure),
+ Cancel(CancelAuthoredWork),
+}
+
+impl AuthoredAtomicCommand {
+ pub fn commit_id(&self) -> AtomicCommitId {
+ let digest = self.digest();
+ let mut hasher = Sha256::new();
+ hash_field(&mut hasher, b"radroots.authored.atomic.id.v2");
+ hash_field(&mut hasher, self.phase_bytes());
+ hash_field(&mut hasher, &self.target_bytes());
+ if let Some(generation) = self.generation() {
+ hasher.update(generation.get().to_be_bytes());
+ }
+ hash_field(&mut hasher, digest.as_bytes());
+ let bytes: [u8; 32] = hasher.finalize().into();
+ let mut id = [0_u8; 16];
+ id.copy_from_slice(&bytes[..16]);
+ AtomicCommitId::new(id).expect("SHA-256 derived commit identity is nonzero")
+ }
+
+ pub fn digest(&self) -> AtomicCommitDigest {
+ let mut hasher = Sha256::new();
+ hash_field(&mut hasher, b"radroots.authored.atomic.digest.v2");
+ hash_field(&mut hasher, self.phase_bytes());
+ hash_field(&mut hasher, &self.target_bytes());
+ match self {
+ Self::Prepare(value) => hash_field(&mut hasher, value.input_digest.as_bytes()),
+ Self::Claim(value) => {
+ hash_field(&mut hasher, value.claim.token());
+ hasher.update(value.claim.generation().get().to_be_bytes());
+ hasher.update(value.claim.row_revision().get().to_be_bytes());
+ }
+ Self::ApplySigned(value) => hash_field(&mut hasher, value.event.raw_json().as_bytes()),
+ Self::ApplyAdmission(value) => {
+ hasher.update([value.state as u8]);
+ hash_failure(&mut hasher, value.failure.as_ref());
+ }
+ Self::ApplyDelivery(value) => hash_delivery(&mut hasher, &value.outcome),
+ Self::ApplyFailure(value) => hash_failure(&mut hasher, Some(&value.failure)),
+ Self::Cancel(value) => hasher.update(value.cancelled_at_unix_ms.to_be_bytes()),
+ }
+ AtomicCommitDigest::new(hasher.finalize().into())
+ }
+
+ pub const fn requested_at_unix_ms(&self) -> u64 {
+ match self {
+ Self::Prepare(value) => value.requested_at_unix_ms,
+ Self::Claim(value) => value.claim.acquired_at_unix_ms(),
+ Self::ApplySigned(value) => value.applied_at_unix_ms,
+ Self::ApplyAdmission(value) => value.applied_at_unix_ms,
+ Self::ApplyDelivery(value) => value.applied_at_unix_ms,
+ Self::ApplyFailure(value) => value.applied_at_unix_ms,
+ Self::Cancel(value) => value.cancelled_at_unix_ms,
+ }
+ }
+
+ fn phase_bytes(&self) -> &'static [u8] {
+ match self {
+ Self::Prepare(_) => b"prepare",
+ Self::Claim(_) => b"claim",
+ Self::ApplySigned(_) => b"signing",
+ Self::ApplyAdmission(_) => b"admission",
+ Self::ApplyDelivery(_) => b"delivery",
+ Self::ApplyFailure(value) => match value.failure.phase() {
+ WorkPhase::Signing => b"signing_failure",
+ WorkPhase::Admission => b"admission_failure",
+ WorkPhase::Delivery => b"delivery_failure",
+ },
+ Self::Cancel(_) => b"cancel",
+ }
+ }
+
+ fn target_bytes(&self) -> [u8; 16] {
+ match self {
+ Self::Prepare(value) => *value.operation.operation_id().as_bytes(),
+ Self::Claim(value) => match &value.target {
+ ClaimAuthoredTarget::ArtifactSigning(id)
+ | ClaimAuthoredTarget::ArtifactAdmission(id) => *id.as_bytes(),
+ ClaimAuthoredTarget::DeliveryPlan(id) => *id.as_bytes(),
+ },
+ Self::ApplySigned(value) => *value.artifact_id.as_bytes(),
+ Self::ApplyAdmission(value) => *value.artifact_id.as_bytes(),
+ Self::ApplyDelivery(value) => *value.plan_id.as_bytes(),
+ Self::ApplyFailure(value) => match &value.target {
+ AuthoredWorkTarget::Artifact(id) => *id.as_bytes(),
+ AuthoredWorkTarget::DeliveryPlan(id) => *id.as_bytes(),
+ },
+ Self::Cancel(value) => match &value.target {
+ CancelAuthoredTarget::ArtifactSigning(id)
+ | CancelAuthoredTarget::ArtifactAdmission(id) => *id.as_bytes(),
+ CancelAuthoredTarget::DeliveryPlan(id) => *id.as_bytes(),
+ },
+ }
+ }
+
+ fn generation(&self) -> Option<NonZeroU64> {
+ match self {
+ Self::ApplySigned(value) => Some(value.fence.generation),
+ Self::Claim(value) => Some(value.claim.generation()),
+ Self::ApplyAdmission(value) => Some(value.fence.generation),
+ Self::ApplyDelivery(value) => Some(value.fence.generation),
+ Self::ApplyFailure(value) => Some(value.fence.generation),
+ Self::Prepare(_) | Self::Cancel(_) => None,
+ }
+ }
+}
+
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub enum AuthoredAtomicOutcome {
+ Prepared {
+ operation: AuthoredOperation,
+ artifacts: Vec<AuthoredArtifact>,
+ delivery_plans: Vec<AuthoredDeliveryPlan>,
+ },
+ Artifact(AuthoredArtifact),
+ DeliveryPlan(AuthoredDeliveryPlan),
+}
+
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct AuthoredAtomicReceipt {
+ commit_id: AtomicCommitId,
+ digest: AtomicCommitDigest,
+ disposition: AtomicCommitDisposition,
+ committed_at_unix_ms: u64,
+ outcome: AuthoredAtomicOutcome,
+}
+
+impl AuthoredAtomicReceipt {
+ pub fn new(
+ command: &AuthoredAtomicCommand,
+ disposition: AtomicCommitDisposition,
+ committed_at_unix_ms: u64,
+ outcome: AuthoredAtomicOutcome,
+ ) -> Result<Self, Error> {
+ if committed_at_unix_ms < command.requested_at_unix_ms() {
+ return Err(Error::AtomicWorkflowMismatch);
+ }
+ Ok(Self {
+ commit_id: command.commit_id(),
+ digest: command.digest(),
+ disposition,
+ committed_at_unix_ms,
+ outcome,
+ })
+ }
+ pub const fn commit_id(&self) -> AtomicCommitId {
+ self.commit_id
+ }
+ pub const fn digest(&self) -> AtomicCommitDigest {
+ self.digest
+ }
+ pub const fn disposition(&self) -> AtomicCommitDisposition {
+ self.disposition
+ }
+ pub const fn committed_at_unix_ms(&self) -> u64 {
+ self.committed_at_unix_ms
+ }
+ pub const fn outcome(&self) -> &AuthoredAtomicOutcome {
+ &self.outcome
+ }
+}
+
+pub trait AuthoredAtomicStorage: Send + Sync {
+ fn execute_authored(
+ &self,
+ command: AuthoredAtomicCommand,
+ ) -> BoxFuture<'_, Result<AuthoredAtomicReceipt, Error>>;
+ fn authored_receipt(
+ &self,
+ commit_id: AtomicCommitId,
+ ) -> BoxFuture<'_, Result<Option<AuthoredAtomicReceipt>, Error>>;
+ fn authored_operation(
+ &self,
+ operation_id: OperationInstanceId,
+ ) -> BoxFuture<'_, Result<Option<AuthoredOperation>, Error>>;
+ fn authored_artifact(
+ &self,
+ artifact_id: AuthoredArtifactId,
+ ) -> BoxFuture<'_, Result<Option<AuthoredArtifact>, Error>>;
+ fn authored_delivery_plan(
+ &self,
+ plan_id: AuthoredDeliveryPlanId,
+ ) -> BoxFuture<'_, Result<Option<AuthoredDeliveryPlan>, Error>>;
+}
+
+fn hash_failure(hasher: &mut Sha256, failure: Option<&WorkFailure>) {
+ if let Some(failure) = failure {
+ hash_field(hasher, failure.code().as_bytes());
+ hasher.update([failure.phase() as u8, failure.class() as u8]);
+ hasher.update(
+ failure
+ .retry_after_unix_ms()
+ .unwrap_or_default()
+ .to_be_bytes(),
+ );
+ if let Some(diagnostic) = failure.diagnostic() {
+ hash_field(hasher, diagnostic.as_bytes());
+ }
+ }
+}
+
+fn hash_delivery(hasher: &mut Sha256, outcome: &DeliveryAttemptOutcome) {
+ let entries = match outcome {
+ DeliveryAttemptOutcome::Receipt(receipt) => receipt.target_receipts(),
+ DeliveryAttemptOutcome::SinkFailure(failure) => {
+ hash_field(hasher, failure.code().as_bytes());
+ hasher.update([failure.retryability() as u8]);
+ failure.partial_evidence()
+ }
+ };
+ for entry in entries {
+ hash_field(hasher, entry.target().fingerprint().as_str().as_bytes());
+ hasher.update([entry.was_attempted() as u8, entry.outcome().kind() as u8]);
+ hasher.update([entry.outcome().retryability() as u8]);
+ if let Some(code) = entry.outcome().code() {
+ hash_field(hasher, code.as_bytes());
+ }
+ if let Some(message) = entry.outcome().message() {
+ hash_field(hasher, message.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 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/src/authored_delivery.rs b/crates/storage/src/authored_delivery.rs
@@ -5,6 +5,8 @@ use radroots_transport::{
DeliveryReceipt, DeliveryRequest, SinkFailure,
outcome::Retryability,
policy::{SatisfactionClass, SatisfactionPolicy, SatisfactionState, evaluate_satisfaction},
+ sink::{DeliveryPayload, DeliveryRequestId},
+ target::TargetSet,
};
use sha2::{Digest, Sha256};
use std::vec::Vec;
@@ -50,6 +52,69 @@ impl From<AuthoredDeliveryPlanId> for [u8; 16] {
}
}
+/// Exact delivery intent persisted before a signed payload exists.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct AuthoredDeliveryIntent {
+ request_id: DeliveryRequestId,
+ target_set: TargetSet,
+ satisfaction: SatisfactionPolicy,
+ deadline_unix_ms: u64,
+}
+
+impl AuthoredDeliveryIntent {
+ pub fn new(
+ request_id: impl Into<String>,
+ target_set: TargetSet,
+ satisfaction: SatisfactionPolicy,
+ deadline_unix_ms: u64,
+ ) -> Result<Self, Error> {
+ if deadline_unix_ms == 0 || satisfaction.validate_for(&target_set).is_err() {
+ return Err(Error::InvalidAuthoredDeliveryPlan);
+ }
+ Ok(Self {
+ request_id: DeliveryRequestId::parse(request_id)
+ .map_err(|_| Error::InvalidAuthoredDeliveryPlan)?,
+ target_set,
+ satisfaction,
+ deadline_unix_ms,
+ })
+ }
+
+ pub fn from_request(request: &DeliveryRequest) -> Self {
+ Self {
+ request_id: request.request_id().clone(),
+ target_set: request.target_set().clone(),
+ satisfaction: request.satisfaction().clone(),
+ deadline_unix_ms: request.deadline_unix_ms(),
+ }
+ }
+
+ pub fn materialize(&self, payload: DeliveryPayload) -> Result<DeliveryRequest, Error> {
+ DeliveryRequest::new(
+ self.request_id.as_str(),
+ payload,
+ self.target_set.clone(),
+ self.satisfaction.clone(),
+ self.deadline_unix_ms,
+ )
+ .map_err(|_| Error::InvalidAuthoredDeliveryPlan)
+ }
+
+ pub const fn request_id(&self) -> &DeliveryRequestId {
+ &self.request_id
+ }
+ pub const fn target_set(&self) -> &TargetSet {
+ &self.target_set
+ }
+ pub const fn satisfaction(&self) -> &SatisfactionPolicy {
+ &self.satisfaction
+ }
+ pub const fn deadline_unix_ms(&self) -> u64 {
+ self.deadline_unix_ms
+ }
+}
+
#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
@@ -146,7 +211,8 @@ pub struct AuthoredDeliveryPlan {
plan_id: AuthoredDeliveryPlanId,
artifact_id: AuthoredArtifactId,
request_digest: [u8; 32],
- request: DeliveryRequest,
+ intent: AuthoredDeliveryIntent,
+ request: Option<DeliveryRequest>,
state: AuthoredDeliveryState,
attempts: Vec<AuthoredDeliveryAttempt>,
attempt_count: u32,
@@ -162,15 +228,16 @@ impl AuthoredDeliveryPlan {
pub fn new(
plan_id: AuthoredDeliveryPlanId,
artifact_id: AuthoredArtifactId,
- request: DeliveryRequest,
+ intent: AuthoredDeliveryIntent,
created_at_unix_ms: u64,
) -> Result<Self, Error> {
- let request_digest = delivery_request_digest(&request);
+ let request_digest = delivery_intent_digest(&intent);
Self::reconstruct(Self {
plan_id,
artifact_id,
request_digest,
- request,
+ intent,
+ request: None,
state: AuthoredDeliveryState::Pending,
attempts: Vec::new(),
attempt_count: 0,
@@ -183,6 +250,19 @@ impl AuthoredDeliveryPlan {
})
}
+ pub fn new_bound(
+ plan_id: AuthoredDeliveryPlanId,
+ artifact_id: AuthoredArtifactId,
+ request: DeliveryRequest,
+ created_at_unix_ms: u64,
+ ) -> Result<Self, Error> {
+ let intent = AuthoredDeliveryIntent::from_request(&request);
+ let mut plan = Self::new(plan_id, artifact_id, intent, created_at_unix_ms)?;
+ plan.request = Some(request);
+ plan.validate()?;
+ Ok(plan)
+ }
+
pub fn reconstruct(value: Self) -> Result<Self, Error> {
value.validate()?;
Ok(value)
@@ -191,7 +271,7 @@ impl AuthoredDeliveryPlan {
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.request_digest != delivery_intent_digest(&self.intent)
|| 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())
@@ -199,10 +279,21 @@ impl AuthoredDeliveryPlan {
{
return Err(Error::InvalidAuthoredDeliveryPlan);
}
+ if self
+ .request
+ .as_ref()
+ .is_some_and(|request| AuthoredDeliveryIntent::from_request(request) != self.intent)
+ || (!self.attempts.is_empty() && self.request.is_none())
+ {
+ 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
+ .request
+ .as_ref()
+ .is_none_or(|request| attempt.outcome.validate_for(request).is_err())
|| self
.evaluate_outcomes(self.attempts[..=index].iter().map(|value| &value.outcome))
.ok()
@@ -270,6 +361,7 @@ impl AuthoredDeliveryPlan {
pub fn claim(&mut self, claim: WorkClaim, now_unix_ms: u64) -> Result<(), Error> {
if self.state.is_terminal()
+ || self.request.is_none()
|| self.claim.is_some()
|| claim.row_revision() != self.revision
|| claim.acquired_at_unix_ms() != now_unix_ms
@@ -294,6 +386,27 @@ impl AuthoredDeliveryPlan {
Ok(())
}
+ pub fn bind_signed_event(
+ &mut self,
+ event: radroots_event::SignedEvent,
+ bound_at_unix_ms: u64,
+ ) -> Result<(), Error> {
+ if self.request.is_some() || self.attempt_count != 0 || self.state.is_terminal() {
+ return Err(Error::InvalidAuthoredDeliveryPlan);
+ }
+ let previous = self.clone();
+ self.request = Some(self.intent.materialize(DeliveryPayload::new(event))?);
+ if let Err(error) = self.advance(bound_at_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],
@@ -304,8 +417,12 @@ impl AuthoredDeliveryPlan {
recorded_at_unix_ms: u64,
) -> Result<(), Error> {
self.require_claim(token, generation, claim_revision, recorded_at_unix_ms)?;
+ let request = self
+ .request
+ .as_ref()
+ .ok_or(Error::InvalidAuthoredDeliveryPlan)?;
receipt
- .validate_for_request(&self.request)
+ .validate_for_request(request)
.map_err(|_| Error::InvalidAuthoredDeliveryPlan)?;
let satisfaction = self.evaluate_with(DeliveryAttemptOutcome::Receipt(receipt.clone()))?;
let (state, last_failure) = match satisfaction {
@@ -349,8 +466,12 @@ impl AuthoredDeliveryPlan {
recorded_at_unix_ms: u64,
) -> Result<(), Error> {
self.require_claim(token, generation, claim_revision, recorded_at_unix_ms)?;
+ let request = self
+ .request
+ .as_ref()
+ .ok_or(Error::InvalidAuthoredDeliveryPlan)?;
failure
- .validate_for_request(&self.request)
+ .validate_for_request(request)
.map_err(|_| Error::InvalidAuthoredDeliveryPlan)?;
let outcome = DeliveryAttemptOutcome::SinkFailure(failure.clone());
let satisfaction = self.evaluate_with(outcome.clone())?;
@@ -489,8 +610,14 @@ impl AuthoredDeliveryPlan {
}
}
evaluate_satisfaction(
- self.request.satisfaction(),
- self.request.target_set(),
+ self.request
+ .as_ref()
+ .ok_or(Error::InvalidAuthoredDeliveryPlan)?
+ .satisfaction(),
+ self.request
+ .as_ref()
+ .ok_or(Error::InvalidAuthoredDeliveryPlan)?
+ .target_set(),
evidence,
)
.map_err(|_| Error::InvalidAuthoredDeliveryPlan)
@@ -534,8 +661,11 @@ impl AuthoredDeliveryPlan {
pub const fn request_digest(&self) -> &[u8; 32] {
&self.request_digest
}
- pub const fn request(&self) -> &DeliveryRequest {
- &self.request
+ pub const fn intent(&self) -> &AuthoredDeliveryIntent {
+ &self.intent
+ }
+ pub const fn request(&self) -> Option<&DeliveryRequest> {
+ self.request.as_ref()
}
pub const fn state(&self) -> AuthoredDeliveryState {
self.state
@@ -566,7 +696,8 @@ struct AuthoredDeliveryPlanWire {
plan_id: AuthoredDeliveryPlanId,
artifact_id: AuthoredArtifactId,
request_digest: [u8; 32],
- request: DeliveryRequest,
+ intent: AuthoredDeliveryIntent,
+ request: Option<DeliveryRequest>,
state: AuthoredDeliveryState,
attempts: Vec<AuthoredDeliveryAttempt>,
attempt_count: u32,
@@ -586,6 +717,7 @@ impl TryFrom<AuthoredDeliveryPlanWire> for AuthoredDeliveryPlan {
plan_id: value.plan_id,
artifact_id: value.artifact_id,
request_digest: value.request_digest,
+ intent: value.intent,
request: value.request,
state: value.state,
attempts: value.attempts,
@@ -607,6 +739,7 @@ impl From<AuthoredDeliveryPlan> for AuthoredDeliveryPlanWire {
plan_id: value.plan_id,
artifact_id: value.artifact_id,
request_digest: value.request_digest,
+ intent: value.intent,
request: value.request,
state: value.state,
attempts: value.attempts,
@@ -621,21 +754,20 @@ impl From<AuthoredDeliveryPlan> for AuthoredDeliveryPlanWire {
}
}
-fn delivery_request_digest(request: &DeliveryRequest) -> [u8; 32] {
+fn delivery_intent_digest(intent: &AuthoredDeliveryIntent) -> [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, intent.request_id().as_str().as_bytes());
+ for target in intent.target_set().targets() {
hash_field(&mut hasher, target.fingerprint().as_str().as_bytes());
}
- let policy = request.satisfaction();
+ let policy = intent.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.update(intent.deadline_unix_ms().to_be_bytes());
hasher.finalize().into()
}
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_atomic;
pub mod authored_delivery;
pub mod backup;
mod error;
@@ -32,6 +33,7 @@ pub trait Storage:
+ private_artifact::PrivateArtifactStore
+ backup::StorageReliability
+ atomic::AtomicStorage
+ + authored_atomic::AuthoredAtomicStorage
{
}
@@ -43,5 +45,6 @@ impl<T> Storage for T where
+ private_artifact::PrivateArtifactStore
+ backup::StorageReliability
+ atomic::AtomicStorage
+ + authored_atomic::AuthoredAtomicStorage
{
}
diff --git a/crates/storage/src/memory.rs b/crates/storage/src/memory.rs
@@ -11,6 +11,12 @@ use crate::{
AtomicCommit, AtomicCommitDisposition, AtomicCommitId, AtomicCommitOutcome,
AtomicCommitReceipt, AtomicStorage, AtomicWorkflow,
},
+ authored::{AdmissionState, FailureClass, WorkFailure, WorkPhase},
+ authored_atomic::{
+ AuthoredAtomicCommand, AuthoredAtomicOutcome, AuthoredAtomicReceipt, AuthoredAtomicStorage,
+ AuthoredWorkTarget, CancelAuthoredTarget, ClaimAuthoredTarget,
+ },
+ authored_delivery::DeliveryAttemptOutcome,
backup::{
BackupId, BackupOperation, BackupPlan, BackupTransition, ReliabilityRevision,
RestoreOperation, RestorePlan, RestoreTransition, StorageReliability,
@@ -67,6 +73,10 @@ struct State {
backups: Vec<BackupOperation>,
restores: Vec<RestoreOperation>,
atomic_receipts: Vec<AtomicCommitReceipt>,
+ authored_operations: Vec<crate::authored::AuthoredOperation>,
+ authored_artifacts: Vec<crate::authored::AuthoredArtifact>,
+ authored_delivery_plans: Vec<crate::authored_delivery::AuthoredDeliveryPlan>,
+ authored_atomic_receipts: Vec<AuthoredAtomicReceipt>,
closed: bool,
}
@@ -93,6 +103,10 @@ impl MemoryStorage {
backups: Vec::new(),
restores: Vec::new(),
atomic_receipts: Vec::new(),
+ authored_operations: Vec::new(),
+ authored_artifacts: Vec::new(),
+ authored_delivery_plans: Vec::new(),
+ authored_atomic_receipts: Vec::new(),
closed: false,
}),
}
@@ -1293,3 +1307,385 @@ impl AtomicStorage for MemoryStorage {
})
}
}
+
+impl AuthoredAtomicStorage for MemoryStorage {
+ fn execute_authored(
+ &self,
+ command: AuthoredAtomicCommand,
+ ) -> BoxFuture<'_, Result<AuthoredAtomicReceipt, Error>> {
+ Box::pin(async move {
+ let mut state = self.state()?;
+ if let Some(existing) = state
+ .authored_atomic_receipts
+ .iter()
+ .find(|receipt| receipt.commit_id() == command.commit_id())
+ {
+ if existing.digest() != command.digest() {
+ return Err(Error::AtomicCommitConflict);
+ }
+ return AuthoredAtomicReceipt::new(
+ &command,
+ AtomicCommitDisposition::Replay,
+ existing.committed_at_unix_ms(),
+ existing.outcome().clone(),
+ );
+ }
+
+ let mut candidate = state.clone();
+ let outcome = match command.clone() {
+ AuthoredAtomicCommand::Prepare(value) => {
+ if candidate.authored_operations.iter().any(|operation| {
+ operation.operation_id() == value.operation().operation_id()
+ }) || value.artifacts().iter().any(|artifact| {
+ candidate
+ .authored_artifacts
+ .iter()
+ .any(|existing| existing.artifact_id() == artifact.artifact_id())
+ }) || value.delivery_plans().iter().any(|plan| {
+ candidate
+ .authored_delivery_plans
+ .iter()
+ .any(|existing| existing.plan_id() == plan.plan_id())
+ }) {
+ return Err(Error::AtomicCommitConflict);
+ }
+ candidate
+ .authored_operations
+ .push(value.operation().clone());
+ candidate
+ .authored_artifacts
+ .extend(value.artifacts().iter().cloned());
+ candidate
+ .authored_delivery_plans
+ .extend(value.delivery_plans().iter().cloned());
+ AuthoredAtomicOutcome::Prepared {
+ operation: value.operation().clone(),
+ artifacts: value.artifacts().to_vec(),
+ delivery_plans: value.delivery_plans().to_vec(),
+ }
+ }
+ AuthoredAtomicCommand::Claim(value) => match value.target() {
+ ClaimAuthoredTarget::ArtifactSigning(artifact_id) => {
+ let artifact = candidate
+ .authored_artifacts
+ .iter_mut()
+ .find(|artifact| artifact.artifact_id() == *artifact_id)
+ .ok_or(Error::InvalidAuthoredArtifact)?;
+ artifact.set_signing_claim(
+ value.claim().clone(),
+ value.claim().acquired_at_unix_ms(),
+ )?;
+ AuthoredAtomicOutcome::Artifact(artifact.clone())
+ }
+ ClaimAuthoredTarget::ArtifactAdmission(artifact_id) => {
+ let artifact = candidate
+ .authored_artifacts
+ .iter_mut()
+ .find(|artifact| artifact.artifact_id() == *artifact_id)
+ .ok_or(Error::InvalidAuthoredArtifact)?;
+ artifact.set_admission_claim(
+ value.claim().clone(),
+ value.claim().acquired_at_unix_ms(),
+ )?;
+ AuthoredAtomicOutcome::Artifact(artifact.clone())
+ }
+ ClaimAuthoredTarget::DeliveryPlan(plan_id) => {
+ let plan = candidate
+ .authored_delivery_plans
+ .iter_mut()
+ .find(|plan| plan.plan_id() == *plan_id)
+ .ok_or(Error::InvalidAuthoredDeliveryPlan)?;
+ plan.claim(value.claim().clone(), value.claim().acquired_at_unix_ms())?;
+ AuthoredAtomicOutcome::DeliveryPlan(plan.clone())
+ }
+ },
+ AuthoredAtomicCommand::ApplySigned(value) => {
+ let artifact = candidate
+ .authored_artifacts
+ .iter_mut()
+ .find(|artifact| artifact.artifact_id() == value.artifact_id())
+ .ok_or(Error::InvalidAuthoredArtifact)?;
+ require_artifact_claim(
+ artifact.signing_claim(),
+ value.fence(),
+ value.applied_at_unix_ms(),
+ )?;
+ artifact.record_signed(value.event().clone(), value.applied_at_unix_ms())?;
+ let artifact = artifact.clone();
+ for plan in candidate
+ .authored_delivery_plans
+ .iter_mut()
+ .filter(|plan| plan.artifact_id() == value.artifact_id())
+ {
+ plan.bind_signed_event(value.event().clone(), value.applied_at_unix_ms())?;
+ }
+ AuthoredAtomicOutcome::Artifact(artifact)
+ }
+ AuthoredAtomicCommand::ApplyAdmission(value) => {
+ let artifact = candidate
+ .authored_artifacts
+ .iter_mut()
+ .find(|artifact| artifact.artifact_id() == value.artifact_id())
+ .ok_or(Error::InvalidAuthoredArtifact)?;
+ require_artifact_claim(
+ artifact.admission_claim(),
+ value.fence(),
+ value.applied_at_unix_ms(),
+ )?;
+ artifact.record_admission(
+ value.state(),
+ value.failure().cloned(),
+ value.retry().cloned(),
+ value.applied_at_unix_ms(),
+ )?;
+ AuthoredAtomicOutcome::Artifact(artifact.clone())
+ }
+ AuthoredAtomicCommand::ApplyDelivery(value) => {
+ let plan = candidate
+ .authored_delivery_plans
+ .iter_mut()
+ .find(|plan| plan.plan_id() == value.plan_id())
+ .ok_or(Error::InvalidAuthoredDeliveryPlan)?;
+ match value.outcome().clone() {
+ DeliveryAttemptOutcome::Receipt(receipt) => plan.apply_receipt(
+ value.fence().token(),
+ value.fence().generation(),
+ value.fence().row_revision(),
+ receipt,
+ value.retry().cloned(),
+ value.applied_at_unix_ms(),
+ )?,
+ DeliveryAttemptOutcome::SinkFailure(failure) => plan.apply_sink_failure(
+ value.fence().token(),
+ value.fence().generation(),
+ value.fence().row_revision(),
+ failure,
+ value.retry().cloned(),
+ value.applied_at_unix_ms(),
+ )?,
+ }
+ AuthoredAtomicOutcome::DeliveryPlan(plan.clone())
+ }
+ AuthoredAtomicCommand::ApplyFailure(value) => match value.target() {
+ AuthoredWorkTarget::Artifact(artifact_id) => {
+ let artifact = candidate
+ .authored_artifacts
+ .iter_mut()
+ .find(|artifact| artifact.artifact_id() == *artifact_id)
+ .ok_or(Error::InvalidAuthoredArtifact)?;
+ match value.failure().phase() {
+ WorkPhase::Signing => {
+ require_artifact_claim(
+ artifact.signing_claim(),
+ value.fence(),
+ value.applied_at_unix_ms(),
+ )?;
+ artifact.record_signing_failure(
+ value.failure().clone(),
+ value.retry().cloned(),
+ value.applied_at_unix_ms(),
+ )?;
+ }
+ WorkPhase::Admission => {
+ require_artifact_claim(
+ artifact.admission_claim(),
+ value.fence(),
+ value.applied_at_unix_ms(),
+ )?;
+ let state = match value.failure().class() {
+ FailureClass::Retryable => AdmissionState::Retryable,
+ FailureClass::Terminal => AdmissionState::Rejected,
+ FailureClass::Indeterminate => {
+ return Err(Error::InvalidAuthoredTransition);
+ }
+ };
+ artifact.record_admission(
+ state,
+ Some(value.failure().clone()),
+ value.retry().cloned(),
+ value.applied_at_unix_ms(),
+ )?;
+ }
+ WorkPhase::Delivery => {
+ return Err(Error::AtomicWorkflowMismatch);
+ }
+ }
+ AuthoredAtomicOutcome::Artifact(artifact.clone())
+ }
+ AuthoredWorkTarget::DeliveryPlan(plan_id) => {
+ let plan = candidate
+ .authored_delivery_plans
+ .iter_mut()
+ .find(|plan| plan.plan_id() == *plan_id)
+ .ok_or(Error::InvalidAuthoredDeliveryPlan)?;
+ if value.failure().phase() != WorkPhase::Delivery
+ || value.failure().class() == FailureClass::Indeterminate
+ {
+ return Err(Error::AtomicWorkflowMismatch);
+ }
+ let retryability = match value.failure().class() {
+ FailureClass::Retryable => {
+ radroots_transport::outcome::Retryability::Retryable
+ }
+ FailureClass::Terminal => {
+ radroots_transport::outcome::Retryability::Terminal
+ }
+ FailureClass::Indeterminate => unreachable!(),
+ };
+ let failure = radroots_transport::SinkFailure::for_request(
+ plan.request().ok_or(Error::InvalidAuthoredDeliveryPlan)?,
+ value.failure().code(),
+ retryability,
+ value.failure().retry_after_unix_ms(),
+ value.failure().diagnostic().map(str::to_owned),
+ Vec::new(),
+ )
+ .map_err(|_| Error::AtomicWorkflowMismatch)?;
+ plan.apply_sink_failure(
+ value.fence().token(),
+ value.fence().generation(),
+ value.fence().row_revision(),
+ failure,
+ value.retry().cloned(),
+ value.applied_at_unix_ms(),
+ )?;
+ AuthoredAtomicOutcome::DeliveryPlan(plan.clone())
+ }
+ },
+ AuthoredAtomicCommand::Cancel(value) => match value.target() {
+ CancelAuthoredTarget::ArtifactSigning(artifact_id) => {
+ let artifact = candidate
+ .authored_artifacts
+ .iter_mut()
+ .find(|artifact| artifact.artifact_id() == *artifact_id)
+ .ok_or(Error::InvalidAuthoredArtifact)?;
+ if artifact.revision() != value.expected_revision() {
+ return Err(Error::InvalidAuthoredTransition);
+ }
+ artifact.cancel_signing(value.cancelled_at_unix_ms())?;
+ AuthoredAtomicOutcome::Artifact(artifact.clone())
+ }
+ CancelAuthoredTarget::ArtifactAdmission(artifact_id) => {
+ let artifact = candidate
+ .authored_artifacts
+ .iter_mut()
+ .find(|artifact| artifact.artifact_id() == *artifact_id)
+ .ok_or(Error::InvalidAuthoredArtifact)?;
+ if artifact.revision() != value.expected_revision() {
+ return Err(Error::InvalidAuthoredTransition);
+ }
+ let failure = WorkFailure::new(
+ "cancelled",
+ WorkPhase::Admission,
+ FailureClass::Terminal,
+ None,
+ None,
+ )?;
+ artifact.record_admission(
+ AdmissionState::Cancelled,
+ Some(failure),
+ None,
+ value.cancelled_at_unix_ms(),
+ )?;
+ AuthoredAtomicOutcome::Artifact(artifact.clone())
+ }
+ CancelAuthoredTarget::DeliveryPlan(plan_id) => {
+ let plan = candidate
+ .authored_delivery_plans
+ .iter_mut()
+ .find(|plan| plan.plan_id() == *plan_id)
+ .ok_or(Error::InvalidAuthoredDeliveryPlan)?;
+ if plan.revision() != value.expected_revision() {
+ return Err(Error::InvalidAuthoredDeliveryPlan);
+ }
+ plan.cancel(value.cancelled_at_unix_ms())?;
+ AuthoredAtomicOutcome::DeliveryPlan(plan.clone())
+ }
+ },
+ };
+ let receipt = AuthoredAtomicReceipt::new(
+ &command,
+ AtomicCommitDisposition::Committed,
+ command.requested_at_unix_ms(),
+ outcome,
+ )?;
+ candidate.authored_atomic_receipts.push(receipt.clone());
+ *state = candidate;
+ Ok(receipt)
+ })
+ }
+
+ fn authored_receipt(
+ &self,
+ commit_id: AtomicCommitId,
+ ) -> BoxFuture<'_, Result<Option<AuthoredAtomicReceipt>, Error>> {
+ Box::pin(async move {
+ Ok(self
+ .state()?
+ .authored_atomic_receipts
+ .iter()
+ .find(|receipt| receipt.commit_id() == commit_id)
+ .cloned())
+ })
+ }
+
+ fn authored_operation(
+ &self,
+ operation_id: OperationInstanceId,
+ ) -> BoxFuture<'_, Result<Option<crate::authored::AuthoredOperation>, Error>> {
+ Box::pin(async move {
+ Ok(self
+ .state()?
+ .authored_operations
+ .iter()
+ .find(|operation| operation.operation_id() == operation_id)
+ .cloned())
+ })
+ }
+
+ fn authored_artifact(
+ &self,
+ artifact_id: crate::authored::AuthoredArtifactId,
+ ) -> BoxFuture<'_, Result<Option<crate::authored::AuthoredArtifact>, Error>> {
+ Box::pin(async move {
+ Ok(self
+ .state()?
+ .authored_artifacts
+ .iter()
+ .find(|artifact| artifact.artifact_id() == artifact_id)
+ .cloned())
+ })
+ }
+
+ fn authored_delivery_plan(
+ &self,
+ plan_id: crate::authored_delivery::AuthoredDeliveryPlanId,
+ ) -> BoxFuture<'_, Result<Option<crate::authored_delivery::AuthoredDeliveryPlan>, Error>> {
+ Box::pin(async move {
+ Ok(self
+ .state()?
+ .authored_delivery_plans
+ .iter()
+ .find(|plan| plan.plan_id() == plan_id)
+ .cloned())
+ })
+ }
+}
+
+fn require_artifact_claim(
+ claim: Option<&crate::authored::WorkClaim>,
+ fence: &crate::authored_atomic::WorkFence,
+ now_unix_ms: u64,
+) -> Result<(), Error> {
+ if !claim.is_some_and(|claim| {
+ claim.matches_fence(
+ fence.token(),
+ fence.generation(),
+ fence.row_revision(),
+ now_unix_ms,
+ )
+ }) {
+ return Err(Error::DeliveryPlanClaimConflict);
+ }
+ Ok(())
+}
diff --git a/crates/storage/tests/authored_atomic.rs b/crates/storage/tests/authored_atomic.rs
@@ -0,0 +1,338 @@
+use core::num::{NonZeroU32, NonZeroU64};
+use futures_executor::block_on;
+use radroots_event::{GenericEventDraft, SignedEvent, wire::v1::Nip01EventWire};
+use radroots_event_codec::authoring::AuthoredEventPlan;
+use radroots_storage::{
+ Error,
+ atomic::{AtomicCommitDigest, AtomicCommitDisposition},
+ authored::{
+ AuthoredArtifact, AuthoredArtifactId, AuthoredOperation, FailureClass, RetrySchedule,
+ SigningState, WorkClaim, WorkFailure, WorkPhase,
+ },
+ authored_atomic::{
+ ApplyDeliveryAttempt, ApplySignedArtifact, AuthoredAtomicCommand, AuthoredAtomicOutcome,
+ AuthoredAtomicStorage, AuthoredWorkTarget, ClaimAuthoredTarget, ClaimAuthoredWork,
+ PrepareAuthoredOperation, WorkFence,
+ },
+ authored_delivery::{
+ AuthoredDeliveryIntent, AuthoredDeliveryPlan, AuthoredDeliveryPlanId,
+ AuthoredDeliveryState, DeliveryAttemptOutcome,
+ },
+ event::SourceGeneration,
+ journal::OperationInstanceId,
+ memory::MemoryStorage,
+};
+use radroots_transport::{
+ DeliveryReceipt, Target, TargetSet,
+ outcome::DeliveryOutcome,
+ policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy},
+ sink::DeliveryTargetReceipt,
+};
+
+const AUTHOR: &str = "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df";
+
+fn authored_plan() -> AuthoredEventPlan {
+ AuthoredEventPlan::from_generic(
+ GenericEventDraft::new(
+ "radroots.social.geochat.v1",
+ 20_000,
+ 1_800_000_100,
+ Vec::new(),
+ "atomic authored plan",
+ AUTHOR,
+ )
+ .expect("draft"),
+ )
+ .expect("plan")
+}
+
+fn signed(plan: &AuthoredEventPlan) -> SignedEvent {
+ let wire = Nip01EventWire {
+ id: plan.expected_event_id().to_hex(),
+ pubkey: plan.author().to_hex(),
+ created_at: plan.created_at(),
+ kind: plan.body().kind(),
+ tags: plan.body().tags().to_vec(),
+ content: plan.body().content().to_owned(),
+ sig: "dd".repeat(64),
+ extra: Default::default(),
+ };
+ let raw = serde_json::to_string(&wire).expect("raw event");
+ SignedEvent::from_wire_verified_id(wire, raw).expect("signed event")
+}
+
+fn ids() -> (
+ OperationInstanceId,
+ AuthoredArtifactId,
+ AuthoredDeliveryPlanId,
+) {
+ (
+ OperationInstanceId::new([1; 16]).expect("operation"),
+ AuthoredArtifactId::new([2; 16]).expect("artifact"),
+ AuthoredDeliveryPlanId::new([3; 16]).expect("delivery"),
+ )
+}
+
+fn intent() -> AuthoredDeliveryIntent {
+ AuthoredDeliveryIntent::new(
+ "atomic-delivery",
+ TargetSet::new(vec![
+ Target::nostr_relay("wss://one.example").expect("one"),
+ Target::nostr_relay("wss://two.example").expect("two"),
+ ])
+ .expect("targets"),
+ SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::any()),
+ 100,
+ )
+ .expect("intent")
+}
+
+fn prepare(input: u8) -> (AuthoredAtomicCommand, AuthoredEventPlan) {
+ let (operation_id, artifact_id, plan_id) = ids();
+ let authored_plan = authored_plan();
+ let artifact = AuthoredArtifact::planned(artifact_id, operation_id, 0, &authored_plan, 10)
+ .expect("artifact");
+ let operation = AuthoredOperation::new(operation_id, vec![artifact_id], 10).expect("operation");
+ let delivery =
+ AuthoredDeliveryPlan::new(plan_id, artifact_id, intent(), 10).expect("delivery plan");
+ let prepare = PrepareAuthoredOperation::new(
+ operation,
+ vec![artifact],
+ vec![delivery],
+ AtomicCommitDigest::new([input; 32]),
+ 10,
+ )
+ .expect("prepare");
+ (AuthoredAtomicCommand::Prepare(prepare), authored_plan)
+}
+
+fn claim(
+ target: ClaimAuthoredTarget,
+ revision: NonZeroU64,
+ token: u8,
+ at: u64,
+) -> (AuthoredAtomicCommand, WorkClaim) {
+ let claim = WorkClaim::new(
+ [token; 16],
+ format!("worker-{token}"),
+ NonZeroU64::new(u64::from(token)).expect("generation"),
+ at,
+ at + 20,
+ revision,
+ )
+ .expect("claim");
+ (
+ AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new(target, claim.clone())),
+ claim,
+ )
+}
+
+#[test]
+fn preparation_is_atomic_deterministic_and_exactly_replayable() {
+ let storage = MemoryStorage::new(SourceGeneration::new([1; 32]).expect("generation"));
+ let (command, _) = prepare(7);
+ assert_eq!(command.commit_id(), command.clone().commit_id());
+ assert_eq!(command.digest(), command.clone().digest());
+
+ let committed = block_on(storage.execute_authored(command.clone())).expect("commit");
+ assert_eq!(committed.disposition(), AtomicCommitDisposition::Committed);
+ let replay = block_on(storage.execute_authored(command.clone())).expect("replay");
+ assert_eq!(replay.disposition(), AtomicCommitDisposition::Replay);
+ assert_eq!(replay.outcome(), committed.outcome());
+ assert_eq!(
+ block_on(storage.authored_receipt(command.commit_id()))
+ .expect("receipt")
+ .expect("stored receipt")
+ .outcome(),
+ committed.outcome()
+ );
+
+ let (conflict, _) = prepare(8);
+ assert_eq!(
+ block_on(storage.execute_authored(conflict)),
+ Err(Error::AtomicCommitConflict)
+ );
+ let operation = block_on(storage.authored_operation(ids().0))
+ .expect("operation query")
+ .expect("operation");
+ assert_eq!(operation.artifact_ids(), &[ids().1]);
+ let delivery = block_on(storage.authored_delivery_plan(ids().2))
+ .expect("delivery query")
+ .expect("delivery");
+ assert!(delivery.request().is_none());
+}
+
+#[test]
+fn signing_atomically_binds_exact_delivery_requests_and_rejects_stale_fences() {
+ let storage = MemoryStorage::new(SourceGeneration::new([1; 32]).expect("generation"));
+ let (prepare, plan) = prepare(7);
+ block_on(storage.execute_authored(prepare)).expect("prepare");
+ let artifact = block_on(storage.authored_artifact(ids().1))
+ .expect("artifact query")
+ .expect("artifact");
+ let (claim_command, active) = claim(
+ ClaimAuthoredTarget::ArtifactSigning(ids().1),
+ artifact.revision(),
+ 4,
+ 11,
+ );
+ block_on(storage.execute_authored(claim_command)).expect("claim");
+
+ let stale = ApplySignedArtifact::new(
+ ids().1,
+ WorkFence::new([8; 16], active.generation(), active.row_revision()).expect("fence"),
+ signed(&plan),
+ 12,
+ )
+ .expect("stale apply");
+ assert_eq!(
+ block_on(storage.execute_authored(AuthoredAtomicCommand::ApplySigned(stale))),
+ Err(Error::DeliveryPlanClaimConflict)
+ );
+ assert_eq!(
+ block_on(storage.authored_artifact(ids().1))
+ .expect("artifact query")
+ .expect("artifact")
+ .signing_state(),
+ SigningState::Planned
+ );
+
+ let apply = ApplySignedArtifact::new(
+ ids().1,
+ WorkFence::new(*active.token(), active.generation(), active.row_revision()).expect("fence"),
+ signed(&plan),
+ 12,
+ )
+ .expect("apply");
+ let receipt = block_on(storage.execute_authored(AuthoredAtomicCommand::ApplySigned(apply)))
+ .expect("signed");
+ assert!(matches!(
+ receipt.outcome(),
+ AuthoredAtomicOutcome::Artifact(_)
+ ));
+ let artifact = block_on(storage.authored_artifact(ids().1))
+ .expect("artifact query")
+ .expect("artifact");
+ assert_eq!(artifact.signing_state(), SigningState::Signed);
+ let delivery = block_on(storage.authored_delivery_plan(ids().2))
+ .expect("delivery query")
+ .expect("delivery");
+ assert_eq!(
+ delivery
+ .request()
+ .expect("bound delivery")
+ .payload()
+ .event()
+ .raw_json(),
+ artifact
+ .signed()
+ .expect("signed artifact")
+ .event()
+ .raw_json()
+ );
+}
+
+#[test]
+fn delivery_attempt_and_work_failure_commands_preserve_atomic_state() {
+ let storage = MemoryStorage::new(SourceGeneration::new([1; 32]).expect("generation"));
+ let (prepare, plan) = prepare(7);
+ block_on(storage.execute_authored(prepare)).expect("prepare");
+ let artifact = block_on(storage.authored_artifact(ids().1))
+ .expect("artifact")
+ .expect("artifact");
+ let (sign_claim, active) = claim(
+ ClaimAuthoredTarget::ArtifactSigning(ids().1),
+ artifact.revision(),
+ 4,
+ 11,
+ );
+ block_on(storage.execute_authored(sign_claim)).expect("claim signing");
+ block_on(
+ storage.execute_authored(AuthoredAtomicCommand::ApplySigned(
+ ApplySignedArtifact::new(
+ ids().1,
+ WorkFence::new(*active.token(), active.generation(), active.row_revision())
+ .expect("fence"),
+ signed(&plan),
+ 12,
+ )
+ .expect("apply"),
+ )),
+ )
+ .expect("sign");
+
+ let delivery = block_on(storage.authored_delivery_plan(ids().2))
+ .expect("delivery")
+ .expect("delivery");
+ let (delivery_claim, active) = claim(
+ ClaimAuthoredTarget::DeliveryPlan(ids().2),
+ delivery.revision(),
+ 5,
+ 13,
+ );
+ block_on(storage.execute_authored(delivery_claim)).expect("claim delivery");
+ let request = block_on(storage.authored_delivery_plan(ids().2))
+ .expect("delivery")
+ .expect("delivery")
+ .request()
+ .expect("bound request")
+ .clone();
+ let receipt = DeliveryReceipt::for_request(
+ &request,
+ request
+ .target_set()
+ .targets()
+ .iter()
+ .cloned()
+ .map(|target| DeliveryTargetReceipt::attempted(target, DeliveryOutcome::accepted()))
+ .collect(),
+ )
+ .expect("receipt");
+ let apply = ApplyDeliveryAttempt::new(
+ ids().2,
+ WorkFence::new(*active.token(), active.generation(), active.row_revision()).expect("fence"),
+ DeliveryAttemptOutcome::Receipt(receipt),
+ None,
+ 14,
+ )
+ .expect("apply delivery");
+ block_on(storage.execute_authored(AuthoredAtomicCommand::ApplyDelivery(apply)))
+ .expect("deliver");
+ assert_eq!(
+ block_on(storage.authored_delivery_plan(ids().2))
+ .expect("delivery")
+ .expect("delivery")
+ .state(),
+ AuthoredDeliveryState::Satisfied
+ );
+
+ let retry_failure = WorkFailure::new(
+ "temporary_signer_failure",
+ WorkPhase::Signing,
+ FailureClass::Retryable,
+ Some(30),
+ None,
+ )
+ .expect("failure");
+ let retry = RetrySchedule::new(NonZeroU32::MIN, 30, retry_failure.clone()).expect("retry");
+ let invalid = radroots_storage::authored_atomic::ApplyWorkFailure::new(
+ AuthoredWorkTarget::Artifact(ids().1),
+ WorkFence::new([9; 16], NonZeroU64::MIN, NonZeroU64::MIN).expect("fence"),
+ retry_failure,
+ Some(retry),
+ 15,
+ )
+ .expect("failure command");
+ let before = block_on(storage.authored_artifact(ids().1))
+ .expect("artifact")
+ .expect("artifact");
+ assert!(
+ block_on(storage.execute_authored(AuthoredAtomicCommand::ApplyFailure(invalid))).is_err()
+ );
+ assert_eq!(
+ block_on(storage.authored_artifact(ids().1))
+ .expect("artifact")
+ .expect("artifact"),
+ before
+ );
+}
diff --git a/crates/storage/tests/authored_delivery.rs b/crates/storage/tests/authored_delivery.rs
@@ -53,7 +53,7 @@ fn request(policy: TargetPolicy) -> DeliveryRequest {
}
fn plan(value: u8, policy: TargetPolicy) -> AuthoredDeliveryPlan {
- AuthoredDeliveryPlan::new(
+ AuthoredDeliveryPlan::new_bound(
AuthoredDeliveryPlanId::new([value; 16]).expect("plan ID"),
AuthoredArtifactId::new([9; 16]).expect("artifact ID"),
request(policy),
@@ -62,6 +62,10 @@ fn plan(value: u8, policy: TargetPolicy) -> AuthoredDeliveryPlan {
.expect("plan")
}
+fn bound_request(plan: &AuthoredDeliveryPlan) -> &DeliveryRequest {
+ plan.request().expect("bound request")
+}
+
fn claim(plan: &AuthoredDeliveryPlan, token: u8, acquired: u64) -> WorkClaim {
WorkClaim::new(
[token; 16],
@@ -118,7 +122,7 @@ fn independent_plans_claim_and_progress_without_cross_blocking() {
.expect("second claim");
let first_receipt = receipt(
- first.request(),
+ bound_request(&first),
vec![DeliveryOutcome::accepted(), DeliveryOutcome::unavailable()],
);
first
@@ -135,7 +139,7 @@ fn independent_plans_claim_and_progress_without_cross_blocking() {
assert!(second.claim_evidence().is_some());
let second_receipt = receipt(
- second.request(),
+ bound_request(&second),
vec![DeliveryOutcome::accepted(), DeliveryOutcome::unavailable()],
);
second
@@ -165,11 +169,11 @@ fn retry_schedule_and_partial_evidence_survive_reconstruction() {
let active = claim(&delivery, 3, 11);
delivery.claim(active.clone(), 11).expect("claim");
let partial = DeliveryTargetReceipt::attempted(
- delivery.request().target_set().targets()[0].clone(),
+ bound_request(&delivery).target_set().targets()[0].clone(),
DeliveryOutcome::accepted(),
);
let sink_failure = SinkFailure::for_request(
- delivery.request(),
+ bound_request(&delivery),
"relay_batch_unavailable",
Retryability::Retryable,
Some(20),
@@ -209,11 +213,11 @@ fn terminal_partial_success_stale_claim_and_invalid_receipt_fail_closed() {
let active = claim(&any, 4, 11);
any.claim(active.clone(), 11).expect("claim");
let partial = DeliveryTargetReceipt::attempted(
- any.request().target_set().targets()[0].clone(),
+ bound_request(&any).target_set().targets()[0].clone(),
DeliveryOutcome::accepted(),
);
let failure = SinkFailure::for_request(
- any.request(),
+ bound_request(&any),
"terminal_batch_failure",
Retryability::Terminal,
None,
@@ -223,10 +227,10 @@ fn terminal_partial_success_stale_claim_and_invalid_receipt_fail_closed() {
.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(),
+ bound_request(&any).payload().clone(),
+ bound_request(&any).target_set().clone(),
+ bound_request(&any).satisfaction().clone(),
+ bound_request(&any).deadline_unix_ms(),
)
.expect("other request");
let invalid_receipt = receipt(
@@ -273,7 +277,7 @@ fn terminal_partial_success_stale_claim_and_invalid_receipt_fail_closed() {
let all_claim = claim(&all, 5, 11);
all.claim(all_claim.clone(), 11).expect("claim");
let terminal = SinkFailure::for_request(
- all.request(),
+ bound_request(&all),
"terminal_batch_failure",
Retryability::Terminal,
None,
@@ -297,7 +301,7 @@ fn terminal_partial_success_stale_claim_and_invalid_receipt_fail_closed() {
fn attempt_limit_is_checked_without_mutating_claimed_state() {
let base = plan(6, TargetPolicy::all());
let pending_receipt = receipt(
- base.request(),
+ bound_request(&base),
vec![
DeliveryOutcome::unavailable(),
DeliveryOutcome::unavailable(),