commit 97ad465d34ab4dcf7b665749a387209d3d807f3d
parent f578069cf894b34695bdf0ed2298e0aa162b5c13
Author: triesap <tyson@radroots.org>
Date: Mon, 14 Sep 2026 13:10:12 +0000
storage: retain late delivery facts and stop intent
- Bind bounded sink results to original durable claim receipts
- Preserve first stop intent and independent scheduling fences
- Add forward SQLite migration and consistent evidence reads
- Verify replay, rollback, provenance and unchanged coverage gates
Diffstat:
18 files changed, 2002 insertions(+), 67 deletions(-)
diff --git a/contracts/api_baselines/radroots_storage.txt b/contracts/api_baselines/radroots_storage.txt
@@ -201,6 +201,7 @@ pub radroots_storage::authored_atomic::AuthoredAtomicCommand::Cancel(radroots_st
pub radroots_storage::authored_atomic::AuthoredAtomicCommand::Claim(radroots_storage::authored_atomic::ClaimAuthoredWork)
pub radroots_storage::authored_atomic::AuthoredAtomicCommand::Prepare(radroots_storage::authored_atomic::PrepareAuthoredOperation)
pub radroots_storage::authored_atomic::AuthoredAtomicCommand::PrepareFromDraft(alloc::boxed::Box<radroots_storage::authored_draft_submission::PrepareFromDraft>)
+pub radroots_storage::authored_atomic::AuthoredAtomicCommand::RecordDelivery(radroots_storage::authored_atomic::RecordDeliveryFact)
pub radroots_storage::authored_atomic::AuthoredAtomicCommand::RecordSigned(radroots_storage::authored_atomic::RecordSignedArtifact)
impl radroots_storage::authored_atomic::AuthoredAtomicCommand
pub fn radroots_storage::authored_atomic::AuthoredAtomicCommand::commit_id(&self) -> radroots_storage::atomic::AtomicCommitId
@@ -286,6 +287,16 @@ pub const fn radroots_storage::authored_atomic::PrepareAuthoredOperation::input_
pub fn radroots_storage::authored_atomic::PrepareAuthoredOperation::new(radroots_storage::authored::AuthoredOperation, alloc::vec::Vec<radroots_storage::authored::AuthoredArtifact>, alloc::vec::Vec<radroots_storage::authored_delivery::AuthoredDeliveryPlan>, radroots_storage::atomic::AtomicCommitDigest, u64) -> core::result::Result<Self, radroots_storage::Error>
pub const fn radroots_storage::authored_atomic::PrepareAuthoredOperation::operation(&self) -> &radroots_storage::authored::AuthoredOperation
pub const fn radroots_storage::authored_atomic::PrepareAuthoredOperation::requested_at_unix_ms(&self) -> u64
+pub struct radroots_storage::authored_atomic::RecordDeliveryFact
+impl radroots_storage::authored_atomic::RecordDeliveryFact
+pub fn radroots_storage::authored_atomic::RecordDeliveryFact::apply_to(&self, &mut radroots_storage::authored_delivery::AuthoredDeliveryPlan, &radroots_storage::authored_atomic::AuthoredAtomicReceipt) -> core::result::Result<(), radroots_storage::Error>
+pub const fn radroots_storage::authored_atomic::RecordDeliveryFact::artifact_id(&self) -> radroots_storage::authored::AuthoredArtifactId
+pub const fn radroots_storage::authored_atomic::RecordDeliveryFact::claim(&self) -> &radroots_storage::authored::WorkClaim
+pub fn radroots_storage::authored_atomic::RecordDeliveryFact::claim_command(&self) -> radroots_storage::authored_atomic::AuthoredAtomicCommand
+pub fn radroots_storage::authored_atomic::RecordDeliveryFact::new(radroots_storage::authored_delivery::AuthoredDeliveryPlanId, radroots_storage::authored::AuthoredArtifactId, radroots_storage::authored::WorkClaim, radroots_storage::authored_delivery::DeliveryAttemptOutcome, u64) -> core::result::Result<Self, radroots_storage::Error>
+pub const fn radroots_storage::authored_atomic::RecordDeliveryFact::observed_at_unix_ms(&self) -> u64
+pub const fn radroots_storage::authored_atomic::RecordDeliveryFact::outcome(&self) -> &radroots_storage::authored_delivery::DeliveryAttemptOutcome
+pub const fn radroots_storage::authored_atomic::RecordDeliveryFact::plan_id(&self) -> radroots_storage::authored_delivery::AuthoredDeliveryPlanId
pub struct radroots_storage::authored_atomic::RecordSignedArtifact
impl radroots_storage::authored_atomic::RecordSignedArtifact
pub fn radroots_storage::authored_atomic::RecordSignedArtifact::apply_to(&self, &mut radroots_storage::authored::AuthoredArtifact, &radroots_storage::authored_atomic::AuthoredAtomicReceipt) -> core::result::Result<(), radroots_storage::Error>
@@ -334,6 +345,11 @@ pub const fn radroots_storage::authored_delivery::AuthoredDeliveryAttempt::outco
pub fn radroots_storage::authored_delivery::AuthoredDeliveryAttempt::reconstruct(core::num::nonzero::NonZeroU32, u64, radroots_storage::authored_delivery::DeliveryAttemptOutcome, radroots_transport::policy::SatisfactionState) -> core::result::Result<Self, radroots_storage::Error>
pub const fn radroots_storage::authored_delivery::AuthoredDeliveryAttempt::recorded_at_unix_ms(&self) -> u64
pub const fn radroots_storage::authored_delivery::AuthoredDeliveryAttempt::satisfaction(&self) -> radroots_transport::policy::SatisfactionState
+pub struct radroots_storage::authored_delivery::AuthoredDeliveryFact
+impl radroots_storage::authored_delivery::AuthoredDeliveryFact
+pub const fn radroots_storage::authored_delivery::AuthoredDeliveryFact::claim(&self) -> &radroots_storage::authored::WorkClaim
+pub const fn radroots_storage::authored_delivery::AuthoredDeliveryFact::observed_at_unix_ms(&self) -> u64
+pub const fn radroots_storage::authored_delivery::AuthoredDeliveryFact::outcome(&self) -> &radroots_storage::authored_delivery::DeliveryAttemptOutcome
pub struct radroots_storage::authored_delivery::AuthoredDeliveryIntent
impl radroots_storage::authored_delivery::AuthoredDeliveryIntent
pub const fn radroots_storage::authored_delivery::AuthoredDeliveryIntent::deadline_unix_ms(&self) -> u64
@@ -364,11 +380,16 @@ pub const fn radroots_storage::authored_delivery::AuthoredDeliveryPlan::plan_id(
pub fn radroots_storage::authored_delivery::AuthoredDeliveryPlan::reconstruct(Self) -> core::result::Result<Self, radroots_storage::Error>
pub const fn radroots_storage::authored_delivery::AuthoredDeliveryPlan::request(&self) -> core::option::Option<&radroots_transport::sink::DeliveryRequest>
pub const fn radroots_storage::authored_delivery::AuthoredDeliveryPlan::request_digest(&self) -> &[u8; 32]
+pub fn radroots_storage::authored_delivery::AuthoredDeliveryPlan::request_stop(&mut self, u64) -> core::result::Result<(), radroots_storage::Error>
pub const fn radroots_storage::authored_delivery::AuthoredDeliveryPlan::retry(&self) -> core::option::Option<&radroots_storage::authored::RetrySchedule>
pub const fn radroots_storage::authored_delivery::AuthoredDeliveryPlan::revision(&self) -> core::num::nonzero::NonZeroU64
pub const fn radroots_storage::authored_delivery::AuthoredDeliveryPlan::state(&self) -> radroots_storage::authored_delivery::AuthoredDeliveryState
pub const fn radroots_storage::authored_delivery::AuthoredDeliveryPlan::updated_at_unix_ms(&self) -> u64
pub fn radroots_storage::authored_delivery::AuthoredDeliveryPlan::validate(&self) -> core::result::Result<(), radroots_storage::Error>
+impl radroots_storage::authored_delivery::AuthoredDeliveryPlan
+pub fn radroots_storage::authored_delivery::AuthoredDeliveryPlan::delivery_facts(&self) -> &[radroots_storage::authored_delivery::AuthoredDeliveryFact]
+pub fn radroots_storage::authored_delivery::AuthoredDeliveryPlan::delivery_satisfaction(&self) -> core::result::Result<radroots_transport::policy::SatisfactionState, radroots_storage::Error>
+pub const fn radroots_storage::authored_delivery::AuthoredDeliveryPlan::stop_requested_at_unix_ms(&self) -> core::option::Option<u64>
pub struct radroots_storage::authored_delivery::AuthoredDeliveryPlanId(_)
impl radroots_storage::authored_delivery::AuthoredDeliveryPlanId
pub const fn radroots_storage::authored_delivery::AuthoredDeliveryPlanId::as_bytes(&self) -> &[u8; 16]
diff --git a/contracts/architecture/decisions/authored_delivery_facts.v1.json b/contracts/architecture/decisions/authored_delivery_facts.v1.json
@@ -0,0 +1,23 @@
+{
+ "schema": "radroots.authored-delivery-facts.v1",
+ "status": "approved",
+ "owners": [
+ "radroots_storage",
+ "radroots_storage_sqlite"
+ ],
+ "command": "authored_atomic::RecordDeliveryFact; AuthoredAtomicCommand::RecordDelivery",
+ "provenance": "The owning transaction retrieves its immutable original delivery Claim receipt. Bind exact plan, artifact, signed request including raw bytes, intent, complete claim and original row identity. Caller-supplied receipts are not storage authority. Expiry and supersession do not invalidate facts.",
+ "scheduling": "Facts are an orthogonal append-only collection with their own observation times and immutable atomic receipts. Recording facts does not alter scheduling state, revision, updated time, live claim, legacy attempts or retry. Ordinary mutations always load the current plan inside the serialized transaction and retain their exact fences. Scheduling callers must reconcile current facts and stop intent before new effects.",
+ "stop": "Retain the first explicit stop time. Pending or retryable plans become Cancelled and lose scheduling authority; settled states retain their outcome and stop intent. The legacy direct cancel API remains strict. Stop does not erase or revoke external acceptance.",
+ "evidence": "One immutable final sink result per complete original claim. Equal replay changes nothing, conflicting results fail closed. Cumulative satisfaction folds legacy attempts and all delivery facts without retroactively rewriting legacy attempt satisfaction. Facts may arrive out of observation order. Receipt time is at least the observation and current scheduling time.",
+ "bounds": "Preserve the 1024 legacy scheduling-attempt limit and 4194304-byte SQLite snapshot limit. Facts have an independent finite 1024-entry limit; no eviction or implicit threshold increase. At capacity, new claims are rejected and a distinct late fact returns the existing typed overflow without changing stored data or claiming success. Exact replay remains admissible.",
+ "migration": {
+ "version": 16,
+ "name": "authored_delivery_facts",
+ "sha256": "f78d35fbe60255f152c4e7a8c23add6677ac3ae74eed7db1f3d8ef4ba0b8f3c3",
+ "compatibility": "All previous SQL, checksums, command hashes and retained receipts remain unchanged. Add a nullable normalized first-stop time and append-only normalized facts. Old Cancelled snapshots infer the stop from their row time; other old snapshots default to no facts or stop. Legacy v10 conversion writes only v11 columns before this forward migration."
+ },
+ "atomicity": "Memory commits a validated candidate; SQLite persists the plan, normalized facts and immutable atomic receipt in the same transaction, returning only after COMMIT. Failure rolls back every write. Each delivery-plan read validates its normalized children within one read snapshot and releases that snapshot before returning.",
+ "consumer": "Sync adoption is separate. These storage mechanics alone do not claim whole-operation cancellation, UI integration or remote rollback.",
+ "identity": "The new logical delivery_fact_v1 command binds plan, artifact, full original claim and normalized outcome details, excluding observation time. The original claim binds the exact request. Receipt replay additionally compares the full typed outcome; invalid request bindings cannot replay as a valid result. All historical command hashes are unchanged."
+}
diff --git a/crates/storage/src/authored_atomic.rs b/crates/storage/src/authored_atomic.rs
@@ -2,6 +2,8 @@
mod signing_evidence;
pub use signing_evidence::RecordSignedArtifact;
+mod delivery_evidence;
+pub use delivery_evidence::RecordDeliveryFact;
use core::num::NonZeroU64;
use radroots_event::SignedEvent;
@@ -132,9 +134,11 @@ impl PrepareAuthoredOperation {
.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())
+ || delivery_plans.iter().any(|plan| {
+ !artifact_ids.contains(&plan.artifact_id())
+ || plan.validate().is_err()
+ || !plan.delivery_facts().is_empty()
+ })
{
return Err(Error::AtomicWorkflowMismatch);
}
@@ -423,6 +427,7 @@ pub enum AuthoredAtomicCommand {
RecordSigned(RecordSignedArtifact),
ApplyAdmission(ApplyAdmissionResult),
ApplyDelivery(ApplyDeliveryAttempt),
+ RecordDelivery(RecordDeliveryFact),
ApplyFailure(ApplyWorkFailure),
Cancel(CancelAuthoredWork),
}
@@ -482,6 +487,17 @@ impl AuthoredAtomicCommand {
hash_failure(&mut hasher, value.failure.as_ref());
}
Self::ApplyDelivery(value) => hash_delivery(&mut hasher, &value.outcome),
+ Self::RecordDelivery(value) => {
+ hash_field(&mut hasher, value.artifact_id().as_bytes());
+ let claim = value.claim();
+ hash_field(&mut hasher, claim.token());
+ hash_field(&mut hasher, claim.owner().as_bytes());
+ hasher.update(claim.generation().get().to_be_bytes());
+ hasher.update(claim.row_revision().get().to_be_bytes());
+ hasher.update(claim.acquired_at_unix_ms().to_be_bytes());
+ hasher.update(claim.expires_at_unix_ms().to_be_bytes());
+ value.hash_outcome(&mut hasher);
+ }
Self::ApplyFailure(value) => hash_failure(&mut hasher, Some(&value.failure)),
Self::Cancel(value) => hasher.update(value.cancelled_at_unix_ms.to_be_bytes()),
}
@@ -497,6 +513,7 @@ impl AuthoredAtomicCommand {
Self::RecordSigned(value) => value.observed_at_unix_ms(),
Self::ApplyAdmission(value) => value.applied_at_unix_ms,
Self::ApplyDelivery(value) => value.applied_at_unix_ms,
+ Self::RecordDelivery(value) => value.observed_at_unix_ms(),
Self::ApplyFailure(value) => value.applied_at_unix_ms,
Self::Cancel(value) => value.cancelled_at_unix_ms,
}
@@ -511,6 +528,7 @@ impl AuthoredAtomicCommand {
Self::RecordSigned(_) => b"signed_fact_v1",
Self::ApplyAdmission(_) => b"admission",
Self::ApplyDelivery(_) => b"delivery",
+ Self::RecordDelivery(_) => b"delivery_fact_v1",
Self::ApplyFailure(value) => match value.failure.phase() {
WorkPhase::Signing => b"signing_failure",
WorkPhase::Admission => b"admission_failure",
@@ -535,6 +553,7 @@ impl AuthoredAtomicCommand {
Self::RecordSigned(value) => *value.artifact_id().as_bytes(),
Self::ApplyAdmission(value) => *value.artifact_id.as_bytes(),
Self::ApplyDelivery(value) => *value.plan_id.as_bytes(),
+ Self::RecordDelivery(value) => *value.plan_id().as_bytes(),
Self::ApplyFailure(value) => match &value.target {
AuthoredWorkTarget::Artifact(id) => *id.as_bytes(),
AuthoredWorkTarget::DeliveryPlan(id) => *id.as_bytes(),
@@ -554,6 +573,7 @@ impl AuthoredAtomicCommand {
Self::Claim(value) => Some(value.claim.generation()),
Self::ApplyAdmission(value) => Some(value.fence.generation),
Self::ApplyDelivery(value) => Some(value.fence.generation),
+ Self::RecordDelivery(value) => Some(value.claim().generation()),
Self::ApplyFailure(value) => Some(value.fence.generation),
Self::Prepare(_) | Self::PrepareFromDraft(_) | Self::Cancel(_) => None,
}
@@ -591,6 +611,11 @@ impl AuthoredAtomicReceipt {
}
match (command, &self.outcome) {
(
+ AuthoredAtomicCommand::RecordDelivery(value),
+ AuthoredAtomicOutcome::DeliveryPlan(plan),
+ ) => value.matches_plan(plan),
+ (AuthoredAtomicCommand::RecordDelivery(_), _) => false,
+ (
AuthoredAtomicCommand::RecordSigned(value),
AuthoredAtomicOutcome::Artifact(artifact),
) => signed_fact_matches(value, artifact),
@@ -622,6 +647,15 @@ impl AuthoredAtomicReceipt {
}
match (command, &outcome) {
(
+ AuthoredAtomicCommand::RecordDelivery(value),
+ AuthoredAtomicOutcome::DeliveryPlan(plan),
+ ) if value.matches_plan(plan) && committed_at_unix_ms >= plan.updated_at_unix_ms() => {
+ plan.validate()?
+ }
+ (AuthoredAtomicCommand::RecordDelivery(_), _) => {
+ return Err(Error::AtomicWorkflowMismatch);
+ }
+ (
AuthoredAtomicCommand::RecordSigned(value),
AuthoredAtomicOutcome::Artifact(artifact),
) if signed_fact_matches(value, artifact) => artifact.validate()?,
diff --git a/crates/storage/src/authored_atomic/delivery_evidence.rs b/crates/storage/src/authored_atomic/delivery_evidence.rs
@@ -0,0 +1,150 @@
+//! Delivery facts authorized by the backend's original immutable claim receipt.
+
+use super::{
+ AuthoredAtomicCommand, AuthoredAtomicOutcome, AuthoredAtomicReceipt, ClaimAuthoredTarget,
+ ClaimAuthoredWork,
+};
+use crate::{
+ Error,
+ authored::{AuthoredArtifactId, WorkClaim},
+ authored_delivery::{AuthoredDeliveryPlan, AuthoredDeliveryPlanId, DeliveryAttemptOutcome},
+};
+
+/// Retains a completed sink result without extending a lease or granting work.
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct RecordDeliveryFact {
+ plan_id: AuthoredDeliveryPlanId,
+ artifact_id: AuthoredArtifactId,
+ claim: WorkClaim,
+ outcome: DeliveryAttemptOutcome,
+ observed_at_unix_ms: u64,
+}
+
+impl RecordDeliveryFact {
+ pub fn new(
+ plan_id: AuthoredDeliveryPlanId,
+ artifact_id: AuthoredArtifactId,
+ claim: WorkClaim,
+ outcome: DeliveryAttemptOutcome,
+ observed_at_unix_ms: u64,
+ ) -> Result<Self, Error> {
+ claim.validate()?;
+ if observed_at_unix_ms < claim.acquired_at_unix_ms() {
+ return Err(Error::AtomicWorkflowMismatch);
+ }
+ Ok(Self {
+ plan_id,
+ artifact_id,
+ claim,
+ outcome,
+ observed_at_unix_ms,
+ })
+ }
+
+ pub const fn plan_id(&self) -> AuthoredDeliveryPlanId {
+ self.plan_id
+ }
+ pub const fn artifact_id(&self) -> AuthoredArtifactId {
+ self.artifact_id
+ }
+ pub const fn claim(&self) -> &WorkClaim {
+ &self.claim
+ }
+ pub const fn outcome(&self) -> &DeliveryAttemptOutcome {
+ &self.outcome
+ }
+ pub const fn observed_at_unix_ms(&self) -> u64 {
+ self.observed_at_unix_ms
+ }
+
+ pub fn claim_command(&self) -> AuthoredAtomicCommand {
+ AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new(
+ ClaimAuthoredTarget::DeliveryPlan(self.plan_id),
+ self.claim.clone(),
+ ))
+ }
+
+ /// The backend must retrieve `original_claim` inside its own transaction.
+ /// Caller-provided receipts must never substitute for persisted provenance.
+ pub fn apply_to(
+ &self,
+ plan: &mut AuthoredDeliveryPlan,
+ original_claim: &AuthoredAtomicReceipt,
+ ) -> Result<(), Error> {
+ let AuthoredAtomicOutcome::DeliveryPlan(original) = original_claim.outcome() else {
+ return Err(Error::AtomicWorkflowMismatch);
+ };
+ original.validate()?;
+ plan.validate()?;
+ if !original_claim.matches_command(&self.claim_command())
+ || original_claim.committed_at_unix_ms() != self.claim.acquired_at_unix_ms()
+ || original.claim_evidence() != Some(&self.claim)
+ || original.plan_id() != self.plan_id
+ || plan.plan_id() != self.plan_id
+ || original.artifact_id() != self.artifact_id
+ || plan.artifact_id() != self.artifact_id
+ || original.request().is_none()
+ || original.request() != plan.request()
+ || original.intent() != plan.intent()
+ || original.created_at_unix_ms() != plan.created_at_unix_ms()
+ || original.revision() > plan.revision()
+ || original.updated_at_unix_ms() > plan.updated_at_unix_ms()
+ {
+ return Err(Error::AtomicWorkflowMismatch);
+ }
+ plan.record_delivery_fact(
+ self.claim.clone(),
+ self.outcome.clone(),
+ self.observed_at_unix_ms,
+ )
+ }
+
+ pub(super) fn hash_outcome(&self, hasher: &mut sha2::Sha256) {
+ use sha2::Digest;
+ let entries = match &self.outcome {
+ DeliveryAttemptOutcome::Receipt(receipt) => {
+ hasher.update([0]);
+ super::hash_field(hasher, receipt.request_id().as_str().as_bytes());
+ receipt.target_receipts()
+ }
+ DeliveryAttemptOutcome::SinkFailure(failure) => {
+ hasher.update([1]);
+ hasher.update([failure.retry_after_unix_ms().is_some() as u8]);
+ hasher.update(
+ failure
+ .retry_after_unix_ms()
+ .unwrap_or_default()
+ .to_be_bytes(),
+ );
+ hash_optional(hasher, failure.message());
+ failure.partial_evidence()
+ }
+ };
+ hasher.update(
+ u64::try_from(entries.len())
+ .unwrap_or(u64::MAX)
+ .to_be_bytes(),
+ );
+ super::hash_delivery(hasher, &self.outcome);
+ for entry in entries {
+ hash_optional(hasher, entry.target().label().map(|label| label.as_str()));
+ }
+ }
+
+ pub(super) fn matches_plan(&self, plan: &AuthoredDeliveryPlan) -> bool {
+ plan.plan_id() == self.plan_id
+ && plan.artifact_id() == self.artifact_id
+ && plan
+ .delivery_facts()
+ .iter()
+ .any(|fact| fact.claim() == &self.claim && fact.outcome() == &self.outcome)
+ }
+}
+
+fn hash_optional(hasher: &mut sha2::Sha256, value: Option<&str>) {
+ use sha2::Digest;
+ hasher.update([value.is_some() as u8]);
+ if let Some(value) = value {
+ super::hash_field(hasher, value.as_bytes());
+ }
+}
diff --git a/crates/storage/src/authored_delivery.rs b/crates/storage/src/authored_delivery.rs
@@ -1,5 +1,8 @@
//! Independent durable delivery-plan, attempt, retry, and evidence models.
+mod facts;
+pub use facts::AuthoredDeliveryFact;
+
use core::num::{NonZeroU32, NonZeroU64};
use radroots_transport::{
DeliveryReceipt, DeliveryRequest, SinkFailure,
@@ -215,6 +218,8 @@ pub struct AuthoredDeliveryPlan {
request: Option<DeliveryRequest>,
state: AuthoredDeliveryState,
attempts: Vec<AuthoredDeliveryAttempt>,
+ delivery_facts: Vec<AuthoredDeliveryFact>,
+ stop_requested_at_unix_ms: Option<u64>,
attempt_count: u32,
retry: Option<RetrySchedule>,
claim: Option<WorkClaim>,
@@ -240,6 +245,8 @@ impl AuthoredDeliveryPlan {
request: None,
state: AuthoredDeliveryState::Pending,
attempts: Vec::new(),
+ delivery_facts: Vec::new(),
+ stop_requested_at_unix_ms: None,
attempt_count: 0,
retry: None,
claim: None,
@@ -269,6 +276,7 @@ impl AuthoredDeliveryPlan {
}
pub fn validate(&self) -> Result<(), Error> {
+ self.validate_facts()?;
if self.created_at_unix_ms == 0
|| self.updated_at_unix_ms < self.created_at_unix_ms
|| self.request_digest != delivery_intent_digest(&self.intent)
@@ -365,6 +373,8 @@ impl AuthoredDeliveryPlan {
|| claim.generation() <= existing.generation()
});
if self.state.is_terminal()
+ || self.stop_requested_at_unix_ms.is_some()
+ || self.delivery_facts.len() >= DELIVERY_PLAN_ATTEMPTS_MAX as usize
|| self.request.is_none()
|| existing_blocks
|| claim.row_revision() != self.revision
@@ -559,11 +569,25 @@ impl AuthoredDeliveryPlan {
if self.state.is_terminal() {
return Err(Error::InvalidAuthoredDeliveryPlan);
}
+ self.request_stop(cancelled_at_unix_ms)
+ }
+
+ /// Retains the first stop intent, including when delivery already settled.
+ pub fn request_stop(&mut self, cancelled_at_unix_ms: u64) -> Result<(), Error> {
+ if cancelled_at_unix_ms < self.created_at_unix_ms {
+ return Err(Error::InvalidAuthoredDeliveryPlan);
+ }
+ if self.stop_requested_at_unix_ms.is_some() {
+ return Ok(());
+ }
let previous = self.clone();
- self.state = AuthoredDeliveryState::Cancelled;
+ self.stop_requested_at_unix_ms = Some(cancelled_at_unix_ms);
+ if !self.state.is_terminal() {
+ self.state = AuthoredDeliveryState::Cancelled;
+ self.last_failure = None;
+ }
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);
@@ -757,6 +781,10 @@ struct AuthoredDeliveryPlanWire {
request: Option<DeliveryRequest>,
state: AuthoredDeliveryState,
attempts: Vec<AuthoredDeliveryAttempt>,
+ #[serde(default)]
+ delivery_facts: Vec<AuthoredDeliveryFact>,
+ #[serde(default)]
+ stop_requested_at_unix_ms: Option<u64>,
attempt_count: u32,
retry: Option<RetrySchedule>,
claim: Option<WorkClaim>,
@@ -778,6 +806,11 @@ impl TryFrom<AuthoredDeliveryPlanWire> for AuthoredDeliveryPlan {
request: value.request,
state: value.state,
attempts: value.attempts,
+ delivery_facts: value.delivery_facts,
+ stop_requested_at_unix_ms: value.stop_requested_at_unix_ms.or_else(|| {
+ (value.state == AuthoredDeliveryState::Cancelled)
+ .then_some(value.updated_at_unix_ms)
+ }),
attempt_count: value.attempt_count,
retry: value.retry,
claim: value.claim,
@@ -800,6 +833,8 @@ impl From<AuthoredDeliveryPlan> for AuthoredDeliveryPlanWire {
request: value.request,
state: value.state,
attempts: value.attempts,
+ delivery_facts: value.delivery_facts,
+ stop_requested_at_unix_ms: value.stop_requested_at_unix_ms,
attempt_count: value.attempt_count,
retry: value.retry,
claim: value.claim,
diff --git a/crates/storage/src/authored_delivery/facts.rs b/crates/storage/src/authored_delivery/facts.rs
@@ -0,0 +1,114 @@
+//! Immutable delivery observations, independent of the scheduling revision.
+
+use super::{
+ AuthoredDeliveryPlan, DELIVERY_PLAN_ATTEMPTS_MAX, DeliveryAttemptOutcome, Error,
+ SatisfactionState, WorkClaim,
+};
+
+/// One final sink result from an exactly identified durable claim.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct AuthoredDeliveryFact {
+ claim: WorkClaim,
+ outcome: DeliveryAttemptOutcome,
+ observed_at_unix_ms: u64,
+}
+
+impl AuthoredDeliveryFact {
+ pub const fn claim(&self) -> &WorkClaim {
+ &self.claim
+ }
+ pub const fn outcome(&self) -> &DeliveryAttemptOutcome {
+ &self.outcome
+ }
+ pub const fn observed_at_unix_ms(&self) -> u64 {
+ self.observed_at_unix_ms
+ }
+}
+
+impl AuthoredDeliveryPlan {
+ pub fn delivery_facts(&self) -> &[AuthoredDeliveryFact] {
+ &self.delivery_facts
+ }
+
+ pub const fn stop_requested_at_unix_ms(&self) -> Option<u64> {
+ self.stop_requested_at_unix_ms
+ }
+
+ /// Cumulative transport evidence, independent of cancellation or retry state.
+ pub fn delivery_satisfaction(&self) -> Result<SatisfactionState, Error> {
+ self.evaluate_outcomes(
+ self.attempts
+ .iter()
+ .map(|attempt| &attempt.outcome)
+ .chain(self.delivery_facts.iter().map(|fact| &fact.outcome)),
+ )
+ }
+
+ pub(super) fn validate_facts(&self) -> Result<(), Error> {
+ if self.delivery_facts.len() > DELIVERY_PLAN_ATTEMPTS_MAX as usize
+ || self.stop_requested_at_unix_ms.is_some_and(|at| {
+ at < self.created_at_unix_ms
+ || at > self.updated_at_unix_ms
+ || !self.state.is_terminal()
+ })
+ {
+ return Err(Error::InvalidAuthoredDeliveryPlan);
+ }
+ let mut claims = std::collections::BTreeSet::new();
+ for fact in &self.delivery_facts {
+ if fact.claim.validate().is_err()
+ || fact.claim.acquired_at_unix_ms() < self.created_at_unix_ms
+ || fact.claim.acquired_at_unix_ms() > self.updated_at_unix_ms
+ || fact.claim.row_revision() >= self.revision
+ || fact.observed_at_unix_ms < fact.claim.acquired_at_unix_ms()
+ || self
+ .request
+ .as_ref()
+ .is_none_or(|request| fact.outcome.validate_for(request).is_err())
+ || !claims.insert((
+ fact.claim.token(),
+ fact.claim.owner(),
+ fact.claim.generation(),
+ fact.claim.row_revision(),
+ fact.claim.acquired_at_unix_ms(),
+ fact.claim.expires_at_unix_ms(),
+ ))
+ {
+ return Err(Error::InvalidAuthoredDeliveryPlan);
+ }
+ }
+ Ok(())
+ }
+
+ /// Called only after the atomic owner has checked its original claim receipt.
+ pub(crate) fn record_delivery_fact(
+ &mut self,
+ claim: WorkClaim,
+ outcome: DeliveryAttemptOutcome,
+ observed_at_unix_ms: u64,
+ ) -> Result<(), Error> {
+ if let Some(prior) = self.delivery_facts.iter().find(|fact| fact.claim == claim) {
+ return if prior.outcome == outcome {
+ Ok(())
+ } else {
+ Err(Error::AtomicWorkflowMismatch)
+ };
+ }
+ if self.delivery_facts.len() >= DELIVERY_PLAN_ATTEMPTS_MAX as usize {
+ return Err(Error::DeliveryAttemptOverflow);
+ }
+ self.delivery_facts.push(AuthoredDeliveryFact {
+ claim,
+ outcome,
+ observed_at_unix_ms,
+ });
+ // Facts have their own observation time and immutable atomic receipt. They
+ // do not mutate scheduling revision/time, fences, attempts or backoff.
+ if let Err(error) = self.validate() {
+ self.delivery_facts.pop();
+ return Err(error);
+ }
+ Ok(())
+ }
+}
diff --git a/crates/storage/src/memory.rs b/crates/storage/src/memory.rs
@@ -1642,6 +1642,21 @@ impl AuthoredAtomicStorage for MemoryStorage {
}
AuthoredAtomicOutcome::Artifact(artifact)
}
+ AuthoredAtomicCommand::RecordDelivery(value) => {
+ let claim_id = value.claim_command().commit_id();
+ let original = candidate
+ .authored_atomic_receipts
+ .iter()
+ .find(|receipt| receipt.commit_id() == claim_id)
+ .ok_or(Error::AtomicWorkflowMismatch)?;
+ let plan = candidate
+ .authored_delivery_plans
+ .iter_mut()
+ .find(|plan| plan.plan_id() == value.plan_id())
+ .ok_or(Error::InvalidAuthoredDeliveryPlan)?;
+ value.apply_to(plan, original)?;
+ AuthoredAtomicOutcome::DeliveryPlan(plan.clone())
+ }
AuthoredAtomicCommand::ApplyDelivery(value) => {
let plan = candidate
.authored_delivery_plans
@@ -1800,13 +1815,19 @@ impl AuthoredAtomicStorage for MemoryStorage {
if plan.revision() != value.expected_revision() {
return Err(Error::InvalidAuthoredDeliveryPlan);
}
- plan.cancel(value.cancelled_at_unix_ms())?;
+ plan.request_stop(value.cancelled_at_unix_ms())?;
AuthoredAtomicOutcome::DeliveryPlan(plan.clone())
}
},
};
let committed_at = match (&command, &outcome) {
(
+ AuthoredAtomicCommand::RecordDelivery(_),
+ AuthoredAtomicOutcome::DeliveryPlan(plan),
+ ) => command
+ .requested_at_unix_ms()
+ .max(plan.updated_at_unix_ms()),
+ (
AuthoredAtomicCommand::RecordSigned(_),
AuthoredAtomicOutcome::Artifact(artifact),
) => command
diff --git a/crates/storage/tests/authored_delivery_facts.rs b/crates/storage/tests/authored_delivery_facts.rs
@@ -0,0 +1,795 @@
+use futures_executor::block_on;
+use radroots_storage::{
+ Error,
+ atomic::{AtomicCommitDigest, AtomicCommitDisposition},
+ authored::WorkClaim,
+ authored_atomic::{
+ ApplyDeliveryAttempt, AuthoredAtomicCommand, AuthoredAtomicOutcome, AuthoredAtomicReceipt,
+ AuthoredAtomicStorage, CancelAuthoredTarget, CancelAuthoredWork, ClaimAuthoredTarget,
+ ClaimAuthoredWork, RecordDeliveryFact, WorkFence,
+ },
+ authored_delivery::{
+ AuthoredDeliveryPlan, AuthoredDeliveryState, DELIVERY_PLAN_ATTEMPTS_MAX,
+ DeliveryAttemptOutcome,
+ },
+ event::SourceGeneration,
+ memory::MemoryStorage,
+};
+use radroots_transport::{
+ DeliveryReceipt, SinkFailure,
+ outcome::{DeliveryOutcome, Retryability},
+ policy::SatisfactionState,
+ sink::DeliveryTargetReceipt,
+};
+use std::num::NonZeroU64;
+
+#[path = "authored_signing/fixture.rs"]
+mod fixture;
+use fixture::*;
+
+fn plan(storage: &MemoryStorage) -> AuthoredDeliveryPlan {
+ block_on(storage.authored_delivery_plan(ids().2))
+ .unwrap()
+ .unwrap()
+}
+
+fn prepared() -> (MemoryStorage, WorkClaim, AuthoredAtomicReceipt) {
+ let storage = MemoryStorage::new(SourceGeneration::new([1; 32]).unwrap());
+ let (command, event) = prepare();
+ block_on(storage.execute_authored(command)).unwrap();
+ let signing = claim(NonZeroU64::MIN, 1, 11);
+ block_on(
+ storage.execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new(
+ ClaimAuthoredTarget::ArtifactSigning(ids().1),
+ signing.clone(),
+ ))),
+ )
+ .unwrap();
+ block_on(storage.execute_authored(record(event, signing, 12))).unwrap();
+ let delivery = claim(plan(&storage).revision(), 2, 13);
+ let receipt = block_on(storage.execute_authored(AuthoredAtomicCommand::Claim(
+ ClaimAuthoredWork::new(ClaimAuthoredTarget::DeliveryPlan(ids().2), delivery.clone()),
+ )))
+ .unwrap();
+ (storage, delivery, receipt)
+}
+
+fn outcome(plan: &AuthoredDeliveryPlan, accepted: bool) -> DeliveryAttemptOutcome {
+ let request = plan.request().unwrap();
+ DeliveryAttemptOutcome::Receipt(
+ DeliveryReceipt::for_request(
+ request,
+ request
+ .target_set()
+ .targets()
+ .iter()
+ .cloned()
+ .map(|target| {
+ DeliveryTargetReceipt::attempted(
+ target,
+ if accepted {
+ DeliveryOutcome::accepted()
+ } else {
+ DeliveryOutcome::unavailable()
+ },
+ )
+ })
+ .collect(),
+ )
+ .unwrap(),
+ )
+}
+
+fn fact(
+ plan: &AuthoredDeliveryPlan,
+ claim: WorkClaim,
+ accepted: bool,
+ at: u64,
+) -> AuthoredAtomicCommand {
+ AuthoredAtomicCommand::RecordDelivery(
+ RecordDeliveryFact::new(
+ plan.plan_id(),
+ plan.artifact_id(),
+ claim,
+ outcome(plan, accepted),
+ at,
+ )
+ .unwrap(),
+ )
+}
+
+fn stop(storage: &MemoryStorage, at: u64) {
+ block_on(
+ storage.execute_authored(AuthoredAtomicCommand::Cancel(
+ CancelAuthoredWork::new(
+ CancelAuthoredTarget::DeliveryPlan(ids().2),
+ plan(storage).revision(),
+ at,
+ )
+ .unwrap(),
+ )),
+ )
+ .unwrap();
+}
+
+#[test]
+fn expired_cancelled_delivery_retains_exact_evidence_and_idempotent_first_time() {
+ let (storage, active, original) = prepared();
+ stop(&storage, 20);
+ let before = plan(&storage);
+ let command = fact(&before, active.clone(), true, 50);
+ let receipt = block_on(storage.execute_authored(command.clone())).unwrap();
+ assert!(receipt.matches_command(&command));
+ let after = plan(&storage);
+ assert_eq!(after.state(), AuthoredDeliveryState::Cancelled);
+ assert_eq!(after.stop_requested_at_unix_ms(), Some(20));
+ assert_eq!(after.request().unwrap().payload().event().raw_json(), RAW);
+ assert_eq!(
+ after.delivery_satisfaction().unwrap(),
+ SatisfactionState::Satisfied
+ );
+ assert_eq!(after.revision(), before.revision());
+ assert_eq!(after.updated_at_unix_ms(), before.updated_at_unix_ms());
+ assert_eq!(after.attempt_count(), 0);
+ assert_eq!(after.delivery_facts()[0].claim(), &active);
+ assert_eq!(after.delivery_facts()[0].observed_at_unix_ms(), 50);
+ let replay_command = fact(&after, active.clone(), true, 90);
+ assert_eq!(replay_command.commit_id(), command.commit_id());
+ let replay = block_on(storage.execute_authored(replay_command)).unwrap();
+ assert_eq!(replay.disposition(), AtomicCommitDisposition::Replay);
+ assert_eq!(replay.committed_at_unix_ms(), 50);
+ assert_eq!(plan(&storage), after);
+ assert_eq!(
+ block_on(storage.authored_receipt(original.commit_id()))
+ .unwrap()
+ .unwrap(),
+ original
+ );
+ stop(&storage, 99);
+ assert_eq!(plan(&storage), after);
+ assert!(block_on(storage.execute_authored(fact(&after, active.clone(), false, 95))).is_err());
+ let fenced = AuthoredAtomicCommand::ApplyDelivery(
+ ApplyDeliveryAttempt::new(
+ ids().2,
+ WorkFence::new(*active.token(), active.generation(), active.row_revision()).unwrap(),
+ outcome(&after, true),
+ None,
+ 50,
+ )
+ .unwrap(),
+ );
+ assert_eq!(
+ block_on(storage.execute_authored(fenced)),
+ Err(Error::DeliveryPlanClaimConflict)
+ );
+ assert!(
+ block_on(
+ storage.execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new(
+ ClaimAuthoredTarget::DeliveryPlan(ids().2),
+ claim(after.revision(), 9, 100)
+ )))
+ )
+ .is_err()
+ );
+}
+
+#[test]
+fn superseded_result_does_not_steal_new_claim_or_regress_acceptance() {
+ let (storage, old, _) = prepared();
+ let newer = claim(plan(&storage).revision(), 3, 40);
+ block_on(
+ storage.execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new(
+ ClaimAuthoredTarget::DeliveryPlan(ids().2),
+ newer.clone(),
+ ))),
+ )
+ .unwrap();
+ let before = plan(&storage);
+ block_on(storage.execute_authored(fact(&before, old, true, 30))).unwrap();
+ let after = plan(&storage);
+ assert_eq!(after.claim_evidence(), Some(&newer));
+ assert_eq!(after.revision(), before.revision());
+ assert_eq!(after.updated_at_unix_ms(), 40);
+ let partial = SinkFailure::for_request(
+ after.request().unwrap(),
+ "sink_lost",
+ Retryability::Retryable,
+ Some(65),
+ Some("connection lost".into()),
+ vec![],
+ )
+ .unwrap();
+ let command = AuthoredAtomicCommand::RecordDelivery(
+ RecordDeliveryFact::new(
+ ids().2,
+ ids().1,
+ newer.clone(),
+ DeliveryAttemptOutcome::SinkFailure(partial.clone()),
+ 45,
+ )
+ .unwrap(),
+ );
+ block_on(storage.execute_authored(command)).unwrap();
+ assert_eq!(
+ plan(&storage).delivery_satisfaction().unwrap(),
+ SatisfactionState::Satisfied
+ );
+ assert_eq!(plan(&storage).claim_evidence(), Some(&newer));
+ let retry = radroots_storage::authored::RetrySchedule::new(
+ std::num::NonZeroU32::MIN,
+ 65,
+ radroots_storage::authored::WorkFailure::new(
+ "sink_lost",
+ radroots_storage::authored::WorkPhase::Delivery,
+ radroots_storage::authored::FailureClass::Retryable,
+ Some(65),
+ Some("connection lost".into()),
+ )
+ .unwrap(),
+ )
+ .unwrap();
+ let fenced = AuthoredAtomicCommand::ApplyDelivery(
+ ApplyDeliveryAttempt::new(
+ ids().2,
+ WorkFence::new(*newer.token(), newer.generation(), newer.row_revision()).unwrap(),
+ DeliveryAttemptOutcome::SinkFailure(partial),
+ Some(retry.clone()),
+ 46,
+ )
+ .unwrap(),
+ );
+ block_on(storage.execute_authored(fenced)).unwrap();
+ let transitioned = plan(&storage);
+ assert_eq!(transitioned.state(), AuthoredDeliveryState::Retryable);
+ assert_eq!(transitioned.retry(), Some(&retry));
+ assert_eq!(transitioned.attempt_count(), 1);
+ assert_eq!(
+ transitioned.delivery_satisfaction().unwrap(),
+ SatisfactionState::Satisfied
+ );
+ stop(&storage, 47);
+ assert_eq!(
+ plan(&storage).delivery_satisfaction().unwrap(),
+ SatisfactionState::Satisfied
+ );
+}
+
+#[test]
+fn backend_provenance_rejects_every_forged_complete_claim_and_wrong_artifact() {
+ let (storage, original, _) = prepared();
+ let before = plan(&storage);
+ for changed in 0..6 {
+ let forged = WorkClaim::new(
+ if changed == 0 {
+ [9; 16]
+ } else {
+ *original.token()
+ },
+ if changed == 1 {
+ "forged"
+ } else {
+ original.owner()
+ },
+ if changed == 2 {
+ NonZeroU64::new(99).unwrap()
+ } else {
+ original.generation()
+ },
+ if changed == 3 {
+ 14
+ } else {
+ original.acquired_at_unix_ms()
+ },
+ if changed == 4 {
+ 90
+ } else {
+ original.expires_at_unix_ms()
+ },
+ if changed == 5 {
+ NonZeroU64::new(99).unwrap()
+ } else {
+ original.row_revision()
+ },
+ )
+ .unwrap();
+ let command = fact(&before, forged, true, 100);
+ assert_eq!(
+ block_on(storage.execute_authored(command.clone())),
+ Err(Error::AtomicWorkflowMismatch)
+ );
+ assert!(
+ block_on(storage.authored_receipt(command.commit_id()))
+ .unwrap()
+ .is_none()
+ );
+ assert_eq!(plan(&storage), before);
+ }
+ let wrong = AuthoredAtomicCommand::RecordDelivery(
+ RecordDeliveryFact::new(
+ ids().2,
+ radroots_storage::authored::AuthoredArtifactId::new([99; 16]).unwrap(),
+ original.clone(),
+ outcome(&before, true),
+ 50,
+ )
+ .unwrap(),
+ );
+ assert_eq!(
+ block_on(storage.execute_authored(wrong)),
+ Err(Error::AtomicWorkflowMismatch)
+ );
+ assert!(
+ RecordDeliveryFact::new(ids().2, ids().1, original, outcome(&before, true), 12).is_err()
+ );
+}
+
+#[test]
+fn late_fact_rejects_rebound_raw_request_and_caller_receipt_with_wrong_time() {
+ let (storage, active, original) = prepared();
+ let before = plan(&storage);
+ let value =
+ RecordDeliveryFact::new(ids().2, ids().1, active.clone(), outcome(&before, true), 50)
+ .unwrap();
+ let request = before
+ .intent()
+ .materialize(radroots_transport::sink::DeliveryPayload::new(event(
+ OTHER_RAW,
+ )))
+ .unwrap();
+ let mut wire = serde_json::to_value(&before).unwrap();
+ wire["request"] = serde_json::to_value(request).unwrap();
+ let mut rebound: AuthoredDeliveryPlan = serde_json::from_value(wire).unwrap();
+ let saved = rebound.clone();
+ assert_eq!(
+ value.apply_to(&mut rebound, &original),
+ Err(Error::AtomicWorkflowMismatch)
+ );
+ assert_eq!(rebound, saved);
+ let wrong = AuthoredAtomicReceipt::from_durable_parts(
+ original.commit_id(),
+ original.digest(),
+ AtomicCommitDisposition::Committed,
+ 14,
+ original.outcome().clone(),
+ )
+ .unwrap();
+ let mut unchanged = before.clone();
+ assert_eq!(
+ value.apply_to(&mut unchanged, &wrong),
+ Err(Error::AtomicWorkflowMismatch)
+ );
+ assert_eq!(unchanged, before);
+}
+
+#[test]
+fn invalid_result_binding_rolls_back_and_changed_failure_details_have_distinct_ids() {
+ let (storage, active, _) = prepared();
+ let before = plan(&storage);
+ let request = before.request().unwrap();
+ let wrong = radroots_transport::DeliveryRequest::new(
+ "wrong-request",
+ request.payload().clone(),
+ request.target_set().clone(),
+ request.satisfaction().clone(),
+ request.deadline_unix_ms(),
+ )
+ .unwrap();
+ let wrong = DeliveryReceipt::for_request(
+ &wrong,
+ wrong
+ .target_set()
+ .targets()
+ .iter()
+ .cloned()
+ .map(|target| DeliveryTargetReceipt::attempted(target, DeliveryOutcome::accepted()))
+ .collect(),
+ )
+ .unwrap();
+ let command = AuthoredAtomicCommand::RecordDelivery(
+ RecordDeliveryFact::new(
+ ids().2,
+ ids().1,
+ active.clone(),
+ DeliveryAttemptOutcome::Receipt(wrong),
+ 50,
+ )
+ .unwrap(),
+ );
+ assert_eq!(
+ block_on(storage.execute_authored(command.clone())),
+ Err(Error::InvalidAuthoredDeliveryPlan)
+ );
+ assert_eq!(plan(&storage), before);
+ assert!(
+ block_on(storage.authored_receipt(command.commit_id()))
+ .unwrap()
+ .is_none()
+ );
+ let mut identities = std::collections::BTreeSet::new();
+ for retry in [None, Some(65), Some(66)] {
+ for message in [
+ None,
+ Some("connection lost".to_owned()),
+ Some("connection closed".to_owned()),
+ ] {
+ let failure = SinkFailure::for_request(
+ request,
+ "sink_lost",
+ Retryability::Retryable,
+ retry,
+ message,
+ vec![],
+ )
+ .unwrap();
+ let command = AuthoredAtomicCommand::RecordDelivery(
+ RecordDeliveryFact::new(
+ ids().2,
+ ids().1,
+ active.clone(),
+ DeliveryAttemptOutcome::SinkFailure(failure),
+ 50,
+ )
+ .unwrap(),
+ );
+ assert!(identities.insert(*command.commit_id().as_bytes()));
+ }
+ }
+ assert_eq!(identities.len(), 9);
+}
+
+#[test]
+fn stop_after_satisfaction_and_legacy_cancelled_snapshot_preserve_intent() {
+ let (storage, active, _) = prepared();
+ let before = plan(&storage);
+ let command = AuthoredAtomicCommand::ApplyDelivery(
+ ApplyDeliveryAttempt::new(
+ ids().2,
+ WorkFence::new(*active.token(), active.generation(), active.row_revision()).unwrap(),
+ outcome(&before, true),
+ None,
+ 14,
+ )
+ .unwrap(),
+ );
+ block_on(storage.execute_authored(command)).unwrap();
+ stop(&storage, 15);
+ let after = plan(&storage);
+ assert_eq!(after.state(), AuthoredDeliveryState::Satisfied);
+ assert_eq!(after.stop_requested_at_unix_ms(), Some(15));
+ assert_eq!(
+ after.delivery_satisfaction().unwrap(),
+ SatisfactionState::Satisfied
+ );
+ let (storage, _, _) = prepared();
+ stop(&storage, 20);
+ let expected = plan(&storage);
+ let mut legacy = serde_json::to_value(&expected).unwrap();
+ legacy.as_object_mut().unwrap().remove("delivery_facts");
+ legacy
+ .as_object_mut()
+ .unwrap()
+ .remove("stop_requested_at_unix_ms");
+ assert_eq!(
+ serde_json::from_value::<AuthoredDeliveryPlan>(legacy).unwrap(),
+ expected
+ );
+ let mut pending = before.clone();
+ assert!(pending.request_stop(9).is_err());
+ assert_eq!(pending, before);
+ assert!(pending.request_stop(12).is_err());
+ assert_eq!(pending, before);
+}
+
+#[test]
+fn receipt_binding_rejects_wrong_outcome_digest_and_observation_time() {
+ let (storage, active, _) = prepared();
+ let before = plan(&storage);
+ let command = fact(&before, active, true, 50);
+ assert!(
+ AuthoredAtomicReceipt::new(
+ &command,
+ AtomicCommitDisposition::Committed,
+ 50,
+ AuthoredAtomicOutcome::DeliveryPlan(before.clone())
+ )
+ .is_err()
+ );
+ let receipt = block_on(storage.execute_authored(command.clone())).unwrap();
+ assert!(
+ AuthoredAtomicReceipt::new(
+ &command,
+ AtomicCommitDisposition::Committed,
+ 49,
+ receipt.outcome().clone()
+ )
+ .is_err()
+ );
+ let wrong = AuthoredAtomicReceipt::from_durable_parts(
+ receipt.commit_id(),
+ AtomicCommitDigest::new([99; 32]),
+ AtomicCommitDisposition::Committed,
+ 50,
+ receipt.outcome().clone(),
+ )
+ .unwrap();
+ assert!(!wrong.matches_command(&command));
+ let wrong = AuthoredAtomicReceipt::from_durable_parts(
+ receipt.commit_id(),
+ receipt.digest(),
+ AtomicCommitDisposition::Committed,
+ 50,
+ AuthoredAtomicOutcome::Artifact(
+ block_on(storage.authored_artifact(ids().1))
+ .unwrap()
+ .unwrap(),
+ ),
+ )
+ .unwrap();
+ assert!(!wrong.matches_command(&command));
+}
+
+#[test]
+fn fact_capacity_and_structural_corruption_fail_without_evicting_evidence() {
+ let (storage, active, original) = prepared();
+ let before = plan(&storage);
+ let command = fact(&before, active.clone(), true, 50);
+ block_on(storage.execute_authored(command)).unwrap();
+ let after = plan(&storage);
+ let base = serde_json::to_value(&after).unwrap();
+ for (path, value) in [
+ ("observed_at_unix_ms", serde_json::json!(1)),
+ ("claim", serde_json::json!(null)),
+ ("outcome", serde_json::json!(null)),
+ ] {
+ let mut corrupt = base.clone();
+ corrupt["delivery_facts"][0][path] = value;
+ assert!(serde_json::from_value::<AuthoredDeliveryPlan>(corrupt).is_err());
+ }
+ let mut duplicate = base.clone();
+ duplicate["delivery_facts"]
+ .as_array_mut()
+ .unwrap()
+ .push(base["delivery_facts"][0].clone());
+ assert!(serde_json::from_value::<AuthoredDeliveryPlan>(duplicate).is_err());
+ let mut full = base.clone();
+ full["revision"] = serde_json::json!(2048);
+ full["claim"] = serde_json::Value::Null;
+ let entries = (1..=DELIVERY_PLAN_ATTEMPTS_MAX)
+ .map(|index| {
+ let mut entry = base["delivery_facts"][0].clone();
+ let claim = WorkClaim::new(
+ [7; 16],
+ "capacity",
+ NonZeroU64::new(u64::from(index)).unwrap(),
+ 13,
+ 33,
+ NonZeroU64::new(u64::from(index)).unwrap(),
+ )
+ .unwrap();
+ entry["claim"] = serde_json::to_value(claim).unwrap();
+ entry
+ })
+ .collect::<Vec<_>>();
+ full["delivery_facts"] = serde_json::to_value(entries).unwrap();
+ let mut full_plan: AuthoredDeliveryPlan = serde_json::from_value(full.clone()).unwrap();
+ let checkpoint = full_plan.clone();
+ let value =
+ RecordDeliveryFact::new(ids().2, ids().1, active, outcome(&before, true), 50).unwrap();
+ assert_eq!(
+ value.apply_to(&mut full_plan, &original),
+ Err(Error::DeliveryAttemptOverflow)
+ );
+ assert_eq!(full_plan, checkpoint);
+ assert!(
+ full_plan
+ .claim(claim(full_plan.revision(), 9, 60), 60)
+ .is_err()
+ );
+ let mut replay = full.clone();
+ replay["delivery_facts"][0] = base["delivery_facts"][0].clone();
+ let mut replay_plan: AuthoredDeliveryPlan = serde_json::from_value(replay).unwrap();
+ let checkpoint = replay_plan.clone();
+ value.apply_to(&mut replay_plan, &original).unwrap();
+ assert_eq!(replay_plan, checkpoint);
+ full["delivery_facts"]
+ .as_array_mut()
+ .unwrap()
+ .push(base["delivery_facts"][0].clone());
+ assert!(serde_json::from_value::<AuthoredDeliveryPlan>(full).is_err());
+}
+
+#[test]
+fn corrupt_original_claim_receipts_and_rebound_current_rows_are_rejected() {
+ let (storage, active, original) = prepared();
+ let before = plan(&storage);
+ let value =
+ RecordDeliveryFact::new(ids().2, ids().1, active, outcome(&before, true), 50).unwrap();
+ let AuthoredAtomicOutcome::DeliveryPlan(original_plan) = original.outcome() else {
+ unreachable!()
+ };
+ for updates in [
+ vec![("plan_id", serde_json::json!(vec![9_u8; 16]))],
+ vec![("artifact_id", serde_json::json!(vec![9_u8; 16]))],
+ vec![("created_at_unix_ms", serde_json::json!(9))],
+ vec![("request", serde_json::Value::Null)],
+ vec![("claim", serde_json::Value::Null)],
+ ] {
+ let mut wire = serde_json::to_value(original_plan).unwrap();
+ for (key, replacement) in updates {
+ wire[key] = replacement;
+ }
+ let altered: AuthoredDeliveryPlan = serde_json::from_value(wire).unwrap();
+ let forged = AuthoredAtomicReceipt::from_durable_parts(
+ original.commit_id(),
+ original.digest(),
+ AtomicCommitDisposition::Committed,
+ original.committed_at_unix_ms(),
+ AuthoredAtomicOutcome::DeliveryPlan(altered),
+ )
+ .unwrap();
+ let mut current = before.clone();
+ assert_eq!(
+ value.apply_to(&mut current, &forged),
+ Err(Error::AtomicWorkflowMismatch)
+ );
+ assert_eq!(current, before);
+ }
+ for updates in [
+ vec![("plan_id", serde_json::json!(vec![9_u8; 16]))],
+ vec![("artifact_id", serde_json::json!(vec![9_u8; 16]))],
+ vec![("created_at_unix_ms", serde_json::json!(9))],
+ vec![
+ ("claim", serde_json::Value::Null),
+ ("revision", serde_json::json!(2)),
+ ],
+ vec![
+ ("claim", serde_json::Value::Null),
+ ("updated_at_unix_ms", serde_json::json!(12)),
+ ],
+ ] {
+ let mut wire = serde_json::to_value(&before).unwrap();
+ for (key, replacement) in updates {
+ wire[key] = replacement;
+ }
+ let mut current: AuthoredDeliveryPlan = serde_json::from_value(wire).unwrap();
+ let unchanged = current.clone();
+ assert_eq!(
+ value.apply_to(&mut current, &original),
+ Err(Error::AtomicWorkflowMismatch)
+ );
+ assert_eq!(current, unchanged);
+ }
+ let wrong_kind = AuthoredAtomicReceipt::from_durable_parts(
+ original.commit_id(),
+ original.digest(),
+ AtomicCommitDisposition::Committed,
+ original.committed_at_unix_ms(),
+ AuthoredAtomicOutcome::Artifact(
+ block_on(storage.authored_artifact(ids().1))
+ .unwrap()
+ .unwrap(),
+ ),
+ )
+ .unwrap();
+ let mut current = before.clone();
+ assert_eq!(
+ value.apply_to(&mut current, &wrong_kind),
+ Err(Error::AtomicWorkflowMismatch)
+ );
+ let wrong_digest = AuthoredAtomicReceipt::from_durable_parts(
+ original.commit_id(),
+ AtomicCommitDigest::new([99; 32]),
+ AtomicCommitDisposition::Committed,
+ original.committed_at_unix_ms(),
+ original.outcome().clone(),
+ )
+ .unwrap();
+ assert_eq!(
+ value.apply_to(&mut current, &wrong_digest),
+ Err(Error::AtomicWorkflowMismatch)
+ );
+ assert_eq!(current, before);
+}
+
+#[test]
+fn forged_fact_snapshots_cannot_bypass_stop_time_claim_or_request_invariants() {
+ let (storage, active, _) = prepared();
+ let before = plan(&storage);
+ block_on(storage.execute_authored(fact(&before, active, true, 50))).unwrap();
+ let current = plan(&storage);
+ let base = serde_json::to_value(¤t).unwrap();
+ for at in [9, 12, 99] {
+ let mut wire = base.clone();
+ wire["stop_requested_at_unix_ms"] = serde_json::json!(at);
+ assert!(serde_json::from_value::<AuthoredDeliveryPlan>(wire).is_err());
+ }
+ for (at, revision) in [(9, 2), (15, 2), (13, 3)] {
+ let forged = WorkClaim::new(
+ [8; 16],
+ "invalid-fact",
+ NonZeroU64::MIN,
+ at,
+ at + 20,
+ NonZeroU64::new(revision).unwrap(),
+ )
+ .unwrap();
+ let mut wire = base.clone();
+ wire["delivery_facts"][0]["claim"] = serde_json::to_value(forged).unwrap();
+ assert!(serde_json::from_value::<AuthoredDeliveryPlan>(wire).is_err());
+ }
+ let mut wire = base;
+ wire["request"] = serde_json::Value::Null;
+ assert!(serde_json::from_value::<AuthoredDeliveryPlan>(wire).is_err());
+}
+
+#[test]
+fn receipt_matching_checks_each_plan_fact_and_monotonic_row_time() {
+ let (storage, active, _) = prepared();
+ let before = plan(&storage);
+ let command = fact(&before, active, true, 50);
+ let receipt = block_on(storage.execute_authored(command.clone())).unwrap();
+ let current = plan(&storage);
+ let base = serde_json::to_value(¤t).unwrap();
+ for key in ["plan_id", "artifact_id"] {
+ let mut wire = base.clone();
+ wire[key] = serde_json::json!(vec![9_u8; 16]);
+ let altered: AuthoredDeliveryPlan = serde_json::from_value(wire).unwrap();
+ let forged = AuthoredAtomicReceipt::from_durable_parts(
+ receipt.commit_id(),
+ receipt.digest(),
+ AtomicCommitDisposition::Committed,
+ 50,
+ AuthoredAtomicOutcome::DeliveryPlan(altered.clone()),
+ )
+ .unwrap();
+ assert!(!forged.matches_command(&command));
+ assert!(
+ AuthoredAtomicReceipt::new(
+ &command,
+ AtomicCommitDisposition::Committed,
+ 50,
+ AuthoredAtomicOutcome::DeliveryPlan(altered)
+ )
+ .is_err()
+ );
+ }
+ for change_claim in [false, true] {
+ let mut wire = base.clone();
+ if change_claim {
+ let forged = WorkClaim::new(
+ [8; 16],
+ "different",
+ NonZeroU64::MIN,
+ 13,
+ 33,
+ NonZeroU64::new(2).unwrap(),
+ )
+ .unwrap();
+ wire["delivery_facts"][0]["claim"] = serde_json::to_value(forged).unwrap();
+ } else {
+ wire["delivery_facts"][0]["outcome"] =
+ serde_json::to_value(outcome(¤t, false)).unwrap();
+ }
+ let altered: AuthoredDeliveryPlan = serde_json::from_value(wire).unwrap();
+ let forged = AuthoredAtomicReceipt::from_durable_parts(
+ receipt.commit_id(),
+ receipt.digest(),
+ AtomicCommitDisposition::Committed,
+ 50,
+ AuthoredAtomicOutcome::DeliveryPlan(altered),
+ )
+ .unwrap();
+ assert!(!forged.matches_command(&command));
+ }
+ let mut stopped = current;
+ stopped.request_stop(100).unwrap();
+ assert!(
+ AuthoredAtomicReceipt::new(
+ &command,
+ AtomicCommitDisposition::Committed,
+ 60,
+ AuthoredAtomicOutcome::DeliveryPlan(stopped)
+ )
+ .is_err()
+ );
+}
diff --git a/crates/storage_sqlite/src/authored.rs b/crates/storage_sqlite/src/authored.rs
@@ -1,4 +1,6 @@
use crate::SqliteStorage;
+#[path = "authored_delivery_facts.rs"]
+mod delivery_facts;
use radroots_storage::{
Error,
atomic::{AtomicCommitDigest, AtomicCommitDisposition, AtomicCommitId},
@@ -205,6 +207,11 @@ async fn commit_outcome(
outcome: AuthoredAtomicOutcome,
) -> Result<AuthoredAtomicReceipt, Error> {
let committed_at = match (command, &outcome) {
+ (AuthoredAtomicCommand::RecordDelivery(_), AuthoredAtomicOutcome::DeliveryPlan(plan)) => {
+ command
+ .requested_at_unix_ms()
+ .max(plan.updated_at_unix_ms())
+ }
(AuthoredAtomicCommand::RecordSigned(_), AuthoredAtomicOutcome::Artifact(artifact)) => {
command
.requested_at_unix_ms()
@@ -406,6 +413,23 @@ async fn execute_command(
persist_artifact(transaction, &artifact).await?;
Ok(AuthoredAtomicOutcome::Artifact(artifact))
}
+ AuthoredAtomicCommand::RecordDelivery(value) => {
+ let row = sqlx::query(
+ "SELECT commit_id, commit_digest, requested_at_unix_ms,
+ committed_at_unix_ms, receipt
+ FROM radroots_runtime_authored_atomic_commits WHERE commit_id = ?",
+ )
+ .bind(value.claim_command().commit_id().as_bytes().as_slice())
+ .fetch_optional(&mut **transaction)
+ .await
+ .map_err(map_backend)?
+ .ok_or(Error::AtomicWorkflowMismatch)?;
+ let original = decode_receipt_row(&row)?;
+ let mut plan = load_plan_tx(transaction, value.plan_id()).await?;
+ value.apply_to(&mut plan, &original)?;
+ persist_plan(transaction, &plan).await?;
+ Ok(AuthoredAtomicOutcome::DeliveryPlan(plan))
+ }
AuthoredAtomicCommand::ApplyDelivery(value) => {
let mut plan = load_plan_tx(transaction, value.plan_id()).await?;
match value.outcome().clone() {
@@ -498,7 +522,7 @@ async fn execute_command(
CancelAuthoredTarget::DeliveryPlan(id) => {
let mut plan = load_plan_tx(transaction, *id).await?;
require_revision(plan.revision().get(), value.expected_revision().get())?;
- plan.cancel(value.cancelled_at_unix_ms())?;
+ plan.request_stop(value.cancelled_at_unix_ms())?;
persist_plan(transaction, &plan).await?;
Ok(AuthoredAtomicOutcome::DeliveryPlan(plan))
}
@@ -672,6 +696,15 @@ pub(crate) async fn persist_plan(
transaction: &mut sqlx::Transaction<'_, Sqlite>,
plan: &AuthoredDeliveryPlan,
) -> Result<(), Error> {
+ persist_plan_v11(transaction, plan).await?;
+ delivery_facts::persist(transaction, plan).await
+}
+
+// Legacy conversion runs before the delivery-facts forward migration.
+pub(crate) async fn persist_plan_v11(
+ transaction: &mut sqlx::Transaction<'_, Sqlite>,
+ plan: &AuthoredDeliveryPlan,
+) -> Result<(), Error> {
let claim = plan.claim_evidence();
sqlx::query(
"INSERT INTO radroots_runtime_authored_delivery_plans (
@@ -787,38 +820,47 @@ async fn load_plan_tx(
transaction: &mut sqlx::Transaction<'_, Sqlite>,
plan_id: AuthoredDeliveryPlanId,
) -> Result<AuthoredDeliveryPlan, Error> {
- let plan =
- sqlx::query("SELECT * FROM radroots_runtime_authored_delivery_plans WHERE plan_id = ?")
- .bind(plan_id.as_bytes().as_slice())
- .fetch_optional(&mut **transaction)
- .await
- .map_err(map_backend)?
- .as_ref()
- .map(decode_plan_row)
- .transpose()?
- .ok_or(Error::InvalidAuthoredDeliveryPlan)?;
- validate_plan_children_tx(transaction, &plan).await?;
- Ok(plan)
+ load_optional_plan_tx(transaction, plan_id)
+ .await?
+ .ok_or(Error::InvalidAuthoredDeliveryPlan)
}
-async fn load_plan_pool(
- storage: &SqliteStorage,
+async fn load_optional_plan_tx(
+ transaction: &mut sqlx::Transaction<'_, Sqlite>,
plan_id: AuthoredDeliveryPlanId,
) -> Result<Option<AuthoredDeliveryPlan>, Error> {
let Some(row) =
sqlx::query("SELECT * FROM radroots_runtime_authored_delivery_plans WHERE plan_id = ?")
.bind(plan_id.as_bytes().as_slice())
- .fetch_optional(storage.pool())
+ .fetch_optional(&mut **transaction)
.await
.map_err(map_backend)?
else {
return Ok(None);
};
let plan = decode_plan_row(&row)?;
- validate_plan_children_pool(storage, &plan).await?;
+ validate_plan_children_tx(transaction, &plan).await?;
Ok(Some(plan))
}
+async fn load_plan_pool(
+ storage: &SqliteStorage,
+ plan_id: AuthoredDeliveryPlanId,
+) -> Result<Option<AuthoredDeliveryPlan>, Error> {
+ // Plan and append-only facts must come from the same read snapshot even
+ // when a late result commits between the individual SELECT statements.
+ let mut transaction = storage.pool().begin().await.map_err(map_backend)?;
+ let result = load_optional_plan_tx(&mut transaction, plan_id).await;
+ let rollback = transaction.rollback().await.map_err(map_backend);
+ match result {
+ Ok(plan) => {
+ rollback?;
+ Ok(plan)
+ }
+ Err(error) => Err(error),
+ }
+}
+
async fn validate_plan_children_tx(
transaction: &mut sqlx::Transaction<'_, Sqlite>,
plan: &AuthoredDeliveryPlan,
@@ -839,30 +881,13 @@ async fn validate_plan_children_tx(
.fetch_all(&mut **transaction)
.await
.map_err(map_backend)?;
- validate_plan_children(plan, &targets, &attempts)
-}
-
-async fn validate_plan_children_pool(
- storage: &SqliteStorage,
- plan: &AuthoredDeliveryPlan,
-) -> Result<(), Error> {
- let targets = sqlx::query(
- "SELECT ordinal, target_fingerprint, target_snapshot
- FROM radroots_runtime_authored_delivery_targets WHERE plan_id = ? ORDER BY ordinal",
- )
- .bind(plan.plan_id().as_bytes().as_slice())
- .fetch_all(storage.pool())
- .await
- .map_err(map_backend)?;
- let attempts = sqlx::query(
- "SELECT attempt, satisfaction, recorded_at_unix_ms, outcome_snapshot
- FROM radroots_runtime_authored_delivery_attempts WHERE plan_id = ? ORDER BY attempt",
- )
- .bind(plan.plan_id().as_bytes().as_slice())
- .fetch_all(storage.pool())
- .await
- .map_err(map_backend)?;
- validate_plan_children(plan, &targets, &attempts)
+ validate_plan_children(plan, &targets, &attempts)?;
+ let facts = sqlx::query(delivery_facts::SELECT)
+ .bind(plan.plan_id().as_bytes().as_slice())
+ .fetch_all(&mut **transaction)
+ .await
+ .map_err(map_backend)?;
+ delivery_facts::validate(plan, &facts)
}
fn validate_plan_children(
@@ -966,6 +991,11 @@ fn decode_plan_row(row: &SqliteRow) -> Result<AuthoredDeliveryPlan, Error> {
|| column::<Vec<u8>>(row, "request_digest")?.as_slice() != value.request_digest()
|| column::<String>(row, "state")? != delivery_state_name(value.state())
|| column::<i64>(row, "attempt_count")? != i64::from(value.attempt_count())
+ || column::<Option<i64>>(row, "stop_requested_at_unix_ms")?
+ != value
+ .stop_requested_at_unix_ms()
+ .map(i64_from_u64)
+ .transpose()?
|| !claim_columns_match(row, "claim", value.claim_evidence())?
|| column::<Option<i64>>(row, "retry_not_before_unix_ms")?
!= value
@@ -1148,6 +1178,7 @@ fn command_target(command: &AuthoredAtomicCommand) -> [u8; 16] {
},
AuthoredAtomicCommand::ApplySigned(value) => *value.artifact_id().as_bytes(),
AuthoredAtomicCommand::RecordSigned(value) => *value.artifact_id().as_bytes(),
+ AuthoredAtomicCommand::RecordDelivery(value) => *value.plan_id().as_bytes(),
AuthoredAtomicCommand::ApplyAdmission(value) => *value.artifact_id().as_bytes(),
AuthoredAtomicCommand::ApplyDelivery(value) => *value.plan_id().as_bytes(),
AuthoredAtomicCommand::ApplyFailure(value) => match value.target() {
@@ -1168,7 +1199,9 @@ fn command_phase(command: &AuthoredAtomicCommand) -> &'static str {
AuthoredAtomicCommand::Claim(_) => "claim",
AuthoredAtomicCommand::ApplySigned(_) | AuthoredAtomicCommand::RecordSigned(_) => "signing",
AuthoredAtomicCommand::ApplyAdmission(_) => "admission",
- AuthoredAtomicCommand::ApplyDelivery(_) => "delivery",
+ AuthoredAtomicCommand::ApplyDelivery(_) | AuthoredAtomicCommand::RecordDelivery(_) => {
+ "delivery"
+ }
AuthoredAtomicCommand::ApplyFailure(value) => match value.failure().phase() {
WorkPhase::Signing => "signing_failure",
WorkPhase::Admission => "admission_failure",
@@ -2607,5 +2640,10 @@ mod signed_fact_tests;
#[cfg(test)]
#[cfg_attr(coverage_nightly, coverage(off))]
+#[path = "authored_delivery_fact_tests.rs"]
+mod delivery_fact_tests;
+
+#[cfg(test)]
+#[cfg_attr(coverage_nightly, coverage(off))]
#[path = "authored_signed_fact_fixture.rs"]
pub(crate) mod signed_fact_fixture;
diff --git a/crates/storage_sqlite/src/authored_delivery_fact_tests.rs b/crates/storage_sqlite/src/authored_delivery_fact_tests.rs
@@ -0,0 +1,328 @@
+use super::signed_fact_fixture::{RAW, claim, ids, record};
+use super::*;
+use crate::OpenMode;
+use radroots_storage::{
+ authored::WorkClaim,
+ authored_atomic::{CancelAuthoredWork, ClaimAuthoredWork, RecordDeliveryFact},
+};
+use radroots_transport::{DeliveryReceipt, outcome::DeliveryOutcome, sink::DeliveryTargetReceipt};
+use tempfile::TempDir;
+
+async fn prepared(temp: &TempDir) -> (SqliteStorage, WorkClaim, AuthoredAtomicReceipt) {
+ let (store, event, signing) = super::signed_fact_tests::prepared(temp).await;
+ store
+ .execute_authored(record(event, signing, 12))
+ .await
+ .unwrap();
+ let plan = store
+ .authored_delivery_plan(ids().2)
+ .await
+ .unwrap()
+ .unwrap();
+ let active = claim(plan.revision(), 5, 13);
+ let original = store
+ .execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new(
+ ClaimAuthoredTarget::DeliveryPlan(ids().2),
+ active.clone(),
+ )))
+ .await
+ .unwrap();
+ (store, active, original)
+}
+
+fn fact(plan: &AuthoredDeliveryPlan, active: WorkClaim) -> AuthoredAtomicCommand {
+ let request = plan.request().unwrap();
+ let receipt = DeliveryReceipt::for_request(
+ request,
+ request
+ .target_set()
+ .targets()
+ .iter()
+ .cloned()
+ .map(|target| DeliveryTargetReceipt::attempted(target, DeliveryOutcome::accepted()))
+ .collect(),
+ )
+ .unwrap();
+ AuthoredAtomicCommand::RecordDelivery(
+ RecordDeliveryFact::new(
+ ids().2,
+ ids().1,
+ active,
+ DeliveryAttemptOutcome::Receipt(receipt),
+ 50,
+ )
+ .unwrap(),
+ )
+}
+
+#[tokio::test]
+async fn stopped_late_delivery_reopens_and_preserves_original_receipt_and_exact_raw() {
+ let temp = TempDir::new().unwrap();
+ let (store, active, original) = prepared(&temp).await;
+ let plan = store
+ .authored_delivery_plan(ids().2)
+ .await
+ .unwrap()
+ .unwrap();
+ let command = fact(&plan, active);
+ store
+ .execute_authored(AuthoredAtomicCommand::Cancel(
+ CancelAuthoredWork::new(
+ CancelAuthoredTarget::DeliveryPlan(ids().2),
+ plan.revision(),
+ 20,
+ )
+ .unwrap(),
+ ))
+ .await
+ .unwrap();
+ let receipt = store.execute_authored(command.clone()).await.unwrap();
+ let retained = store
+ .authored_delivery_plan(ids().2)
+ .await
+ .unwrap()
+ .unwrap();
+ assert_eq!(retained.state(), AuthoredDeliveryState::Cancelled);
+ assert_eq!(retained.stop_requested_at_unix_ms(), Some(20));
+ assert_eq!(
+ retained.delivery_satisfaction().unwrap(),
+ SatisfactionState::Satisfied
+ );
+ assert_eq!(
+ retained.request().unwrap().payload().event().raw_json(),
+ RAW
+ );
+ assert_eq!(retained.delivery_facts().len(), 1);
+ for sql in [
+ "UPDATE radroots_runtime_authored_delivery_facts SET observed_at_unix_ms = 99",
+ "DELETE FROM radroots_runtime_authored_delivery_facts",
+ "UPDATE radroots_runtime_authored_delivery_plans SET stop_requested_at_unix_ms = NULL",
+ ] {
+ assert!(sqlx::query(sql).execute(store.pool()).await.is_err());
+ }
+ store.close().await.unwrap();
+ let store = super::signed_fact_tests::open(&temp, OpenMode::ReadWriteExisting).await;
+ assert_eq!(
+ store
+ .authored_delivery_plan(ids().2)
+ .await
+ .unwrap()
+ .unwrap(),
+ retained
+ );
+ assert_eq!(
+ store
+ .authored_receipt(original.commit_id())
+ .await
+ .unwrap()
+ .unwrap(),
+ original
+ );
+ let replay = store.execute_authored(command).await.unwrap();
+ assert_eq!(replay.disposition(), AtomicCommitDisposition::Replay);
+ assert_eq!(replay.outcome(), receipt.outcome());
+ assert_eq!(
+ store
+ .authored_delivery_plan(ids().2)
+ .await
+ .unwrap()
+ .unwrap(),
+ retained
+ );
+ store.close().await.unwrap();
+}
+
+#[tokio::test]
+async fn real_delivery_commit_failure_rolls_back_facts_snapshot_and_receipt() {
+ let temp = TempDir::new().unwrap();
+ let (store, active, _) = prepared(&temp).await;
+ let before = store
+ .authored_delivery_plan(ids().2)
+ .await
+ .unwrap()
+ .unwrap();
+ let command = fact(&before, active);
+ sqlx::query("CREATE TABLE delivery_fact_commit_fault (parent BLOB REFERENCES radroots_runtime_authored_operations(operation_id) DEFERRABLE INITIALLY DEFERRED)").execute(store.pool()).await.unwrap();
+ sqlx::query("CREATE TRIGGER delivery_fact_commit_fault_trigger AFTER INSERT ON radroots_runtime_authored_atomic_commits WHEN NEW.phase IN ('delivery', 'cancel') BEGIN INSERT INTO delivery_fact_commit_fault VALUES (x'99999999999999999999999999999999'); END").execute(store.pool()).await.unwrap();
+ let stop = AuthoredAtomicCommand::Cancel(
+ CancelAuthoredWork::new(
+ CancelAuthoredTarget::DeliveryPlan(ids().2),
+ before.revision(),
+ 20,
+ )
+ .unwrap(),
+ );
+ assert!(store.execute_authored(stop.clone()).await.is_err());
+ assert_eq!(
+ store
+ .authored_delivery_plan(ids().2)
+ .await
+ .unwrap()
+ .unwrap(),
+ before
+ );
+ assert!(
+ store
+ .authored_receipt(stop.commit_id())
+ .await
+ .unwrap()
+ .is_none()
+ );
+ assert!(store.execute_authored(command.clone()).await.is_err());
+ assert_eq!(
+ store
+ .authored_delivery_plan(ids().2)
+ .await
+ .unwrap()
+ .unwrap(),
+ before
+ );
+ assert!(
+ store
+ .authored_receipt(command.commit_id())
+ .await
+ .unwrap()
+ .is_none()
+ );
+ assert_eq!(
+ sqlx::query_scalar::<_, i64>(
+ "SELECT COUNT(*) FROM radroots_runtime_authored_delivery_facts"
+ )
+ .fetch_one(store.pool())
+ .await
+ .unwrap(),
+ 0
+ );
+ assert_eq!(
+ sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM delivery_fact_commit_fault")
+ .fetch_one(store.pool())
+ .await
+ .unwrap(),
+ 0
+ );
+ sqlx::query("DROP TRIGGER delivery_fact_commit_fault_trigger")
+ .execute(store.pool())
+ .await
+ .unwrap();
+ sqlx::query("DROP TABLE delivery_fact_commit_fault")
+ .execute(store.pool())
+ .await
+ .unwrap();
+ store.close().await.unwrap();
+ let store = super::signed_fact_tests::open(&temp, OpenMode::ReadWriteExisting).await;
+ assert_eq!(
+ store
+ .authored_delivery_plan(ids().2)
+ .await
+ .unwrap()
+ .unwrap(),
+ before
+ );
+ store.execute_authored(command.clone()).await.unwrap();
+ assert_eq!(
+ store
+ .authored_delivery_plan(ids().2)
+ .await
+ .unwrap()
+ .unwrap()
+ .delivery_satisfaction()
+ .unwrap(),
+ SatisfactionState::Satisfied
+ );
+ store.close().await.unwrap();
+}
+
+#[tokio::test]
+async fn forged_claim_and_normalized_fact_corruption_fail_closed() {
+ let temp = TempDir::new().unwrap();
+ let (store, active, _) = prepared(&temp).await;
+ let before = store
+ .authored_delivery_plan(ids().2)
+ .await
+ .unwrap()
+ .unwrap();
+ let forged = WorkClaim::new(
+ *active.token(),
+ "forged-owner",
+ active.generation(),
+ active.acquired_at_unix_ms(),
+ active.expires_at_unix_ms(),
+ active.row_revision(),
+ )
+ .unwrap();
+ assert_eq!(
+ store.execute_authored(fact(&before, forged)).await,
+ Err(Error::AtomicWorkflowMismatch)
+ );
+ assert_eq!(
+ store
+ .authored_delivery_plan(ids().2)
+ .await
+ .unwrap()
+ .unwrap(),
+ before
+ );
+ store.execute_authored(fact(&before, active)).await.unwrap();
+ sqlx::query("DROP TRIGGER radroots_runtime_authored_delivery_facts_update_guard")
+ .execute(store.pool())
+ .await
+ .unwrap();
+ sqlx::query("UPDATE radroots_runtime_authored_delivery_facts SET observed_at_unix_ms = 99")
+ .execute(store.pool())
+ .await
+ .unwrap();
+ assert_eq!(
+ store.authored_delivery_plan(ids().2).await,
+ Err(Error::InvalidAuthoredDeliveryPlan)
+ );
+ store.close().await.unwrap();
+}
+
+#[tokio::test]
+async fn delivery_read_snapshot_remains_consistent_across_a_late_commit() {
+ let temp = TempDir::new().unwrap();
+ let (store, active, _) = prepared(&temp).await;
+ let before = store
+ .authored_delivery_plan(ids().2)
+ .await
+ .unwrap()
+ .unwrap();
+ let mut read = store.pool().begin().await.unwrap();
+ let frozen = load_optional_plan_tx(&mut read, ids().2)
+ .await
+ .unwrap()
+ .unwrap();
+ assert_eq!(frozen, before);
+ store.execute_authored(fact(&before, active)).await.unwrap();
+ validate_plan_children_tx(&mut read, &frozen).await.unwrap();
+ assert_eq!(
+ load_optional_plan_tx(&mut read, ids().2)
+ .await
+ .unwrap()
+ .unwrap(),
+ frozen
+ );
+ read.rollback().await.unwrap();
+ let current = store
+ .authored_delivery_plan(ids().2)
+ .await
+ .unwrap()
+ .unwrap();
+ assert_eq!(current.delivery_facts().len(), 1);
+ assert_eq!(current.revision(), before.revision());
+ let mut stale = store.pool().begin().await.unwrap();
+ assert_eq!(
+ delivery_facts::persist(&mut stale, &before).await,
+ Err(Error::InvalidAuthoredDeliveryPlan)
+ );
+ stale.rollback().await.unwrap();
+ assert_eq!(
+ store
+ .authored_delivery_plan(ids().2)
+ .await
+ .unwrap()
+ .unwrap(),
+ current
+ );
+ store.close().await.unwrap();
+}
diff --git a/crates/storage_sqlite/src/authored_delivery_facts.rs b/crates/storage_sqlite/src/authored_delivery_facts.rs
@@ -0,0 +1,79 @@
+//! Normalized, append-only delivery facts in the authored transaction.
+
+use crate::authored::{column, decode_snapshot, encode_snapshot, i64_from_u64, map_backend};
+use radroots_storage::{
+ Error,
+ authored_atomic::{AuthoredAtomicCommand, ClaimAuthoredTarget, ClaimAuthoredWork},
+ authored_delivery::{AuthoredDeliveryFact, AuthoredDeliveryPlan},
+};
+use sqlx::{Sqlite, sqlite::SqliteRow};
+
+pub(crate) const SELECT: &str = "SELECT ordinal, claim_id, observed_at_unix_ms, fact_snapshot FROM radroots_runtime_authored_delivery_facts WHERE plan_id = ? ORDER BY ordinal";
+
+fn claim_id(
+ plan: &AuthoredDeliveryPlan,
+ fact: &AuthoredDeliveryFact,
+) -> radroots_storage::atomic::AtomicCommitId {
+ AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new(
+ ClaimAuthoredTarget::DeliveryPlan(plan.plan_id()),
+ fact.claim().clone(),
+ ))
+ .commit_id()
+}
+
+pub(crate) async fn persist(
+ transaction: &mut sqlx::Transaction<'_, Sqlite>,
+ plan: &AuthoredDeliveryPlan,
+) -> Result<(), Error> {
+ sqlx::query("UPDATE radroots_runtime_authored_delivery_plans SET stop_requested_at_unix_ms = ? WHERE plan_id = ?")
+ .bind(plan.stop_requested_at_unix_ms().map(i64_from_u64).transpose()?)
+ .bind(plan.plan_id().as_bytes().as_slice())
+ .execute(&mut **transaction).await.map_err(map_backend)?;
+ let existing = sqlx::query(SELECT)
+ .bind(plan.plan_id().as_bytes().as_slice())
+ .fetch_all(&mut **transaction)
+ .await
+ .map_err(map_backend)?;
+ if existing.len() > plan.delivery_facts().len() {
+ return Err(Error::InvalidAuthoredDeliveryPlan);
+ }
+ validate_prefix(plan, &existing)?;
+ for (ordinal, fact) in plan
+ .delivery_facts()
+ .iter()
+ .enumerate()
+ .skip(existing.len())
+ {
+ sqlx::query("INSERT INTO radroots_runtime_authored_delivery_facts (plan_id, ordinal, claim_id, observed_at_unix_ms, fact_snapshot) VALUES (?, ?, ?, ?, ?)")
+ .bind(plan.plan_id().as_bytes().as_slice())
+ .bind(i64::try_from(ordinal).map_err(|_| Error::InvalidAuthoredDeliveryPlan)?)
+ .bind(claim_id(plan, fact).as_bytes().as_slice())
+ .bind(i64_from_u64(fact.observed_at_unix_ms())?)
+ .bind(encode_snapshot(fact)?)
+ .execute(&mut **transaction).await.map_err(map_backend)?;
+ }
+ Ok(())
+}
+
+pub(crate) fn validate(plan: &AuthoredDeliveryPlan, rows: &[SqliteRow]) -> Result<(), Error> {
+ if rows.len() != plan.delivery_facts().len() {
+ return Err(Error::InvalidAuthoredDeliveryPlan);
+ }
+ validate_prefix(plan, rows)
+}
+
+fn validate_prefix(plan: &AuthoredDeliveryPlan, rows: &[SqliteRow]) -> Result<(), Error> {
+ for (ordinal, (row, expected)) in rows.iter().zip(plan.delivery_facts()).enumerate() {
+ let decoded = decode_snapshot::<AuthoredDeliveryFact>(column(row, "fact_snapshot")?)?;
+ if column::<i64>(row, "ordinal")?
+ != i64::try_from(ordinal).map_err(|_| Error::InvalidAuthoredDeliveryPlan)?
+ || column::<Vec<u8>>(row, "claim_id")?.as_slice() != claim_id(plan, expected).as_bytes()
+ || column::<i64>(row, "observed_at_unix_ms")?
+ != i64_from_u64(expected.observed_at_unix_ms())?
+ || decoded != *expected
+ {
+ return Err(Error::InvalidAuthoredDeliveryPlan);
+ }
+ }
+ Ok(())
+}
diff --git a/crates/storage_sqlite/src/authored_signed_fact_tests.rs b/crates/storage_sqlite/src/authored_signed_fact_tests.rs
@@ -12,7 +12,7 @@ use tempfile::TempDir;
use super::signed_fact_fixture as fixture;
use fixture::*;
-async fn open(temp: &TempDir, mode: OpenMode) -> SqliteStorage {
+pub(super) async fn open(temp: &TempDir, mode: OpenMode) -> SqliteStorage {
let options = OpenOptions::new(Paths::from_directory(temp.path()).unwrap(), mode);
let options = if matches!(mode, OpenMode::Create) {
options
@@ -24,7 +24,9 @@ async fn open(temp: &TempDir, mode: OpenMode) -> SqliteStorage {
SqliteStorage::open(options).await.unwrap()
}
-async fn prepared(temp: &TempDir) -> (SqliteStorage, radroots_event::SignedEvent, WorkClaim) {
+pub(super) async fn prepared(
+ temp: &TempDir,
+) -> (SqliteStorage, radroots_event::SignedEvent, WorkClaim) {
let store = open(temp, OpenMode::Create).await;
let (preparation, event) = prepare();
store.execute_authored(preparation).await.unwrap();
diff --git a/crates/storage_sqlite/src/migration.rs b/crates/storage_sqlite/src/migration.rs
@@ -371,7 +371,7 @@ async fn metadata(
fn validate_plan(plan: &MigrationPlan) -> Result<(), Error> {
let valid = plan.minimum_version > 0
&& plan.minimum_version <= plan.current_version
- && plan.current_version <= 15
+ && plan.current_version <= 16
&& plan.steps.len() == usize::try_from(plan.current_version).unwrap_or(usize::MAX)
&& plan
.steps
@@ -491,6 +491,7 @@ const fn set_user_version_sql(version: u32) -> Option<&'static str> {
13 => Some("PRAGMA user_version = 13"),
14 => Some("PRAGMA user_version = 14"),
15 => Some("PRAGMA user_version = 15"),
+ 16 => Some("PRAGMA user_version = 16"),
_ => None,
}
}
@@ -717,7 +718,9 @@ mod tests {
.execute(&mut newer)
.await
.expect("application id");
- sqlx::raw_sql("PRAGMA user_version = 16")
+ let newer_version = 17;
+ assert_eq!(newer_version, runtime::CURRENT_VERSION + 1);
+ sqlx::raw_sql("PRAGMA user_version = 17")
.execute(&mut newer)
.await
.expect("newer version");
@@ -726,10 +729,13 @@ mod tests {
Err(Error::SchemaTooNew {
database: RUNTIME_DATABASE,
supported: runtime::CURRENT_VERSION,
- actual: 16,
- })
+ actual,
+ }) if actual == newer_version
));
- assert_eq!(pragma(&mut newer, "user_version").await, 16);
+ assert_eq!(
+ pragma(&mut newer, "user_version").await,
+ i64::from(newer_version)
+ );
let mut wrong_identity = connection().await;
establish_runtime_version(&mut wrong_identity, 1).await;
@@ -868,3 +874,8 @@ mod draft_query_tests;
#[cfg_attr(coverage_nightly, coverage(off))]
#[path = "migration_signed_facts_tests.rs"]
mod signed_facts_tests;
+
+#[cfg(test)]
+#[cfg_attr(coverage_nightly, coverage(off))]
+#[path = "migration_delivery_facts_tests.rs"]
+mod delivery_facts_tests;
diff --git a/crates/storage_sqlite/src/migration/authored_v10.rs b/crates/storage_sqlite/src/migration/authored_v10.rs
@@ -312,7 +312,7 @@ pub(crate) async fn apply(
.await
.map_err(|_| metadata_error())?;
persist_imported_artifact(transaction, &artifact).await?;
- authored::persist_plan(transaction, &plan)
+ authored::persist_plan_v11(transaction, &plan)
.await
.map_err(|_| metadata_error())?;
}
diff --git a/crates/storage_sqlite/src/migration/runtime/0016_authored_delivery_facts.up.sql b/crates/storage_sqlite/src/migration/runtime/0016_authored_delivery_facts.up.sql
@@ -0,0 +1,43 @@
+-- Preserve scheduling state and revision while retaining independent effects.
+ALTER TABLE radroots_runtime_authored_delivery_plans
+ADD COLUMN stop_requested_at_unix_ms INTEGER CHECK (
+ stop_requested_at_unix_ms IS NULL OR (
+ stop_requested_at_unix_ms >= created_at_unix_ms
+ AND stop_requested_at_unix_ms <= updated_at_unix_ms
+ AND state IN ('satisfied', 'exhausted', 'failed_terminal', 'cancelled')
+ AND claim_token IS NULL
+ )
+);
+
+UPDATE radroots_runtime_authored_delivery_plans
+SET stop_requested_at_unix_ms = updated_at_unix_ms WHERE state = 'cancelled';
+
+CREATE TRIGGER radroots_runtime_authored_delivery_stop_guard
+BEFORE UPDATE ON radroots_runtime_authored_delivery_plans
+WHEN OLD.stop_requested_at_unix_ms IS NOT NULL
+ AND NEW.stop_requested_at_unix_ms IS NOT OLD.stop_requested_at_unix_ms
+BEGIN
+ SELECT RAISE(ABORT, 'delivery stop intent is immutable');
+END;
+
+CREATE TABLE radroots_runtime_authored_delivery_facts (
+ plan_id BLOB NOT NULL CHECK (length(plan_id) = 16) REFERENCES radroots_runtime_authored_delivery_plans(plan_id),
+ ordinal INTEGER NOT NULL CHECK (ordinal >= 0 AND ordinal < 1024),
+ claim_id BLOB NOT NULL CHECK (length(claim_id) = 16) REFERENCES radroots_runtime_authored_atomic_commits(commit_id),
+ observed_at_unix_ms INTEGER NOT NULL CHECK (observed_at_unix_ms > 0),
+ fact_snapshot BLOB NOT NULL CHECK (length(fact_snapshot) BETWEEN 2 AND 4194304),
+ PRIMARY KEY (plan_id, ordinal),
+ UNIQUE (plan_id, claim_id)
+) STRICT;
+
+CREATE TRIGGER radroots_runtime_authored_delivery_facts_update_guard
+BEFORE UPDATE ON radroots_runtime_authored_delivery_facts
+BEGIN
+ SELECT RAISE(ABORT, 'delivery facts are immutable');
+END;
+
+CREATE TRIGGER radroots_runtime_authored_delivery_facts_delete_guard
+BEFORE DELETE ON radroots_runtime_authored_delivery_facts
+BEGIN
+ SELECT RAISE(ABORT, 'delivery facts are immutable');
+END;
diff --git a/crates/storage_sqlite/src/migration/runtime/mod.rs b/crates/storage_sqlite/src/migration/runtime/mod.rs
@@ -6,7 +6,7 @@
/// Lowest runtime schema version this package can recognize.
pub const MINIMUM_VERSION: u32 = 1;
/// Current runtime schema version created by this package.
-pub const CURRENT_VERSION: u32 = 15;
+pub const CURRENT_VERSION: u32 = 16;
const RUNTIME_V1_SQL: &str = include_str!("0001_runtime.up.sql");
const CANONICAL_EVENT_STORAGE_V2_SQL: &str = include_str!("0002_canonical_event_storage.up.sql");
@@ -28,6 +28,8 @@ const AUTHORED_DRAFT_QUERY_V14_SQL: &str =
include_str!("0014_authored_draft_query_metadata.up.sql");
const AUTHORED_SIGNED_FACTS_V15_SQL: &str = include_str!("0015_authored_signed_facts.up.sql");
+const AUTHORED_DELIVERY_FACTS_V16_SQL: &str = include_str!("0016_authored_delivery_facts.up.sql");
+
/// Stable, non-SQL description of one forward runtime migration.
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct MigrationDescriptor {
@@ -716,6 +718,89 @@ const RUNTIME_V15_OBJECTS: &[&str] = &[
"radroots_runtime_source_generations_sequence_guard",
];
+const RUNTIME_V16_OBJECTS: &[&str] = &[
+ "radroots_runtime_atomic_commits",
+ "radroots_runtime_authored_artifacts",
+ "radroots_runtime_authored_artifacts_admission_ready_idx",
+ "radroots_runtime_authored_artifacts_signed_fact_guard",
+ "radroots_runtime_authored_artifacts_signing_ready_idx",
+ "radroots_runtime_authored_atomic_commits",
+ "radroots_runtime_authored_atomic_commits_delete_guard",
+ "radroots_runtime_authored_atomic_commits_update_guard",
+ "radroots_runtime_authored_delivery_attempts",
+ "radroots_runtime_authored_delivery_facts",
+ "radroots_runtime_authored_delivery_facts_delete_guard",
+ "radroots_runtime_authored_delivery_facts_update_guard",
+ "radroots_runtime_authored_delivery_plans",
+ "radroots_runtime_authored_delivery_ready_idx",
+ "radroots_runtime_authored_delivery_stop_guard",
+ "radroots_runtime_authored_delivery_targets",
+ "radroots_runtime_authored_draft_author_head_idx",
+ "radroots_runtime_authored_draft_revisions",
+ "radroots_runtime_authored_draft_revisions_delete_guard",
+ "radroots_runtime_authored_draft_revisions_update_guard",
+ "radroots_runtime_authored_draft_scope_head_idx",
+ "radroots_runtime_authored_migration_evidence",
+ "radroots_runtime_authored_migration_evidence_delete_guard",
+ "radroots_runtime_authored_migration_evidence_update_guard",
+ "radroots_runtime_authored_operations",
+ "radroots_runtime_delivery_evidence",
+ "radroots_runtime_delivery_evidence_item_idx",
+ "radroots_runtime_event_index_checkpoints",
+ "radroots_runtime_event_index_manifests",
+ "radroots_runtime_event_index_shards",
+ "radroots_runtime_event_provenance",
+ "radroots_runtime_event_provenance_observed_idx",
+ "radroots_runtime_events",
+ "radroots_runtime_events_admission_idx",
+ "radroots_runtime_events_contract_metadata_guard",
+ "radroots_runtime_events_contract_metadata_insert_guard",
+ "radroots_runtime_events_delete_guard",
+ "radroots_runtime_events_event_id_idx",
+ "radroots_runtime_events_raw_update_guard",
+ "radroots_runtime_journal_idempotency_idx",
+ "radroots_runtime_journal_operations",
+ "radroots_runtime_journal_recovery_idx",
+ "radroots_runtime_legacy_event_staging",
+ "radroots_runtime_legacy_event_staging_delete_guard",
+ "radroots_runtime_legacy_event_staging_insert_guard",
+ "radroots_runtime_legacy_event_staging_update_guard",
+ "radroots_runtime_legacy_import_commit_delete_guard",
+ "radroots_runtime_legacy_import_commit_update_guard",
+ "radroots_runtime_legacy_import_commits",
+ "radroots_runtime_legacy_import_delete_guard",
+ "radroots_runtime_legacy_import_identity_guard",
+ "radroots_runtime_legacy_import_member_delete_guard",
+ "radroots_runtime_legacy_import_member_identity_guard",
+ "radroots_runtime_legacy_import_member_state_guard",
+ "radroots_runtime_legacy_import_members",
+ "radroots_runtime_legacy_import_state_guard",
+ "radroots_runtime_legacy_import_state_idx",
+ "radroots_runtime_legacy_imports",
+ "radroots_runtime_legacy_outbox_staging",
+ "radroots_runtime_legacy_outbox_staging_delete_guard",
+ "radroots_runtime_legacy_outbox_staging_insert_guard",
+ "radroots_runtime_legacy_outbox_staging_parent_idx",
+ "radroots_runtime_legacy_outbox_staging_update_guard",
+ "radroots_runtime_outbox_items",
+ "radroots_runtime_outbox_operation_idx",
+ "radroots_runtime_outbox_ready_idx",
+ "radroots_runtime_outbox_targets",
+ "radroots_runtime_projection_checkpoints",
+ "radroots_runtime_projection_documents",
+ "radroots_runtime_projection_invalidations",
+ "radroots_runtime_projection_rebuilds",
+ "radroots_runtime_projection_rebuilds_stage_idx",
+ "radroots_runtime_projection_snapshots",
+ "radroots_runtime_projection_snapshots_created_idx",
+ "radroots_runtime_projection_snapshots_update_guard",
+ "radroots_runtime_source_generations",
+ "radroots_runtime_source_generations_active_idx",
+ "radroots_runtime_source_generations_delete_guard",
+ "radroots_runtime_source_generations_identity_guard",
+ "radroots_runtime_source_generations_sequence_guard",
+];
+
/// Ordered, immutable runtime migration plan.
pub const MIGRATIONS: &[MigrationDescriptor] = &[
MigrationDescriptor {
@@ -808,6 +893,12 @@ pub const MIGRATIONS: &[MigrationDescriptor] = &[
up_sha256: "d7ed7312b36f6ae633d2018d71d1bc0117097523ff46ecfaca8a99ca739aa026",
owned_objects: RUNTIME_V15_OBJECTS,
},
+ MigrationDescriptor {
+ version: 16,
+ name: "authored_delivery_facts",
+ up_sha256: "f78d35fbe60255f152c4e7a8c23add6677ac3ae74eed7db1f3d8ef4ba0b8f3c3",
+ owned_objects: RUNTIME_V16_OBJECTS,
+ },
];
pub(crate) const fn migration_sql(version: u32) -> Option<&'static str> {
@@ -827,6 +918,7 @@ pub(crate) const fn migration_sql(version: u32) -> Option<&'static str> {
13 => Some(AUTHORED_DRAFT_REVISIONS_V13_SQL),
14 => Some(AUTHORED_DRAFT_QUERY_V14_SQL),
15 => Some(AUTHORED_SIGNED_FACTS_V15_SQL),
+ 16 => Some(AUTHORED_DELIVERY_FACTS_V16_SQL),
_ => None,
}
}
@@ -870,8 +962,8 @@ mod tests {
fn migration_plan_matches_governed_snapshot() {
let snapshot = toml::from_str::<PlanSnapshot>(PLAN_SNAPSHOT).expect("valid snapshot");
assert_eq!(MINIMUM_VERSION, 1);
- assert_eq!(CURRENT_VERSION, 15);
- assert_eq!(MIGRATIONS.len(), 15);
+ assert_eq!(CURRENT_VERSION, 16);
+ assert_eq!(MIGRATIONS.len(), 16);
let migration = MIGRATIONS[8];
assert_eq!(snapshot.schema_version, 1);
assert_eq!(snapshot.database, "runtime.sqlite");
diff --git a/crates/storage_sqlite/src/migration_delivery_facts_tests.rs b/crates/storage_sqlite/src/migration_delivery_facts_tests.rs
@@ -0,0 +1,142 @@
+use super::{
+ tests::{connection, establish_runtime_version, pragma},
+ *,
+};
+
+const FAILING_V16: &str = concat!(
+ include_str!("migration/runtime/0016_authored_delivery_facts.up.sql"),
+ "\nINSERT INTO missing_fixture_table VALUES (1);"
+);
+
+fn plan(current: u32, fail: bool) -> MigrationPlan {
+ MigrationPlan {
+ database: RUNTIME_DATABASE,
+ application_id: RUNTIME_APPLICATION_ID,
+ set_application_id_sql: SET_RUNTIME_APPLICATION_ID,
+ minimum_version: runtime::MINIMUM_VERSION,
+ current_version: current,
+ steps: runtime::MIGRATIONS
+ .iter()
+ .take(current as usize)
+ .map(|step| MigrationStep {
+ version: step.version(),
+ sql: if fail && step.version() == 16 {
+ FAILING_V16
+ } else {
+ runtime::migration_sql(step.version()).unwrap()
+ },
+ owned_objects: step.owned_objects(),
+ })
+ .collect(),
+ }
+}
+
+#[tokio::test]
+async fn delivery_fact_upgrade_preserves_legacy_snapshots_receipts_and_cancelled_stop() {
+ let mut connection = connection().await;
+ establish_runtime_version(&mut connection, 15).await;
+ let old = super::signed_facts_tests::seed(&mut connection).await;
+ let bytes: Vec<u8> =
+ sqlx::query_scalar("SELECT snapshot FROM radroots_runtime_authored_delivery_plans")
+ .fetch_one(&mut connection)
+ .await
+ .unwrap();
+ let mut wire: serde_json::Value = serde_json::from_slice(&bytes).unwrap();
+ wire["state"] = serde_json::json!("cancelled");
+ wire.as_object_mut().unwrap().remove("delivery_facts");
+ wire.as_object_mut()
+ .unwrap()
+ .remove("stop_requested_at_unix_ms");
+ let historical = serde_json::to_vec(&wire).unwrap();
+ sqlx::query(
+ "UPDATE radroots_runtime_authored_delivery_plans SET state = 'cancelled', snapshot = ?",
+ )
+ .bind(&historical)
+ .execute(&mut connection)
+ .await
+ .unwrap();
+ assert!(matches!(
+ migrate_runtime(&mut connection, OpenMode::ReadOnly).await,
+ Err(Error::SchemaMigrationRequired {
+ actual: 15,
+ current: 16,
+ ..
+ })
+ ));
+ assert_eq!(
+ migrate_runtime(&mut connection, OpenMode::ReadWriteExisting)
+ .await
+ .unwrap()
+ .applied(),
+ 1
+ );
+ assert_eq!(pragma(&mut connection, "user_version").await, 16);
+ super::signed_facts_tests::retained(&mut connection, &old).await;
+ let (actual, stop): (Vec<u8>, Option<i64>) = sqlx::query_as(
+ "SELECT snapshot, stop_requested_at_unix_ms FROM radroots_runtime_authored_delivery_plans",
+ )
+ .fetch_one(&mut connection)
+ .await
+ .unwrap();
+ assert_eq!(actual, historical);
+ assert_eq!(stop, Some(10));
+ let decoded: radroots_storage::authored_delivery::AuthoredDeliveryPlan =
+ serde_json::from_slice(&actual).unwrap();
+ assert_eq!(decoded.stop_requested_at_unix_ms(), Some(10));
+ assert!(decoded.delivery_facts().is_empty());
+ for mode in [OpenMode::ReadOnly, OpenMode::ReadWriteExisting] {
+ assert!(matches!(
+ migrate(&mut connection, mode, &plan(15, false)).await,
+ Err(Error::SchemaTooNew {
+ actual: 16,
+ supported: 15,
+ ..
+ })
+ ));
+ }
+ connection.close().await.unwrap();
+}
+
+#[tokio::test]
+async fn failed_delivery_fact_migration_rolls_back_every_new_schema_object() {
+ let mut connection = connection().await;
+ establish_runtime_version(&mut connection, 15).await;
+ let old = super::signed_facts_tests::seed(&mut connection).await;
+ assert!(
+ migrate(
+ &mut connection,
+ OpenMode::ReadWriteExisting,
+ &plan(16, true)
+ )
+ .await
+ .is_err()
+ );
+ assert_eq!(pragma(&mut connection, "user_version").await, 15);
+ super::signed_facts_tests::retained(&mut connection, &old).await;
+ assert!(
+ sqlx::query(
+ "SELECT stop_requested_at_unix_ms FROM radroots_runtime_authored_delivery_plans"
+ )
+ .fetch_all(&mut connection)
+ .await
+ .is_err()
+ );
+ assert_eq!(sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM sqlite_schema WHERE name LIKE 'radroots_runtime_authored_delivery_facts%' OR name = 'radroots_runtime_authored_delivery_stop_guard'").fetch_one(&mut connection).await.unwrap(), 0);
+ migrate_runtime(&mut connection, OpenMode::ReadWriteExisting)
+ .await
+ .unwrap();
+ super::signed_facts_tests::retained(&mut connection, &old).await;
+ connection.close().await.unwrap();
+}
+
+#[test]
+fn delivery_fact_decision_binds_unchanged_historical_migrations_and_exact_successor() {
+ let decision: serde_json::Value = serde_json::from_str(include_str!(
+ "../../../contracts/architecture/decisions/authored_delivery_facts.v1.json"
+ ))
+ .unwrap();
+ let migration = runtime::MIGRATIONS[15];
+ assert_eq!(decision["migration"]["version"], migration.version());
+ assert_eq!(decision["migration"]["name"], migration.name());
+ assert_eq!(decision["migration"]["sha256"], migration.up_sha256());
+}
diff --git a/crates/storage_sqlite/src/migration_signed_facts_tests.rs b/crates/storage_sqlite/src/migration_signed_facts_tests.rs
@@ -36,7 +36,7 @@ fn plan(current: u32, fail: bool) -> MigrationPlan {
}
}
-async fn seed(connection: &mut SqliteConnection) -> (Vec<u8>, Vec<u8>, Vec<u8>) {
+pub(super) async fn seed(connection: &mut SqliteConnection) -> (Vec<u8>, Vec<u8>, Vec<u8>) {
let (command, _) = signed_fact_fixture::prepare();
let AuthoredAtomicCommand::Prepare(prepared) = &command else {
unreachable!()
@@ -83,7 +83,10 @@ async fn seed(connection: &mut SqliteConnection) -> (Vec<u8>, Vec<u8>, Vec<u8>)
)
}
-async fn retained(connection: &mut SqliteConnection, expected: &(Vec<u8>, Vec<u8>, Vec<u8>)) {
+pub(super) async fn retained(
+ connection: &mut SqliteConnection,
+ expected: &(Vec<u8>, Vec<u8>, Vec<u8>),
+) {
let artifact: Vec<u8> =
sqlx::query_scalar("SELECT snapshot FROM radroots_runtime_authored_artifacts")
.fetch_one(&mut *connection)
@@ -103,7 +106,7 @@ async fn v15_preserves_v14_rows_and_receipts_and_prior_schema_policy_fails_close
establish_runtime_version(&mut connection, 14).await;
let old = seed(&mut connection).await;
assert!(matches!(
- migrate_runtime(&mut connection, OpenMode::ReadOnly).await,
+ migrate(&mut connection, OpenMode::ReadOnly, &plan(15, false)).await,
Err(Error::SchemaMigrationRequired {
actual: 14,
current: 15,
@@ -111,10 +114,14 @@ async fn v15_preserves_v14_rows_and_receipts_and_prior_schema_policy_fails_close
})
));
assert_eq!(
- migrate_runtime(&mut connection, OpenMode::ReadWriteExisting)
- .await
- .unwrap()
- .applied(),
+ migrate(
+ &mut connection,
+ OpenMode::ReadWriteExisting,
+ &plan(15, false)
+ )
+ .await
+ .unwrap()
+ .applied(),
1
);
assert_eq!(pragma(&mut connection, "user_version").await, 15);
@@ -136,7 +143,7 @@ async fn v15_preserves_v14_rows_and_receipts_and_prior_schema_policy_fails_close
));
}
assert_eq!(
- migrate_runtime(&mut connection, OpenMode::ReadOnly)
+ migrate(&mut connection, OpenMode::ReadOnly, &plan(15, false))
.await
.unwrap()
.applied(),