commit eff47745c38195d2e48674d6b98236b714dcb880 parent 97ad465d34ab4dcf7b665749a387209d3d807f3d Author: triesap <tyson@radroots.org> Date: Mon, 14 Sep 2026 15:38:45 +0000 storage: reconcile delivery facts with bounded claim history - Index original issued claims and retain explicit legacy uncertainty - Reconcile exact facts under current authority without duplicate attempts - Preserve immutable provenance through forward migration and rollback - Verify legacy binding, capacity, corruption and unchanged coverage gates Diffstat:
22 files changed, 3070 insertions(+), 18 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::ReconcileDelivery(radroots_storage::authored_atomic::ReconcileDeliveryFacts) 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 @@ -287,6 +288,13 @@ 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::ReconcileDeliveryFacts +impl radroots_storage::authored_atomic::ReconcileDeliveryFacts +pub fn radroots_storage::authored_atomic::ReconcileDeliveryFacts::apply_to(&self, &radroots_storage::authored_delivery::AuthoredDeliveryHistory) -> core::result::Result<radroots_storage::authored_delivery::AuthoredDeliveryPlan, radroots_storage::Error> +pub fn radroots_storage::authored_atomic::ReconcileDeliveryFacts::new(&radroots_storage::authored_delivery::AuthoredDeliveryPlan, core::option::Option<radroots_storage::authored_atomic::WorkFence>, core::option::Option<radroots_storage::authored::RetrySchedule>, u64) -> core::result::Result<Self, radroots_storage::Error> +pub const fn radroots_storage::authored_atomic::ReconcileDeliveryFacts::plan_id(&self) -> radroots_storage::authored_delivery::AuthoredDeliveryPlanId +pub const fn radroots_storage::authored_atomic::ReconcileDeliveryFacts::reconciled_at_unix_ms(&self) -> u64 +pub const fn radroots_storage::authored_atomic::ReconcileDeliveryFacts::revision(&self) -> core::num::nonzero::NonZeroU64 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> @@ -315,12 +323,14 @@ pub const fn radroots_storage::authored_atomic::WorkFence::row_revision(&self) - pub const fn radroots_storage::authored_atomic::WorkFence::token(&self) -> &[u8; 16] pub trait radroots_storage::authored_atomic::AuthoredAtomicStorage: core::marker::Send + core::marker::Sync pub fn radroots_storage::authored_atomic::AuthoredAtomicStorage::authored_artifact(&self, radroots_storage::authored::AuthoredArtifactId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::authored::AuthoredArtifact>, radroots_storage::Error>> +pub fn radroots_storage::authored_atomic::AuthoredAtomicStorage::authored_delivery_history(&self, radroots_storage::authored_delivery::AuthoredDeliveryPlanId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::authored_delivery::AuthoredDeliveryHistory>, radroots_storage::Error>> pub fn radroots_storage::authored_atomic::AuthoredAtomicStorage::authored_delivery_plan(&self, radroots_storage::authored_delivery::AuthoredDeliveryPlanId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::authored_delivery::AuthoredDeliveryPlan>, radroots_storage::Error>> pub fn radroots_storage::authored_atomic::AuthoredAtomicStorage::authored_operation(&self, radroots_storage::journal::OperationInstanceId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::authored::AuthoredOperation>, radroots_storage::Error>> pub fn radroots_storage::authored_atomic::AuthoredAtomicStorage::authored_receipt(&self, radroots_storage::atomic::AtomicCommitId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::authored_atomic::AuthoredAtomicReceipt>, radroots_storage::Error>> pub fn radroots_storage::authored_atomic::AuthoredAtomicStorage::execute_authored(&self, radroots_storage::authored_atomic::AuthoredAtomicCommand) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_storage::authored_atomic::AuthoredAtomicReceipt, radroots_storage::Error>> impl radroots_storage::authored_atomic::AuthoredAtomicStorage for radroots_storage::memory::MemoryStorage pub fn radroots_storage::memory::MemoryStorage::authored_artifact(&self, radroots_storage::authored::AuthoredArtifactId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::authored::AuthoredArtifact>, radroots_storage::Error>> +pub fn radroots_storage::memory::MemoryStorage::authored_delivery_history(&self, radroots_storage::authored_delivery::AuthoredDeliveryPlanId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::authored_delivery::AuthoredDeliveryHistory>, radroots_storage::Error>> pub fn radroots_storage::memory::MemoryStorage::authored_delivery_plan(&self, radroots_storage::authored_delivery::AuthoredDeliveryPlanId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::authored_delivery::AuthoredDeliveryPlan>, radroots_storage::Error>> pub fn radroots_storage::memory::MemoryStorage::authored_operation(&self, radroots_storage::journal::OperationInstanceId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::authored::AuthoredOperation>, radroots_storage::Error>> pub fn radroots_storage::memory::MemoryStorage::authored_receipt(&self, radroots_storage::atomic::AtomicCommitId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::authored_atomic::AuthoredAtomicReceipt>, radroots_storage::Error>> @@ -341,15 +351,33 @@ pub radroots_storage::authored_delivery::DeliveryAttemptOutcome::SinkFailure(rad pub struct radroots_storage::authored_delivery::AuthoredDeliveryAttempt impl radroots_storage::authored_delivery::AuthoredDeliveryAttempt pub const fn radroots_storage::authored_delivery::AuthoredDeliveryAttempt::attempt(&self) -> core::num::nonzero::NonZeroU32 +pub const fn radroots_storage::authored_delivery::AuthoredDeliveryAttempt::claim_evidence(&self) -> core::option::Option<&radroots_storage::authored::WorkClaim> pub const fn radroots_storage::authored_delivery::AuthoredDeliveryAttempt::outcome(&self) -> &radroots_storage::authored_delivery::DeliveryAttemptOutcome 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::AuthoredDeliveryClaim +impl radroots_storage::authored_delivery::AuthoredDeliveryClaim +pub const fn radroots_storage::authored_delivery::AuthoredDeliveryClaim::claim(&self) -> &radroots_storage::authored::WorkClaim +pub const fn radroots_storage::authored_delivery::AuthoredDeliveryClaim::prior_attempt_count(&self) -> u32 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::AuthoredDeliveryHistory +impl radroots_storage::authored_delivery::AuthoredDeliveryHistory +pub fn radroots_storage::authored_delivery::AuthoredDeliveryHistory::claims(&self) -> &[radroots_storage::authored_delivery::AuthoredDeliveryClaim] +pub fn radroots_storage::authored_delivery::AuthoredDeliveryHistory::has_unresolved_claims(&self) -> bool +pub const fn radroots_storage::authored_delivery::AuthoredDeliveryHistory::is_complete(&self) -> bool +pub const fn radroots_storage::authored_delivery::AuthoredDeliveryHistory::is_truncated(&self) -> bool +pub fn radroots_storage::authored_delivery::AuthoredDeliveryHistory::mark_truncated(&mut self) +pub fn radroots_storage::authored_delivery::AuthoredDeliveryHistory::new(radroots_storage::authored_delivery::AuthoredDeliveryPlan, core::option::Option<&radroots_storage::authored_atomic::AuthoredAtomicReceipt>) -> core::result::Result<Self, radroots_storage::Error> +pub const fn radroots_storage::authored_delivery::AuthoredDeliveryHistory::plan(&self) -> &radroots_storage::authored_delivery::AuthoredDeliveryPlan +pub fn radroots_storage::authored_delivery::AuthoredDeliveryHistory::proves_no_issued_attempt(&self) -> bool +pub fn radroots_storage::authored_delivery::AuthoredDeliveryHistory::push_claim(&mut self, &radroots_storage::authored_atomic::AuthoredAtomicReceipt) -> core::result::Result<(), radroots_storage::Error> +pub fn radroots_storage::authored_delivery::AuthoredDeliveryHistory::require_pending_fact_provenance(&self) -> core::result::Result<(), radroots_storage::Error> +pub fn radroots_storage::authored_delivery::AuthoredDeliveryHistory::validate(&self) -> core::result::Result<(), radroots_storage::Error> 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 @@ -390,6 +418,8 @@ 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> +impl radroots_storage::authored_delivery::AuthoredDeliveryPlan +pub fn radroots_storage::authored_delivery::AuthoredDeliveryPlan::pending_delivery_facts(&self) -> impl core::iter::traits::iterator::Iterator<Item = &radroots_storage::authored_delivery::AuthoredDeliveryFact> 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] @@ -919,6 +949,7 @@ pub fn radroots_storage::memory::MemoryStorage::commit(&self, radroots_storage:: pub fn radroots_storage::memory::MemoryStorage::receipt(&self, radroots_storage::atomic::AtomicCommitId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::atomic::AtomicCommitReceipt>, radroots_storage::Error>> impl radroots_storage::authored_atomic::AuthoredAtomicStorage for radroots_storage::memory::MemoryStorage pub fn radroots_storage::memory::MemoryStorage::authored_artifact(&self, radroots_storage::authored::AuthoredArtifactId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::authored::AuthoredArtifact>, radroots_storage::Error>> +pub fn radroots_storage::memory::MemoryStorage::authored_delivery_history(&self, radroots_storage::authored_delivery::AuthoredDeliveryPlanId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::authored_delivery::AuthoredDeliveryHistory>, radroots_storage::Error>> pub fn radroots_storage::memory::MemoryStorage::authored_delivery_plan(&self, radroots_storage::authored_delivery::AuthoredDeliveryPlanId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::authored_delivery::AuthoredDeliveryPlan>, radroots_storage::Error>> pub fn radroots_storage::memory::MemoryStorage::authored_operation(&self, radroots_storage::journal::OperationInstanceId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::authored::AuthoredOperation>, radroots_storage::Error>> pub fn radroots_storage::memory::MemoryStorage::authored_receipt(&self, radroots_storage::atomic::AtomicCommitId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::authored_atomic::AuthoredAtomicReceipt>, radroots_storage::Error>> diff --git a/contracts/api_baselines/radroots_storage_sqlite.txt b/contracts/api_baselines/radroots_storage_sqlite.txt @@ -551,6 +551,7 @@ pub fn radroots_storage_sqlite::SqliteStorage::commit(&self, radroots_storage::a pub fn radroots_storage_sqlite::SqliteStorage::receipt(&self, radroots_storage::atomic::AtomicCommitId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::atomic::AtomicCommitReceipt>, radroots_storage::error::Error>> impl radroots_storage::authored_atomic::AuthoredAtomicStorage for radroots_storage_sqlite::SqliteStorage pub fn radroots_storage_sqlite::SqliteStorage::authored_artifact(&self, radroots_storage::authored::AuthoredArtifactId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::authored::AuthoredArtifact>, radroots_storage::error::Error>> +pub fn radroots_storage_sqlite::SqliteStorage::authored_delivery_history(&self, radroots_storage::authored_delivery::AuthoredDeliveryPlanId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::authored_delivery::history::AuthoredDeliveryHistory>, radroots_storage::error::Error>> pub fn radroots_storage_sqlite::SqliteStorage::authored_delivery_plan(&self, radroots_storage::authored_delivery::AuthoredDeliveryPlanId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::authored_delivery::AuthoredDeliveryPlan>, radroots_storage::error::Error>> pub fn radroots_storage_sqlite::SqliteStorage::authored_operation(&self, radroots_storage::journal::OperationInstanceId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::authored::AuthoredOperation>, radroots_storage::error::Error>> pub fn radroots_storage_sqlite::SqliteStorage::authored_receipt(&self, radroots_storage::atomic::AtomicCommitId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::authored_atomic::AuthoredAtomicReceipt>, radroots_storage::error::Error>> diff --git a/contracts/architecture/decisions/authored_delivery_reconciliation.v1.json b/contracts/architecture/decisions/authored_delivery_reconciliation.v1.json @@ -0,0 +1,21 @@ +{ + "schema": "radroots.authored-delivery-reconciliation.v1", + "status": "approved", + "owners": [ + "radroots_storage", + "radroots_storage_sqlite" + ], + "history": "AuthoredDeliveryHistory is a consistent bounded projection of the current plan and backend-owned immutable preparation/Claim receipts. A verified initial preparation without prior attempts or claims permits a scoped no-issued-attempt proof. Missing legacy provenance or truncated history remains explicitly uncertain. Validate exact original plan/artifact/intent/request/claim binding; malformed provenance fails closed.", + "history_bounds": "Retain all historical receipts. Return at most1024 issued claims and an explicit truncation flag; fetch bounded indexed IDs before decoding one bounded receipt at a time. New claims fail closed once the per-plan durable claim bound is reached. No history eviction or increased1024 attempt/fact or4194304-byte snapshot limit.", + "reconciliation": "A distinct ReconcileDeliveryFacts atomic command compares the current scheduling revision and complete retained fact-set digest. It requires the exact currently valid fence, or explicit current observation that no competing lease remains valid. A late worker with an expired/superseded fence cannot reconcile scheduling. Explicit stop and terminal scheduling states reject new reconciliation. RecordDeliveryFact remains independent of scheduling authority.", + "idempotence": "Each reconciled scheduling attempt carries the complete originating fact claim. Match it to the immutable raw fact and admit each claim once. Reconciliation appends only pending raw facts in stored order, preserving legacy prefix satisfaction and monotonic record time. New command identities bind all current-authority/fact/retry inputs. Every historical command hash and receipt remains unchanged. A matching legacy fenced application is identified by its original claim-receipt attempt ordinal and valid claim interval and receives an exact marker without increasing its attempt count or rewriting its raw result/time. Conflicting legacy results fail reconciliation without discarding retained facts.", + "policy_boundary": "Storage owns validated state transitions and atomicity. Sync supplies retry policy and current host time, records raw valid sink results before clock-dependent scheduling, and performs fresh reconciliation before any further effect. Reconciliation itself never invokes transport or a signer and never fabricates another external attempt.", + "compatibility": "Old attempt snapshots default to no exact claim marker and retain their bytes. A forward runtime migration adds an indexed immutable delivery-claim projection and immutable fact-to-attempt reconciliation rows. Backfill from original typed delivery Claim receipts; preserve historical SQL/checksums and legacy v10 conversion ordering. Missing original preparation remains explicit legacy uncertainty.", + "atomicity": "Memory publishes a validated candidate. SQLite commits plan, normalized attempt provenance, claim projection and immutable command receipt together, returning only after COMMIT. Consistent reads release their snapshot before returning. Stale revisions, changed fact sets, conflicting markers, bounds, migration faults and COMMIT faults cannot partially change durable state.", + "consumer": "Shared Sync adoption and application/native integration are separate checkpoints. No remote rollback, physical-device, release-publication or deployment claim is made.", + "migration": { + "version": 17, + "name": "authored_delivery_reconciliation", + "sha256": "01463e7effeb368dd0577bfed79f566f7bb2af76d03e8e76898cc53332e2e573" + } +} diff --git a/crates/storage/src/authored_atomic.rs b/crates/storage/src/authored_atomic.rs @@ -4,6 +4,8 @@ mod signing_evidence; pub use signing_evidence::RecordSignedArtifact; mod delivery_evidence; pub use delivery_evidence::RecordDeliveryFact; +mod delivery_reconciliation; +pub use delivery_reconciliation::ReconcileDeliveryFacts; use core::num::NonZeroU64; use radroots_event::SignedEvent; @@ -428,6 +430,7 @@ pub enum AuthoredAtomicCommand { ApplyAdmission(ApplyAdmissionResult), ApplyDelivery(ApplyDeliveryAttempt), RecordDelivery(RecordDeliveryFact), + ReconcileDelivery(ReconcileDeliveryFacts), ApplyFailure(ApplyWorkFailure), Cancel(CancelAuthoredWork), } @@ -498,6 +501,7 @@ impl AuthoredAtomicCommand { hasher.update(claim.expires_at_unix_ms().to_be_bytes()); value.hash_outcome(&mut hasher); } + Self::ReconcileDelivery(value) => value.hash_into(&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()), } @@ -514,6 +518,7 @@ impl AuthoredAtomicCommand { 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::ReconcileDelivery(value) => value.reconciled_at_unix_ms(), Self::ApplyFailure(value) => value.applied_at_unix_ms, Self::Cancel(value) => value.cancelled_at_unix_ms, } @@ -529,6 +534,7 @@ impl AuthoredAtomicCommand { Self::ApplyAdmission(_) => b"admission", Self::ApplyDelivery(_) => b"delivery", Self::RecordDelivery(_) => b"delivery_fact_v1", + Self::ReconcileDelivery(_) => b"delivery_reconciliation_v1", Self::ApplyFailure(value) => match value.failure.phase() { WorkPhase::Signing => b"signing_failure", WorkPhase::Admission => b"admission_failure", @@ -554,6 +560,7 @@ impl AuthoredAtomicCommand { 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::ReconcileDelivery(value) => *value.plan_id().as_bytes(), Self::ApplyFailure(value) => match &value.target { AuthoredWorkTarget::Artifact(id) => *id.as_bytes(), AuthoredWorkTarget::DeliveryPlan(id) => *id.as_bytes(), @@ -575,7 +582,10 @@ impl AuthoredAtomicCommand { 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, + Self::Prepare(_) + | Self::PrepareFromDraft(_) + | Self::Cancel(_) + | Self::ReconcileDelivery(_) => None, } } } @@ -611,6 +621,11 @@ impl AuthoredAtomicReceipt { } match (command, &self.outcome) { ( + AuthoredAtomicCommand::ReconcileDelivery(value), + AuthoredAtomicOutcome::DeliveryPlan(plan), + ) => value.matches_plan(plan), + (AuthoredAtomicCommand::ReconcileDelivery(_), _) => false, + ( AuthoredAtomicCommand::RecordDelivery(value), AuthoredAtomicOutcome::DeliveryPlan(plan), ) => value.matches_plan(plan), @@ -647,6 +662,15 @@ impl AuthoredAtomicReceipt { } match (command, &outcome) { ( + AuthoredAtomicCommand::ReconcileDelivery(value), + AuthoredAtomicOutcome::DeliveryPlan(plan), + ) if value.matches_plan(plan) && committed_at_unix_ms >= plan.updated_at_unix_ms() => { + plan.validate()? + } + (AuthoredAtomicCommand::ReconcileDelivery(_), _) => { + return Err(Error::AtomicWorkflowMismatch); + } + ( AuthoredAtomicCommand::RecordDelivery(value), AuthoredAtomicOutcome::DeliveryPlan(plan), ) if value.matches_plan(plan) && committed_at_unix_ms >= plan.updated_at_unix_ms() => { @@ -763,6 +787,14 @@ impl AuthoredAtomicOutcome { } pub trait AuthoredAtomicStorage: Send + Sync { + /// Read exact issued-claim provenance; unsupported backends fail closed. + fn authored_delivery_history( + &self, + _plan_id: AuthoredDeliveryPlanId, + ) -> BoxFuture<'_, Result<Option<crate::authored_delivery::AuthoredDeliveryHistory>, Error>> + { + Box::pin(async { Err(Error::BackendUnavailable) }) + } fn execute_authored( &self, command: AuthoredAtomicCommand, diff --git a/crates/storage/src/authored_atomic/delivery_reconciliation.rs b/crates/storage/src/authored_atomic/delivery_reconciliation.rs @@ -0,0 +1,120 @@ +//! Current-authority command for reconciling already retained delivery facts. + +use super::{AuthoredAtomicCommand, RecordDeliveryFact, WorkFence}; +use crate::{ + Error, + authored::RetrySchedule, + authored_delivery::{AuthoredDeliveryHistory, AuthoredDeliveryPlan, AuthoredDeliveryPlanId}, +}; +use core::num::NonZeroU64; +use sha2::{Digest, Sha256}; + +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct ReconcileDeliveryFacts { + plan_id: AuthoredDeliveryPlanId, + revision: NonZeroU64, + facts_digest: [u8; 32], + fence: Option<WorkFence>, + retry: Option<RetrySchedule>, + reconciled_at_unix_ms: u64, +} + +impl ReconcileDeliveryFacts { + pub fn new( + plan: &AuthoredDeliveryPlan, + fence: Option<WorkFence>, + retry: Option<RetrySchedule>, + reconciled_at_unix_ms: u64, + ) -> Result<Self, Error> { + plan.validate()?; + if reconciled_at_unix_ms == 0 || plan.pending_delivery_facts().next().is_none() { + return Err(Error::AtomicWorkflowMismatch); + } + Ok(Self { + plan_id: plan.plan_id(), + revision: plan.revision(), + facts_digest: facts_digest(plan)?, + fence, + retry, + reconciled_at_unix_ms, + }) + } + pub const fn plan_id(&self) -> AuthoredDeliveryPlanId { + self.plan_id + } + pub const fn revision(&self) -> NonZeroU64 { + self.revision + } + pub const fn reconciled_at_unix_ms(&self) -> u64 { + self.reconciled_at_unix_ms + } + /// Backends supply their own consistent original-claim history. + pub fn apply_to( + &self, + history: &AuthoredDeliveryHistory, + ) -> Result<AuthoredDeliveryPlan, Error> { + history.require_pending_fact_provenance()?; + let mut plan = history.plan().clone(); + if plan.plan_id() != self.plan_id + || plan.revision() != self.revision + || facts_digest(&plan)? != self.facts_digest + { + return Err(Error::DeliveryPlanClaimConflict); + } + plan.reconcile_delivery_facts( + self.fence.as_ref(), + self.retry.clone(), + self.reconciled_at_unix_ms, + history.claims(), + )?; + Ok(plan) + } + + pub(super) fn matches_plan(&self, plan: &AuthoredDeliveryPlan) -> bool { + plan.plan_id() == self.plan_id + && self.revision.get().checked_add(1) == Some(plan.revision().get()) + && plan.updated_at_unix_ms() == self.reconciled_at_unix_ms + && plan.claim_evidence().is_none() + && plan.stop_requested_at_unix_ms().is_none() + && plan.pending_delivery_facts().next().is_none() + && facts_digest(plan).ok() == Some(self.facts_digest) + && plan.retry() == self.retry.as_ref() + } + + pub(super) fn hash_into(&self, hasher: &mut Sha256) { + hasher.update(self.revision.get().to_be_bytes()); + hasher.update(self.facts_digest); + hasher.update(self.reconciled_at_unix_ms.to_be_bytes()); + hasher.update([self.fence.is_some() as u8]); + if let Some(fence) = &self.fence { + hasher.update(fence.token()); + hasher.update(fence.generation().get().to_be_bytes()); + hasher.update(fence.row_revision().get().to_be_bytes()); + } + hasher.update([self.retry.is_some() as u8]); + if let Some(retry) = &self.retry { + hasher.update(retry.attempt().get().to_be_bytes()); + hasher.update(retry.not_before_unix_ms().to_be_bytes()); + super::hash_failure(hasher, Some(retry.failure())); + } + } +} + +fn facts_digest(plan: &AuthoredDeliveryPlan) -> Result<[u8; 32], Error> { + let mut hasher = Sha256::new(); + super::hash_field(&mut hasher, b"radroots.authored.delivery-facts.v1"); + hasher.update(plan.plan_id().as_bytes()); + hasher.update(plan.artifact_id().as_bytes()); + for fact in plan.delivery_facts() { + let command = AuthoredAtomicCommand::RecordDelivery(RecordDeliveryFact::new( + plan.plan_id(), + plan.artifact_id(), + fact.claim().clone(), + fact.outcome().clone(), + fact.observed_at_unix_ms(), + )?); + hasher.update(command.digest().as_bytes()); + hasher.update(fact.observed_at_unix_ms().to_be_bytes()); + } + Ok(hasher.finalize().into()) +} diff --git a/crates/storage/src/authored_delivery.rs b/crates/storage/src/authored_delivery.rs @@ -2,6 +2,9 @@ mod facts; pub use facts::AuthoredDeliveryFact; +mod history; +pub use history::{AuthoredDeliveryClaim, AuthoredDeliveryHistory}; +mod reconciliation; use core::num::{NonZeroU32, NonZeroU64}; use radroots_transport::{ @@ -167,6 +170,11 @@ pub struct AuthoredDeliveryAttempt { recorded_at_unix_ms: u64, outcome: DeliveryAttemptOutcome, satisfaction: SatisfactionState, + #[cfg_attr( + feature = "serde", + serde(default, skip_serializing_if = "Option::is_none") + )] + claim: Option<WorkClaim>, } impl AuthoredDeliveryAttempt { @@ -184,6 +192,7 @@ impl AuthoredDeliveryAttempt { recorded_at_unix_ms, outcome, satisfaction, + claim: None, }) } @@ -199,6 +208,10 @@ impl AuthoredDeliveryAttempt { pub const fn satisfaction(&self) -> SatisfactionState { self.satisfaction } + /// Exact fact provenance for new reconciliation; absent on historical attempts. + pub const fn claim_evidence(&self) -> Option<&WorkClaim> { + self.claim.as_ref() + } } #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] diff --git a/crates/storage/src/authored_delivery/facts.rs b/crates/storage/src/authored_delivery/facts.rs @@ -78,6 +78,20 @@ impl AuthoredDeliveryPlan { return Err(Error::InvalidAuthoredDeliveryPlan); } } + let mut reconciled = Vec::new(); + for attempt in &self.attempts { + if let Some(claim) = attempt.claim_evidence() { + if reconciled.contains(&claim) + || !self + .delivery_facts + .iter() + .any(|fact| fact.claim() == claim && fact.outcome() == attempt.outcome()) + { + return Err(Error::InvalidAuthoredDeliveryPlan); + } + reconciled.push(claim); + } + } Ok(()) } diff --git a/crates/storage/src/authored_delivery/history.rs b/crates/storage/src/authored_delivery/history.rs @@ -0,0 +1,199 @@ +//! Bounded projections of backend-owned delivery claim receipts. + +use super::{AuthoredDeliveryPlan, DELIVERY_PLAN_ATTEMPTS_MAX, WorkClaim}; +use crate::{ + Error, + authored_atomic::{ + AuthoredAtomicCommand, AuthoredAtomicOutcome, AuthoredAtomicReceipt, ClaimAuthoredTarget, + ClaimAuthoredWork, + }, +}; + +/// One issued claim and its original scheduling-attempt boundary. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct AuthoredDeliveryClaim { + claim: WorkClaim, + prior_attempt_count: u32, +} + +impl AuthoredDeliveryClaim { + pub const fn claim(&self) -> &WorkClaim { + &self.claim + } + pub const fn prior_attempt_count(&self) -> u32 { + self.prior_attempt_count + } +} + +/// A consistent, bounded read of a plan and its immutable issued claims. +/// +/// Backends must supply their own retained receipts from the same read snapshot. +/// A truncated or unproven legacy history never proves absence of an attempt. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct AuthoredDeliveryHistory { + plan: AuthoredDeliveryPlan, + claims: Vec<AuthoredDeliveryClaim>, + initial_boundary_proven: bool, + truncated: bool, +} + +impl AuthoredDeliveryHistory { + pub fn new( + plan: AuthoredDeliveryPlan, + preparation: Option<&AuthoredAtomicReceipt>, + ) -> Result<Self, Error> { + plan.validate()?; + let initial_boundary_proven = match preparation.map(AuthoredAtomicReceipt::outcome) { + Some(AuthoredAtomicOutcome::Prepared { delivery_plans, .. }) => { + Self::initial_boundary(&plan, delivery_plans)? + } + Some(AuthoredAtomicOutcome::Submitted(value)) => { + value.validate()?; + Self::initial_boundary(&plan, value.preparation().delivery_plans())? + } + Some(_) => return Err(Error::AtomicWorkflowMismatch), + None => false, + }; + Ok(Self { + plan, + claims: Vec::new(), + initial_boundary_proven, + truncated: false, + }) + } + + fn initial_boundary( + plan: &AuthoredDeliveryPlan, + originals: &[AuthoredDeliveryPlan], + ) -> Result<bool, Error> { + let original = originals + .iter() + .find(|original| original.plan_id() == plan.plan_id()) + .ok_or(Error::AtomicWorkflowMismatch)?; + original.validate()?; + if original.artifact_id() != plan.artifact_id() + || 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); + } + Ok(original.claim_evidence().is_none() && original.attempts().is_empty()) + } + + pub fn push_claim(&mut self, receipt: &AuthoredAtomicReceipt) -> Result<(), Error> { + if self.claims.len() >= DELIVERY_PLAN_ATTEMPTS_MAX as usize { + return Err(Error::DeliveryAttemptOverflow); + } + let AuthoredAtomicOutcome::DeliveryPlan(original) = receipt.outcome() else { + return Err(Error::AtomicWorkflowMismatch); + }; + original.validate()?; + let claim = original + .claim_evidence() + .ok_or(Error::AtomicWorkflowMismatch)?; + let command = AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( + ClaimAuthoredTarget::DeliveryPlan(self.plan.plan_id()), + claim.clone(), + )); + if !receipt.matches_command(&command) + || receipt.committed_at_unix_ms() != claim.acquired_at_unix_ms() + || original.plan_id() != self.plan.plan_id() + || original.artifact_id() != self.plan.artifact_id() + || original.request().is_none() + || original.request() != self.plan.request() + || original.intent() != self.plan.intent() + || original.created_at_unix_ms() != self.plan.created_at_unix_ms() + || original.revision() > self.plan.revision() + || original.updated_at_unix_ms() > self.plan.updated_at_unix_ms() + || self.claims.iter().any(|entry| entry.claim == *claim) + { + return Err(Error::AtomicWorkflowMismatch); + } + self.claims.push(AuthoredDeliveryClaim { + claim: claim.clone(), + prior_attempt_count: original.attempt_count(), + }); + Ok(()) + } + + /// Records that additional retained claims exceeded this bounded read. + pub fn mark_truncated(&mut self) { + self.truncated = true; + } + pub const fn plan(&self) -> &AuthoredDeliveryPlan { + &self.plan + } + pub fn claims(&self) -> &[AuthoredDeliveryClaim] { + &self.claims + } + pub const fn is_complete(&self) -> bool { + self.initial_boundary_proven && !self.truncated + } + pub const fn is_truncated(&self) -> bool { + self.truncated + } + pub fn validate(&self) -> Result<(), Error> { + self.plan.validate()?; + if self.is_complete() + && (self + .plan + .claim_evidence() + .is_some_and(|claim| !self.claims.iter().any(|entry| entry.claim() == claim)) + || self.plan.delivery_facts().iter().any(|fact| { + !self + .claims + .iter() + .any(|entry| entry.claim() == fact.claim()) + })) + { + return Err(Error::AtomicWorkflowMismatch); + } + Ok(()) + } + /// Require original backend claim provenance before scheduling reconciliation. + pub fn require_pending_fact_provenance(&self) -> Result<(), Error> { + self.validate()?; + if self.is_truncated() { + return Err(Error::DeliveryAttemptOverflow); + } + if self.plan.pending_delivery_facts().any(|fact| { + !self + .claims + .iter() + .any(|entry| entry.claim() == fact.claim()) + }) { + return Err(Error::AtomicWorkflowMismatch); + } + Ok(()) + } + pub fn proves_no_issued_attempt(&self) -> bool { + self.is_complete() + && self.claims.is_empty() + && self.plan.attempts().is_empty() + && self.plan.delivery_facts().is_empty() + && self.plan.claim_evidence().is_none() + } + + pub fn has_unresolved_claims(&self) -> bool { + !self.is_complete() + || self.claims.iter().any(|entry| { + !self + .plan + .delivery_facts() + .iter() + .any(|fact| fact.claim() == &entry.claim) + && !self.plan.attempts().iter().any(|attempt| { + // Historical fenced applications have no explicit claim + // marker. Their original attempt boundary and valid lease + // interval together identify the retained application. + attempt.claim_evidence().is_none() + && entry.prior_attempt_count.checked_add(1) + == Some(attempt.attempt().get()) + && attempt.recorded_at_unix_ms() >= entry.claim.acquired_at_unix_ms() + && attempt.recorded_at_unix_ms() < entry.claim.expires_at_unix_ms() + }) + }) + } +} diff --git a/crates/storage/src/authored_delivery/reconciliation.rs b/crates/storage/src/authored_delivery/reconciliation.rs @@ -0,0 +1,190 @@ +//! Reconcile immutable facts only under current scheduling authority. + +use super::{ + AuthoredDeliveryAttempt, AuthoredDeliveryClaim, AuthoredDeliveryFact, AuthoredDeliveryPlan, + AuthoredDeliveryState, DELIVERY_PLAN_ATTEMPTS_MAX, DeliveryAttemptOutcome, Error, FailureClass, + NonZeroU32, RetrySchedule, Retryability, SatisfactionState, WorkFailure, WorkPhase, +}; +use crate::authored_atomic::WorkFence; + +impl AuthoredDeliveryPlan { + pub fn pending_delivery_facts(&self) -> impl Iterator<Item = &AuthoredDeliveryFact> { + self.delivery_facts.iter().filter(|fact| { + !self + .attempts + .iter() + .any(|attempt| attempt.claim_evidence() == Some(fact.claim())) + }) + } + + pub(crate) fn reconcile_delivery_facts( + &mut self, + fence: Option<&WorkFence>, + retry: Option<RetrySchedule>, + at: u64, + claims: &[AuthoredDeliveryClaim], + ) -> Result<(), Error> { + self.validate()?; + if self.stop_requested_at_unix_ms.is_some() || self.state.is_terminal() { + return Err(Error::InvalidAuthoredTransition); + } + match fence { + Some(fence) => { + self.require_claim(fence.token(), fence.generation(), fence.row_revision(), at)? + } + None if self + .claim + .as_ref() + .is_some_and(|claim| at < claim.expires_at_unix_ms()) => + { + return Err(Error::DeliveryPlanClaimConflict); + } + None => {} + } + let pending: Vec<_> = self.pending_delivery_facts().cloned().collect(); + if pending.is_empty() + || pending.iter().any(|fact| at < fact.observed_at_unix_ms()) + || at < self.updated_at_unix_ms + { + return Err(Error::AtomicWorkflowMismatch); + } + let mut candidate = self.clone(); + for fact in &pending { + // A legacy fenced application may already have persisted this exact + // result. Both its original ordinal and valid lease interval must + // match: reconciling another fact can clear a live lease early. + let original = claims + .iter() + .find(|entry| entry.claim() == fact.claim()) + .ok_or(Error::AtomicWorkflowMismatch)?; + if let Some(attempt) = candidate.attempts.iter_mut().find(|attempt| { + attempt.claim_evidence().is_none() + && original.prior_attempt_count().checked_add(1) + == Some(attempt.attempt().get()) + && attempt.recorded_at_unix_ms() >= fact.claim().acquired_at_unix_ms() + && attempt.recorded_at_unix_ms() < fact.claim().expires_at_unix_ms() + }) { + if attempt.outcome() != fact.outcome() { + return Err(Error::AtomicWorkflowMismatch); + } + attempt.claim = Some(fact.claim().clone()); + continue; + } + if candidate.attempts.len() >= DELIVERY_PLAN_ATTEMPTS_MAX as usize { + return Err(Error::DeliveryAttemptOverflow); + } + let satisfaction = candidate.evaluate_with(fact.outcome().clone())?; + let next = u32::try_from(candidate.attempts.len() + 1) + .ok() + .and_then(NonZeroU32::new) + .ok_or(Error::DeliveryAttemptOverflow)?; + candidate.attempts.push(AuthoredDeliveryAttempt { + attempt: next, + recorded_at_unix_ms: at, + outcome: fact.outcome().clone(), + satisfaction, + claim: Some(fact.claim().clone()), + }); + } + candidate.attempt_count = + u32::try_from(candidate.attempts.len()).map_err(|_| Error::DeliveryAttemptOverflow)?; + let last = candidate + .attempts + .last() + .ok_or(Error::AtomicWorkflowMismatch)?; + let (state, failure) = reconciliation_state( + last.outcome(), + last.satisfaction(), + candidate.attempt_count, + retry.as_ref(), + )?; + candidate.state = state; + candidate.last_failure = failure; + candidate.retry = retry; + candidate.claim = None; + candidate.advance(at)?; + candidate.validate()?; + *self = candidate; + Ok(()) + } +} + +fn reconciliation_state( + outcome: &DeliveryAttemptOutcome, + satisfaction: SatisfactionState, + attempt_count: u32, + retry: Option<&RetrySchedule>, +) -> Result<(AuthoredDeliveryState, Option<WorkFailure>), Error> { + let failure = match outcome { + DeliveryAttemptOutcome::Receipt(_) => None, + DeliveryAttemptOutcome::SinkFailure(failure) => Some(WorkFailure::new( + failure.code(), + WorkPhase::Delivery, + if failure.retryability() == Retryability::Retryable { + FailureClass::Retryable + } else { + FailureClass::Terminal + }, + failure.retry_after_unix_ms(), + failure.message().map(str::to_owned), + )?), + }; + match satisfaction { + SatisfactionState::Satisfied if retry.is_none() => { + return Ok((AuthoredDeliveryState::Satisfied, None)); + } + SatisfactionState::Exhausted if retry.is_none() => { + // Exhausted acceptance is settled even when the last adapter failure + // was retryable; retain only a compatible terminal diagnostic. + return Ok(( + AuthoredDeliveryState::Exhausted, + failure.filter(|failure| failure.class() == FailureClass::Terminal), + )); + } + SatisfactionState::Satisfied | SatisfactionState::Exhausted => { + return Err(Error::InvalidRetrySchedule); + } + SatisfactionState::Pending => {} + } + if attempt_count == DELIVERY_PLAN_ATTEMPTS_MAX { + if retry.is_some() { + return Err(Error::InvalidRetrySchedule); + } + return Ok(( + AuthoredDeliveryState::Exhausted, + Some(WorkFailure::new( + "delivery_attempt_limit", + WorkPhase::Delivery, + FailureClass::Terminal, + None, + None, + )?), + )); + } + if failure + .as_ref() + .is_some_and(|failure| failure.class() == FailureClass::Terminal) + { + if retry.is_some() { + return Err(Error::InvalidRetrySchedule); + } + return Ok((AuthoredDeliveryState::FailedTerminal, failure)); + } + let retry = retry.ok_or(Error::InvalidRetrySchedule)?; + if retry.attempt().get() != attempt_count + || retry.failure().phase() != WorkPhase::Delivery + || failure.as_ref().is_some_and(|failure| { + retry.failure().code() != failure.code() + || retry.failure().diagnostic() != failure.diagnostic() + || failure + .retry_after_unix_ms() + .is_some_and(|at| retry.not_before_unix_ms() < at) + }) + { + return Err(Error::InvalidRetrySchedule); + } + Ok(( + AuthoredDeliveryState::Retryable, + Some(retry.failure().clone()), + )) +} diff --git a/crates/storage/src/memory.rs b/crates/storage/src/memory.rs @@ -1,5 +1,8 @@ //! Deterministic in-memory reference storage backend. +#[path = "memory_delivery_history.rs"] +mod delivery_history; + use radroots_event::EventId; use radroots_protocol::runtime::v1::OperationId; use radroots_transport::{BoxFuture, source::EventProvenance}; @@ -86,6 +89,10 @@ struct State { authored_artifacts: Vec<crate::authored::AuthoredArtifact>, authored_delivery_plans: Vec<crate::authored_delivery::AuthoredDeliveryPlan>, authored_atomic_receipts: Vec<AuthoredAtomicReceipt>, + authored_delivery_history: std::collections::BTreeMap< + crate::authored_delivery::AuthoredDeliveryPlanId, + delivery_history::Entry, + >, authored_drafts: Vec<AuthoredDraft>, closed: bool, } @@ -120,6 +127,7 @@ impl MemoryStorage { authored_artifacts: Vec::new(), authored_delivery_plans: Vec::new(), authored_atomic_receipts: Vec::new(), + authored_delivery_history: std::collections::BTreeMap::new(), authored_drafts: Vec::new(), closed: false, }), @@ -1457,6 +1465,16 @@ fn prepare_authored_memory( } impl AuthoredAtomicStorage for MemoryStorage { + fn authored_delivery_history( + &self, + plan_id: crate::authored_delivery::AuthoredDeliveryPlanId, + ) -> BoxFuture<'_, Result<Option<crate::authored_delivery::AuthoredDeliveryHistory>, Error>> + { + Box::pin(async move { + let state = self.state()?; + delivery_history::history(&state, plan_id) + }) + } fn execute_authored( &self, command: AuthoredAtomicCommand, @@ -1513,6 +1531,8 @@ impl AuthoredAtomicStorage for MemoryStorage { let prepared = prepare_authored_memory(&mut candidate, value.preparation().clone())?; candidate.authored_drafts.push(value.intent().clone()); + let ordinary_index = candidate.authored_atomic_receipts.len(); + delivery_history::register(&mut candidate, &ordinary, ordinary_index)?; candidate .authored_atomic_receipts .push(AuthoredAtomicReceipt::new( @@ -1657,6 +1677,10 @@ impl AuthoredAtomicStorage for MemoryStorage { value.apply_to(plan, original)?; AuthoredAtomicOutcome::DeliveryPlan(plan.clone()) } + AuthoredAtomicCommand::ReconcileDelivery(value) => { + let plan = delivery_history::reconcile(&mut candidate, &value)?; + AuthoredAtomicOutcome::DeliveryPlan(plan) + } AuthoredAtomicCommand::ApplyDelivery(value) => { let plan = candidate .authored_delivery_plans @@ -1841,6 +1865,8 @@ impl AuthoredAtomicStorage for MemoryStorage { committed_at, outcome, )?; + let receipt_index = candidate.authored_atomic_receipts.len(); + delivery_history::register(&mut candidate, &command, receipt_index)?; candidate.authored_atomic_receipts.push(receipt.clone()); *state = candidate; Ok(receipt) diff --git a/crates/storage/src/memory_delivery_history.rs b/crates/storage/src/memory_delivery_history.rs @@ -0,0 +1,118 @@ +//! Instance-owned indexes into the reference backend's immutable receipts. + +use super::State; +use crate::{ + Error, + authored_atomic::{AuthoredAtomicCommand, ClaimAuthoredTarget, ReconcileDeliveryFacts}, + authored_delivery::{ + AuthoredDeliveryHistory, AuthoredDeliveryPlan, AuthoredDeliveryPlanId, + DELIVERY_PLAN_ATTEMPTS_MAX, + }, +}; + +#[derive(Clone)] +pub(super) struct Entry { + preparation: Option<usize>, + claims: Vec<usize>, +} + +pub(super) fn register( + state: &mut State, + command: &AuthoredAtomicCommand, + receipt_index: usize, +) -> Result<(), Error> { + match command { + AuthoredAtomicCommand::Prepare(value) => { + for plan in value.delivery_plans() { + if state + .authored_delivery_history + .insert( + plan.plan_id(), + Entry { + preparation: Some(receipt_index), + claims: Vec::new(), + }, + ) + .is_some() + { + return Err(Error::AtomicCommitConflict); + } + } + } + AuthoredAtomicCommand::Claim(value) => { + if let ClaimAuthoredTarget::DeliveryPlan(plan_id) = value.target() { + let entry = state + .authored_delivery_history + .entry(*plan_id) + .or_insert(Entry { + preparation: None, + claims: Vec::new(), + }); + if entry.claims.len() >= DELIVERY_PLAN_ATTEMPTS_MAX as usize { + return Err(Error::DeliveryAttemptOverflow); + } + entry.claims.push(receipt_index); + } + } + _ => {} + } + Ok(()) +} + +pub(super) fn history( + state: &State, + plan_id: AuthoredDeliveryPlanId, +) -> Result<Option<AuthoredDeliveryHistory>, Error> { + let Some(plan) = state + .authored_delivery_plans + .iter() + .find(|plan| plan.plan_id() == plan_id) + else { + return Ok(None); + }; + let entry = state.authored_delivery_history.get(&plan_id); + let preparation = entry + .and_then(|entry| entry.preparation) + .map(|index| { + state + .authored_atomic_receipts + .get(index) + .ok_or(Error::AtomicWorkflowMismatch) + }) + .transpose()?; + let mut history = AuthoredDeliveryHistory::new(plan.clone(), preparation)?; + if let Some(entry) = entry { + if entry.claims.len() > DELIVERY_PLAN_ATTEMPTS_MAX as usize { + history.mark_truncated(); + } + for index in entry + .claims + .iter() + .take(DELIVERY_PLAN_ATTEMPTS_MAX as usize) + { + history.push_claim( + state + .authored_atomic_receipts + .get(*index) + .ok_or(Error::AtomicWorkflowMismatch)?, + )?; + } + } + history.validate()?; + Ok(Some(history)) +} + +pub(super) fn reconcile( + state: &mut State, + command: &ReconcileDeliveryFacts, +) -> Result<AuthoredDeliveryPlan, Error> { + let history = history(state, command.plan_id())?.ok_or(Error::InvalidAuthoredDeliveryPlan)?; + let plan = command.apply_to(&history)?; + let stored = state + .authored_delivery_plans + .iter_mut() + .find(|value| value.plan_id() == command.plan_id()) + .ok_or(Error::InvalidAuthoredDeliveryPlan)?; + *stored = plan.clone(); + Ok(plan) +} diff --git a/crates/storage/tests/authored_delivery/reconciliation_tests.rs b/crates/storage/tests/authored_delivery/reconciliation_tests.rs @@ -0,0 +1,1039 @@ +use super::*; +use radroots_storage::{ + authored::{FailureClass, RetrySchedule, WorkFailure, WorkPhase}, + authored_atomic::ReconcileDeliveryFacts, + authored_delivery::AuthoredDeliveryHistory, +}; +use std::num::NonZeroU32; + +fn history(storage: &MemoryStorage) -> AuthoredDeliveryHistory { + block_on(storage.authored_delivery_history(ids().2)) + .unwrap() + .unwrap() +} + +fn fence(claim: &WorkClaim) -> WorkFence { + WorkFence::new(*claim.token(), claim.generation(), claim.row_revision()).unwrap() +} + +fn reconcile( + plan: &AuthoredDeliveryPlan, + authority: Option<WorkFence>, + retry: Option<RetrySchedule>, + at: u64, +) -> AuthoredAtomicCommand { + AuthoredAtomicCommand::ReconcileDelivery( + ReconcileDeliveryFacts::new(plan, authority, retry, at).unwrap(), + ) +} + +fn retry(attempt: u32, at: u64) -> RetrySchedule { + RetrySchedule::new( + NonZeroU32::new(attempt).unwrap(), + at, + WorkFailure::new( + "delivery_pending", + WorkPhase::Delivery, + FailureClass::Retryable, + Some(at), + None, + ) + .unwrap(), + ) + .unwrap() +} + +#[test] +fn history_distinguishes_no_issued_work_from_stopped_unresolved_work() { + let untouched = MemoryStorage::new(SourceGeneration::new([71; 32]).unwrap()); + assert!( + block_on(untouched.authored_delivery_history(ids().2)) + .unwrap() + .is_none() + ); + block_on(untouched.execute_authored(prepare().0)).unwrap(); + let before = history(&untouched); + assert!(before.is_complete()); + assert!(before.proves_no_issued_attempt()); + assert!(!before.has_unresolved_claims()); + stop(&untouched, 20); + assert!(history(&untouched).proves_no_issued_attempt()); + + let (storage, active, _) = prepared(); + let issued = history(&storage); + assert_eq!(issued.claims().len(), 1); + assert_eq!(issued.claims()[0].claim(), &active); + assert_eq!(issued.claims()[0].prior_attempt_count(), 0); + assert!(issued.has_unresolved_claims()); + assert!(!issued.proves_no_issued_attempt()); + stop(&storage, 20); + assert!(history(&storage).has_unresolved_claims()); + block_on(storage.execute_authored(fact(&plan(&storage), active, true, 50))).unwrap(); + let observed = history(&storage); + assert!(!observed.has_unresolved_claims()); + assert!(!observed.proves_no_issued_attempt()); + assert_eq!(observed.plan().state(), AuthoredDeliveryState::Cancelled); + assert_eq!( + observed.plan().delivery_satisfaction().unwrap(), + SatisfactionState::Satisfied + ); +} + +#[test] +fn exact_current_fence_reconciles_once_and_preserves_raw_fact_provenance() { + let (storage, active, original) = prepared(); + block_on(storage.execute_authored(fact(&plan(&storage), active.clone(), true, 14))).unwrap(); + let before = plan(&storage); + let command = reconcile(&before, Some(fence(&active)), None, 15); + let receipt = block_on(storage.execute_authored(command.clone())).unwrap(); + let after = plan(&storage); + assert!(receipt.matches_command(&command)); + assert_eq!(after.state(), AuthoredDeliveryState::Satisfied); + assert_eq!(after.attempt_count(), 1); + assert_eq!(after.attempts()[0].claim_evidence(), Some(&active)); + assert_eq!( + after.attempts()[0].outcome(), + before.delivery_facts()[0].outcome() + ); + assert_eq!(after.attempts()[0].recorded_at_unix_ms(), 15); + assert_eq!(after.delivery_facts(), before.delivery_facts()); + assert_eq!(after.pending_delivery_facts().count(), 0); + assert!(!history(&storage).has_unresolved_claims()); + let replay = block_on(storage.execute_authored(command)).unwrap(); + assert_eq!(replay.disposition(), AtomicCommitDisposition::Replay); + assert_eq!(plan(&storage), after); + assert_eq!( + block_on(storage.authored_receipt(original.commit_id())) + .unwrap() + .unwrap(), + original + ); + assert_eq!( + ReconcileDeliveryFacts::new(&after, None, None, 16), + Err(Error::AtomicWorkflowMismatch) + ); +} + +#[test] +fn distinct_current_reconciliation_survives_legacy_generation_identity_reuse() { + let (storage, first, _) = prepared(); + let first_plan = plan(&storage); + block_on(storage.execute_authored(fact(&first_plan, first.clone(), false, 14))).unwrap(); + let first_command = reconcile(&plan(&storage), Some(fence(&first)), Some(retry(1, 18)), 15); + block_on(storage.execute_authored(first_command.clone())).unwrap(); + let second = claim(plan(&storage).revision(), 2, 20); + block_on( + storage.execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( + ClaimAuthoredTarget::DeliveryPlan(ids().2), + second.clone(), + ))), + ) + .unwrap(); + block_on(storage.execute_authored(fact(&plan(&storage), second.clone(), false, 21))).unwrap(); + let second_command = reconcile( + &plan(&storage), + Some(fence(&second)), + Some(retry(2, 25)), + 22, + ); + assert_ne!(first_command.commit_id(), second_command.commit_id()); + let legacy = |claim: &WorkClaim, at| { + AuthoredAtomicCommand::ApplyDelivery( + ApplyDeliveryAttempt::new(ids().2, fence(claim), outcome(&first_plan, false), None, at) + .unwrap(), + ) + }; + assert_eq!( + legacy(&first, 15).commit_id(), + legacy(&second, 22).commit_id() + ); + block_on(storage.execute_authored(second_command)).unwrap(); + let after = plan(&storage); + assert_eq!(after.state(), AuthoredDeliveryState::Retryable); + assert_eq!(after.attempt_count(), 2); + assert_eq!(after.retry().unwrap().not_before_unix_ms(), 25); + assert_eq!(after.attempts()[0].claim_evidence(), Some(&first)); + assert_eq!(after.attempts()[1].claim_evidence(), Some(&second)); + assert!(!history(&storage).has_unresolved_claims()); +} + +#[test] +fn late_fact_cannot_take_a_newer_lease_and_its_marker_cannot_resolve_that_lease() { + 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(); + block_on(storage.execute_authored(fact(&plan(&storage), old.clone(), true, 50))).unwrap(); + let before = plan(&storage); + for authority in [None, Some(fence(&old))] { + let command = reconcile(&before, authority, None, 51); + assert_eq!( + block_on(storage.execute_authored(command)), + Err(Error::DeliveryPlanClaimConflict) + ); + assert_eq!(plan(&storage), before); + } + block_on(storage.execute_authored(reconcile(&before, Some(fence(&newer)), None, 51))).unwrap(); + let after = history(&storage); + assert_eq!(after.plan().attempts()[0].claim_evidence(), Some(&old)); + assert!( + after.has_unresolved_claims(), + "the new worker has no final result yet" + ); + assert_eq!( + after.plan().delivery_satisfaction().unwrap(), + SatisfactionState::Satisfied + ); +} + +#[test] +fn fresh_observer_can_reconcile_after_expiry_but_expired_worker_cannot() { + let (storage, old, _) = prepared(); + block_on(storage.execute_authored(fact(&plan(&storage), old.clone(), true, 50))).unwrap(); + let before = plan(&storage); + assert_eq!( + block_on(storage.execute_authored(reconcile(&before, Some(fence(&old)), None, 51))), + Err(Error::DeliveryPlanClaimConflict) + ); + block_on(storage.execute_authored(reconcile(&before, None, None, 51))).unwrap(); + assert_eq!(plan(&storage).state(), AuthoredDeliveryState::Satisfied); + assert_eq!(plan(&storage).attempts()[0].claim_evidence(), Some(&old)); +} + +#[test] +fn exact_fact_set_and_revision_races_fail_without_partial_scheduling() { + 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(); + block_on(storage.execute_authored(fact(&plan(&storage), old.clone(), true, 50))).unwrap(); + let before = plan(&storage); + let stale = reconcile(&before, Some(fence(&newer)), None, 55); + block_on(storage.execute_authored(fact(&before, newer.clone(), false, 51))).unwrap(); + let current = plan(&storage); + assert_eq!(current.revision(), before.revision()); + assert_eq!( + block_on(storage.execute_authored(stale)), + Err(Error::DeliveryPlanClaimConflict) + ); + assert_eq!(plan(&storage), current); + block_on(storage.execute_authored(reconcile(¤t, Some(fence(&newer)), None, 55))).unwrap(); + let settled = plan(&storage); + assert_eq!(settled.attempt_count(), 2); + assert_eq!(settled.state(), AuthoredDeliveryState::Satisfied); + assert_eq!(settled.attempts()[0].claim_evidence(), Some(&old)); + assert_eq!(settled.attempts()[1].claim_evidence(), Some(&newer)); + assert!( + settled + .attempts() + .iter() + .all(|attempt| attempt.recorded_at_unix_ms() == 55) + ); + + let (stopped, active, _) = prepared(); + block_on(stopped.execute_authored(fact(&plan(&stopped), active, true, 50))).unwrap(); + let stale = reconcile(&plan(&stopped), None, None, 51); + stop(&stopped, 52); + let before = plan(&stopped); + assert_eq!( + block_on(stopped.execute_authored(stale)), + Err(Error::DeliveryPlanClaimConflict) + ); + assert_eq!( + block_on(stopped.execute_authored(reconcile(&before, None, None, 53))), + Err(Error::InvalidAuthoredTransition) + ); + assert_eq!(plan(&stopped), before); +} + +#[test] +fn historical_application_is_resolved_without_rewriting_its_snapshot_shape() { + let (storage, active, _) = prepared(); + let command = AuthoredAtomicCommand::ApplyDelivery( + ApplyDeliveryAttempt::new( + ids().2, + fence(&active), + outcome(&plan(&storage), true), + None, + 14, + ) + .unwrap(), + ); + block_on(storage.execute_authored(command)).unwrap(); + let legacy = history(&storage); + assert!(!legacy.has_unresolved_claims()); + let attempt = &legacy.plan().attempts()[0]; + assert!(attempt.claim_evidence().is_none()); + assert!( + !serde_json::to_value(attempt) + .unwrap() + .as_object() + .unwrap() + .contains_key("claim") + ); + let restored: AuthoredDeliveryPlan = + serde_json::from_slice(&serde_json::to_vec(legacy.plan()).unwrap()).unwrap(); + assert_eq!(restored, *legacy.plan()); +} + +#[test] +fn missing_truncated_or_malformed_history_never_proves_absence() { + let storage = MemoryStorage::new(SourceGeneration::new([72; 32]).unwrap()); + let original = block_on(storage.execute_authored(prepare().0)).unwrap(); + let initial = plan(&storage); + let unknown = AuthoredDeliveryHistory::new(initial.clone(), None).unwrap(); + assert!(!unknown.is_complete()); + assert!(!unknown.proves_no_issued_attempt()); + assert!(unknown.has_unresolved_claims()); + let mut bounded = AuthoredDeliveryHistory::new(initial, Some(&original)).unwrap(); + bounded.mark_truncated(); + assert!(bounded.is_truncated()); + assert!(!bounded.proves_no_issued_attempt()); + assert!(bounded.has_unresolved_claims()); + assert_eq!( + bounded.require_pending_fact_provenance(), + Err(Error::DeliveryAttemptOverflow) + ); + assert_eq!( + bounded.push_claim(&original), + Err(Error::AtomicWorkflowMismatch) + ); + + let (storage, _, issued) = prepared(); + let mut duplicate = history(&storage); + assert_eq!( + duplicate.push_claim(&issued), + Err(Error::AtomicWorkflowMismatch) + ); + assert!(AuthoredDeliveryHistory::new(plan(&storage), Some(&issued)).is_err()); + let incomplete = AuthoredDeliveryHistory::new(plan(&storage), Some(&original)).unwrap(); + assert_eq!(incomplete.validate(), Err(Error::AtomicWorkflowMismatch)); +} + +#[test] +fn forged_or_duplicate_reconciliation_markers_fail_snapshot_validation() { + let (storage, active, _) = prepared(); + block_on(storage.execute_authored(fact(&plan(&storage), active.clone(), true, 14))).unwrap(); + block_on(storage.execute_authored(reconcile(&plan(&storage), Some(fence(&active)), None, 15))) + .unwrap(); + let valid = serde_json::to_value(plan(&storage)).unwrap(); + let mut missing_fact = valid.clone(); + missing_fact["delivery_facts"] = serde_json::json!([]); + assert!(serde_json::from_value::<AuthoredDeliveryPlan>(missing_fact).is_err()); + let mut forged = valid.clone(); + forged["attempts"][0]["claim"] = serde_json::to_value(claim(NonZeroU64::MIN, 9, 13)).unwrap(); + assert!(serde_json::from_value::<AuthoredDeliveryPlan>(forged).is_err()); + let mut duplicate = valid; + let mut second = duplicate["attempts"][0].clone(); + second["attempt"] = serde_json::json!(2); + duplicate["attempts"].as_array_mut().unwrap().push(second); + duplicate["attempt_count"] = serde_json::json!(2); + assert!(serde_json::from_value::<AuthoredDeliveryPlan>(duplicate).is_err()); +} + +#[test] +fn issued_claim_limit_rejects_new_work_atomically_and_keeps_original_replay() { + let (storage, _, first) = prepared(); + for index in 1..DELIVERY_PLAN_ATTEMPTS_MAX { + let at = 40 + u64::from(index) * 21; + let active = WorkClaim::new( + [7; 16], + "bounded-worker", + NonZeroU64::new(u64::from(index) + 2).unwrap(), + at, + at + 20, + plan(&storage).revision(), + ) + .unwrap(); + block_on( + storage.execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( + ClaimAuthoredTarget::DeliveryPlan(ids().2), + active, + ))), + ) + .unwrap(); + } + let before = plan(&storage); + let at = 40 + u64::from(DELIVERY_PLAN_ATTEMPTS_MAX) * 21; + let active = WorkClaim::new( + [8; 16], + "overflow-worker", + NonZeroU64::new(2048).unwrap(), + at, + at + 20, + before.revision(), + ) + .unwrap(); + let command = AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( + ClaimAuthoredTarget::DeliveryPlan(ids().2), + active, + )); + assert_eq!( + block_on(storage.execute_authored(command.clone())), + Err(Error::DeliveryAttemptOverflow) + ); + assert_eq!(plan(&storage), before); + assert!( + block_on(storage.authored_receipt(command.commit_id())) + .unwrap() + .is_none() + ); + let history = history(&storage); + assert_eq!(history.claims().len(), DELIVERY_PLAN_ATTEMPTS_MAX as usize); + assert!(history.is_complete()); + assert!(history.has_unresolved_claims()); + assert_eq!( + history.clone().push_claim(&first), + Err(Error::DeliveryAttemptOverflow) + ); + let AuthoredAtomicOutcome::DeliveryPlan(original) = first.outcome() else { + unreachable!() + }; + let original = AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( + ClaimAuthoredTarget::DeliveryPlan(ids().2), + original.claim_evidence().unwrap().clone(), + )); + assert_eq!( + block_on(storage.execute_authored(original)) + .unwrap() + .disposition(), + AtomicCommitDisposition::Replay + ); + assert_eq!(plan(&storage), before); +} + +#[test] +fn reconciliation_rejects_future_facts_zero_time_and_invalid_retry_without_mutation() { + for accepted in [false, true] { + let (storage, active, _) = prepared(); + block_on(storage.execute_authored(fact(&plan(&storage), active.clone(), accepted, 30))) + .unwrap(); + let before = plan(&storage); + assert_eq!( + ReconcileDeliveryFacts::new(&before, None, None, 0), + Err(Error::AtomicWorkflowMismatch) + ); + assert_eq!( + block_on(storage.execute_authored(reconcile(&before, Some(fence(&active)), None, 29))), + Err(Error::AtomicWorkflowMismatch) + ); + let invalid = if accepted { Some(retry(1, 60)) } else { None }; + assert_eq!( + block_on(storage.execute_authored(reconcile(&before, None, invalid, 50))), + Err(Error::InvalidRetrySchedule) + ); + if !accepted { + assert_eq!( + block_on(storage.execute_authored(reconcile( + &before, + None, + Some(retry(2, 60)), + 50 + ))), + Err(Error::InvalidRetrySchedule) + ); + } + assert_eq!(plan(&storage), before); + } +} + +#[test] +fn raw_sink_failure_controls_retry_diagnostic_and_provider_backoff() { + let (storage, active, _) = prepared(); + let failure = SinkFailure::for_request( + plan(&storage).request().unwrap(), + "sink_lost", + Retryability::Retryable, + Some(65), + Some("connection lost".into()), + vec![], + ) + .unwrap(); + block_on( + storage.execute_authored(AuthoredAtomicCommand::RecordDelivery( + RecordDeliveryFact::new( + ids().2, + ids().1, + active, + DeliveryAttemptOutcome::SinkFailure(failure.clone()), + 50, + ) + .unwrap(), + )), + ) + .unwrap(); + let before = plan(&storage); + for (code, diagnostic, at) in [ + ("other", Some("connection lost"), 65), + ("sink_lost", None, 65), + ("sink_lost", Some("connection lost"), 64), + ] { + let retry = RetrySchedule::new( + NonZeroU32::MIN, + at, + WorkFailure::new( + code, + WorkPhase::Delivery, + FailureClass::Retryable, + None, + diagnostic.map(str::to_owned), + ) + .unwrap(), + ) + .unwrap(); + assert_eq!( + block_on(storage.execute_authored(reconcile(&before, None, Some(retry), 51))), + Err(Error::InvalidRetrySchedule) + ); + assert_eq!(plan(&storage), before); + } + let retry = RetrySchedule::new( + NonZeroU32::MIN, + 65, + WorkFailure::new( + "sink_lost", + WorkPhase::Delivery, + FailureClass::Retryable, + Some(65), + Some("connection lost".into()), + ) + .unwrap(), + ) + .unwrap(); + block_on(storage.execute_authored(reconcile(&before, None, Some(retry.clone()), 51))).unwrap(); + let after = plan(&storage); + assert_eq!(after.state(), AuthoredDeliveryState::Retryable); + assert_eq!(after.retry(), Some(&retry)); + assert_eq!( + after.attempts()[0].outcome(), + &DeliveryAttemptOutcome::SinkFailure(failure) + ); +} + +#[test] +fn matching_late_fact_marks_legacy_attempt_without_inventing_another_attempt() { + for accepted in [false, true] { + let (storage, active, _) = prepared(); + let command = AuthoredAtomicCommand::ApplyDelivery( + ApplyDeliveryAttempt::new( + ids().2, + fence(&active), + outcome(&plan(&storage), false), + Some(retry(1, 18)), + 14, + ) + .unwrap(), + ); + let original = block_on(storage.execute_authored(command)).unwrap(); + block_on(storage.execute_authored(fact(&plan(&storage), active.clone(), accepted, 50))) + .unwrap(); + let before = plan(&storage); + let command = reconcile(&before, None, Some(retry(1, 60)), 51); + if accepted { + assert_eq!( + block_on(storage.execute_authored(command)), + Err(Error::AtomicWorkflowMismatch) + ); + assert_eq!(plan(&storage), before); + // Conflicting raw facts remain factual evidence, not a replacement + // for the immutable earlier attempt or authority for another effect. + assert_eq!( + before.delivery_satisfaction().unwrap(), + SatisfactionState::Satisfied + ); + } else { + block_on(storage.execute_authored(command.clone())).unwrap(); + let after = plan(&storage); + assert_eq!(after.attempt_count(), 1); + assert_eq!(after.attempts()[0].recorded_at_unix_ms(), 14); + assert_eq!(after.attempts()[0].claim_evidence(), Some(&active)); + assert_eq!(after.pending_delivery_facts().count(), 0); + assert_eq!( + block_on(storage.execute_authored(command)) + .unwrap() + .disposition(), + AtomicCommitDisposition::Replay + ); + } + assert_eq!( + block_on(storage.authored_receipt(original.commit_id())) + .unwrap() + .unwrap(), + original + ); + } +} + +#[test] +fn partial_acceptance_and_terminal_evidence_settle_without_retry_authority() { + for (retryability, partial, expected) in [ + ( + Retryability::Terminal, + None, + AuthoredDeliveryState::FailedTerminal, + ), + ( + Retryability::Retryable, + Some(true), + AuthoredDeliveryState::Satisfied, + ), + ( + Retryability::Retryable, + Some(false), + AuthoredDeliveryState::Exhausted, + ), + ( + Retryability::Terminal, + Some(false), + AuthoredDeliveryState::Exhausted, + ), + ] { + let (storage, active, _) = prepared(); + let before = plan(&storage); + let evidence = partial + .map(|accepted| { + DeliveryTargetReceipt::attempted( + before.request().unwrap().target_set().targets()[0].clone(), + if accepted { + DeliveryOutcome::accepted() + } else { + DeliveryOutcome::rejected() + }, + ) + }) + .into_iter() + .collect(); + let failure = SinkFailure::for_request( + before.request().unwrap(), + "sink_lost", + retryability, + None, + None, + evidence, + ) + .unwrap(); + block_on( + storage.execute_authored(AuthoredAtomicCommand::RecordDelivery( + RecordDeliveryFact::new( + ids().2, + ids().1, + active, + DeliveryAttemptOutcome::SinkFailure(failure), + 50, + ) + .unwrap(), + )), + ) + .unwrap(); + let before = plan(&storage); + assert_eq!( + block_on(storage.execute_authored(reconcile(&before, None, Some(retry(1, 60)), 51))), + Err(Error::InvalidRetrySchedule) + ); + assert_eq!(plan(&storage), before); + block_on(storage.execute_authored(reconcile(&before, None, None, 51))).unwrap(); + assert_eq!(plan(&storage).state(), expected); + assert!(plan(&storage).retry().is_none()); + } +} + +#[test] +fn earlier_unresolved_claim_cannot_adopt_a_new_workers_legacy_attempt() { + let (storage, first, _) = prepared(); + let second = claim(plan(&storage).revision(), 3, 40); + block_on( + storage.execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( + ClaimAuthoredTarget::DeliveryPlan(ids().2), + second.clone(), + ))), + ) + .unwrap(); + block_on(storage.execute_authored(fact(&plan(&storage), first, false, 41))).unwrap(); + block_on(storage.execute_authored(reconcile( + &plan(&storage), + Some(fence(&second)), + Some(retry(1, 43)), + 42, + ))) + .unwrap(); + let third = claim(plan(&storage).revision(), 4, 44); + block_on( + storage.execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( + ClaimAuthoredTarget::DeliveryPlan(ids().2), + third.clone(), + ))), + ) + .unwrap(); + let legacy = AuthoredAtomicCommand::ApplyDelivery( + ApplyDeliveryAttempt::new( + ids().2, + fence(&third), + outcome(&plan(&storage), false), + Some(retry(2, 47)), + 45, + ) + .unwrap(), + ); + block_on(storage.execute_authored(legacy)).unwrap(); + assert!(history(&storage).has_unresolved_claims()); + block_on(storage.execute_authored(fact(&plan(&storage), second.clone(), false, 50))).unwrap(); + block_on(storage.execute_authored(reconcile(&plan(&storage), None, Some(retry(3, 65)), 61))) + .unwrap(); + let after = plan(&storage); + assert_eq!(after.attempt_count(), 3); + assert!(after.attempts()[1].claim_evidence().is_none()); + assert_eq!(after.attempts()[2].claim_evidence(), Some(&second)); + assert!(!history(&storage).has_unresolved_claims()); +} + +fn restored_receipt( + original: &AuthoredAtomicReceipt, + plan: AuthoredDeliveryPlan, + at: u64, +) -> AuthoredAtomicReceipt { + AuthoredAtomicReceipt::from_durable_parts( + original.commit_id(), + original.digest(), + AtomicCommitDisposition::Committed, + at, + AuthoredAtomicOutcome::DeliveryPlan(plan), + ) + .unwrap() +} + +#[test] +fn original_claim_history_rejects_each_mismatched_binding() { + let (storage, _, original) = prepared(); + let current = plan(&storage); + for field in ["plan_id", "artifact_id", "request", "created_at_unix_ms"] { + let mut wire = serde_json::to_value(¤t).unwrap(); + match field { + "plan_id" => { + wire[field] = serde_json::to_value( + radroots_storage::authored_delivery::AuthoredDeliveryPlanId::new([99; 16]) + .unwrap(), + ) + .unwrap() + } + "artifact_id" => { + wire[field] = serde_json::to_value( + radroots_storage::authored::AuthoredArtifactId::new([99; 16]).unwrap(), + ) + .unwrap() + } + "request" => wire[field] = serde_json::Value::Null, + _ => wire[field] = serde_json::json!(9), + } + let altered = serde_json::from_value(wire).unwrap(); + let forged = restored_receipt(&original, altered, 13); + let mut history = AuthoredDeliveryHistory::new(current.clone(), None).unwrap(); + assert_eq!( + history.push_claim(&forged), + Err(Error::AtomicWorkflowMismatch), + "{field}" + ); + assert!(history.claims().is_empty()); + } + let mut rebound = serde_json::to_value(¤t).unwrap(); + rebound["request"] = serde_json::to_value( + current + .intent() + .materialize(radroots_transport::sink::DeliveryPayload::new(event( + OTHER_RAW, + ))) + .unwrap(), + ) + .unwrap(); + let mut history = + AuthoredDeliveryHistory::new(serde_json::from_value(rebound).unwrap(), None).unwrap(); + assert_eq!( + history.push_claim(&original), + Err(Error::AtomicWorkflowMismatch) + ); + let mut history = AuthoredDeliveryHistory::new(current.clone(), None).unwrap(); + assert_eq!( + history.push_claim(&restored_receipt(&original, current.clone(), 14)), + Err(Error::AtomicWorkflowMismatch) + ); + let corrupt_digest = AuthoredAtomicReceipt::from_durable_parts( + original.commit_id(), + AtomicCommitDigest::new([99; 32]), + AtomicCommitDisposition::Committed, + 13, + original.outcome().clone(), + ) + .unwrap(); + assert_eq!( + history.push_claim(&corrupt_digest), + Err(Error::AtomicWorkflowMismatch) + ); + // A structurally valid future original cannot explain the current row. + for (revision, at) in [ + (current.revision().get() + 1, 13), + (current.revision().get(), 14), + ] { + let mut wire = serde_json::to_value(¤t).unwrap(); + let active = claim(NonZeroU64::new(revision - 1).unwrap(), 8, at); + wire["revision"] = serde_json::json!(revision); + wire["updated_at_unix_ms"] = serde_json::json!(at); + wire["claim"] = serde_json::to_value(&active).unwrap(); + let future = serde_json::from_value(wire).unwrap(); + let command = AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( + ClaimAuthoredTarget::DeliveryPlan(ids().2), + active, + )); + let receipt = AuthoredAtomicReceipt::new( + &command, + AtomicCommitDisposition::Committed, + at, + AuthoredAtomicOutcome::DeliveryPlan(future), + ) + .unwrap(); + assert_eq!( + history.push_claim(&receipt), + Err(Error::AtomicWorkflowMismatch) + ); + } +} + +#[test] +fn preparation_history_requires_exact_initial_identity_and_monotonic_row() { + let storage = MemoryStorage::new(SourceGeneration::new([73; 32]).unwrap()); + let original = block_on(storage.execute_authored(prepare().0)).unwrap(); + let initial = plan(&storage); + for field in ["plan_id", "artifact_id", "created_at_unix_ms"] { + let mut wire = serde_json::to_value(&initial).unwrap(); + match field { + "plan_id" => { + wire[field] = serde_json::to_value( + radroots_storage::authored_delivery::AuthoredDeliveryPlanId::new([99; 16]) + .unwrap(), + ) + .unwrap() + } + "artifact_id" => { + wire[field] = serde_json::to_value( + radroots_storage::authored::AuthoredArtifactId::new([99; 16]).unwrap(), + ) + .unwrap() + } + _ => wire[field] = serde_json::json!(9), + } + assert!( + AuthoredDeliveryHistory::new(serde_json::from_value(wire).unwrap(), Some(&original)) + .is_err(), + "{field}" + ); + } + let intent = radroots_storage::authored_delivery::AuthoredDeliveryIntent::new( + "different-intent", + initial.intent().target_set().clone(), + initial.intent().satisfaction().clone(), + 100, + ) + .unwrap(); + let changed = AuthoredDeliveryPlan::new(ids().2, ids().1, intent, 10).unwrap(); + assert!(AuthoredDeliveryHistory::new(changed, Some(&original)).is_err()); + for field in ["revision", "updated_at_unix_ms"] { + let mut wire = serde_json::to_value(&initial).unwrap(); + wire[field] = serde_json::json!(12); + let altered = serde_json::from_value(wire).unwrap(); + let AuthoredAtomicOutcome::Prepared { + operation, + artifacts, + .. + } = original.outcome() + else { + unreachable!() + }; + let receipt = AuthoredAtomicReceipt::from_durable_parts( + original.commit_id(), + original.digest(), + AtomicCommitDisposition::Committed, + 12, + AuthoredAtomicOutcome::Prepared { + operation: operation.clone(), + artifacts: artifacts.clone(), + delivery_plans: vec![altered], + }, + ) + .unwrap(); + assert!( + AuthoredDeliveryHistory::new(initial.clone(), Some(&receipt)).is_err(), + "{field}" + ); + } +} + +#[test] +fn reconciliation_receipts_cannot_substitute_a_different_plan_or_result() { + let (storage, active, _) = prepared(); + block_on(storage.execute_authored(fact(&plan(&storage), active.clone(), false, 14))).unwrap(); + let command = reconcile( + &plan(&storage), + Some(fence(&active)), + Some(retry(1, 20)), + 15, + ); + let receipt = block_on(storage.execute_authored(command.clone())).unwrap(); + let after = plan(&storage); + for field in [ + "plan_id", + "revision", + "updated_at_unix_ms", + "claim", + "stop", + "pending", + "facts", + "retry", + ] { + let mut wire = serde_json::to_value(&after).unwrap(); + match field { + "plan_id" => { + wire[field] = serde_json::to_value( + radroots_storage::authored_delivery::AuthoredDeliveryPlanId::new([99; 16]) + .unwrap(), + ) + .unwrap() + } + "revision" => wire[field] = serde_json::json!(after.revision().get() + 1), + "updated_at_unix_ms" => wire[field] = serde_json::json!(16), + "claim" => { + wire[field] = serde_json::to_value(claim( + NonZeroU64::new(after.revision().get() - 1).unwrap(), + 9, + 15, + )) + .unwrap(); + } + "stop" => { + wire["stop_requested_at_unix_ms"] = serde_json::json!(15); + wire["state"] = serde_json::json!("cancelled"); + wire["retry"] = serde_json::Value::Null; + wire["last_failure"] = serde_json::Value::Null; + } + "pending" => wire["attempts"][0]["claim"] = serde_json::Value::Null, + "facts" => wire["delivery_facts"][0]["observed_at_unix_ms"] = serde_json::json!(15), + _ => { + let other = retry(1, 21); + wire["retry"] = serde_json::to_value(&other).unwrap(); + wire["last_failure"] = serde_json::to_value(other.failure()).unwrap(); + } + } + let altered: AuthoredDeliveryPlan = serde_json::from_value(wire).unwrap(); + let wrong = restored_receipt(&receipt, altered.clone(), 15); + assert!(!wrong.matches_command(&command), "{field}"); + assert_eq!( + AuthoredAtomicReceipt::new( + &command, + AtomicCommitDisposition::Committed, + 15, + AuthoredAtomicOutcome::DeliveryPlan(altered) + ), + Err(Error::AtomicWorkflowMismatch), + "{field}" + ); + } + let artifact = block_on(storage.authored_artifact(ids().1)) + .unwrap() + .unwrap(); + let wrong = AuthoredAtomicReceipt::from_durable_parts( + receipt.commit_id(), + receipt.digest(), + AtomicCommitDisposition::Committed, + 15, + AuthoredAtomicOutcome::Artifact(artifact.clone()), + ) + .unwrap(); + assert!(!wrong.matches_command(&command)); + assert!( + AuthoredAtomicReceipt::new( + &command, + AtomicCommitDisposition::Committed, + 15, + AuthoredAtomicOutcome::Artifact(artifact) + ) + .is_err() + ); +} + +#[test] +fn missing_fact_claims_and_prepopulated_preparation_remain_uncertain() { + let (storage, active, _) = prepared(); + let original = block_on(storage.authored_receipt(prepare().0.commit_id())) + .unwrap() + .unwrap(); + let AuthoredAtomicOutcome::Prepared { + operation, + artifacts, + .. + } = original.outcome() + else { + unreachable!() + }; + let prepopulated = AuthoredAtomicReceipt::from_durable_parts( + original.commit_id(), + original.digest(), + AtomicCommitDisposition::Committed, + 13, + AuthoredAtomicOutcome::Prepared { + operation: operation.clone(), + artifacts: artifacts.clone(), + delivery_plans: vec![plan(&storage)], + }, + ) + .unwrap(); + let legacy = AuthoredDeliveryHistory::new(plan(&storage), Some(&prepopulated)).unwrap(); + assert!(!legacy.is_complete()); + assert!(!legacy.proves_no_issued_attempt()); + assert!(legacy.has_unresolved_claims()); + block_on(storage.execute_authored(fact(&plan(&storage), active, true, 50))).unwrap(); + stop(&storage, 51); + let unknown = AuthoredDeliveryHistory::new(plan(&storage), None).unwrap(); + assert_eq!( + unknown.require_pending_fact_provenance(), + Err(Error::AtomicWorkflowMismatch) + ); + let incomplete = AuthoredDeliveryHistory::new(plan(&storage), Some(&original)).unwrap(); + assert_eq!(incomplete.validate(), Err(Error::AtomicWorkflowMismatch)); +} + +#[test] +fn later_legacy_attempt_outside_old_lease_cannot_resolve_old_claim() { + let (storage, _, _) = 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(); + block_on( + storage.execute_authored(AuthoredAtomicCommand::ApplyDelivery( + ApplyDeliveryAttempt::new( + ids().2, + fence(&newer), + outcome(&plan(&storage), true), + None, + 41, + ) + .unwrap(), + )), + ) + .unwrap(); + let history = history(&storage); + assert_eq!(history.plan().attempt_count(), 1); + assert_eq!(history.plan().state(), AuthoredDeliveryState::Satisfied); + assert!(history.has_unresolved_claims()); +} diff --git a/crates/storage/tests/authored_delivery_facts.rs b/crates/storage/tests/authored_delivery_facts.rs @@ -27,6 +27,9 @@ use std::num::NonZeroU64; mod fixture; use fixture::*; +#[path = "authored_delivery/reconciliation_tests.rs"] +mod reconciliation; + fn plan(storage: &MemoryStorage) -> AuthoredDeliveryPlan { block_on(storage.authored_delivery_plan(ids().2)) .unwrap() diff --git a/crates/storage_sqlite/src/authored.rs b/crates/storage_sqlite/src/authored.rs @@ -1,6 +1,8 @@ use crate::SqliteStorage; #[path = "authored_delivery_facts.rs"] mod delivery_facts; +#[path = "authored_delivery_reconciliation.rs"] +mod delivery_reconciliation; use radroots_storage::{ Error, atomic::{AtomicCommitDigest, AtomicCommitDisposition, AtomicCommitId}, @@ -30,6 +32,26 @@ struct ReceiptSnapshot { } impl AuthoredAtomicStorage for SqliteStorage { + fn authored_delivery_history( + &self, + plan_id: AuthoredDeliveryPlanId, + ) -> BoxFuture< + '_, + Result<Option<radroots_storage::authored_delivery::AuthoredDeliveryHistory>, Error>, + > { + Box::pin(async move { + let mut transaction = self.pool().begin().await.map_err(map_backend)?; + let result = delivery_reconciliation::history(&mut transaction, plan_id).await; + let rollback = transaction.rollback().await.map_err(map_backend); + match result { + Ok(history) => { + rollback?; + Ok(history) + } + Err(error) => Err(error), + } + }) + } fn execute_authored( &self, command: AuthoredAtomicCommand, @@ -243,6 +265,7 @@ async fn commit_outcome( .execute(&mut **transaction) .await .map_err(map_backend)?; + delivery_reconciliation::record_claim(transaction, command).await?; Ok(receipt) } @@ -430,6 +453,10 @@ async fn execute_command( persist_plan(transaction, &plan).await?; Ok(AuthoredAtomicOutcome::DeliveryPlan(plan)) } + AuthoredAtomicCommand::ReconcileDelivery(value) => { + let plan = delivery_reconciliation::reconcile(transaction, &value).await?; + Ok(AuthoredAtomicOutcome::DeliveryPlan(plan)) + } AuthoredAtomicCommand::ApplyDelivery(value) => { let mut plan = load_plan_tx(transaction, value.plan_id()).await?; match value.outcome().clone() { @@ -697,7 +724,8 @@ pub(crate) async fn persist_plan( plan: &AuthoredDeliveryPlan, ) -> Result<(), Error> { persist_plan_v11(transaction, plan).await?; - delivery_facts::persist(transaction, plan).await + delivery_facts::persist(transaction, plan).await?; + delivery_reconciliation::persist(transaction, plan).await } // Legacy conversion runs before the delivery-facts forward migration. @@ -887,7 +915,8 @@ async fn validate_plan_children_tx( .fetch_all(&mut **transaction) .await .map_err(map_backend)?; - delivery_facts::validate(plan, &facts) + delivery_facts::validate(plan, &facts)?; + delivery_reconciliation::validate(transaction, plan).await } fn validate_plan_children( @@ -1179,6 +1208,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::ReconcileDelivery(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() { @@ -1199,9 +1229,9 @@ fn command_phase(command: &AuthoredAtomicCommand) -> &'static str { AuthoredAtomicCommand::Claim(_) => "claim", AuthoredAtomicCommand::ApplySigned(_) | AuthoredAtomicCommand::RecordSigned(_) => "signing", AuthoredAtomicCommand::ApplyAdmission(_) => "admission", - AuthoredAtomicCommand::ApplyDelivery(_) | AuthoredAtomicCommand::RecordDelivery(_) => { - "delivery" - } + AuthoredAtomicCommand::ApplyDelivery(_) + | AuthoredAtomicCommand::RecordDelivery(_) + | AuthoredAtomicCommand::ReconcileDelivery(_) => "delivery", AuthoredAtomicCommand::ApplyFailure(value) => match value.failure().phase() { WorkPhase::Signing => "signing_failure", WorkPhase::Admission => "admission_failure", @@ -2645,5 +2675,10 @@ mod delivery_fact_tests; #[cfg(test)] #[cfg_attr(coverage_nightly, coverage(off))] +#[path = "authored_delivery_reconciliation_tests.rs"] +mod delivery_reconciliation_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 @@ -30,7 +30,7 @@ async fn prepared(temp: &TempDir) -> (SqliteStorage, WorkClaim, AuthoredAtomicRe (store, active, original) } -fn fact(plan: &AuthoredDeliveryPlan, active: WorkClaim) -> AuthoredAtomicCommand { +pub(super) fn fact(plan: &AuthoredDeliveryPlan, active: WorkClaim) -> AuthoredAtomicCommand { let request = plan.request().unwrap(); let receipt = DeliveryReceipt::for_request( request, diff --git a/crates/storage_sqlite/src/authored_delivery_reconciliation.rs b/crates/storage_sqlite/src/authored_delivery_reconciliation.rs @@ -0,0 +1,187 @@ +//! Indexed issued claims and immutable fact-to-attempt provenance. + +use super::{ + column, decode_receipt_row, load_artifact_tx, load_optional_plan_tx, map_backend, persist_plan, +}; +use radroots_storage::{ + Error, + atomic::AtomicCommitId, + authored_atomic::{ + AuthoredAtomicCommand, AuthoredAtomicReceipt, ClaimAuthoredTarget, ClaimAuthoredWork, + ReconcileDeliveryFacts, + }, + authored_delivery::{ + AuthoredDeliveryHistory, AuthoredDeliveryPlan, AuthoredDeliveryPlanId, + DELIVERY_PLAN_ATTEMPTS_MAX, + }, +}; +use sqlx::{Sqlite, sqlite::SqliteRow}; + +const MARKERS: &str = "SELECT attempt, claim_id FROM radroots_runtime_authored_delivery_reconciliations WHERE plan_id = ? ORDER BY attempt LIMIT 1025"; + +async fn receipt( + transaction: &mut sqlx::Transaction<'_, Sqlite>, + id: &[u8], +) -> Result<AuthoredAtomicReceipt, Error> { + 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(id).fetch_one(&mut **transaction).await.map_err(map_backend)?; + decode_receipt_row(&row) +} + +pub(super) async fn history( + transaction: &mut sqlx::Transaction<'_, Sqlite>, + plan_id: AuthoredDeliveryPlanId, +) -> Result<Option<AuthoredDeliveryHistory>, Error> { + let Some(plan) = load_optional_plan_tx(transaction, plan_id).await? else { + return Ok(None); + }; + let operation_id = load_artifact_tx(transaction, plan.artifact_id()) + .await? + .operation_id(); + let preparations: Vec<Vec<u8>> = sqlx::query_scalar("SELECT commit_id FROM radroots_runtime_authored_atomic_commits WHERE target_id = ? AND phase = 'prepare' AND json_type(CAST(receipt AS TEXT), '$.outcome.prepared') = 'object' ORDER BY commit_id LIMIT 2") + .bind(operation_id.as_bytes().as_slice()).fetch_all(&mut **transaction).await.map_err(map_backend)?; + if preparations.len() > 1 { + return Err(Error::AtomicWorkflowMismatch); + } + let original = match preparations.first() { + Some(id) => Some(receipt(transaction, id).await?), + None => None, + }; + let mut history = AuthoredDeliveryHistory::new(plan, original.as_ref())?; + drop(original); + // Fetch only bounded IDs, then decode one bounded original receipt at a time. + let ids: Vec<Vec<u8>> = sqlx::query_scalar("SELECT claim_id FROM radroots_runtime_authored_delivery_claims WHERE plan_id = ? ORDER BY claim_id LIMIT 1025") + .bind(plan_id.as_bytes().as_slice()).fetch_all(&mut **transaction).await.map_err(map_backend)?; + if ids.len() > DELIVERY_PLAN_ATTEMPTS_MAX as usize { + history.mark_truncated(); + } + for id in ids.iter().take(DELIVERY_PLAN_ATTEMPTS_MAX as usize) { + history.push_claim(&receipt(transaction, id).await?)?; + } + history.validate()?; + Ok(Some(history)) +} + +pub(super) async fn record_claim( + transaction: &mut sqlx::Transaction<'_, Sqlite>, + command: &AuthoredAtomicCommand, +) -> Result<(), Error> { + let AuthoredAtomicCommand::Claim(value) = command else { + return Ok(()); + }; + let ClaimAuthoredTarget::DeliveryPlan(plan_id) = value.target() else { + return Ok(()); + }; + let count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM (SELECT 1 FROM radroots_runtime_authored_delivery_claims WHERE plan_id = ? LIMIT 1024)") + .bind(plan_id.as_bytes().as_slice()).fetch_one(&mut **transaction).await.map_err(map_backend)?; + if count >= i64::from(DELIVERY_PLAN_ATTEMPTS_MAX) { + return Err(Error::DeliveryAttemptOverflow); + } + sqlx::query( + "INSERT INTO radroots_runtime_authored_delivery_claims (plan_id, claim_id) VALUES (?, ?)", + ) + .bind(plan_id.as_bytes().as_slice()) + .bind(command.commit_id().as_bytes().as_slice()) + .execute(&mut **transaction) + .await + .map_err(map_backend)?; + Ok(()) +} + +pub(super) async fn reconcile( + transaction: &mut sqlx::Transaction<'_, Sqlite>, + command: &ReconcileDeliveryFacts, +) -> Result<AuthoredDeliveryPlan, Error> { + let history = history(transaction, command.plan_id()) + .await? + .ok_or(Error::InvalidAuthoredDeliveryPlan)?; + let plan = command.apply_to(&history)?; + persist_plan(transaction, &plan).await?; + Ok(plan) +} + +fn marker_id( + plan_id: AuthoredDeliveryPlanId, + claim: &radroots_storage::authored::WorkClaim, +) -> AtomicCommitId { + AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( + ClaimAuthoredTarget::DeliveryPlan(plan_id), + claim.clone(), + )) + .commit_id() +} + +fn validate_existing( + plan: &AuthoredDeliveryPlan, + rows: &[SqliteRow], +) -> Result<std::collections::BTreeSet<u32>, Error> { + let mut existing = std::collections::BTreeSet::new(); + for row in rows { + let ordinal = u32::try_from(column::<i64>(row, "attempt")?) + .map_err(|_| Error::InvalidAuthoredDeliveryPlan)?; + let index = ordinal + .checked_sub(1) + .and_then(|index| usize::try_from(index).ok()) + .ok_or(Error::InvalidAuthoredDeliveryPlan)?; + let claim = plan + .attempts() + .get(index) + .and_then(|attempt| attempt.claim_evidence()) + .ok_or(Error::InvalidAuthoredDeliveryPlan)?; + if !existing.insert(ordinal) + || column::<Vec<u8>>(row, "claim_id")?.as_slice() + != marker_id(plan.plan_id(), claim).as_bytes() + { + return Err(Error::InvalidAuthoredDeliveryPlan); + } + } + Ok(existing) +} + +pub(super) async fn persist( + transaction: &mut sqlx::Transaction<'_, Sqlite>, + plan: &AuthoredDeliveryPlan, +) -> Result<(), Error> { + let existing = sqlx::query(MARKERS) + .bind(plan.plan_id().as_bytes().as_slice()) + .fetch_all(&mut **transaction) + .await + .map_err(map_backend)?; + let existing = validate_existing(plan, &existing)?; + for (attempt, claim) in plan + .attempts() + .iter() + .filter_map(|attempt| { + attempt + .claim_evidence() + .map(|claim| (attempt.attempt(), claim)) + }) + .filter(|(attempt, _)| !existing.contains(&attempt.get())) + { + sqlx::query("INSERT INTO radroots_runtime_authored_delivery_reconciliations (plan_id, attempt, claim_id) VALUES (?, ?, ?)") + .bind(plan.plan_id().as_bytes().as_slice()).bind(i64::from(attempt.get())) + .bind(marker_id(plan.plan_id(), claim).as_bytes().as_slice()) + .execute(&mut **transaction).await.map_err(map_backend)?; + } + Ok(()) +} + +pub(super) async fn validate( + transaction: &mut sqlx::Transaction<'_, Sqlite>, + plan: &AuthoredDeliveryPlan, +) -> Result<(), Error> { + let rows = sqlx::query(MARKERS) + .bind(plan.plan_id().as_bytes().as_slice()) + .fetch_all(&mut **transaction) + .await + .map_err(map_backend)?; + let count = plan + .attempts() + .iter() + .filter(|attempt| attempt.claim_evidence().is_some()) + .count(); + if rows.len() != count { + return Err(Error::InvalidAuthoredDeliveryPlan); + } + validate_existing(plan, &rows).map(|_| ()) +} diff --git a/crates/storage_sqlite/src/authored_delivery_reconciliation_tests.rs b/crates/storage_sqlite/src/authored_delivery_reconciliation_tests.rs @@ -0,0 +1,637 @@ +use super::signed_fact_fixture::{claim, ids, record}; +use super::*; +use crate::OpenMode; +use radroots_storage::{ + authored::WorkClaim, + authored_atomic::{ClaimAuthoredWork, ReconcileDeliveryFacts}, +}; +use tempfile::TempDir; + +async fn signed(temp: &TempDir) -> SqliteStorage { + let (store, event, signing) = signed_fact_tests::prepared(temp).await; + store + .execute_authored(record(event, signing, 12)) + .await + .unwrap(); + store +} + +async fn plan(store: &SqliteStorage) -> AuthoredDeliveryPlan { + store + .authored_delivery_plan(ids().2) + .await + .unwrap() + .unwrap() +} + +async fn issue(store: &SqliteStorage) -> (WorkClaim, AuthoredAtomicReceipt) { + let active = claim(plan(store).await.revision(), 5, 13); + let receipt = store + .execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( + ClaimAuthoredTarget::DeliveryPlan(ids().2), + active.clone(), + ))) + .await + .unwrap(); + (active, receipt) +} + +fn command(plan: &AuthoredDeliveryPlan) -> AuthoredAtomicCommand { + AuthoredAtomicCommand::ReconcileDelivery( + ReconcileDeliveryFacts::new(plan, None, None, 51).unwrap(), + ) +} + +#[tokio::test] +async fn claim_and_reconciliation_provenance_reopen_exactly_and_survive_stop() { + let temp = TempDir::new().unwrap(); + let store = signed(&temp).await; + assert!( + store + .authored_delivery_history(ids().2) + .await + .unwrap() + .unwrap() + .proves_no_issued_attempt() + ); + let (active, issued) = issue(&store).await; + let history = store + .authored_delivery_history(ids().2) + .await + .unwrap() + .unwrap(); + assert_eq!(history.claims().len(), 1); + assert!(history.has_unresolved_claims()); + store + .execute_authored(delivery_fact_tests::fact( + &plan(&store).await, + active.clone(), + )) + .await + .unwrap(); + let before = plan(&store).await; + let command = command(&before); + let committed = store.execute_authored(command.clone()).await.unwrap(); + let after = plan(&store).await; + assert_eq!(after.state(), AuthoredDeliveryState::Satisfied); + assert_eq!(after.attempt_count(), 1); + assert_eq!(after.attempts()[0].claim_evidence(), Some(&active)); + assert_eq!(after.delivery_facts(), before.delivery_facts()); + assert_eq!( + store + .authored_receipt(issued.commit_id()) + .await + .unwrap() + .unwrap(), + issued + ); + assert!( + !store + .authored_delivery_history(ids().2) + .await + .unwrap() + .unwrap() + .has_unresolved_claims() + ); + for sql in [ + "DELETE FROM radroots_runtime_authored_delivery_claims", + "UPDATE radroots_runtime_authored_delivery_claims SET claim_id = claim_id", + "DELETE FROM radroots_runtime_authored_delivery_reconciliations", + "UPDATE radroots_runtime_authored_delivery_reconciliations SET attempt = 2", + ] { + assert!(sqlx::query(sql).execute(store.pool()).await.is_err()); + } + store.close().await.unwrap(); + let store = signed_fact_tests::open(&temp, OpenMode::ReadWriteExisting).await; + assert_eq!(plan(&store).await, after); + let replay = store.execute_authored(command).await.unwrap(); + assert_eq!(replay.disposition(), AtomicCommitDisposition::Replay); + assert_eq!(replay.outcome(), committed.outcome()); + store + .execute_authored(AuthoredAtomicCommand::Cancel( + radroots_storage::authored_atomic::CancelAuthoredWork::new( + CancelAuthoredTarget::DeliveryPlan(ids().2), + after.revision(), + 60, + ) + .unwrap(), + )) + .await + .unwrap(); + let stopped = plan(&store).await; + assert_eq!(stopped.stop_requested_at_unix_ms(), Some(60)); + assert_eq!(stopped.attempts(), after.attempts()); + assert_eq!(stopped.delivery_facts(), after.delivery_facts()); + store.close().await.unwrap(); +} + +#[tokio::test] +async fn reconciliation_commit_failure_rolls_back_marker_plan_and_receipt() { + let temp = TempDir::new().unwrap(); + let store = signed(&temp).await; + let (active, _) = issue(&store).await; + store + .execute_authored(delivery_fact_tests::fact(&plan(&store).await, active)) + .await + .unwrap(); + let before = plan(&store).await; + let before_history = store + .authored_delivery_history(ids().2) + .await + .unwrap() + .unwrap(); + let command = command(&before); + sqlx::query("CREATE TABLE reconciliation_commit_fault (parent BLOB REFERENCES radroots_runtime_authored_operations(operation_id) DEFERRABLE INITIALLY DEFERRED)").execute(store.pool()).await.unwrap(); + sqlx::query("CREATE TRIGGER reconciliation_commit_fault_trigger AFTER INSERT ON radroots_runtime_authored_delivery_reconciliations BEGIN INSERT INTO reconciliation_commit_fault VALUES (x'99999999999999999999999999999999'); END").execute(store.pool()).await.unwrap(); + assert!(store.execute_authored(command.clone()).await.is_err()); + assert_eq!(plan(&store).await, before); + assert_eq!( + store + .authored_delivery_history(ids().2) + .await + .unwrap() + .unwrap(), + before_history + ); + assert!( + store + .authored_receipt(command.commit_id()) + .await + .unwrap() + .is_none() + ); + assert_eq!( + sqlx::query_scalar::<_, i64>( + "SELECT COUNT(*) FROM radroots_runtime_authored_delivery_reconciliations" + ) + .fetch_one(store.pool()) + .await + .unwrap(), + 0 + ); + sqlx::query("DROP TRIGGER reconciliation_commit_fault_trigger") + .execute(store.pool()) + .await + .unwrap(); + store.execute_authored(command).await.unwrap(); + assert_eq!(plan(&store).await.attempt_count(), 1); + // The existing writer may replace these rows only within its transaction; + // a committed missing attempt must never orphan retained reconciliation. + assert!( + sqlx::query("DELETE FROM radroots_runtime_authored_delivery_attempts") + .execute(store.pool()) + .await + .is_err() + ); + assert_eq!(plan(&store).await.attempt_count(), 1); + store.close().await.unwrap(); +} + +#[tokio::test] +async fn original_claim_commit_failure_preserves_no_issued_attempt_proof() { + let temp = TempDir::new().unwrap(); + let store = signed(&temp).await; + let before = plan(&store).await; + let active = claim(before.revision(), 5, 13); + let command = AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( + ClaimAuthoredTarget::DeliveryPlan(ids().2), + active, + )); + sqlx::query("CREATE TABLE claim_index_commit_fault (parent BLOB REFERENCES radroots_runtime_authored_operations(operation_id) DEFERRABLE INITIALLY DEFERRED)").execute(store.pool()).await.unwrap(); + sqlx::query("CREATE TRIGGER claim_index_commit_fault_trigger AFTER INSERT ON radroots_runtime_authored_delivery_claims BEGIN INSERT INTO claim_index_commit_fault VALUES (x'99999999999999999999999999999999'); END").execute(store.pool()).await.unwrap(); + assert!(store.execute_authored(command.clone()).await.is_err()); + assert_eq!(plan(&store).await, before); + assert!( + store + .authored_receipt(command.commit_id()) + .await + .unwrap() + .is_none() + ); + assert!( + store + .authored_delivery_history(ids().2) + .await + .unwrap() + .unwrap() + .proves_no_issued_attempt() + ); + sqlx::query("DROP TRIGGER claim_index_commit_fault_trigger") + .execute(store.pool()) + .await + .unwrap(); + store.execute_authored(command).await.unwrap(); + let issued = store + .authored_delivery_history(ids().2) + .await + .unwrap() + .unwrap(); + assert_eq!(issued.claims().len(), 1); + assert!(issued.has_unresolved_claims()); + assert!(!issued.proves_no_issued_attempt()); + store.close().await.unwrap(); +} + +#[tokio::test] +async fn history_read_snapshot_cannot_mix_old_plan_with_new_claim_or_reconciliation() { + let temp = TempDir::new().unwrap(); + let store = signed(&temp).await; + let mut read = store.pool().begin().await.unwrap(); + let frozen = delivery_reconciliation::history(&mut read, ids().2) + .await + .unwrap() + .unwrap(); + assert!(frozen.proves_no_issued_attempt()); + let (active, _) = issue(&store).await; + assert_eq!( + delivery_reconciliation::history(&mut read, ids().2) + .await + .unwrap() + .unwrap(), + frozen + ); + read.rollback().await.unwrap(); + let mut read = store.pool().begin().await.unwrap(); + let issued = delivery_reconciliation::history(&mut read, ids().2) + .await + .unwrap() + .unwrap(); + store + .execute_authored(delivery_fact_tests::fact(&plan(&store).await, active)) + .await + .unwrap(); + store + .execute_authored(command(&plan(&store).await)) + .await + .unwrap(); + assert_eq!( + delivery_reconciliation::history(&mut read, ids().2) + .await + .unwrap() + .unwrap(), + issued + ); + read.rollback().await.unwrap(); + let after = store + .authored_delivery_history(ids().2) + .await + .unwrap() + .unwrap(); + assert!(!after.has_unresolved_claims()); + assert_eq!(after.plan().attempt_count(), 1); + let mut stale = store.pool().begin().await.unwrap(); + assert_eq!( + delivery_reconciliation::persist(&mut stale, issued.plan()).await, + Err(Error::InvalidAuthoredDeliveryPlan) + ); + stale.rollback().await.unwrap(); + assert_eq!( + store + .authored_delivery_history(ids().2) + .await + .unwrap() + .unwrap(), + after + ); + store.close().await.unwrap(); +} + +#[tokio::test] +async fn oversized_historical_claims_remain_retained_but_cannot_authorize_new_work() { + let temp = TempDir::new().unwrap(); + let store = signed(&temp).await; + let mut current = plan(&store).await; + let mut transaction = store.pool().begin().await.unwrap(); + // Historical databases may exceed the new admission bound. Construct their + // typed original receipts without invoking the now-bounded command path. + for index in 0..1025u64 { + let at = 13 + index * 21; + let active = WorkClaim::new( + [7; 16], + "historical-worker", + std::num::NonZeroU64::new(index + 5).unwrap(), + at, + at + 20, + current.revision(), + ) + .unwrap(); + current.claim(active.clone(), at).unwrap(); + let command = AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( + ClaimAuthoredTarget::DeliveryPlan(ids().2), + active, + )); + let receipt = AuthoredAtomicReceipt::new( + &command, + AtomicCommitDisposition::Committed, + at, + AuthoredAtomicOutcome::DeliveryPlan(current.clone()), + ) + .unwrap(); + sqlx::query("INSERT INTO radroots_runtime_authored_atomic_commits (commit_id, commit_digest, phase, target_id, requested_at_unix_ms, committed_at_unix_ms, receipt) VALUES (?, ?, 'claim', ?, ?, ?, ?)") + .bind(receipt.commit_id().as_bytes().as_slice()).bind(receipt.digest().as_bytes().as_slice()).bind(ids().2.as_bytes().as_slice()).bind(at as i64).bind(at as i64) + .bind(serde_json::to_vec(&serde_json::json!({"outcome": receipt.outcome()})).unwrap()).execute(&mut *transaction).await.unwrap(); + sqlx::query("INSERT INTO radroots_runtime_authored_delivery_claims (plan_id, claim_id) VALUES (?, ?)") + .bind(ids().2.as_bytes().as_slice()).bind(receipt.commit_id().as_bytes().as_slice()).execute(&mut *transaction).await.unwrap(); + } + persist_plan(&mut transaction, ¤t).await.unwrap(); + transaction.commit().await.unwrap(); + let history = store + .authored_delivery_history(ids().2) + .await + .unwrap() + .unwrap(); + assert_eq!(history.claims().len(), 1024); + assert!(history.is_truncated()); + assert!(!history.is_complete()); + assert!(!history.proves_no_issued_attempt()); + assert!(history.has_unresolved_claims()); + assert_eq!( + history.require_pending_fact_provenance(), + Err(Error::DeliveryAttemptOverflow) + ); + let at = 13 + 1025 * 21; + let active = WorkClaim::new( + [8; 16], + "new-worker", + std::num::NonZeroU64::new(2048).unwrap(), + at, + at + 20, + current.revision(), + ) + .unwrap(); + let command = AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( + ClaimAuthoredTarget::DeliveryPlan(ids().2), + active, + )); + assert_eq!( + store.execute_authored(command.clone()).await, + Err(Error::DeliveryAttemptOverflow) + ); + assert_eq!(plan(&store).await, current); + assert!( + store + .authored_receipt(command.commit_id()) + .await + .unwrap() + .is_none() + ); + assert_eq!( + sqlx::query_scalar::<_, i64>( + "SELECT COUNT(*) FROM radroots_runtime_authored_delivery_claims" + ) + .fetch_one(store.pool()) + .await + .unwrap(), + 1025 + ); + store.close().await.unwrap(); +} + +#[tokio::test] +async fn forged_indexed_claim_receipt_fails_closed_instead_of_proving_absence() { + let temp = TempDir::new().unwrap(); + let store = signed(&temp).await; + let (_, issued) = issue(&store).await; + let before = plan(&store).await; + sqlx::query("INSERT INTO radroots_runtime_authored_atomic_commits (commit_id, commit_digest, phase, target_id, requested_at_unix_ms, committed_at_unix_ms, receipt) VALUES (?, ?, 'claim', ?, 13, 13, ?)") + .bind([9u8; 16].as_slice()).bind([8u8; 32].as_slice()).bind(ids().2.as_bytes().as_slice()) + .bind(serde_json::to_vec(&serde_json::json!({"outcome": issued.outcome()})).unwrap()).execute(store.pool()).await.unwrap(); + sqlx::query( + "INSERT INTO radroots_runtime_authored_delivery_claims (plan_id, claim_id) VALUES (?, ?)", + ) + .bind(ids().2.as_bytes().as_slice()) + .bind([9u8; 16].as_slice()) + .execute(store.pool()) + .await + .unwrap(); + assert!(store.authored_delivery_history(ids().2).await.is_err()); + assert_eq!(plan(&store).await, before); + store.close().await.unwrap(); +} + +#[tokio::test] +async fn late_legacy_marker_can_precede_existing_marker_without_rewriting_it() { + use radroots_storage::{ + authored::{FailureClass, RetrySchedule, WorkFailure, WorkPhase}, + authored_atomic::{ApplyDeliveryAttempt, RecordDeliveryFact, WorkFence}, + authored_delivery::DeliveryAttemptOutcome, + }; + use radroots_transport::{ + DeliveryReceipt, outcome::DeliveryOutcome, sink::DeliveryTargetReceipt, + }; + let retry = |attempt, at| { + RetrySchedule::new( + std::num::NonZeroU32::new(attempt).unwrap(), + at, + WorkFailure::new( + "delivery_pending", + WorkPhase::Delivery, + FailureClass::Retryable, + Some(at), + None, + ) + .unwrap(), + ) + .unwrap() + }; + let fence = |claim: &WorkClaim| { + WorkFence::new(*claim.token(), claim.generation(), claim.row_revision()).unwrap() + }; + let temp = TempDir::new().unwrap(); + let store = signed(&temp).await; + let (first, _) = issue(&store).await; + let current = plan(&store).await; + let request = current.request().unwrap(); + let outcome = DeliveryAttemptOutcome::Receipt( + DeliveryReceipt::for_request( + request, + request + .target_set() + .targets() + .iter() + .cloned() + .map(|target| { + DeliveryTargetReceipt::attempted(target, DeliveryOutcome::unavailable()) + }) + .collect(), + ) + .unwrap(), + ); + store + .execute_authored(AuthoredAtomicCommand::ApplyDelivery( + ApplyDeliveryAttempt::new( + ids().2, + fence(&first), + outcome.clone(), + Some(retry(1, 18)), + 14, + ) + .unwrap(), + )) + .await + .unwrap(); + let second = claim(plan(&store).await.revision(), 6, 20); + store + .execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( + ClaimAuthoredTarget::DeliveryPlan(ids().2), + second.clone(), + ))) + .await + .unwrap(); + store + .execute_authored(AuthoredAtomicCommand::RecordDelivery( + RecordDeliveryFact::new(ids().2, ids().1, second.clone(), outcome.clone(), 21).unwrap(), + )) + .await + .unwrap(); + store + .execute_authored(AuthoredAtomicCommand::ReconcileDelivery( + ReconcileDeliveryFacts::new( + &plan(&store).await, + Some(fence(&second)), + Some(retry(2, 25)), + 22, + ) + .unwrap(), + )) + .await + .unwrap(); + let marker_before: Vec<u8> = sqlx::query_scalar( + "SELECT claim_id FROM radroots_runtime_authored_delivery_reconciliations WHERE attempt = 2", + ) + .fetch_one(store.pool()) + .await + .unwrap(); + store + .execute_authored(AuthoredAtomicCommand::RecordDelivery( + RecordDeliveryFact::new(ids().2, ids().1, first.clone(), outcome, 50).unwrap(), + )) + .await + .unwrap(); + store + .execute_authored(AuthoredAtomicCommand::ReconcileDelivery( + ReconcileDeliveryFacts::new(&plan(&store).await, None, Some(retry(2, 60)), 51).unwrap(), + )) + .await + .unwrap(); + let after = plan(&store).await; + assert_eq!(after.attempt_count(), 2); + assert_eq!(after.attempts()[0].recorded_at_unix_ms(), 14); + assert_eq!(after.attempts()[0].claim_evidence(), Some(&first)); + assert_eq!(after.attempts()[1].recorded_at_unix_ms(), 22); + assert_eq!(after.attempts()[1].claim_evidence(), Some(&second)); + assert_eq!(sqlx::query_scalar::<_, Vec<u8>>("SELECT claim_id FROM radroots_runtime_authored_delivery_reconciliations WHERE attempt = 2").fetch_one(store.pool()).await.unwrap(), marker_before); + store.close().await.unwrap(); + let store = signed_fact_tests::open(&temp, OpenMode::ReadWriteExisting).await; + assert_eq!(plan(&store).await, after); + assert!( + !store + .authored_delivery_history(ids().2) + .await + .unwrap() + .unwrap() + .has_unresolved_claims() + ); + store.close().await.unwrap(); +} + +#[tokio::test] +async fn history_rejects_ambiguous_preparation_and_distinguishes_missing_plan() { + let temp = TempDir::new().unwrap(); + let store = signed(&temp).await; + assert!( + store + .authored_delivery_history(AuthoredDeliveryPlanId::new([99; 16]).unwrap()) + .await + .unwrap() + .is_none() + ); + let before = plan(&store).await; + sqlx::query("INSERT INTO radroots_runtime_authored_atomic_commits (commit_id, commit_digest, phase, target_id, requested_at_unix_ms, committed_at_unix_ms, receipt) SELECT ?, commit_digest, phase, target_id, requested_at_unix_ms, committed_at_unix_ms, receipt FROM radroots_runtime_authored_atomic_commits WHERE phase = 'prepare'") + .bind([9u8; 16].as_slice()).execute(store.pool()).await.unwrap(); + assert_eq!( + store.authored_delivery_history(ids().2).await, + Err(Error::AtomicWorkflowMismatch) + ); + assert_eq!(plan(&store).await, before); + store.close().await.unwrap(); +} + +#[tokio::test] +async fn missing_or_rebound_normalized_marker_cannot_return_a_valid_plan() { + for missing in [false, true] { + let temp = TempDir::new().unwrap(); + let store = signed(&temp).await; + let (first, _) = issue(&store).await; + let second = claim(plan(&store).await.revision(), 6, 40); + let second_command = AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( + ClaimAuthoredTarget::DeliveryPlan(ids().2), + second, + )); + store + .execute_authored(second_command.clone()) + .await + .unwrap(); + store + .execute_authored(delivery_fact_tests::fact(&plan(&store).await, first)) + .await + .unwrap(); + store + .execute_authored(AuthoredAtomicCommand::ReconcileDelivery( + ReconcileDeliveryFacts::new(&plan(&store).await, None, None, 61).unwrap(), + )) + .await + .unwrap(); + let before: Vec<u8> = + sqlx::query_scalar("SELECT snapshot FROM radroots_runtime_authored_delivery_plans") + .fetch_one(store.pool()) + .await + .unwrap(); + // Isolated corruption fixtures deliberately remove the applicable guard. + // The normal command path cannot erase or rebind these immutable rows. + if missing { + sqlx::query( + "DROP TRIGGER radroots_runtime_authored_delivery_reconciliations_delete_guard", + ) + .execute(store.pool()) + .await + .unwrap(); + sqlx::query("DELETE FROM radroots_runtime_authored_delivery_reconciliations") + .execute(store.pool()) + .await + .unwrap(); + } else { + sqlx::query( + "DROP TRIGGER radroots_runtime_authored_delivery_reconciliations_update_guard", + ) + .execute(store.pool()) + .await + .unwrap(); + sqlx::query( + "UPDATE radroots_runtime_authored_delivery_reconciliations SET claim_id = ?", + ) + .bind(second_command.commit_id().as_bytes().as_slice()) + .execute(store.pool()) + .await + .unwrap(); + } + assert_eq!( + store.authored_delivery_plan(ids().2).await, + Err(Error::InvalidAuthoredDeliveryPlan) + ); + assert_eq!( + store.authored_delivery_history(ids().2).await, + Err(Error::InvalidAuthoredDeliveryPlan) + ); + assert_eq!( + sqlx::query_scalar::<_, Vec<u8>>( + "SELECT snapshot FROM radroots_runtime_authored_delivery_plans" + ) + .fetch_one(store.pool()) + .await + .unwrap(), + before + ); + store.close().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 <= 16 + && plan.current_version <= 17 && plan.steps.len() == usize::try_from(plan.current_version).unwrap_or(usize::MAX) && plan .steps @@ -492,6 +492,7 @@ const fn set_user_version_sql(version: u32) -> Option<&'static str> { 14 => Some("PRAGMA user_version = 14"), 15 => Some("PRAGMA user_version = 15"), 16 => Some("PRAGMA user_version = 16"), + 17 => Some("PRAGMA user_version = 17"), _ => None, } } @@ -718,9 +719,9 @@ mod tests { .execute(&mut newer) .await .expect("application id"); - let newer_version = 17; + let newer_version = 18; assert_eq!(newer_version, runtime::CURRENT_VERSION + 1); - sqlx::raw_sql("PRAGMA user_version = 17") + sqlx::raw_sql("PRAGMA user_version = 18") .execute(&mut newer) .await .expect("newer version"); @@ -879,3 +880,8 @@ mod signed_facts_tests; #[cfg_attr(coverage_nightly, coverage(off))] #[path = "migration_delivery_facts_tests.rs"] mod delivery_facts_tests; + +#[cfg(test)] +#[cfg_attr(coverage_nightly, coverage(off))] +#[path = "migration_delivery_reconciliation_tests.rs"] +mod delivery_reconciliation_tests; diff --git a/crates/storage_sqlite/src/migration/runtime/0017_authored_delivery_reconciliation.up.sql b/crates/storage_sqlite/src/migration/runtime/0017_authored_delivery_reconciliation.up.sql @@ -0,0 +1,65 @@ +-- Project original issued delivery claims without rewriting their receipts. +-- Reject ambiguous claim namespaces before absence of a delivery projection +-- can be interpreted as evidence that no work was issued. +SELECT json(CASE + WHEN json_type(CAST(receipt AS TEXT), '$.outcome.delivery_plan') = 'object' + AND json_type(CAST(receipt AS TEXT), '$.outcome.artifact') IS NULL THEN 'null' + WHEN json_type(CAST(receipt AS TEXT), '$.outcome.artifact') = 'object' + AND json_type(CAST(receipt AS TEXT), '$.outcome.delivery_plan') IS NULL THEN 'null' + ELSE 'invalid claim outcome' +END) +FROM radroots_runtime_authored_atomic_commits WHERE phase = 'claim'; + +CREATE INDEX radroots_runtime_authored_atomic_target_phase_idx +ON radroots_runtime_authored_atomic_commits(target_id, phase, commit_id); + +CREATE TABLE radroots_runtime_authored_delivery_claims ( + plan_id BLOB NOT NULL REFERENCES radroots_runtime_authored_delivery_plans(plan_id), + claim_id BLOB NOT NULL REFERENCES radroots_runtime_authored_atomic_commits(commit_id), + PRIMARY KEY (plan_id, claim_id) +) STRICT, WITHOUT ROWID; + +INSERT INTO radroots_runtime_authored_delivery_claims (plan_id, claim_id) +SELECT target_id, commit_id FROM radroots_runtime_authored_atomic_commits +WHERE phase = 'claim' + AND json_type(CAST(receipt AS TEXT), '$.outcome.delivery_plan') = 'object'; + +CREATE TRIGGER radroots_runtime_authored_delivery_claims_update_guard +BEFORE UPDATE ON radroots_runtime_authored_delivery_claims +BEGIN + SELECT RAISE(ABORT, 'issued delivery claims are immutable'); +END; + +CREATE TRIGGER radroots_runtime_authored_delivery_claims_delete_guard +BEFORE DELETE ON radroots_runtime_authored_delivery_claims +BEGIN + SELECT RAISE(ABORT, 'issued delivery claims are retained'); +END; + +CREATE TABLE radroots_runtime_authored_delivery_reconciliations ( + plan_id BLOB NOT NULL, + attempt INTEGER NOT NULL CHECK (attempt BETWEEN 1 AND 1024), + claim_id BLOB NOT NULL, + PRIMARY KEY (plan_id, attempt), + UNIQUE (plan_id, claim_id), + FOREIGN KEY (plan_id, claim_id) + REFERENCES radroots_runtime_authored_delivery_claims(plan_id, claim_id), + -- The compatible normalized writer replaces attempt rows inside one + -- transaction. Require their exact final keys at COMMIT, without cascading + -- deletion or changing this immutable provenance during that replacement. + FOREIGN KEY (plan_id, attempt) + REFERENCES radroots_runtime_authored_delivery_attempts(plan_id, attempt) + DEFERRABLE INITIALLY DEFERRED +) STRICT, WITHOUT ROWID; + +CREATE TRIGGER radroots_runtime_authored_delivery_reconciliations_update_guard +BEFORE UPDATE ON radroots_runtime_authored_delivery_reconciliations +BEGIN + SELECT RAISE(ABORT, 'delivery reconciliation provenance is immutable'); +END; + +CREATE TRIGGER radroots_runtime_authored_delivery_reconciliations_delete_guard +BEFORE DELETE ON radroots_runtime_authored_delivery_reconciliations +BEGIN + SELECT RAISE(ABORT, 'delivery reconciliation provenance is retained'); +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 = 16; +pub const CURRENT_VERSION: u32 = 17; 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"); @@ -30,6 +30,9 @@ const AUTHORED_SIGNED_FACTS_V15_SQL: &str = include_str!("0015_authored_signed_f const AUTHORED_DELIVERY_FACTS_V16_SQL: &str = include_str!("0016_authored_delivery_facts.up.sql"); +const AUTHORED_DELIVERY_RECONCILIATION_V17_SQL: &str = + include_str!("0017_authored_delivery_reconciliation.up.sql"); + /// Stable, non-SQL description of one forward runtime migration. #[derive(Clone, Copy, Debug, Eq, PartialEq)] pub struct MigrationDescriptor { @@ -801,6 +804,96 @@ const RUNTIME_V16_OBJECTS: &[&str] = &[ "radroots_runtime_source_generations_sequence_guard", ]; +const RUNTIME_V17_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_atomic_target_phase_idx", + "radroots_runtime_authored_delivery_attempts", + "radroots_runtime_authored_delivery_claims", + "radroots_runtime_authored_delivery_claims_delete_guard", + "radroots_runtime_authored_delivery_claims_update_guard", + "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_reconciliations", + "radroots_runtime_authored_delivery_reconciliations_delete_guard", + "radroots_runtime_authored_delivery_reconciliations_update_guard", + "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 { @@ -899,6 +992,12 @@ pub const MIGRATIONS: &[MigrationDescriptor] = &[ up_sha256: "f78d35fbe60255f152c4e7a8c23add6677ac3ae74eed7db1f3d8ef4ba0b8f3c3", owned_objects: RUNTIME_V16_OBJECTS, }, + MigrationDescriptor { + version: 17, + name: "authored_delivery_reconciliation", + up_sha256: "01463e7effeb368dd0577bfed79f566f7bb2af76d03e8e76898cc53332e2e573", + owned_objects: RUNTIME_V17_OBJECTS, + }, ]; pub(crate) const fn migration_sql(version: u32) -> Option<&'static str> { @@ -919,6 +1018,7 @@ pub(crate) const fn migration_sql(version: u32) -> Option<&'static str> { 14 => Some(AUTHORED_DRAFT_QUERY_V14_SQL), 15 => Some(AUTHORED_SIGNED_FACTS_V15_SQL), 16 => Some(AUTHORED_DELIVERY_FACTS_V16_SQL), + 17 => Some(AUTHORED_DELIVERY_RECONCILIATION_V17_SQL), _ => None, } } @@ -962,8 +1062,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, 16); - assert_eq!(MIGRATIONS.len(), 16); + assert_eq!(CURRENT_VERSION, 17); + assert_eq!(MIGRATIONS.len(), 17); 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 @@ -56,7 +56,7 @@ async fn delivery_fact_upgrade_preserves_legacy_snapshots_receipts_and_cancelled .await .unwrap(); assert!(matches!( - migrate_runtime(&mut connection, OpenMode::ReadOnly).await, + migrate(&mut connection, OpenMode::ReadOnly, &plan(16, false)).await, Err(Error::SchemaMigrationRequired { actual: 15, current: 16, @@ -64,10 +64,14 @@ async fn delivery_fact_upgrade_preserves_legacy_snapshots_receipts_and_cancelled }) )); assert_eq!( - migrate_runtime(&mut connection, OpenMode::ReadWriteExisting) - .await - .unwrap() - .applied(), + migrate( + &mut connection, + OpenMode::ReadWriteExisting, + &plan(16, false) + ) + .await + .unwrap() + .applied(), 1 ); assert_eq!(pragma(&mut connection, "user_version").await, 16); diff --git a/crates/storage_sqlite/src/migration_delivery_reconciliation_tests.rs b/crates/storage_sqlite/src/migration_delivery_reconciliation_tests.rs @@ -0,0 +1,211 @@ +use super::{ + tests::{connection, establish_runtime_version, pragma}, + *, +}; +use crate::authored::signed_fact_fixture::{claim, prepare}; +use radroots_storage::{ + atomic::AtomicCommitDisposition, + authored_atomic::{ + AuthoredAtomicCommand, AuthoredAtomicOutcome, AuthoredAtomicReceipt, ClaimAuthoredTarget, + ClaimAuthoredWork, + }, +}; + +const FAILING_V17: &str = concat!( + include_str!("migration/runtime/0017_authored_delivery_reconciliation.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() == 17 { + FAILING_V17 + } else { + runtime::migration_sql(step.version()).unwrap() + }, + owned_objects: step.owned_objects(), + }) + .collect(), + } +} + +async fn seed(connection: &mut SqliteConnection) { + super::signed_facts_tests::seed(connection).await; + let (command, event) = prepare(); + let AuthoredAtomicCommand::Prepare(prepared) = command else { + unreachable!() + }; + let mut delivery = prepared.delivery_plans()[0].clone(); + delivery.bind_signed_event(event, 12).unwrap(); + for (token, at) in [(5, 13), (6, 40)] { + let active = claim(delivery.revision(), token, at); + delivery.claim(active.clone(), at).unwrap(); + let command = AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( + ClaimAuthoredTarget::DeliveryPlan(delivery.plan_id()), + active, + )); + let receipt = AuthoredAtomicReceipt::new( + &command, + AtomicCommitDisposition::Committed, + at, + AuthoredAtomicOutcome::DeliveryPlan(delivery.clone()), + ) + .unwrap(); + sqlx::query("INSERT INTO radroots_runtime_authored_atomic_commits (commit_id, commit_digest, phase, target_id, requested_at_unix_ms, committed_at_unix_ms, receipt) VALUES (?, ?, 'claim', ?, ?, ?, ?)") + .bind(receipt.commit_id().as_bytes().as_slice()).bind(receipt.digest().as_bytes().as_slice()) + .bind(delivery.plan_id().as_bytes().as_slice()).bind(at as i64).bind(at as i64) + .bind(serde_json::to_vec(&serde_json::json!({"outcome": receipt.outcome()})).unwrap()) + .execute(&mut *connection).await.unwrap(); + } + delivery.request_stop(61).unwrap(); + let mut transaction = connection.begin().await.unwrap(); + crate::authored::persist_plan_v11(&mut transaction, &delivery) + .await + .unwrap(); + transaction.commit().await.unwrap(); +} + +async fn snapshots(connection: &mut SqliteConnection) -> (Vec<u8>, Vec<(Vec<u8>, Vec<u8>)>) { + let plan = sqlx::query_scalar("SELECT snapshot FROM radroots_runtime_authored_delivery_plans") + .fetch_one(&mut *connection) + .await + .unwrap(); + let receipts = sqlx::query_as("SELECT commit_id, receipt FROM radroots_runtime_authored_atomic_commits ORDER BY commit_id").fetch_all(&mut *connection).await.unwrap(); + (plan, receipts) +} + +#[tokio::test] +async fn v17_backfills_stopped_issued_history_without_rewriting_old_receipts() { + let mut connection = connection().await; + establish_runtime_version(&mut connection, 16).await; + seed(&mut connection).await; + let before = snapshots(&mut connection).await; + assert!(matches!( + migrate(&mut connection, OpenMode::ReadOnly, &plan(17, false)).await, + Err(Error::SchemaMigrationRequired { + actual: 16, + current: 17, + .. + }) + )); + assert_eq!( + migrate( + &mut connection, + OpenMode::ReadWriteExisting, + &plan(17, false) + ) + .await + .unwrap() + .applied(), + 1 + ); + assert_eq!(snapshots(&mut connection).await, before); + let projected: Vec<Vec<u8>> = sqlx::query_scalar( + "SELECT claim_id FROM radroots_runtime_authored_delivery_claims ORDER BY claim_id", + ) + .fetch_all(&mut connection) + .await + .unwrap(); + let original: Vec<Vec<u8>> = sqlx::query_scalar("SELECT commit_id FROM radroots_runtime_authored_atomic_commits WHERE phase = 'claim' ORDER BY commit_id").fetch_all(&mut connection).await.unwrap(); + assert_eq!(projected.len(), 2); + assert_eq!(projected, original); + for mode in [OpenMode::ReadOnly, OpenMode::ReadWriteExisting] { + assert!(matches!( + migrate(&mut connection, mode, &plan(16, false)).await, + Err(Error::SchemaTooNew { + actual: 17, + supported: 16, + .. + }) + )); + } + assert_eq!( + migrate(&mut connection, OpenMode::ReadOnly, &plan(17, false)) + .await + .unwrap() + .applied(), + 0 + ); + connection.close().await.unwrap(); +} + +#[tokio::test] +async fn failed_v17_rolls_back_all_objects_and_preserves_original_claim_bytes() { + let mut connection = connection().await; + establish_runtime_version(&mut connection, 16).await; + seed(&mut connection).await; + let before = snapshots(&mut connection).await; + assert!( + migrate( + &mut connection, + OpenMode::ReadWriteExisting, + &plan(17, true) + ) + .await + .is_err() + ); + assert_eq!(pragma(&mut connection, "user_version").await, 16); + assert_eq!(snapshots(&mut connection).await, before); + assert_eq!(sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM sqlite_schema WHERE name LIKE 'radroots_runtime_authored_delivery_claims%' OR name LIKE 'radroots_runtime_authored_delivery_reconciliations%' OR name = 'radroots_runtime_authored_atomic_target_phase_idx'").fetch_one(&mut connection).await.unwrap(), 0); + migrate( + &mut connection, + OpenMode::ReadWriteExisting, + &plan(17, false), + ) + .await + .unwrap(); + assert_eq!(snapshots(&mut connection).await, before); + connection.close().await.unwrap(); +} + +#[tokio::test] +async fn malformed_claim_namespaces_cannot_silently_disappear_during_backfill() { + for wire in [ + "{}", + "{!", + r#"{"outcome":{"delivery_plan":null}}"#, + r#"{"outcome":{"delivery_plan":{},"artifact":{}}}"#, + ] { + let mut connection = connection().await; + establish_runtime_version(&mut connection, 16).await; + super::signed_facts_tests::seed(&mut connection).await; + sqlx::query("INSERT INTO radroots_runtime_authored_atomic_commits (commit_id, commit_digest, phase, target_id, requested_at_unix_ms, committed_at_unix_ms, receipt) VALUES (?, ?, 'claim', ?, 13, 13, ?)") + .bind([9u8; 16].as_slice()).bind([8u8; 32].as_slice()).bind([3u8; 16].as_slice()).bind(wire.as_bytes()).execute(&mut connection).await.unwrap(); + let before = snapshots(&mut connection).await; + assert!( + migrate( + &mut connection, + OpenMode::ReadWriteExisting, + &plan(17, false) + ) + .await + .is_err(), + "{wire}" + ); + assert_eq!(pragma(&mut connection, "user_version").await, 16); + assert_eq!(snapshots(&mut connection).await, before); + connection.close().await.unwrap(); + } +} + +#[test] +fn delivery_reconciliation_decision_binds_exact_forward_migration() { + let decision: serde_json::Value = serde_json::from_str(include_str!( + "../../../contracts/architecture/decisions/authored_delivery_reconciliation.v1.json" + )) + .unwrap(); + let migration = runtime::MIGRATIONS[16]; + assert_eq!(decision["migration"]["version"], migration.version()); + assert_eq!(decision["migration"]["name"], migration.name()); + assert_eq!(decision["migration"]["sha256"], migration.up_sha256()); +}