commit 27c373feaa38c3b8a4a5cb2e65544f459baabc60
parent e9e390200b5d32818fcdbbe408dce32ea939b72b
Author: triesap <tyson@radroots.org>
Date: Fri, 7 Aug 2026 13:54:26 +0000
feat(mobile): persist phase 1 drafts and outbox
- add immutable integrity-bound draft revisions to memory and SQLite storage
- queue all five Add flows offline with frozen relay and media prerequisites
- preserve signed bytes and delivery evidence during phase-aware cancellation
- verify crash recovery conflicts bounds and release coverage gates
Diffstat:
18 files changed, 3523 insertions(+), 18 deletions(-)
diff --git a/Cargo.lock b/Cargo.lock
@@ -3573,7 +3573,9 @@ dependencies = [
"radroots_event_codec",
"radroots_identity",
"radroots_sdk",
+ "radroots_signing",
"radroots_storage",
+ "radroots_sync",
"radroots_transport",
"serde",
"serde_json",
diff --git a/crates/mobile_core/Cargo.toml b/crates/mobile_core/Cargo.toml
@@ -27,11 +27,18 @@ mobile-social = [
]
[dependencies]
+radroots_blossom = { workspace = true, default-features = false, features = [
+ "serde",
+ "std",
+] }
radroots_sdk = { workspace = true, features = ["sqlite"] }
radroots_event = { workspace = true, default-features = false, features = ["std"] }
radroots_event_codec = { workspace = true, default-features = false, features = ["json", "std"] }
radroots_identity = { workspace = true, default-features = false, features = ["std"] }
+radroots_signing = { workspace = true, default-features = false, features = ["std"] }
radroots_storage = { workspace = true, default-features = false }
+radroots_sync = { workspace = true, default-features = false }
+radroots_transport = { workspace = true, default-features = false, features = ["std"] }
chrono = { workspace = true }
hex = { workspace = true }
serde = { workspace = true, features = ["derive"] }
@@ -41,11 +48,6 @@ thiserror = { workspace = true }
[dev-dependencies]
nostr = { workspace = true, features = ["std"] }
-radroots_blossom = { workspace = true, default-features = false, features = [
- "serde",
- "std",
-] }
radroots_sdk = { workspace = true, features = ["memory", "sqlite"] }
-radroots_transport = { workspace = true, default-features = false, features = ["std"] }
tempfile = { workspace = true }
tokio = { workspace = true, features = ["macros", "rt"] }
diff --git a/crates/mobile_core/src/runtime/product_surface.rs b/crates/mobile_core/src/runtime/product_surface.rs
@@ -9,6 +9,8 @@ mod context;
mod cursor;
mod identity;
mod model;
+#[cfg(feature = "mobile-social")]
+mod outbox;
mod projection;
mod ranking;
mod today;
@@ -30,6 +32,11 @@ pub use model::{
SearchResult, SearchResultType, SupportingProfile, ThreadEntry, ThreadReference, TodayCard,
TodayCardType, TodayPage,
};
+#[cfg(feature = "mobile-social")]
+pub use outbox::{
+ Phase1CancellationPolicy, Phase1DraftError, Phase1DraftStatus, Phase1MediaPrerequisite,
+ Phase1MediaStage, Phase1OutboxState, Phase1QueuePolicy, Phase1RelaySatisfaction,
+};
pub use projection::{ProductEventClassification, ProductEventExclusion, classify_admitted_event};
pub use ranking::{RankError, TODAY_RANK_SCHEMA_VERSION, TimeRelevance, TodayRank, TodayRankInput};
pub use today::{
diff --git a/crates/mobile_core/src/runtime/product_surface/outbox.rs b/crates/mobile_core/src/runtime/product_surface/outbox.rs
@@ -0,0 +1,1432 @@
+use std::collections::BTreeSet;
+
+use radroots_blossom::{BlobUrl, MediaType};
+use radroots_event::contract::AuthorRole;
+use radroots_event_codec::authoring::PlanWireV1;
+use radroots_identity::PublicKey;
+use radroots_signing::{Actor, actor::ActorSource, request::CancellationPolicy};
+use radroots_storage::{
+ authored::{AdmissionState, SigningState},
+ authored_delivery::{AuthoredDeliveryState, DeliveryAttemptOutcome},
+ authored_draft::{
+ AuthoredDraft, AuthoredDraftId, AuthoredDraftRevision, AuthoredDraftStage,
+ AuthoredDraftStore,
+ },
+ journal::{IdempotencyKey, OperationInstanceId},
+};
+use radroots_sync::{PushRequest, PushStatus, policy::SyncId};
+use radroots_transport::{
+ Target, TargetSet,
+ outcome::DeliveryOutcomeKind,
+ policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy},
+};
+use serde::{Deserialize, Serialize};
+use sha2::{Digest, Sha256};
+use thiserror::Error;
+
+use super::{
+ AddCommandType, CardId, CardSourceIdentity, LocalAuthorOverlay, LocalNetwork, Phase1AddCommand,
+ TodayCardType, TodayError,
+};
+use crate::runtime::RadrootsRuntime;
+
+const DRAFT_PAYLOAD_SCHEMA: &str = "radroots.mobile.phase1-draft.v1";
+const DRAFT_SCHEMA_VERSION: u16 = 1;
+const DRAFT_MEDIA_MAX: usize = 20;
+const DRAFT_LOCAL_REFERENCE_MAX_BYTES: usize = 4_096;
+const DRAFT_FAILURE_CODE_MAX_BYTES: usize = 96;
+const DRAFT_OPERATION_DOMAIN: &[u8] = b"radroots.mobile.phase1-draft-operation.v1\0";
+
+/// Durable state of one media prerequisite referenced by an Add command.
+#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
+#[serde(rename_all = "snake_case")]
+pub enum Phase1MediaStage {
+ Pending,
+ Preparing,
+ Uploading,
+ Verified,
+ Failed,
+ Orphaned,
+}
+
+/// Exact local and remote identity of one Add media prerequisite.
+#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
+#[serde(deny_unknown_fields)]
+pub struct Phase1MediaPrerequisite {
+ local_reference: String,
+ url: String,
+ sha256: String,
+ media_type: String,
+ byte_size: u64,
+ stage: Phase1MediaStage,
+ failure_code: Option<String>,
+}
+
+impl Phase1MediaPrerequisite {
+ pub fn new(
+ local_reference: impl Into<String>,
+ url: impl Into<String>,
+ sha256: impl Into<String>,
+ media_type: impl Into<String>,
+ byte_size: u64,
+ stage: Phase1MediaStage,
+ ) -> Result<Self, Phase1DraftError> {
+ let value = Self {
+ local_reference: local_reference.into(),
+ url: url.into(),
+ sha256: sha256.into(),
+ media_type: media_type.into(),
+ byte_size,
+ stage,
+ failure_code: None,
+ };
+ value.validate()?;
+ Ok(value)
+ }
+
+ pub fn with_failure_code(mut self, code: impl Into<String>) -> Result<Self, Phase1DraftError> {
+ self.stage = Phase1MediaStage::Failed;
+ self.failure_code = Some(code.into());
+ self.validate()?;
+ Ok(self)
+ }
+
+ fn validate(&self) -> Result<(), Phase1DraftError> {
+ let blob = BlobUrl::parse(self.url.as_str()).map_err(|_| Phase1DraftError::InvalidMedia)?;
+ let hash = blob.hash_path().hash().to_string();
+ if self.local_reference.is_empty()
+ || self.local_reference.len() > DRAFT_LOCAL_REFERENCE_MAX_BYTES
+ || self.local_reference != self.local_reference.trim()
+ || self.local_reference.chars().any(char::is_control)
+ || self.sha256 != hash
+ || MediaType::parse(self.media_type.as_str()).is_err()
+ || self.byte_size == 0
+ || self.failure_code.as_ref().is_some_and(|code| {
+ code.is_empty()
+ || code.len() > DRAFT_FAILURE_CODE_MAX_BYTES
+ || code != code.trim()
+ || code.chars().any(char::is_control)
+ })
+ || (self.stage == Phase1MediaStage::Failed) != self.failure_code.is_some()
+ {
+ return Err(Phase1DraftError::InvalidMedia);
+ }
+ Ok(())
+ }
+
+ pub fn url(&self) -> &str {
+ self.url.as_str()
+ }
+ pub const fn stage(&self) -> Phase1MediaStage {
+ self.stage
+ }
+}
+
+/// Closed delivery-satisfaction profiles available to Phase 1 Add.
+#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
+#[serde(rename_all = "snake_case")]
+pub enum Phase1RelaySatisfaction {
+ AnyAccepted,
+ AllAccepted,
+ AnyDelivered,
+ AllDelivered,
+}
+
+/// Stable product-level cancellation policy persisted with queue intent.
+#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
+#[serde(rename_all = "snake_case")]
+pub enum Phase1CancellationPolicy {
+ PreservePublishedRequest,
+ LocalCooperative,
+}
+
+/// Exact relay and deadline intent frozen before an operation is prepared.
+#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
+#[serde(deny_unknown_fields)]
+pub struct Phase1QueuePolicy {
+ relay_urls: Vec<String>,
+ satisfaction: Phase1RelaySatisfaction,
+ delivery_deadline_unix_ms: u64,
+ cancellation: Phase1CancellationPolicy,
+}
+
+impl Phase1QueuePolicy {
+ pub fn new(
+ relay_urls: Vec<String>,
+ satisfaction: Phase1RelaySatisfaction,
+ delivery_deadline_unix_ms: u64,
+ cancellation: Phase1CancellationPolicy,
+ ) -> Result<Self, Phase1DraftError> {
+ let value = Self {
+ relay_urls,
+ satisfaction,
+ delivery_deadline_unix_ms,
+ cancellation,
+ };
+ value.materialize()?;
+ Ok(value)
+ }
+
+ fn materialize(
+ &self,
+ ) -> Result<(TargetSet, SatisfactionPolicy, CancellationPolicy), Phase1DraftError> {
+ if self.delivery_deadline_unix_ms == 0 || self.relay_urls.is_empty() {
+ return Err(Phase1DraftError::InvalidQueuePolicy);
+ }
+ let mut canonical = BTreeSet::new();
+ let targets = self
+ .relay_urls
+ .iter()
+ .map(|url| {
+ let target =
+ Target::nostr_relay(url).map_err(|_| Phase1DraftError::InvalidQueuePolicy)?;
+ if target.uri().as_str() != url || !canonical.insert(url.as_str()) {
+ return Err(Phase1DraftError::InvalidQueuePolicy);
+ }
+ Ok(target)
+ })
+ .collect::<Result<Vec<_>, _>>()?;
+ let targets = TargetSet::new(targets).map_err(|_| Phase1DraftError::InvalidQueuePolicy)?;
+ let (class, target_policy) = match self.satisfaction {
+ Phase1RelaySatisfaction::AnyAccepted => {
+ (SatisfactionClass::Accepted, TargetPolicy::any())
+ }
+ Phase1RelaySatisfaction::AllAccepted => {
+ (SatisfactionClass::Accepted, TargetPolicy::all())
+ }
+ Phase1RelaySatisfaction::AnyDelivered => {
+ (SatisfactionClass::Delivered, TargetPolicy::any())
+ }
+ Phase1RelaySatisfaction::AllDelivered => {
+ (SatisfactionClass::Delivered, TargetPolicy::all())
+ }
+ };
+ let cancellation = match self.cancellation {
+ Phase1CancellationPolicy::PreservePublishedRequest => {
+ CancellationPolicy::PreservePublishedRequest
+ }
+ Phase1CancellationPolicy::LocalCooperative => CancellationPolicy::LocalCooperative,
+ };
+ Ok((
+ targets,
+ SatisfactionPolicy::new(class, target_policy),
+ cancellation,
+ ))
+ }
+}
+
+/// Honest aggregate state for a local draft and its durable authored operation.
+#[derive(Clone, Copy, Debug, Eq, PartialEq)]
+pub enum Phase1OutboxState {
+ Draft,
+ MediaPreparing,
+ MediaUploading,
+ ReadyToSign,
+ Signing,
+ Signed,
+ Queued,
+ Delivering,
+ PartiallyDelivered,
+ Retryable,
+ Terminal,
+ Cancelled,
+ Complete,
+}
+
+impl Phase1OutboxState {
+ pub const fn label(self) -> &'static str {
+ match self {
+ Self::Draft => "draft",
+ Self::MediaPreparing => "media_preparing",
+ Self::MediaUploading => "media_uploading",
+ Self::ReadyToSign => "ready_to_sign",
+ Self::Signing => "signing",
+ Self::Signed => "signed",
+ Self::Queued => "queued",
+ Self::Delivering => "delivering",
+ Self::PartiallyDelivered => "partially_delivered",
+ Self::Retryable => "retryable",
+ Self::Terminal => "terminal",
+ Self::Cancelled => "cancelled",
+ Self::Complete => "complete",
+ }
+ }
+}
+
+/// Current reconstructable product view of one Add draft.
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct Phase1DraftStatus {
+ draft: AuthoredDraft,
+ command_type: AddCommandType,
+ media: Vec<Phase1MediaPrerequisite>,
+ state: Phase1OutboxState,
+ card_id: CardId,
+ push: Option<PushStatus>,
+}
+
+impl Phase1DraftStatus {
+ pub const fn draft(&self) -> &AuthoredDraft {
+ &self.draft
+ }
+ pub const fn command_type(&self) -> AddCommandType {
+ self.command_type
+ }
+ pub fn media(&self) -> &[Phase1MediaPrerequisite] {
+ self.media.as_slice()
+ }
+ pub const fn state(&self) -> Phase1OutboxState {
+ self.state
+ }
+ pub const fn card_id(&self) -> CardId {
+ self.card_id
+ }
+ pub const fn push(&self) -> Option<&PushStatus> {
+ self.push.as_ref()
+ }
+}
+
+#[derive(Clone, Debug, Error, Eq, PartialEq)]
+pub enum Phase1DraftError {
+ #[error("authenticated draft identity is unavailable")]
+ IdentityUnavailable,
+ #[error("phase 1 draft input is invalid")]
+ InvalidDraft,
+ #[error("phase 1 media prerequisite is invalid")]
+ InvalidMedia,
+ #[error("phase 1 queue policy is invalid")]
+ InvalidQueuePolicy,
+ #[error("phase 1 draft revision conflicts with durable state")]
+ RevisionConflict,
+ #[error("phase 1 draft is not found")]
+ NotFound,
+ #[error("phase 1 draft is terminal")]
+ Terminal,
+ #[error("phase 1 draft media is not ready")]
+ MediaNotReady,
+ #[error("phase 1 authored operation is unavailable")]
+ OperationUnavailable,
+ #[error("phase 1 draft persistence failed")]
+ Storage,
+ #[error("phase 1 draft payload is corrupt")]
+ Corrupt,
+ #[error("phase 1 authored operation failed")]
+ Operation,
+ #[error("phase 1 Today overlay failed")]
+ Overlay,
+}
+
+#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
+#[serde(deny_unknown_fields)]
+struct Phase1DraftPayload {
+ schema_version: u16,
+ command_type: AddCommandType,
+ plan_wire_json: Vec<u8>,
+ media: Vec<Phase1MediaPrerequisite>,
+ queue: Option<Phase1QueuePolicy>,
+}
+
+impl Phase1DraftPayload {
+ fn new(
+ command: &Phase1AddCommand,
+ plan_wire_json: Vec<u8>,
+ media: Vec<Phase1MediaPrerequisite>,
+ ) -> Result<Self, Phase1DraftError> {
+ let value = Self {
+ schema_version: DRAFT_SCHEMA_VERSION,
+ command_type: command.command_type(),
+ plan_wire_json,
+ media,
+ queue: None,
+ };
+ value.validate()?;
+ Ok(value)
+ }
+
+ fn validate(&self) -> Result<(), Phase1DraftError> {
+ if self.schema_version != DRAFT_SCHEMA_VERSION || self.media.len() > DRAFT_MEDIA_MAX {
+ return Err(Phase1DraftError::Corrupt);
+ }
+ let integrity = PlanWireV1::from_json(self.plan_wire_json.as_slice())
+ .map_err(|_| Phase1DraftError::Corrupt)?;
+ let plan = integrity.plan();
+ let expected_media = media_urls(plan.body().tags())?;
+ let actual_media = self
+ .media
+ .iter()
+ .map(|media| {
+ media.validate()?;
+ Ok(media.url.as_str())
+ })
+ .collect::<Result<BTreeSet<_>, Phase1DraftError>>()?;
+ if expected_media != actual_media || actual_media.len() != self.media.len() {
+ return Err(Phase1DraftError::InvalidMedia);
+ }
+ if let Some(queue) = &self.queue {
+ queue.materialize()?;
+ }
+ Ok(())
+ }
+
+ fn decode(draft: &AuthoredDraft) -> Result<Self, Phase1DraftError> {
+ if draft.payload_schema() != DRAFT_PAYLOAD_SCHEMA {
+ return Err(Phase1DraftError::Corrupt);
+ }
+ let value = serde_json::from_slice::<Self>(draft.payload())
+ .map_err(|_| Phase1DraftError::Corrupt)?;
+ value.validate()?;
+ let plan = PlanWireV1::from_json(value.plan_wire_json.as_slice())
+ .map_err(|_| Phase1DraftError::Corrupt)?;
+ if plan.plan().author().as_bytes() != draft.author() {
+ return Err(Phase1DraftError::Corrupt);
+ }
+ Ok(value)
+ }
+
+ fn encode(&self) -> Result<Vec<u8>, Phase1DraftError> {
+ self.validate()?;
+ serde_json::to_vec(self).map_err(|_| Phase1DraftError::InvalidDraft)
+ }
+}
+
+impl RadrootsRuntime {
+ /// Creates or replaces the editable content of one immutable-revision draft.
+ #[allow(clippy::too_many_arguments)]
+ pub async fn phase1_save_draft(
+ &self,
+ draft_id: [u8; 16],
+ command: Phase1AddCommand,
+ authored_at_unix_s: u64,
+ media: Vec<Phase1MediaPrerequisite>,
+ expected_revision: Option<u64>,
+ persisted_at_unix_ms: u64,
+ ) -> Result<Phase1DraftStatus, Phase1DraftError> {
+ let author = self.draft_author()?;
+ let draft_id =
+ AuthoredDraftId::new(draft_id).map_err(|_| Phase1DraftError::InvalidDraft)?;
+ let plan = command
+ .authored_plan(authored_at_unix_s, hex::encode(author))
+ .map_err(|_| Phase1DraftError::InvalidDraft)?;
+ let wire = PlanWireV1::from_plan(&plan)
+ .to_json()
+ .map_err(|_| Phase1DraftError::InvalidDraft)?;
+ let payload = Phase1DraftPayload::new(&command, wire, media)?;
+ let bytes = payload.encode()?;
+ let storage = self.storage()?;
+ let expected = expected_revision
+ .map(AuthoredDraftRevision::new)
+ .transpose()
+ .map_err(|_| Phase1DraftError::RevisionConflict)?;
+ let draft = if let Some(expected) = expected {
+ let head = storage
+ .authored_draft_head(draft_id)
+ .await
+ .map_err(|_| Phase1DraftError::Storage)?
+ .ok_or(Phase1DraftError::NotFound)?;
+ if head.revision() != expected
+ || head.stage().is_terminal()
+ || matches!(
+ head.stage(),
+ AuthoredDraftStage::ReadyToSign | AuthoredDraftStage::Queued
+ )
+ {
+ return Err(Phase1DraftError::RevisionConflict);
+ }
+ head.successor(
+ bytes,
+ draft_stage_for_media(&payload.media),
+ None,
+ persisted_at_unix_ms,
+ )
+ .map_err(|_| Phase1DraftError::RevisionConflict)?
+ } else {
+ AuthoredDraft::initial(
+ draft_id,
+ author,
+ DRAFT_PAYLOAD_SCHEMA,
+ bytes,
+ draft_stage_for_media(&payload.media),
+ None,
+ persisted_at_unix_ms,
+ )
+ .map_err(|_| Phase1DraftError::InvalidDraft)?
+ };
+ let receipt = storage
+ .append_authored_draft(draft, expected)
+ .await
+ .map_err(map_draft_storage_error)?;
+ self.draft_status_from(receipt.draft().clone()).await
+ }
+
+ /// Advances one media prerequisite without mutating any prior revision.
+ pub async fn phase1_update_draft_media(
+ &self,
+ draft_id: [u8; 16],
+ expected_revision: u64,
+ url: &str,
+ stage: Phase1MediaStage,
+ failure_code: Option<String>,
+ updated_at_unix_ms: u64,
+ ) -> Result<Phase1DraftStatus, Phase1DraftError> {
+ let draft_id =
+ AuthoredDraftId::new(draft_id).map_err(|_| Phase1DraftError::InvalidDraft)?;
+ let expected = AuthoredDraftRevision::new(expected_revision)
+ .map_err(|_| Phase1DraftError::RevisionConflict)?;
+ let storage = self.storage()?;
+ let head = storage
+ .authored_draft_head(draft_id)
+ .await
+ .map_err(|_| Phase1DraftError::Storage)?
+ .ok_or(Phase1DraftError::NotFound)?;
+ if head.revision() != expected
+ || head.stage().is_terminal()
+ || matches!(
+ head.stage(),
+ AuthoredDraftStage::ReadyToSign | AuthoredDraftStage::Queued
+ )
+ {
+ return Err(Phase1DraftError::RevisionConflict);
+ }
+ let mut payload = Phase1DraftPayload::decode(&head)?;
+ let media = payload
+ .media
+ .iter_mut()
+ .find(|media| media.url == url)
+ .ok_or(Phase1DraftError::InvalidMedia)?;
+ if !valid_media_transition(media.stage, stage) {
+ return Err(Phase1DraftError::InvalidMedia);
+ }
+ media.stage = stage;
+ media.failure_code = failure_code;
+ media.validate()?;
+ let next_stage = match stage {
+ Phase1MediaStage::Pending | Phase1MediaStage::Preparing => {
+ AuthoredDraftStage::MediaPreparing
+ }
+ Phase1MediaStage::Uploading
+ | Phase1MediaStage::Verified
+ | Phase1MediaStage::Failed
+ | Phase1MediaStage::Orphaned => AuthoredDraftStage::MediaUploading,
+ };
+ let next = head
+ .successor(payload.encode()?, next_stage, None, updated_at_unix_ms)
+ .map_err(|_| Phase1DraftError::RevisionConflict)?;
+ let receipt = storage
+ .append_authored_draft(next, Some(expected))
+ .await
+ .map_err(map_draft_storage_error)?;
+ self.draft_status_from(receipt.draft().clone()).await
+ }
+
+ /// Freezes queue intent and atomically prepares the canonical outbox before
+ /// returning `queued`. Network connectivity is neither read nor required.
+ pub async fn phase1_queue_draft(
+ &self,
+ draft_id: [u8; 16],
+ expected_revision: u64,
+ policy: Phase1QueuePolicy,
+ queued_at_unix_ms: u64,
+ ) -> Result<Phase1DraftStatus, Phase1DraftError> {
+ let draft_id =
+ AuthoredDraftId::new(draft_id).map_err(|_| Phase1DraftError::InvalidDraft)?;
+ let expected = AuthoredDraftRevision::new(expected_revision)
+ .map_err(|_| Phase1DraftError::RevisionConflict)?;
+ let storage = self.storage()?;
+ let head = storage
+ .authored_draft_head(draft_id)
+ .await
+ .map_err(|_| Phase1DraftError::Storage)?
+ .ok_or(Phase1DraftError::NotFound)?;
+ if head.revision() != expected {
+ return Err(Phase1DraftError::RevisionConflict);
+ }
+ let mut payload = Phase1DraftPayload::decode(&head)?;
+ let ready = match head.stage() {
+ AuthoredDraftStage::ReadyToSign => {
+ if payload.queue.as_ref() != Some(&policy) {
+ return Err(Phase1DraftError::RevisionConflict);
+ }
+ head
+ }
+ AuthoredDraftStage::Queued => {
+ if payload.queue.as_ref() != Some(&policy) {
+ return Err(Phase1DraftError::RevisionConflict);
+ }
+ return self.draft_status_from(head).await;
+ }
+ AuthoredDraftStage::Cancelled => return Err(Phase1DraftError::Terminal),
+ AuthoredDraftStage::Draft
+ | AuthoredDraftStage::MediaPreparing
+ | AuthoredDraftStage::MediaUploading => {
+ if payload
+ .media
+ .iter()
+ .any(|media| media.stage != Phase1MediaStage::Verified)
+ {
+ return Err(Phase1DraftError::MediaNotReady);
+ }
+ policy.materialize()?;
+ payload.queue = Some(policy);
+ let bytes = payload.encode()?;
+ let operation_id = operation_id(draft_id, bytes.as_slice())?;
+ let operation_id = OperationInstanceId::new(*operation_id.as_bytes())
+ .map_err(|_| Phase1DraftError::InvalidDraft)?;
+ let ready = head
+ .successor(
+ bytes,
+ AuthoredDraftStage::ReadyToSign,
+ Some(operation_id),
+ queued_at_unix_ms,
+ )
+ .map_err(|_| Phase1DraftError::RevisionConflict)?;
+ storage
+ .append_authored_draft(ready.clone(), Some(expected))
+ .await
+ .map_err(map_draft_storage_error)?;
+ ready
+ }
+ };
+ self.finish_queue(ready, queued_at_unix_ms).await
+ }
+
+ /// Resumes a queue transition interrupted after its durable ready-to-sign
+ /// checkpoint, including the crash window after outbox preparation.
+ pub async fn phase1_recover_draft_queue(
+ &self,
+ draft_id: [u8; 16],
+ recovered_at_unix_ms: u64,
+ ) -> Result<Phase1DraftStatus, Phase1DraftError> {
+ let draft_id =
+ AuthoredDraftId::new(draft_id).map_err(|_| Phase1DraftError::InvalidDraft)?;
+ let head = self
+ .storage()?
+ .authored_draft_head(draft_id)
+ .await
+ .map_err(|_| Phase1DraftError::Storage)?
+ .ok_or(Phase1DraftError::NotFound)?;
+ match head.stage() {
+ AuthoredDraftStage::ReadyToSign => self.finish_queue(head, recovered_at_unix_ms).await,
+ AuthoredDraftStage::Queued | AuthoredDraftStage::Cancelled => {
+ self.draft_status_from(head).await
+ }
+ AuthoredDraftStage::Draft
+ | AuthoredDraftStage::MediaPreparing
+ | AuthoredDraftStage::MediaUploading => Err(Phase1DraftError::InvalidDraft),
+ }
+ }
+
+ /// Returns durable draft state composed with canonical authored-operation state.
+ pub async fn phase1_draft_status(
+ &self,
+ draft_id: [u8; 16],
+ ) -> Result<Phase1DraftStatus, Phase1DraftError> {
+ let draft_id =
+ AuthoredDraftId::new(draft_id).map_err(|_| Phase1DraftError::InvalidDraft)?;
+ let head = self
+ .storage()?
+ .authored_draft_head(draft_id)
+ .await
+ .map_err(|_| Phase1DraftError::Storage)?
+ .ok_or(Phase1DraftError::NotFound)?;
+ self.draft_status_from(head).await
+ }
+
+ /// Lists the newest immutable revision of each draft for the active author.
+ pub async fn phase1_draft_heads(
+ &self,
+ limit: u16,
+ ) -> Result<Vec<Phase1DraftStatus>, Phase1DraftError> {
+ let drafts = self
+ .storage()?
+ .authored_draft_heads(self.draft_author()?, limit)
+ .await
+ .map_err(map_draft_storage_error)?;
+ let mut statuses = Vec::with_capacity(drafts.len());
+ for draft in drafts {
+ statuses.push(self.draft_status_from(draft).await?);
+ }
+ Ok(statuses)
+ }
+
+ /// Cancels still-pending authored work and records uploaded-but-unreferenced
+ /// media as possible orphans without deleting any evidence.
+ pub async fn phase1_cancel_draft(
+ &self,
+ draft_id: [u8; 16],
+ expected_revision: u64,
+ cancelled_at_unix_ms: u64,
+ ) -> Result<Phase1DraftStatus, Phase1DraftError> {
+ let draft_id =
+ AuthoredDraftId::new(draft_id).map_err(|_| Phase1DraftError::InvalidDraft)?;
+ let expected = AuthoredDraftRevision::new(expected_revision)
+ .map_err(|_| Phase1DraftError::RevisionConflict)?;
+ let storage = self.storage()?;
+ let head = storage
+ .authored_draft_head(draft_id)
+ .await
+ .map_err(|_| Phase1DraftError::Storage)?
+ .ok_or(Phase1DraftError::NotFound)?;
+ if head.stage() == AuthoredDraftStage::Cancelled {
+ return self.draft_status_from(head).await;
+ }
+ if head.revision() != expected {
+ return Err(Phase1DraftError::RevisionConflict);
+ }
+ let mut payload = Phase1DraftPayload::decode(&head)?;
+ let push = self.push_status_for(&head).await?;
+ if let Some(status) = &push {
+ self.sync()?
+ .cancel_push(sync_id_for(&head)?)
+ .await
+ .map_err(|_| Phase1DraftError::Operation)?;
+ if status.artifact().signed().is_none() {
+ mark_possible_orphans(&mut payload.media);
+ }
+ } else {
+ mark_possible_orphans(&mut payload.media);
+ }
+ let next = head
+ .successor(
+ payload.encode()?,
+ AuthoredDraftStage::Cancelled,
+ head.operation_id(),
+ cancelled_at_unix_ms,
+ )
+ .map_err(|_| Phase1DraftError::RevisionConflict)?;
+ let receipt = storage
+ .append_authored_draft(next, Some(expected))
+ .await
+ .map_err(map_draft_storage_error)?;
+ self.draft_status_from(receipt.draft().clone()).await
+ }
+
+ /// Applies the current durable operation state to an already-projected
+ /// active-author card. The overlay remains local and never changes event truth.
+ pub async fn phase1_apply_draft_overlay(
+ &self,
+ context: &LocalNetwork,
+ draft_id: [u8; 16],
+ ) -> Result<Phase1DraftStatus, Phase1DraftError> {
+ let status = self.phase1_draft_status(draft_id).await?;
+ let operation_id = status
+ .draft
+ .operation_id()
+ .map(|id| hex::encode(id.as_bytes()))
+ .ok_or(Phase1DraftError::OperationUnavailable)?;
+ self.phase1_set_local_author_overlay(
+ context,
+ status.card_id,
+ Some(LocalAuthorOverlay {
+ operation_id,
+ state: status.state.label().to_owned(),
+ }),
+ )
+ .await
+ .map_err(map_overlay_error)?;
+ Ok(status)
+ }
+
+ async fn finish_queue(
+ &self,
+ ready: AuthoredDraft,
+ queued_at_unix_ms: u64,
+ ) -> Result<Phase1DraftStatus, Phase1DraftError> {
+ let request = push_request(&ready)?;
+ self.sync()?
+ .prepare_push(request)
+ .await
+ .map_err(|_| Phase1DraftError::Operation)?;
+ let queued = ready
+ .successor(
+ ready.payload().to_vec(),
+ AuthoredDraftStage::Queued,
+ ready.operation_id(),
+ queued_at_unix_ms.max(ready.updated_at_unix_ms()),
+ )
+ .map_err(|_| Phase1DraftError::RevisionConflict)?;
+ let receipt = self
+ .storage()?
+ .append_authored_draft(queued, Some(ready.revision()))
+ .await
+ .map_err(map_draft_storage_error)?;
+ self.draft_status_from(receipt.draft().clone()).await
+ }
+
+ async fn draft_status_from(
+ &self,
+ draft: AuthoredDraft,
+ ) -> Result<Phase1DraftStatus, Phase1DraftError> {
+ if draft.author() != &self.draft_author()? {
+ return Err(Phase1DraftError::Corrupt);
+ }
+ let payload = Phase1DraftPayload::decode(&draft)?;
+ let integrity = PlanWireV1::from_json(payload.plan_wire_json.as_slice())
+ .map_err(|_| Phase1DraftError::Corrupt)?;
+ let card_id = card_id(payload.command_type, integrity.plan())?;
+ let push = self.push_status_for(&draft).await?;
+ if draft.stage() == AuthoredDraftStage::Queued && push.is_none() {
+ return Err(Phase1DraftError::Corrupt);
+ }
+ let state = aggregate_state(&draft, push.as_ref());
+ Ok(Phase1DraftStatus {
+ draft,
+ command_type: payload.command_type,
+ media: payload.media,
+ state,
+ card_id,
+ push,
+ })
+ }
+
+ async fn push_status_for(
+ &self,
+ draft: &AuthoredDraft,
+ ) -> Result<Option<PushStatus>, Phase1DraftError> {
+ let Some(_) = draft.operation_id() else {
+ return Ok(None);
+ };
+ self.sync()?
+ .push_status(sync_id_for(draft)?)
+ .await
+ .map_err(|_| Phase1DraftError::Operation)
+ }
+
+ fn storage(&self) -> Result<&dyn AuthoredDraftStore, Phase1DraftError> {
+ self.client
+ .storage()
+ .map(|storage| storage as &dyn AuthoredDraftStore)
+ .map_err(|_| Phase1DraftError::Storage)
+ }
+
+ fn sync(&self) -> Result<radroots_sdk::sync::Operations<'_>, Phase1DraftError> {
+ self.client
+ .sync()
+ .map_err(|_| Phase1DraftError::OperationUnavailable)?
+ .ok_or(Phase1DraftError::OperationUnavailable)
+ }
+
+ fn draft_author(&self) -> Result<[u8; 32], Phase1DraftError> {
+ self.store_public_key
+ .map(|key| *key.as_bytes())
+ .ok_or(Phase1DraftError::IdentityUnavailable)
+ }
+}
+
+fn push_request(draft: &AuthoredDraft) -> Result<PushRequest, Phase1DraftError> {
+ let payload = Phase1DraftPayload::decode(draft)?;
+ let policy = payload.queue.ok_or(Phase1DraftError::InvalidQueuePolicy)?;
+ let (targets, satisfaction, cancellation) = policy.materialize()?;
+ let plan = PlanWireV1::from_json(payload.plan_wire_json.as_slice())
+ .map_err(|_| Phase1DraftError::Corrupt)?
+ .into_plan();
+ let public_key = PublicKey::from_bytes(*draft.author())
+ .map_err(|_| Phase1DraftError::IdentityUnavailable)?;
+ let actor = Actor::new(public_key, ActorSource::ExplicitPublicKey, AuthorRole::ALL)
+ .map_err(|_| Phase1DraftError::IdentityUnavailable)?;
+ let sync_id = sync_id_for(draft)?;
+ let idempotency = IdempotencyKey::parse(format!(
+ "phase1-draft-{}",
+ hex::encode(draft.draft_id().as_bytes())
+ ))
+ .map_err(|_| Phase1DraftError::InvalidDraft)?;
+ PushRequest::new(
+ sync_id,
+ idempotency,
+ actor,
+ plan,
+ targets,
+ satisfaction,
+ policy.delivery_deadline_unix_ms,
+ cancellation,
+ )
+ .map_err(|_| Phase1DraftError::InvalidQueuePolicy)
+}
+
+fn operation_id(
+ draft_id: AuthoredDraftId,
+ ready_payload: &[u8],
+) -> Result<SyncId, Phase1DraftError> {
+ let mut hasher = Sha256::new();
+ hasher.update(DRAFT_OPERATION_DOMAIN);
+ hasher.update(draft_id.as_bytes());
+ hasher.update(
+ u64::try_from(ready_payload.len())
+ .map_err(|_| Phase1DraftError::InvalidDraft)?
+ .to_be_bytes(),
+ );
+ hasher.update(ready_payload);
+ let digest: [u8; 32] = hasher.finalize().into();
+ let mut value = [0_u8; 16];
+ value.copy_from_slice(&digest[..16]);
+ if value.iter().all(|byte| *byte == 0) {
+ value[15] = 1;
+ }
+ SyncId::new(value).map_err(|_| Phase1DraftError::InvalidDraft)
+}
+
+fn sync_id_for(draft: &AuthoredDraft) -> Result<SyncId, Phase1DraftError> {
+ draft
+ .operation_id()
+ .ok_or(Phase1DraftError::OperationUnavailable)
+ .and_then(|id| SyncId::new(*id.as_bytes()).map_err(|_| Phase1DraftError::Corrupt))
+}
+
+fn media_urls(tags: &[Vec<String>]) -> Result<BTreeSet<&str>, Phase1DraftError> {
+ let mut urls = BTreeSet::new();
+ for tag in tags {
+ let candidate = match tag.first().map(String::as_str) {
+ Some("imeta") => tag
+ .iter()
+ .skip(1)
+ .find_map(|value| value.strip_prefix("url ")),
+ Some("image") => tag.get(1).map(String::as_str),
+ _ => None,
+ };
+ if let Some(candidate) = candidate
+ && BlobUrl::parse(candidate).is_ok()
+ && !urls.insert(candidate)
+ {
+ return Err(Phase1DraftError::InvalidMedia);
+ }
+ }
+ Ok(urls)
+}
+
+fn draft_stage_for_media(media: &[Phase1MediaPrerequisite]) -> AuthoredDraftStage {
+ if media.is_empty() {
+ AuthoredDraftStage::Draft
+ } else if media
+ .iter()
+ .any(|media| media.stage == Phase1MediaStage::Uploading)
+ {
+ AuthoredDraftStage::MediaUploading
+ } else {
+ AuthoredDraftStage::MediaPreparing
+ }
+}
+
+fn mark_possible_orphans(media: &mut [Phase1MediaPrerequisite]) {
+ for media in media {
+ if media.stage == Phase1MediaStage::Verified {
+ media.stage = Phase1MediaStage::Orphaned;
+ }
+ }
+}
+
+const fn valid_media_transition(previous: Phase1MediaStage, next: Phase1MediaStage) -> bool {
+ match previous {
+ Phase1MediaStage::Pending => matches!(
+ next,
+ Phase1MediaStage::Pending
+ | Phase1MediaStage::Preparing
+ | Phase1MediaStage::Uploading
+ | Phase1MediaStage::Failed
+ ),
+ Phase1MediaStage::Preparing => matches!(
+ next,
+ Phase1MediaStage::Preparing | Phase1MediaStage::Uploading | Phase1MediaStage::Failed
+ ),
+ Phase1MediaStage::Uploading => matches!(
+ next,
+ Phase1MediaStage::Uploading | Phase1MediaStage::Verified | Phase1MediaStage::Failed
+ ),
+ Phase1MediaStage::Failed => matches!(
+ next,
+ Phase1MediaStage::Preparing | Phase1MediaStage::Uploading | Phase1MediaStage::Failed
+ ),
+ Phase1MediaStage::Verified => matches!(next, Phase1MediaStage::Verified),
+ Phase1MediaStage::Orphaned => matches!(next, Phase1MediaStage::Orphaned),
+ }
+}
+
+fn aggregate_state(draft: &AuthoredDraft, push: Option<&PushStatus>) -> Phase1OutboxState {
+ if draft.stage() == AuthoredDraftStage::Cancelled {
+ return Phase1OutboxState::Cancelled;
+ }
+ let Some(push) = push else {
+ return match draft.stage() {
+ AuthoredDraftStage::Draft => Phase1OutboxState::Draft,
+ AuthoredDraftStage::MediaPreparing => Phase1OutboxState::MediaPreparing,
+ AuthoredDraftStage::MediaUploading => Phase1OutboxState::MediaUploading,
+ AuthoredDraftStage::ReadyToSign => Phase1OutboxState::ReadyToSign,
+ AuthoredDraftStage::Queued => Phase1OutboxState::Queued,
+ AuthoredDraftStage::Cancelled => Phase1OutboxState::Cancelled,
+ };
+ };
+ if push.settlement().is_successful() {
+ return Phase1OutboxState::Complete;
+ }
+ if push.settlement().has_failures() {
+ return if push.settlement().retryable() != 0 || push.settlement().delivery_retryable() != 0
+ {
+ Phase1OutboxState::Retryable
+ } else if push.settlement().cancelled() != 0 || push.settlement().delivery_cancelled() != 0
+ {
+ Phase1OutboxState::Cancelled
+ } else {
+ Phase1OutboxState::Terminal
+ };
+ }
+ match push.artifact().signing_state() {
+ SigningState::Planned if push.artifact().signing_claim().is_some() => {
+ Phase1OutboxState::Signing
+ }
+ SigningState::Planned => Phase1OutboxState::Queued,
+ SigningState::Retryable => Phase1OutboxState::Retryable,
+ SigningState::Indeterminate | SigningState::FailedTerminal => Phase1OutboxState::Terminal,
+ SigningState::Cancelled => Phase1OutboxState::Cancelled,
+ SigningState::Signed => match push.artifact().admission_state() {
+ AdmissionState::Pending => Phase1OutboxState::Signed,
+ AdmissionState::Retryable => Phase1OutboxState::Retryable,
+ AdmissionState::Rejected => Phase1OutboxState::Terminal,
+ AdmissionState::Cancelled => Phase1OutboxState::Cancelled,
+ AdmissionState::Inserted | AdmissionState::Duplicate => match push
+ .delivery_plan()
+ .state()
+ {
+ AuthoredDeliveryState::Pending
+ if push.delivery_plan().claim_evidence().is_some() =>
+ {
+ Phase1OutboxState::Delivering
+ }
+ AuthoredDeliveryState::Pending if push.delivery_plan().attempts().is_empty() => {
+ Phase1OutboxState::Queued
+ }
+ AuthoredDeliveryState::Pending => Phase1OutboxState::PartiallyDelivered,
+ AuthoredDeliveryState::Retryable if has_delivery_success(push) => {
+ Phase1OutboxState::PartiallyDelivered
+ }
+ AuthoredDeliveryState::Retryable => Phase1OutboxState::Retryable,
+ AuthoredDeliveryState::Satisfied => Phase1OutboxState::Complete,
+ AuthoredDeliveryState::Exhausted | AuthoredDeliveryState::FailedTerminal => {
+ Phase1OutboxState::Terminal
+ }
+ AuthoredDeliveryState::Cancelled => Phase1OutboxState::Cancelled,
+ },
+ },
+ }
+}
+
+fn has_delivery_success(push: &PushStatus) -> bool {
+ push.delivery_plan().attempts().iter().any(|attempt| {
+ let evidence = match attempt.outcome() {
+ DeliveryAttemptOutcome::Receipt(receipt) => receipt.target_receipts(),
+ DeliveryAttemptOutcome::SinkFailure(failure) => failure.partial_evidence(),
+ };
+ evidence.iter().any(|receipt| {
+ matches!(
+ receipt.outcome().kind(),
+ DeliveryOutcomeKind::Accepted | DeliveryOutcomeKind::Delivered
+ )
+ })
+ })
+}
+
+fn card_id(
+ command_type: AddCommandType,
+ plan: &radroots_event_codec::authoring::AuthoredEventPlan,
+) -> Result<CardId, Phase1DraftError> {
+ let card_type = match command_type {
+ AddCommandType::CreateUpdate => TodayCardType::Update,
+ AddCommandType::CreatePhotoUpdate => TodayCardType::PhotoUpdate,
+ AddCommandType::CreateAsk => TodayCardType::Ask,
+ AddCommandType::CreateEvent => TodayCardType::Event,
+ AddCommandType::CreateFoodAvailability => TodayCardType::FoodAvailability,
+ };
+ let source = if (30_000..40_000).contains(&plan.body().kind()) {
+ let identifier = plan
+ .body()
+ .tags()
+ .iter()
+ .find(|tag| tag.first().map(String::as_str) == Some("d") && tag.len() == 2)
+ .and_then(|tag| tag.get(1))
+ .ok_or(Phase1DraftError::Corrupt)?;
+ CardSourceIdentity::address(
+ plan.body().kind(),
+ plan.author().to_hex(),
+ identifier.clone(),
+ )
+ .map_err(|_| Phase1DraftError::Corrupt)?
+ } else {
+ CardSourceIdentity::Event(*plan.expected_event_id())
+ };
+ Ok(CardId::derive(card_type, &source))
+}
+
+fn map_draft_storage_error(error: radroots_storage::Error) -> Phase1DraftError {
+ match error {
+ radroots_storage::Error::DraftRevisionConflict => Phase1DraftError::RevisionConflict,
+ radroots_storage::Error::DraftNotFound => Phase1DraftError::NotFound,
+ radroots_storage::Error::CorruptAuthoredDraft => Phase1DraftError::Corrupt,
+ _ => Phase1DraftError::Storage,
+ }
+}
+
+fn map_overlay_error(_: TodayError) -> Phase1DraftError {
+ Phase1DraftError::Overlay
+}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+ use crate::runtime::product_surface::{
+ CANONICAL_ADD_COMMAND_TYPES, CreateAsk, CreateEvent, CreateFoodAvailability,
+ CreatePhotoUpdate, CreateUpdate,
+ };
+ use radroots_blossom::{BlobDescriptor, Sha256 as BlossomSha256};
+ use radroots_event::{
+ calendar::{AuthoredCalendarDateEvent, CalendarDate},
+ food::availability::{
+ FoodAvailabilityDetails, FoodAvailabilityDetailsParts, FoodAvailabilityStatus,
+ FoodContent, FoodCurrency, FoodIdentifier, FoodPrice, FoodPublishedAt, FoodText,
+ FoodUnit,
+ },
+ media::AuthoredImage,
+ post::{AuthoredPostImage, PostImageDimensions},
+ };
+ use radroots_sdk::ClientBuilder;
+
+ const AUTHOR: &str = "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa";
+
+ fn runtime() -> RadrootsRuntime {
+ RadrootsRuntime::from_client_builder(
+ ClientBuilder::memory_default(),
+ Some(PublicKey::from_hex(AUTHOR).unwrap()),
+ )
+ .unwrap()
+ }
+
+ fn policy() -> Phase1QueuePolicy {
+ Phase1QueuePolicy::new(
+ vec![
+ "wss://relay-one.example".to_owned(),
+ "wss://relay-two.example".to_owned(),
+ ],
+ Phase1RelaySatisfaction::AllAccepted,
+ 2_000_000_000_000,
+ Phase1CancellationPolicy::LocalCooperative,
+ )
+ .unwrap()
+ }
+
+ #[tokio::test]
+ async fn media_free_draft_queues_offline_and_recovers_exactly() {
+ let runtime = runtime();
+ let id = [7; 16];
+ let draft = runtime
+ .phase1_save_draft(
+ id,
+ Phase1AddCommand::CreateUpdate(CreateUpdate::new("Harvest").unwrap()),
+ 1_900_000_000,
+ Vec::new(),
+ None,
+ 10,
+ )
+ .await
+ .unwrap();
+ assert_eq!(draft.state(), Phase1OutboxState::Draft);
+ let queued = runtime
+ .phase1_queue_draft(id, 1, policy(), 11)
+ .await
+ .unwrap();
+ assert_eq!(queued.state(), Phase1OutboxState::Queued);
+ assert_eq!(queued.draft().revision().get(), 3);
+ assert!(queued.push().is_some());
+ let recovered = runtime.phase1_recover_draft_queue(id, 12).await.unwrap();
+ assert_eq!(recovered.draft(), queued.draft());
+ assert_eq!(recovered.push(), queued.push());
+ }
+
+ #[tokio::test]
+ async fn queue_recovery_closes_both_preparation_crash_windows() {
+ let runtime = runtime();
+ for (id_byte, prepare_before_recovery) in [(10, false), (11, true)] {
+ let id = [id_byte; 16];
+ let saved = runtime
+ .phase1_save_draft(
+ id,
+ Phase1AddCommand::CreateUpdate(CreateUpdate::new("Recover").unwrap()),
+ 1_900_000_010,
+ Vec::new(),
+ None,
+ 40,
+ )
+ .await
+ .unwrap();
+ let mut payload = Phase1DraftPayload::decode(saved.draft()).unwrap();
+ payload.queue = Some(policy());
+ let bytes = payload.encode().unwrap();
+ let draft_id = saved.draft().draft_id();
+ let operation = operation_id(draft_id, bytes.as_slice()).unwrap();
+ let operation = OperationInstanceId::new(*operation.as_bytes()).unwrap();
+ let ready = saved
+ .draft()
+ .successor(bytes, AuthoredDraftStage::ReadyToSign, Some(operation), 41)
+ .unwrap();
+ runtime
+ .storage()
+ .unwrap()
+ .append_authored_draft(ready.clone(), Some(saved.draft().revision()))
+ .await
+ .unwrap();
+ if prepare_before_recovery {
+ runtime
+ .sync()
+ .unwrap()
+ .prepare_push(push_request(&ready).unwrap())
+ .await
+ .unwrap();
+ }
+ let recovered = runtime.phase1_recover_draft_queue(id, 42).await.unwrap();
+ assert_eq!(recovered.draft().stage(), AuthoredDraftStage::Queued);
+ assert_eq!(recovered.draft().revision().get(), 3);
+ assert!(recovered.push().is_some());
+ }
+ }
+
+ #[tokio::test]
+ async fn all_five_add_flows_queue_without_network_access() {
+ let runtime = runtime();
+ let (photo, media) = photo_command();
+ let commands = [
+ (
+ Phase1AddCommand::CreateUpdate(CreateUpdate::new("Harvest update").unwrap()),
+ Vec::new(),
+ ),
+ (photo, vec![media]),
+ (
+ Phase1AddCommand::CreateAsk(CreateAsk::new("Who has basil?", Vec::new()).unwrap()),
+ Vec::new(),
+ ),
+ (
+ Phase1AddCommand::CreateEvent(CreateEvent::date(
+ AuthoredCalendarDateEvent::new(
+ "market-day",
+ "Saturday Market",
+ CalendarDate::parse("2026-08-08").unwrap(),
+ )
+ .unwrap(),
+ )),
+ Vec::new(),
+ ),
+ (
+ Phase1AddCommand::CreateFoodAvailability(CreateFoodAvailability::new(food())),
+ Vec::new(),
+ ),
+ ];
+ for (index, (command, media)) in commands.into_iter().enumerate() {
+ let mut id = [30; 16];
+ id[15] = u8::try_from(index + 1).unwrap();
+ let saved = runtime
+ .phase1_save_draft(
+ id,
+ command,
+ 1_784_347_200,
+ media,
+ None,
+ 100 + u64::try_from(index).unwrap() * 10,
+ )
+ .await
+ .unwrap();
+ let queued = runtime
+ .phase1_queue_draft(
+ id,
+ saved.draft().revision().get(),
+ policy(),
+ 101 + u64::try_from(index).unwrap() * 10,
+ )
+ .await
+ .unwrap();
+ assert_eq!(queued.state(), Phase1OutboxState::Queued);
+ assert_eq!(queued.command_type(), CANONICAL_ADD_COMMAND_TYPES[index]);
+ }
+ }
+
+ #[tokio::test]
+ async fn cancellation_is_terminal_and_preserves_operation_evidence() {
+ let runtime = runtime();
+ let id = [8; 16];
+ runtime
+ .phase1_save_draft(
+ id,
+ Phase1AddCommand::CreateUpdate(CreateUpdate::new("Cancelled").unwrap()),
+ 1_900_000_001,
+ Vec::new(),
+ None,
+ 20,
+ )
+ .await
+ .unwrap();
+ let queued = runtime
+ .phase1_queue_draft(id, 1, policy(), 21)
+ .await
+ .unwrap();
+ let cancelled = runtime
+ .phase1_cancel_draft(id, queued.draft().revision().get(), 22)
+ .await
+ .unwrap();
+ assert_eq!(cancelled.state(), Phase1OutboxState::Cancelled);
+ let push = cancelled.push().expect("retained push evidence");
+ assert_eq!(push.artifact().signing_state(), SigningState::Cancelled);
+ assert_eq!(
+ push.delivery_plan().state(),
+ AuthoredDeliveryState::Cancelled
+ );
+ assert!(push.settlement().is_settled());
+ }
+
+ #[tokio::test]
+ async fn media_phase_revisions_gate_queue_and_record_possible_orphans() {
+ let runtime = runtime();
+ let id = [9; 16];
+ let (command, mut media) = photo_command();
+ media.stage = Phase1MediaStage::Pending;
+ media.validate().unwrap();
+ let saved = runtime
+ .phase1_save_draft(id, command, 1_784_347_200, vec![media], None, 30)
+ .await
+ .unwrap();
+ assert_eq!(saved.state(), Phase1OutboxState::MediaPreparing);
+ assert_eq!(
+ runtime
+ .phase1_queue_draft(id, 1, policy(), 31)
+ .await
+ .unwrap_err(),
+ Phase1DraftError::MediaNotReady
+ );
+ let preparing = runtime
+ .phase1_update_draft_media(
+ id,
+ 1,
+ saved.media()[0].url(),
+ Phase1MediaStage::Preparing,
+ None,
+ 31,
+ )
+ .await
+ .unwrap();
+ let uploading = runtime
+ .phase1_update_draft_media(
+ id,
+ 2,
+ preparing.media()[0].url(),
+ Phase1MediaStage::Uploading,
+ None,
+ 32,
+ )
+ .await
+ .unwrap();
+ let verified = runtime
+ .phase1_update_draft_media(
+ id,
+ 3,
+ uploading.media()[0].url(),
+ Phase1MediaStage::Verified,
+ None,
+ 33,
+ )
+ .await
+ .unwrap();
+ assert_eq!(verified.media()[0].stage(), Phase1MediaStage::Verified);
+ assert_eq!(
+ runtime
+ .phase1_update_draft_media(
+ id,
+ 4,
+ verified.media()[0].url(),
+ Phase1MediaStage::Uploading,
+ None,
+ 34,
+ )
+ .await
+ .unwrap_err(),
+ Phase1DraftError::InvalidMedia
+ );
+ let queued = runtime
+ .phase1_queue_draft(id, 4, policy(), 34)
+ .await
+ .unwrap();
+ let cancelled = runtime
+ .phase1_cancel_draft(id, queued.draft().revision().get(), 35)
+ .await
+ .unwrap();
+ assert_eq!(cancelled.media()[0].stage(), Phase1MediaStage::Orphaned);
+ assert_eq!(runtime.phase1_draft_heads(10).await.unwrap().len(), 1);
+ }
+
+ #[test]
+ fn queue_policy_rejects_duplicates_and_noncanonical_relays() {
+ assert!(
+ Phase1QueuePolicy::new(
+ vec!["wss://relay.example".into(), "wss://relay.example".into()],
+ Phase1RelaySatisfaction::AnyAccepted,
+ 1,
+ Phase1CancellationPolicy::LocalCooperative,
+ )
+ .is_err()
+ );
+ assert!(
+ Phase1QueuePolicy::new(
+ vec!["WSS://relay.example".into()],
+ Phase1RelaySatisfaction::AnyAccepted,
+ 1,
+ Phase1CancellationPolicy::LocalCooperative,
+ )
+ .is_err()
+ );
+ }
+
+ fn photo_command() -> (Phase1AddCommand, Phase1MediaPrerequisite) {
+ let bytes = b"harvest-photo";
+ let hash = BlossomSha256::digest(bytes);
+ let url = format!("https://media.example/{hash}.webp");
+ let media_type = MediaType::parse("image/webp").unwrap();
+ let descriptor = BlobDescriptor::new(
+ BlobUrl::parse(url.as_str()).unwrap(),
+ hash,
+ bytes.len() as u64,
+ media_type.clone(),
+ 1_784_347_100,
+ )
+ .unwrap()
+ .approve_reference()
+ .unwrap()
+ .verify_bytes(bytes, &media_type)
+ .unwrap();
+ let image = AuthoredPostImage::new(
+ AuthoredImage::try_from(descriptor).unwrap(),
+ PostImageDimensions::new(1200, 900).unwrap(),
+ "Harvest",
+ )
+ .unwrap();
+ let command = Phase1AddCommand::CreatePhotoUpdate(
+ CreatePhotoUpdate::new(format!("Harvest photo {url}"), vec![image]).unwrap(),
+ );
+ let prerequisite = Phase1MediaPrerequisite::new(
+ "protected://draft/photo-1",
+ url,
+ hash.to_hex(),
+ "image/webp",
+ bytes.len() as u64,
+ Phase1MediaStage::Verified,
+ )
+ .unwrap();
+ (command, prerequisite)
+ }
+
+ fn food() -> FoodAvailabilityDetails {
+ FoodAvailabilityDetails::new(FoodAvailabilityDetailsParts {
+ content: FoodContent::new("Carrots available this week.").unwrap(),
+ identifier: FoodIdentifier::parse("nantes-carrots").unwrap(),
+ title: FoodText::new("Nantes Carrots").unwrap(),
+ summary: FoodText::new("Fresh bunches").unwrap(),
+ published_at: FoodPublishedAt::new(1_784_347_100).unwrap(),
+ location: FoodText::new("Central Saanich, BC").unwrap(),
+ price: FoodPrice::new("3", FoodCurrency::parse("CAD").unwrap(), FoodUnit::Pound)
+ .unwrap(),
+ quantity: None,
+ status: FoodAvailabilityStatus::Active,
+ images: Vec::new(),
+ })
+ .unwrap()
+ }
+}
diff --git a/crates/sdk/src/sync.rs b/crates/sdk/src/sync.rs
@@ -12,8 +12,8 @@ use radroots_sync::{
policy::Error,
projection::{Reducer, RefreshReceipt, RefreshRequest},
push::{
- AdmissionRunReceipt, DeliveryExecutionReceipt, PushPreparation, PushStatus,
- SigningRunReceipt,
+ AdmissionRunReceipt, DeliveryExecutionReceipt, PushCancellationReceipt, PushPreparation,
+ PushStatus, SigningRunReceipt,
},
};
#[cfg(feature = "sync")]
@@ -164,6 +164,14 @@ impl<'a> Operations<'a> {
self.engine.push_status(operation_id).await
}
+ /// Cancels every remaining local phase while preserving durable evidence.
+ pub async fn cancel_push(
+ &self,
+ operation_id: radroots_sync::policy::SyncId,
+ ) -> Result<PushCancellationReceipt, Error> {
+ self.engine.cancel_push(operation_id).await
+ }
+
/// Runs one bounded signing phase for an exactly prepared operation.
pub async fn sign_prepared(&self, request: PushRequest) -> Result<SigningRunReceipt, Error> {
self.engine.sign_prepared(request).await
diff --git a/crates/storage/src/authored_draft.rs b/crates/storage/src/authored_draft.rs
@@ -0,0 +1,867 @@
+//! Immutable authored-draft revisions stored before outbound side effects.
+
+use core::num::NonZeroU64;
+use radroots_transport::BoxFuture;
+use sha2::{Digest, Sha256};
+use std::{string::String, vec::Vec};
+
+use crate::{Error, journal::OperationInstanceId};
+
+pub const AUTHORED_DRAFT_PAYLOAD_MAX_BYTES: usize = 4 * 1024 * 1024;
+pub const AUTHORED_DRAFT_SCHEMA_MAX_BYTES: usize = 128;
+pub const AUTHORED_DRAFT_QUERY_LIMIT_MAX: u16 = 256;
+
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[cfg_attr(feature = "serde", serde(try_from = "[u8; 16]", into = "[u8; 16]"))]
+#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
+pub struct AuthoredDraftId([u8; 16]);
+
+impl AuthoredDraftId {
+ pub const fn new(value: [u8; 16]) -> Result<Self, Error> {
+ if bytes_are_zero(&value) {
+ Err(Error::InvalidAuthoredDraft)
+ } else {
+ Ok(Self(value))
+ }
+ }
+
+ pub const fn as_bytes(&self) -> &[u8; 16] {
+ &self.0
+ }
+}
+
+impl TryFrom<[u8; 16]> for AuthoredDraftId {
+ type Error = Error;
+
+ fn try_from(value: [u8; 16]) -> Result<Self, Self::Error> {
+ Self::new(value)
+ }
+}
+
+impl From<AuthoredDraftId> for [u8; 16] {
+ fn from(value: AuthoredDraftId) -> Self {
+ value.0
+ }
+}
+
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[cfg_attr(feature = "serde", serde(try_from = "u64", into = "u64"))]
+#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
+pub struct AuthoredDraftRevision(NonZeroU64);
+
+impl AuthoredDraftRevision {
+ pub const INITIAL: Self = Self(NonZeroU64::MIN);
+
+ pub const fn new(value: u64) -> Result<Self, Error> {
+ match NonZeroU64::new(value) {
+ Some(value) => Ok(Self(value)),
+ None => Err(Error::InvalidAuthoredDraft),
+ }
+ }
+
+ pub const fn get(self) -> u64 {
+ self.0.get()
+ }
+
+ pub fn next(self) -> Result<Self, Error> {
+ self.get()
+ .checked_add(1)
+ .ok_or(Error::InvalidAuthoredDraft)
+ .and_then(Self::new)
+ }
+}
+
+impl TryFrom<u64> for AuthoredDraftRevision {
+ type Error = Error;
+
+ fn try_from(value: u64) -> Result<Self, Self::Error> {
+ Self::new(value)
+ }
+}
+
+impl From<AuthoredDraftRevision> for u64 {
+ fn from(value: AuthoredDraftRevision) -> Self {
+ value.get()
+ }
+}
+
+/// Product-independent persistence phases for one local authored draft.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
+#[derive(Clone, Copy, Debug, Eq, PartialEq)]
+pub enum AuthoredDraftStage {
+ Draft,
+ MediaPreparing,
+ MediaUploading,
+ ReadyToSign,
+ Queued,
+ Cancelled,
+}
+
+impl AuthoredDraftStage {
+ pub const fn is_terminal(self) -> bool {
+ matches!(self, Self::Cancelled)
+ }
+}
+
+/// One immutable, integrity-bound draft revision.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[cfg_attr(
+ feature = "serde",
+ serde(
+ try_from = "AuthoredDraftRevisionWire",
+ into = "AuthoredDraftRevisionWire"
+ )
+)]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct AuthoredDraft {
+ draft_id: AuthoredDraftId,
+ revision: AuthoredDraftRevision,
+ author: [u8; 32],
+ payload_schema: String,
+ payload: Vec<u8>,
+ payload_sha256: [u8; 32],
+ stage: AuthoredDraftStage,
+ operation_id: Option<OperationInstanceId>,
+ created_at_unix_ms: u64,
+ updated_at_unix_ms: u64,
+}
+
+#[cfg(feature = "serde")]
+#[derive(serde::Serialize, serde::Deserialize)]
+struct AuthoredDraftRevisionWire {
+ draft_id: AuthoredDraftId,
+ revision: AuthoredDraftRevision,
+ author: [u8; 32],
+ payload_schema: String,
+ payload: Vec<u8>,
+ payload_sha256: [u8; 32],
+ stage: AuthoredDraftStage,
+ operation_id: Option<OperationInstanceId>,
+ created_at_unix_ms: u64,
+ updated_at_unix_ms: u64,
+}
+
+#[cfg(feature = "serde")]
+impl TryFrom<AuthoredDraftRevisionWire> for AuthoredDraft {
+ type Error = Error;
+
+ fn try_from(value: AuthoredDraftRevisionWire) -> Result<Self, Self::Error> {
+ Self::reconstruct(
+ value.draft_id,
+ value.revision,
+ value.author,
+ value.payload_schema,
+ value.payload,
+ value.payload_sha256,
+ value.stage,
+ value.operation_id,
+ value.created_at_unix_ms,
+ value.updated_at_unix_ms,
+ )
+ }
+}
+
+#[cfg(feature = "serde")]
+impl From<AuthoredDraft> for AuthoredDraftRevisionWire {
+ fn from(value: AuthoredDraft) -> Self {
+ Self {
+ draft_id: value.draft_id,
+ revision: value.revision,
+ author: value.author,
+ payload_schema: value.payload_schema,
+ payload: value.payload,
+ payload_sha256: value.payload_sha256,
+ stage: value.stage,
+ operation_id: value.operation_id,
+ created_at_unix_ms: value.created_at_unix_ms,
+ updated_at_unix_ms: value.updated_at_unix_ms,
+ }
+ }
+}
+
+impl AuthoredDraft {
+ #[allow(clippy::too_many_arguments)]
+ pub fn initial(
+ draft_id: AuthoredDraftId,
+ author: [u8; 32],
+ payload_schema: impl Into<String>,
+ payload: Vec<u8>,
+ stage: AuthoredDraftStage,
+ operation_id: Option<OperationInstanceId>,
+ created_at_unix_ms: u64,
+ ) -> Result<Self, Error> {
+ if matches!(
+ stage,
+ AuthoredDraftStage::ReadyToSign
+ | AuthoredDraftStage::Queued
+ | AuthoredDraftStage::Cancelled
+ ) {
+ return Err(Error::InvalidAuthoredDraft);
+ }
+ let payload_sha256 = Sha256::digest(payload.as_slice()).into();
+ Self::reconstruct(
+ draft_id,
+ AuthoredDraftRevision::INITIAL,
+ author,
+ payload_schema.into(),
+ payload,
+ payload_sha256,
+ stage,
+ operation_id,
+ created_at_unix_ms,
+ created_at_unix_ms,
+ )
+ }
+
+ pub fn successor(
+ &self,
+ payload: Vec<u8>,
+ stage: AuthoredDraftStage,
+ operation_id: Option<OperationInstanceId>,
+ updated_at_unix_ms: u64,
+ ) -> Result<Self, Error> {
+ let payload_sha256 = Sha256::digest(payload.as_slice()).into();
+ let next = Self::reconstruct(
+ self.draft_id,
+ self.revision.next()?,
+ self.author,
+ self.payload_schema.clone(),
+ payload,
+ payload_sha256,
+ stage,
+ operation_id,
+ self.created_at_unix_ms,
+ updated_at_unix_ms,
+ )?;
+ next.validate_successor_of(self)?;
+ Ok(next)
+ }
+
+ #[allow(clippy::too_many_arguments)]
+ pub fn reconstruct(
+ draft_id: AuthoredDraftId,
+ revision: AuthoredDraftRevision,
+ author: [u8; 32],
+ payload_schema: impl Into<String>,
+ payload: Vec<u8>,
+ payload_sha256: [u8; 32],
+ stage: AuthoredDraftStage,
+ operation_id: Option<OperationInstanceId>,
+ created_at_unix_ms: u64,
+ updated_at_unix_ms: u64,
+ ) -> Result<Self, Error> {
+ let value = Self {
+ draft_id,
+ revision,
+ author,
+ payload_schema: payload_schema.into(),
+ payload,
+ payload_sha256,
+ stage,
+ operation_id,
+ created_at_unix_ms,
+ updated_at_unix_ms,
+ };
+ value.validate()?;
+ Ok(value)
+ }
+
+ pub fn validate(&self) -> Result<(), Error> {
+ let schema = self.payload_schema.as_str();
+ let requires_operation = matches!(
+ self.stage,
+ AuthoredDraftStage::ReadyToSign | AuthoredDraftStage::Queued
+ );
+ if bytes_are_zero(&self.author)
+ || schema.is_empty()
+ || schema.len() > AUTHORED_DRAFT_SCHEMA_MAX_BYTES
+ || schema != schema.trim()
+ || schema.chars().any(char::is_control)
+ || self.payload.is_empty()
+ || self.payload.len() > AUTHORED_DRAFT_PAYLOAD_MAX_BYTES
+ || Sha256::digest(self.payload.as_slice()).as_slice() != self.payload_sha256
+ || self.created_at_unix_ms == 0
+ || self.updated_at_unix_ms < self.created_at_unix_ms
+ || (requires_operation && self.operation_id.is_none())
+ || (!requires_operation
+ && self.stage != AuthoredDraftStage::Cancelled
+ && self.operation_id.is_some())
+ {
+ return Err(Error::InvalidAuthoredDraft);
+ }
+ Ok(())
+ }
+
+ pub fn validate_successor_of(&self, previous: &Self) -> Result<(), Error> {
+ let identity_matches = self.draft_id == previous.draft_id
+ && self.author == previous.author
+ && self.payload_schema == previous.payload_schema
+ && self.created_at_unix_ms == previous.created_at_unix_ms
+ && self.revision == previous.revision.next()?
+ && self.updated_at_unix_ms >= previous.updated_at_unix_ms;
+ let operation_matches = match (previous.operation_id, self.operation_id) {
+ (Some(previous), Some(next)) => previous == next,
+ (None, _) => true,
+ (Some(_), None) => false,
+ };
+ let frozen_queue_payload = !matches!(
+ (previous.stage, self.stage),
+ (
+ AuthoredDraftStage::ReadyToSign,
+ AuthoredDraftStage::ReadyToSign | AuthoredDraftStage::Queued
+ ) | (AuthoredDraftStage::Queued, AuthoredDraftStage::Queued)
+ ) || self.payload_sha256 == previous.payload_sha256;
+ let stage_allowed = match previous.stage {
+ AuthoredDraftStage::Draft => matches!(
+ self.stage,
+ AuthoredDraftStage::Draft
+ | AuthoredDraftStage::MediaPreparing
+ | AuthoredDraftStage::ReadyToSign
+ | AuthoredDraftStage::Cancelled
+ ),
+ AuthoredDraftStage::MediaPreparing => matches!(
+ self.stage,
+ AuthoredDraftStage::MediaPreparing
+ | AuthoredDraftStage::MediaUploading
+ | AuthoredDraftStage::ReadyToSign
+ | AuthoredDraftStage::Cancelled
+ ),
+ AuthoredDraftStage::MediaUploading => matches!(
+ self.stage,
+ AuthoredDraftStage::MediaPreparing
+ | AuthoredDraftStage::MediaUploading
+ | AuthoredDraftStage::ReadyToSign
+ | AuthoredDraftStage::Cancelled
+ ),
+ AuthoredDraftStage::ReadyToSign => matches!(
+ self.stage,
+ AuthoredDraftStage::ReadyToSign
+ | AuthoredDraftStage::Queued
+ | AuthoredDraftStage::Cancelled
+ ),
+ AuthoredDraftStage::Queued => matches!(
+ self.stage,
+ AuthoredDraftStage::Queued | AuthoredDraftStage::Cancelled
+ ),
+ AuthoredDraftStage::Cancelled => false,
+ };
+ if !identity_matches || !operation_matches || !frozen_queue_payload || !stage_allowed {
+ return Err(Error::DraftRevisionConflict);
+ }
+ Ok(())
+ }
+
+ pub const fn draft_id(&self) -> AuthoredDraftId {
+ self.draft_id
+ }
+ pub const fn revision(&self) -> AuthoredDraftRevision {
+ self.revision
+ }
+ pub const fn author(&self) -> &[u8; 32] {
+ &self.author
+ }
+ pub fn payload_schema(&self) -> &str {
+ self.payload_schema.as_str()
+ }
+ pub fn payload(&self) -> &[u8] {
+ self.payload.as_slice()
+ }
+ pub const fn payload_sha256(&self) -> &[u8; 32] {
+ &self.payload_sha256
+ }
+ pub const fn stage(&self) -> AuthoredDraftStage {
+ self.stage
+ }
+ pub const fn operation_id(&self) -> Option<OperationInstanceId> {
+ self.operation_id
+ }
+ pub const fn created_at_unix_ms(&self) -> u64 {
+ self.created_at_unix_ms
+ }
+ pub const fn updated_at_unix_ms(&self) -> u64 {
+ self.updated_at_unix_ms
+ }
+}
+
+#[derive(Clone, Copy, Debug, Eq, PartialEq)]
+pub enum DraftAppendDisposition {
+ Inserted,
+ Replay,
+}
+
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct DraftAppendReceipt {
+ draft: AuthoredDraft,
+ disposition: DraftAppendDisposition,
+}
+
+impl DraftAppendReceipt {
+ pub const fn new(draft: AuthoredDraft, disposition: DraftAppendDisposition) -> Self {
+ Self { draft, disposition }
+ }
+ pub const fn draft(&self) -> &AuthoredDraft {
+ &self.draft
+ }
+ pub const fn disposition(&self) -> DraftAppendDisposition {
+ self.disposition
+ }
+}
+
+pub trait AuthoredDraftStore: Send + Sync {
+ fn append_authored_draft(
+ &self,
+ draft: AuthoredDraft,
+ expected_head: Option<AuthoredDraftRevision>,
+ ) -> BoxFuture<'_, Result<DraftAppendReceipt, Error>>;
+
+ fn authored_draft_head(
+ &self,
+ draft_id: AuthoredDraftId,
+ ) -> BoxFuture<'_, Result<Option<AuthoredDraft>, Error>>;
+
+ fn authored_draft_revision(
+ &self,
+ draft_id: AuthoredDraftId,
+ revision: AuthoredDraftRevision,
+ ) -> BoxFuture<'_, Result<Option<AuthoredDraft>, Error>>;
+
+ fn authored_draft_heads(
+ &self,
+ author: [u8; 32],
+ limit: u16,
+ ) -> BoxFuture<'_, Result<Vec<AuthoredDraft>, Error>>;
+}
+
+const fn bytes_are_zero<const N: usize>(bytes: &[u8; N]) -> bool {
+ let mut index = 0;
+ while index < bytes.len() {
+ if bytes[index] != 0 {
+ return false;
+ }
+ index += 1;
+ }
+ true
+}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+
+ fn draft() -> AuthoredDraft {
+ AuthoredDraft::initial(
+ AuthoredDraftId::new([1; 16]).unwrap(),
+ [2; 32],
+ "radroots.phase1-draft.v1",
+ b"draft".to_vec(),
+ AuthoredDraftStage::Draft,
+ None,
+ 10,
+ )
+ .unwrap()
+ }
+
+ #[test]
+ fn revisions_are_immutable_and_phase_ordered() {
+ let first = draft();
+ let media = first
+ .successor(
+ b"media".to_vec(),
+ AuthoredDraftStage::MediaPreparing,
+ None,
+ 11,
+ )
+ .unwrap();
+ let operation = OperationInstanceId::new([3; 16]).unwrap();
+ let ready = media
+ .successor(
+ b"ready".to_vec(),
+ AuthoredDraftStage::ReadyToSign,
+ Some(operation),
+ 12,
+ )
+ .unwrap();
+ assert!(
+ ready
+ .successor(
+ b"changed".to_vec(),
+ AuthoredDraftStage::Queued,
+ Some(operation),
+ 13,
+ )
+ .is_err()
+ );
+ let queued = ready
+ .successor(
+ b"ready".to_vec(),
+ AuthoredDraftStage::Queued,
+ Some(operation),
+ 13,
+ )
+ .unwrap();
+ assert_eq!(queued.revision().get(), 4);
+ assert!(
+ queued
+ .successor(b"x".to_vec(), AuthoredDraftStage::Draft, None, 14)
+ .is_err()
+ );
+ }
+
+ #[test]
+ fn reconstruction_rejects_tampering_and_invalid_contracts() {
+ let value = draft();
+ assert!(AuthoredDraftId::new([0; 16]).is_err());
+ assert!(AuthoredDraftRevision::new(0).is_err());
+ assert!(
+ AuthoredDraft::reconstruct(
+ value.draft_id(),
+ value.revision(),
+ *value.author(),
+ value.payload_schema(),
+ value.payload().to_vec(),
+ [9; 32],
+ value.stage(),
+ None,
+ value.created_at_unix_ms(),
+ value.updated_at_unix_ms(),
+ )
+ .is_err()
+ );
+ assert!(
+ AuthoredDraft::initial(
+ value.draft_id(),
+ [0; 32],
+ "schema",
+ Vec::new(),
+ AuthoredDraftStage::ReadyToSign,
+ None,
+ 0,
+ )
+ .is_err()
+ );
+ }
+
+ #[test]
+ fn validation_rejects_every_independent_invalid_field() {
+ let value = draft();
+ let digest = |payload: &[u8]| -> [u8; 32] { Sha256::digest(payload).into() };
+ let reconstruct = |author: [u8; 32],
+ schema: String,
+ payload: Vec<u8>,
+ payload_sha256: [u8; 32],
+ stage: AuthoredDraftStage,
+ operation_id: Option<OperationInstanceId>,
+ created_at_unix_ms: u64,
+ updated_at_unix_ms: u64| {
+ AuthoredDraft::reconstruct(
+ value.draft_id(),
+ value.revision(),
+ author,
+ schema,
+ payload,
+ payload_sha256,
+ stage,
+ operation_id,
+ created_at_unix_ms,
+ updated_at_unix_ms,
+ )
+ };
+
+ let operation = OperationInstanceId::new([3; 16]).unwrap();
+ let valid_payload = b"draft".to_vec();
+ let valid_digest = digest(&valid_payload);
+ let invalid = [
+ reconstruct(
+ [0; 32],
+ "schema".into(),
+ valid_payload.clone(),
+ valid_digest,
+ AuthoredDraftStage::Draft,
+ None,
+ 10,
+ 10,
+ ),
+ reconstruct(
+ [2; 32],
+ String::new(),
+ valid_payload.clone(),
+ valid_digest,
+ AuthoredDraftStage::Draft,
+ None,
+ 10,
+ 10,
+ ),
+ reconstruct(
+ [2; 32],
+ "x".repeat(AUTHORED_DRAFT_SCHEMA_MAX_BYTES + 1),
+ valid_payload.clone(),
+ valid_digest,
+ AuthoredDraftStage::Draft,
+ None,
+ 10,
+ 10,
+ ),
+ reconstruct(
+ [2; 32],
+ " schema".into(),
+ valid_payload.clone(),
+ valid_digest,
+ AuthoredDraftStage::Draft,
+ None,
+ 10,
+ 10,
+ ),
+ reconstruct(
+ [2; 32],
+ "bad\nschema".into(),
+ valid_payload.clone(),
+ valid_digest,
+ AuthoredDraftStage::Draft,
+ None,
+ 10,
+ 10,
+ ),
+ reconstruct(
+ [2; 32],
+ "schema".into(),
+ Vec::new(),
+ digest(&[]),
+ AuthoredDraftStage::Draft,
+ None,
+ 10,
+ 10,
+ ),
+ reconstruct(
+ [2; 32],
+ "schema".into(),
+ valid_payload.clone(),
+ [9; 32],
+ AuthoredDraftStage::Draft,
+ None,
+ 10,
+ 10,
+ ),
+ reconstruct(
+ [2; 32],
+ "schema".into(),
+ valid_payload.clone(),
+ valid_digest,
+ AuthoredDraftStage::Draft,
+ None,
+ 0,
+ 10,
+ ),
+ reconstruct(
+ [2; 32],
+ "schema".into(),
+ valid_payload.clone(),
+ valid_digest,
+ AuthoredDraftStage::Draft,
+ None,
+ 10,
+ 9,
+ ),
+ reconstruct(
+ [2; 32],
+ "schema".into(),
+ valid_payload.clone(),
+ valid_digest,
+ AuthoredDraftStage::ReadyToSign,
+ None,
+ 10,
+ 10,
+ ),
+ reconstruct(
+ [2; 32],
+ "schema".into(),
+ valid_payload,
+ valid_digest,
+ AuthoredDraftStage::Draft,
+ Some(operation),
+ 10,
+ 10,
+ ),
+ ];
+ assert!(invalid.into_iter().all(|result| result.is_err()));
+
+ let oversized = vec![0; AUTHORED_DRAFT_PAYLOAD_MAX_BYTES + 1];
+ assert!(
+ reconstruct(
+ [2; 32],
+ "schema".into(),
+ oversized.clone(),
+ digest(&oversized),
+ AuthoredDraftStage::Draft,
+ None,
+ 10,
+ 10,
+ )
+ .is_err()
+ );
+ assert!(
+ AuthoredDraftRevision::new(u64::MAX)
+ .unwrap()
+ .next()
+ .is_err()
+ );
+ assert!(AuthoredDraftStage::Cancelled.is_terminal());
+ assert!(!AuthoredDraftStage::Draft.is_terminal());
+ for forbidden_initial_stage in [
+ AuthoredDraftStage::ReadyToSign,
+ AuthoredDraftStage::Queued,
+ AuthoredDraftStage::Cancelled,
+ ] {
+ assert!(
+ AuthoredDraft::initial(
+ AuthoredDraftId::new([4; 16]).unwrap(),
+ [5; 32],
+ "schema",
+ b"payload".to_vec(),
+ forbidden_initial_stage,
+ None,
+ 10,
+ )
+ .is_err()
+ );
+ }
+ }
+
+ #[test]
+ fn every_stage_pair_obeys_the_transition_matrix() {
+ let stages = [
+ AuthoredDraftStage::Draft,
+ AuthoredDraftStage::MediaPreparing,
+ AuthoredDraftStage::MediaUploading,
+ AuthoredDraftStage::ReadyToSign,
+ AuthoredDraftStage::Queued,
+ AuthoredDraftStage::Cancelled,
+ ];
+ let operation = OperationInstanceId::new([3; 16]).unwrap();
+
+ for previous_stage in stages {
+ for next_stage in stages {
+ let previous_operation = matches!(
+ previous_stage,
+ AuthoredDraftStage::ReadyToSign | AuthoredDraftStage::Queued
+ )
+ .then_some(operation);
+ let previous = AuthoredDraft::reconstruct(
+ AuthoredDraftId::new([1; 16]).unwrap(),
+ AuthoredDraftRevision::INITIAL,
+ [2; 32],
+ "schema",
+ b"payload".to_vec(),
+ Sha256::digest(b"payload").into(),
+ previous_stage,
+ previous_operation,
+ 10,
+ 10,
+ )
+ .unwrap();
+ let next_operation = match next_stage {
+ AuthoredDraftStage::ReadyToSign | AuthoredDraftStage::Queued => Some(operation),
+ AuthoredDraftStage::Cancelled => previous_operation,
+ _ => None,
+ };
+ let payload = if matches!(
+ (previous_stage, next_stage),
+ (
+ AuthoredDraftStage::ReadyToSign,
+ AuthoredDraftStage::ReadyToSign | AuthoredDraftStage::Queued
+ ) | (AuthoredDraftStage::Queued, AuthoredDraftStage::Queued)
+ ) {
+ b"payload".to_vec()
+ } else {
+ b"next".to_vec()
+ };
+ let allowed = match previous_stage {
+ AuthoredDraftStage::Draft => matches!(
+ next_stage,
+ AuthoredDraftStage::Draft
+ | AuthoredDraftStage::MediaPreparing
+ | AuthoredDraftStage::ReadyToSign
+ | AuthoredDraftStage::Cancelled
+ ),
+ AuthoredDraftStage::MediaPreparing => matches!(
+ next_stage,
+ AuthoredDraftStage::MediaPreparing
+ | AuthoredDraftStage::MediaUploading
+ | AuthoredDraftStage::ReadyToSign
+ | AuthoredDraftStage::Cancelled
+ ),
+ AuthoredDraftStage::MediaUploading => matches!(
+ next_stage,
+ AuthoredDraftStage::MediaPreparing
+ | AuthoredDraftStage::MediaUploading
+ | AuthoredDraftStage::ReadyToSign
+ | AuthoredDraftStage::Cancelled
+ ),
+ AuthoredDraftStage::ReadyToSign => matches!(
+ next_stage,
+ AuthoredDraftStage::ReadyToSign
+ | AuthoredDraftStage::Queued
+ | AuthoredDraftStage::Cancelled
+ ),
+ AuthoredDraftStage::Queued => matches!(
+ next_stage,
+ AuthoredDraftStage::Queued | AuthoredDraftStage::Cancelled
+ ),
+ AuthoredDraftStage::Cancelled => false,
+ };
+ assert_eq!(
+ previous
+ .successor(payload, next_stage, next_operation, 11)
+ .is_ok(),
+ allowed,
+ "{previous_stage:?} -> {next_stage:?}"
+ );
+ }
+ }
+
+ let previous = AuthoredDraft::reconstruct(
+ AuthoredDraftId::new([1; 16]).unwrap(),
+ AuthoredDraftRevision::INITIAL,
+ [2; 32],
+ "schema",
+ b"payload".to_vec(),
+ Sha256::digest(b"payload").into(),
+ AuthoredDraftStage::ReadyToSign,
+ Some(operation),
+ 10,
+ 10,
+ )
+ .unwrap();
+ assert!(
+ previous
+ .successor(
+ b"payload".to_vec(),
+ AuthoredDraftStage::ReadyToSign,
+ Some(OperationInstanceId::new([4; 16]).unwrap()),
+ 11,
+ )
+ .is_err()
+ );
+ assert!(
+ previous
+ .successor(b"payload".to_vec(), AuthoredDraftStage::Cancelled, None, 11,)
+ .is_err()
+ );
+ let rebound_identity = AuthoredDraft::reconstruct(
+ AuthoredDraftId::new([9; 16]).unwrap(),
+ previous.revision().next().unwrap(),
+ *previous.author(),
+ previous.payload_schema(),
+ previous.payload().to_vec(),
+ *previous.payload_sha256(),
+ AuthoredDraftStage::ReadyToSign,
+ previous.operation_id(),
+ previous.created_at_unix_ms(),
+ 11,
+ )
+ .unwrap();
+ assert_eq!(
+ rebound_identity.validate_successor_of(&previous),
+ Err(Error::DraftRevisionConflict)
+ );
+ }
+}
diff --git a/crates/storage/src/error.rs b/crates/storage/src/error.rs
@@ -10,6 +10,10 @@ pub enum Error {
InvalidAuthoredArtifact,
InvalidAuthoredTransition,
InvalidAuthoredDeliveryPlan,
+ InvalidAuthoredDraft,
+ DraftRevisionConflict,
+ DraftNotFound,
+ CorruptAuthoredDraft,
DeliveryPlanClaimConflict,
DeliveryAttemptOverflow,
InvalidWorkClaim,
@@ -137,6 +141,10 @@ impl fmt::Display for Error {
Self::InvalidAuthoredArtifact => "storage authored artifact is invalid",
Self::InvalidAuthoredTransition => "storage authored transition is invalid",
Self::InvalidAuthoredDeliveryPlan => "storage authored delivery plan is invalid",
+ Self::InvalidAuthoredDraft => "storage authored draft is invalid",
+ Self::DraftRevisionConflict => "storage authored draft revision conflicts",
+ Self::DraftNotFound => "storage authored draft was not found",
+ Self::CorruptAuthoredDraft => "storage authored draft is corrupt",
Self::DeliveryPlanClaimConflict => "storage delivery plan claim conflicts",
Self::DeliveryAttemptOverflow => "storage delivery attempt counter overflowed",
Self::InvalidWorkClaim => "storage authored work claim is invalid",
diff --git a/crates/storage/src/lib.rs b/crates/storage/src/lib.rs
@@ -5,6 +5,7 @@ pub mod atomic;
pub mod authored;
pub mod authored_atomic;
pub mod authored_delivery;
+pub mod authored_draft;
pub mod backup;
mod error;
pub mod event;
@@ -34,6 +35,7 @@ pub trait Storage:
+ backup::StorageReliability
+ atomic::AtomicStorage
+ authored_atomic::AuthoredAtomicStorage
+ + authored_draft::AuthoredDraftStore
{
}
@@ -46,5 +48,6 @@ impl<T> Storage for T where
+ backup::StorageReliability
+ atomic::AtomicStorage
+ authored_atomic::AuthoredAtomicStorage
+ + authored_draft::AuthoredDraftStore
{
}
diff --git a/crates/storage/src/memory.rs b/crates/storage/src/memory.rs
@@ -17,6 +17,10 @@ use crate::{
AuthoredWorkTarget, CancelAuthoredTarget, ClaimAuthoredTarget,
},
authored_delivery::DeliveryAttemptOutcome,
+ authored_draft::{
+ AUTHORED_DRAFT_QUERY_LIMIT_MAX, AuthoredDraft, AuthoredDraftId, AuthoredDraftRevision,
+ AuthoredDraftStore, DraftAppendDisposition, DraftAppendReceipt,
+ },
backup::{
BackupId, BackupOperation, BackupPlan, BackupTransition, ReliabilityRevision,
RestoreOperation, RestorePlan, RestoreTransition, StorageReliability,
@@ -82,6 +86,7 @@ struct State {
authored_artifacts: Vec<crate::authored::AuthoredArtifact>,
authored_delivery_plans: Vec<crate::authored_delivery::AuthoredDeliveryPlan>,
authored_atomic_receipts: Vec<AuthoredAtomicReceipt>,
+ authored_drafts: Vec<AuthoredDraft>,
closed: bool,
}
@@ -115,6 +120,7 @@ impl MemoryStorage {
authored_artifacts: Vec::new(),
authored_delivery_plans: Vec::new(),
authored_atomic_receipts: Vec::new(),
+ authored_drafts: Vec::new(),
closed: false,
}),
}
@@ -1776,6 +1782,117 @@ impl AuthoredAtomicStorage for MemoryStorage {
}
}
+impl AuthoredDraftStore for MemoryStorage {
+ fn append_authored_draft(
+ &self,
+ draft: AuthoredDraft,
+ expected_head: Option<AuthoredDraftRevision>,
+ ) -> BoxFuture<'_, Result<DraftAppendReceipt, Error>> {
+ Box::pin(async move {
+ draft.validate()?;
+ let mut state = self.state()?;
+ if let Some(existing) = state.authored_drafts.iter().find(|existing| {
+ existing.draft_id() == draft.draft_id() && existing.revision() == draft.revision()
+ }) {
+ return if existing == &draft {
+ Ok(DraftAppendReceipt::new(
+ existing.clone(),
+ DraftAppendDisposition::Replay,
+ ))
+ } else {
+ Err(Error::DraftRevisionConflict)
+ };
+ }
+ let head = state
+ .authored_drafts
+ .iter()
+ .filter(|existing| existing.draft_id() == draft.draft_id())
+ .max_by_key(|existing| existing.revision());
+ match (head, expected_head) {
+ (None, None) if draft.revision() == AuthoredDraftRevision::INITIAL => {}
+ (Some(previous), Some(expected)) if previous.revision() == expected => {
+ draft.validate_successor_of(previous)?;
+ }
+ _ => return Err(Error::DraftRevisionConflict),
+ }
+ state.authored_drafts.push(draft.clone());
+ Ok(DraftAppendReceipt::new(
+ draft,
+ DraftAppendDisposition::Inserted,
+ ))
+ })
+ }
+
+ fn authored_draft_head(
+ &self,
+ draft_id: AuthoredDraftId,
+ ) -> BoxFuture<'_, Result<Option<AuthoredDraft>, Error>> {
+ Box::pin(async move {
+ Ok(self
+ .state()?
+ .authored_drafts
+ .iter()
+ .filter(|draft| draft.draft_id() == draft_id)
+ .max_by_key(|draft| draft.revision())
+ .cloned())
+ })
+ }
+
+ fn authored_draft_revision(
+ &self,
+ draft_id: AuthoredDraftId,
+ revision: AuthoredDraftRevision,
+ ) -> BoxFuture<'_, Result<Option<AuthoredDraft>, Error>> {
+ Box::pin(async move {
+ Ok(self
+ .state()?
+ .authored_drafts
+ .iter()
+ .find(|draft| draft.draft_id() == draft_id && draft.revision() == revision)
+ .cloned())
+ })
+ }
+
+ fn authored_draft_heads(
+ &self,
+ author: [u8; 32],
+ limit: u16,
+ ) -> BoxFuture<'_, Result<Vec<AuthoredDraft>, Error>> {
+ Box::pin(async move {
+ if author.iter().all(|byte| *byte == 0)
+ || limit == 0
+ || limit > AUTHORED_DRAFT_QUERY_LIMIT_MAX
+ {
+ return Err(Error::InvalidAuthoredDraft);
+ }
+ let state = self.state()?;
+ let mut heads = state
+ .authored_drafts
+ .iter()
+ .filter(|draft| draft.author() == &author)
+ .fold(Vec::<AuthoredDraft>::new(), |mut heads, draft| {
+ match heads
+ .iter_mut()
+ .find(|head| head.draft_id() == draft.draft_id())
+ {
+ Some(head) if draft.revision() > head.revision() => *head = draft.clone(),
+ None => heads.push(draft.clone()),
+ Some(_) => {}
+ }
+ heads
+ });
+ heads.sort_by(|left, right| {
+ right
+ .updated_at_unix_ms()
+ .cmp(&left.updated_at_unix_ms())
+ .then_with(|| left.draft_id().cmp(&right.draft_id()))
+ });
+ heads.truncate(usize::from(limit));
+ Ok(heads)
+ })
+ }
+}
+
fn require_artifact_claim(
claim: Option<&crate::authored::WorkClaim>,
fence: &crate::authored_atomic::WorkFence,
diff --git a/crates/storage/tests/authored_draft.rs b/crates/storage/tests/authored_draft.rs
@@ -0,0 +1,143 @@
+use futures_executor::block_on;
+use radroots_storage::{
+ Error,
+ authored_draft::{
+ AUTHORED_DRAFT_QUERY_LIMIT_MAX, AuthoredDraft, AuthoredDraftId, AuthoredDraftRevision,
+ AuthoredDraftStage, AuthoredDraftStore, DraftAppendDisposition,
+ },
+ memory::MemoryStorage,
+};
+
+fn draft(id: u8, author: u8, created_at_unix_ms: u64) -> AuthoredDraft {
+ AuthoredDraft::initial(
+ AuthoredDraftId::new([id; 16]).expect("draft ID"),
+ [author; 32],
+ "radroots.phase1-draft.v1",
+ vec![id],
+ AuthoredDraftStage::Draft,
+ None,
+ created_at_unix_ms,
+ )
+ .expect("draft")
+}
+
+#[test]
+fn memory_draft_store_replays_conflicts_and_queries_exact_heads() {
+ fn accepts_dyn(_: &dyn AuthoredDraftStore) {}
+
+ let store = MemoryStorage::default();
+ accepts_dyn(&store);
+ let first = draft(1, 9, 10);
+ assert_eq!(
+ block_on(store.append_authored_draft(first.clone(), None))
+ .expect("insert")
+ .disposition(),
+ DraftAppendDisposition::Inserted
+ );
+ let replay = block_on(store.append_authored_draft(first.clone(), None)).expect("replay");
+ assert_eq!(replay.disposition(), DraftAppendDisposition::Replay);
+ assert_eq!(replay.draft(), &first);
+
+ let conflicting = AuthoredDraft::initial(
+ first.draft_id(),
+ *first.author(),
+ first.payload_schema(),
+ b"conflict".to_vec(),
+ AuthoredDraftStage::Draft,
+ None,
+ first.created_at_unix_ms(),
+ )
+ .expect("conflicting revision");
+ assert_eq!(
+ block_on(store.append_authored_draft(conflicting, None)),
+ Err(Error::DraftRevisionConflict)
+ );
+
+ let second = first
+ .successor(
+ b"second".to_vec(),
+ AuthoredDraftStage::MediaPreparing,
+ None,
+ 11,
+ )
+ .expect("successor");
+ assert_eq!(
+ block_on(store.append_authored_draft(second.clone(), None)),
+ Err(Error::DraftRevisionConflict)
+ );
+ assert_eq!(
+ block_on(store.append_authored_draft(
+ second.clone(),
+ Some(AuthoredDraftRevision::new(2).expect("revision")),
+ )),
+ Err(Error::DraftRevisionConflict)
+ );
+ block_on(store.append_authored_draft(second.clone(), Some(AuthoredDraftRevision::INITIAL)))
+ .expect("successor insert");
+
+ assert_eq!(
+ block_on(store.authored_draft_head(first.draft_id()))
+ .expect("head")
+ .expect("stored head"),
+ second
+ );
+ assert_eq!(
+ block_on(store.authored_draft_revision(first.draft_id(), first.revision()))
+ .expect("revision")
+ .expect("stored revision"),
+ first
+ );
+ assert!(
+ block_on(store.authored_draft_revision(
+ AuthoredDraftId::new([8; 16]).expect("missing ID"),
+ AuthoredDraftRevision::INITIAL,
+ ))
+ .expect("missing revision query")
+ .is_none()
+ );
+}
+
+#[test]
+fn memory_draft_head_query_is_bounded_filtered_sorted_and_truncated() {
+ let store = MemoryStorage::default();
+ let first = draft(1, 9, 10);
+ let second = draft(2, 9, 20);
+ let other_author = draft(3, 8, 30);
+ for value in [&first, &second, &other_author] {
+ block_on(store.append_authored_draft(value.clone(), None)).expect("insert");
+ }
+ let first_head = first
+ .successor(
+ b"new head".to_vec(),
+ AuthoredDraftStage::MediaPreparing,
+ None,
+ 40,
+ )
+ .expect("new head");
+ block_on(store.append_authored_draft(first_head.clone(), Some(first.revision())))
+ .expect("head insert");
+
+ let heads = block_on(store.authored_draft_heads([9; 32], 2)).expect("heads");
+ assert_eq!(heads, vec![first_head.clone(), second]);
+ assert_eq!(
+ block_on(store.authored_draft_heads([9; 32], 1)).expect("truncated heads"),
+ vec![first_head]
+ );
+ assert!(
+ block_on(store.authored_draft_heads([7; 32], 2))
+ .expect("unmatched author")
+ .is_empty()
+ );
+ assert_eq!(
+ block_on(store.authored_draft_heads([0; 32], 1)),
+ Err(Error::InvalidAuthoredDraft)
+ );
+ assert_eq!(
+ block_on(store.authored_draft_heads([9; 32], 0)),
+ Err(Error::InvalidAuthoredDraft)
+ );
+ assert_eq!(
+ block_on(store.authored_draft_heads([9; 32], AUTHORED_DRAFT_QUERY_LIMIT_MAX + 1,)),
+ Err(Error::InvalidAuthoredDraft)
+ );
+}
diff --git a/crates/storage_sqlite/src/authored_draft.rs b/crates/storage_sqlite/src/authored_draft.rs
@@ -0,0 +1,627 @@
+use crate::SqliteStorage;
+use radroots_storage::{
+ Error,
+ authored_draft::{
+ AUTHORED_DRAFT_QUERY_LIMIT_MAX, AuthoredDraft, AuthoredDraftId, AuthoredDraftRevision,
+ AuthoredDraftStage, AuthoredDraftStore, DraftAppendDisposition, DraftAppendReceipt,
+ },
+ event::BoxFuture,
+};
+use sqlx::{Row, sqlite::SqliteRow};
+
+const SNAPSHOT_MAX_BYTES: usize = 16 * 1024 * 1024;
+
+impl AuthoredDraftStore for SqliteStorage {
+ fn append_authored_draft(
+ &self,
+ draft: AuthoredDraft,
+ expected_head: Option<AuthoredDraftRevision>,
+ ) -> BoxFuture<'_, Result<DraftAppendReceipt, Error>> {
+ Box::pin(async move {
+ draft.validate()?;
+ if self.event_mode() == radroots_storage::status::EventStoreMode::ReadOnly {
+ return Err(Error::BackendUnavailable);
+ }
+ let mut transaction = self
+ .pool()
+ .begin_with("BEGIN IMMEDIATE")
+ .await
+ .map_err(map_backend)?;
+ if let Some(row) = sqlx::query(
+ "SELECT * FROM radroots_runtime_authored_draft_revisions
+ WHERE draft_id = ? AND revision = ?",
+ )
+ .bind(draft.draft_id().as_bytes().as_slice())
+ .bind(i64_from_u64(draft.revision().get())?)
+ .fetch_optional(&mut *transaction)
+ .await
+ .map_err(map_backend)?
+ {
+ let existing = decode_row(&row)?;
+ transaction.rollback().await.map_err(map_backend)?;
+ return if existing == draft {
+ Ok(DraftAppendReceipt::new(
+ existing,
+ DraftAppendDisposition::Replay,
+ ))
+ } else {
+ Err(Error::DraftRevisionConflict)
+ };
+ }
+
+ let head = sqlx::query(
+ "SELECT * FROM radroots_runtime_authored_draft_revisions
+ WHERE draft_id = ? ORDER BY revision DESC LIMIT 1",
+ )
+ .bind(draft.draft_id().as_bytes().as_slice())
+ .fetch_optional(&mut *transaction)
+ .await
+ .map_err(map_backend)?
+ .as_ref()
+ .map(decode_row)
+ .transpose()?;
+ match (head.as_ref(), expected_head) {
+ (None, None) if draft.revision() == AuthoredDraftRevision::INITIAL => {}
+ (Some(previous), Some(expected)) if previous.revision() == expected => {
+ draft.validate_successor_of(previous)?;
+ }
+ _ => {
+ let _ = transaction.rollback().await;
+ return Err(Error::DraftRevisionConflict);
+ }
+ }
+
+ let snapshot = encode_snapshot(&draft)?;
+ sqlx::query(
+ "INSERT INTO radroots_runtime_authored_draft_revisions (
+ draft_id, revision, author, stage, operation_id, payload_sha256,
+ created_at_unix_ms, updated_at_unix_ms, snapshot
+ ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)",
+ )
+ .bind(draft.draft_id().as_bytes().as_slice())
+ .bind(i64_from_u64(draft.revision().get())?)
+ .bind(draft.author().as_slice())
+ .bind(stage_code(draft.stage()))
+ .bind(draft.operation_id().map(|id| id.as_bytes().to_vec()))
+ .bind(draft.payload_sha256().as_slice())
+ .bind(i64_from_u64(draft.created_at_unix_ms())?)
+ .bind(i64_from_u64(draft.updated_at_unix_ms())?)
+ .bind(snapshot)
+ .execute(&mut *transaction)
+ .await
+ .map_err(map_backend)?;
+ transaction.commit().await.map_err(map_backend)?;
+ Ok(DraftAppendReceipt::new(
+ draft,
+ DraftAppendDisposition::Inserted,
+ ))
+ })
+ }
+
+ fn authored_draft_head(
+ &self,
+ draft_id: AuthoredDraftId,
+ ) -> BoxFuture<'_, Result<Option<AuthoredDraft>, Error>> {
+ Box::pin(async move {
+ sqlx::query(
+ "SELECT * FROM radroots_runtime_authored_draft_revisions
+ WHERE draft_id = ? ORDER BY revision DESC LIMIT 1",
+ )
+ .bind(draft_id.as_bytes().as_slice())
+ .fetch_optional(self.pool())
+ .await
+ .map_err(map_backend)?
+ .as_ref()
+ .map(decode_row)
+ .transpose()
+ })
+ }
+
+ fn authored_draft_revision(
+ &self,
+ draft_id: AuthoredDraftId,
+ revision: AuthoredDraftRevision,
+ ) -> BoxFuture<'_, Result<Option<AuthoredDraft>, Error>> {
+ Box::pin(async move {
+ sqlx::query(
+ "SELECT * FROM radroots_runtime_authored_draft_revisions
+ WHERE draft_id = ? AND revision = ?",
+ )
+ .bind(draft_id.as_bytes().as_slice())
+ .bind(i64_from_u64(revision.get())?)
+ .fetch_optional(self.pool())
+ .await
+ .map_err(map_backend)?
+ .as_ref()
+ .map(decode_row)
+ .transpose()
+ })
+ }
+
+ fn authored_draft_heads(
+ &self,
+ author: [u8; 32],
+ limit: u16,
+ ) -> BoxFuture<'_, Result<Vec<AuthoredDraft>, Error>> {
+ Box::pin(async move {
+ if author.iter().all(|byte| *byte == 0)
+ || limit == 0
+ || limit > AUTHORED_DRAFT_QUERY_LIMIT_MAX
+ {
+ return Err(Error::InvalidAuthoredDraft);
+ }
+ sqlx::query(
+ "SELECT revisions.*
+ FROM radroots_runtime_authored_draft_revisions AS revisions
+ WHERE revisions.author = ?
+ AND revisions.revision = (
+ SELECT MAX(head.revision)
+ FROM radroots_runtime_authored_draft_revisions AS head
+ WHERE head.draft_id = revisions.draft_id
+ )
+ ORDER BY revisions.updated_at_unix_ms DESC, revisions.draft_id
+ LIMIT ?",
+ )
+ .bind(author.as_slice())
+ .bind(i64::from(limit))
+ .fetch_all(self.pool())
+ .await
+ .map_err(map_backend)?
+ .iter()
+ .map(decode_row)
+ .collect()
+ })
+ }
+}
+
+fn encode_snapshot(draft: &AuthoredDraft) -> Result<Vec<u8>, Error> {
+ let snapshot = serde_json::to_vec(draft).map_err(|_| Error::InvalidAuthoredDraft)?;
+ if snapshot.is_empty() || snapshot.len() > SNAPSHOT_MAX_BYTES {
+ return Err(Error::InvalidAuthoredDraft);
+ }
+ Ok(snapshot)
+}
+
+fn decode_row(row: &SqliteRow) -> Result<AuthoredDraft, Error> {
+ let snapshot = row
+ .try_get::<Vec<u8>, _>("snapshot")
+ .map_err(|_| Error::CorruptAuthoredDraft)?;
+ if snapshot.is_empty() || snapshot.len() > SNAPSHOT_MAX_BYTES {
+ return Err(Error::CorruptAuthoredDraft);
+ }
+ let draft = serde_json::from_slice::<AuthoredDraft>(snapshot.as_slice())
+ .map_err(|_| Error::CorruptAuthoredDraft)?;
+ let draft_id = fixed::<16>(row, "draft_id")?;
+ let revision = u64_from_i64(
+ row.try_get::<i64, _>("revision")
+ .map_err(|_| Error::CorruptAuthoredDraft)?,
+ )?;
+ let author = fixed::<32>(row, "author")?;
+ let stage = row
+ .try_get::<i64, _>("stage")
+ .map_err(|_| Error::CorruptAuthoredDraft)?;
+ let operation_id = row
+ .try_get::<Option<Vec<u8>>, _>("operation_id")
+ .map_err(|_| Error::CorruptAuthoredDraft)?;
+ let payload_sha256 = fixed::<32>(row, "payload_sha256")?;
+ let created = u64_from_i64(
+ row.try_get::<i64, _>("created_at_unix_ms")
+ .map_err(|_| Error::CorruptAuthoredDraft)?,
+ )?;
+ let updated = u64_from_i64(
+ row.try_get::<i64, _>("updated_at_unix_ms")
+ .map_err(|_| Error::CorruptAuthoredDraft)?,
+ )?;
+ let operation_matches = match (operation_id, draft.operation_id()) {
+ (None, None) => true,
+ (Some(raw), Some(expected)) => raw.as_slice() == expected.as_bytes(),
+ _ => false,
+ };
+ if draft.draft_id().as_bytes() != &draft_id
+ || draft.revision().get() != revision
+ || draft.author() != &author
+ || stage_code(draft.stage()) != stage
+ || !operation_matches
+ || draft.payload_sha256() != &payload_sha256
+ || draft.created_at_unix_ms() != created
+ || draft.updated_at_unix_ms() != updated
+ {
+ return Err(Error::CorruptAuthoredDraft);
+ }
+ Ok(draft)
+}
+
+fn fixed<const N: usize>(row: &SqliteRow, column: &str) -> Result<[u8; N], Error> {
+ row.try_get::<Vec<u8>, _>(column)
+ .map_err(|_| Error::CorruptAuthoredDraft)?
+ .try_into()
+ .map_err(|_| Error::CorruptAuthoredDraft)
+}
+
+const fn stage_code(stage: AuthoredDraftStage) -> i64 {
+ match stage {
+ AuthoredDraftStage::Draft => 0,
+ AuthoredDraftStage::MediaPreparing => 1,
+ AuthoredDraftStage::MediaUploading => 2,
+ AuthoredDraftStage::ReadyToSign => 3,
+ AuthoredDraftStage::Queued => 4,
+ AuthoredDraftStage::Cancelled => 5,
+ }
+}
+
+fn i64_from_u64(value: u64) -> Result<i64, Error> {
+ i64::try_from(value).map_err(|_| Error::InvalidAuthoredDraft)
+}
+
+fn u64_from_i64(value: i64) -> Result<u64, Error> {
+ u64::try_from(value).map_err(|_| Error::CorruptAuthoredDraft)
+}
+
+fn map_backend(_: sqlx::Error) -> Error {
+ Error::BackendUnavailable
+}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+ use crate::{OpenMode, OpenOptions, Paths};
+ use radroots_storage::event::SourceGeneration;
+ use sha2::Digest;
+ use tempfile::TempDir;
+
+ fn draft(id: u8, at: u64) -> AuthoredDraft {
+ AuthoredDraft::initial(
+ AuthoredDraftId::new([id; 16]).unwrap(),
+ [7; 32],
+ "radroots.phase1-draft.v1",
+ vec![id],
+ AuthoredDraftStage::Draft,
+ None,
+ at,
+ )
+ .unwrap()
+ }
+
+ async fn open_store(temp: &TempDir) -> SqliteStorage {
+ let paths = Paths::from_directory(temp.path()).unwrap();
+ SqliteStorage::open(
+ OpenOptions::new(paths, OpenMode::Create)
+ .with_source_generation(SourceGeneration::new([9; 32]).unwrap(), 9)
+ .unwrap(),
+ )
+ .await
+ .unwrap()
+ }
+
+ #[tokio::test]
+ async fn revisions_survive_reopen_and_replay_exactly() {
+ let temp = TempDir::new().unwrap();
+ let store = open_store(&temp).await;
+ let first = draft(1, 10);
+ assert_eq!(
+ store
+ .append_authored_draft(first.clone(), None)
+ .await
+ .unwrap()
+ .disposition(),
+ DraftAppendDisposition::Inserted
+ );
+ assert_eq!(
+ store
+ .append_authored_draft(first.clone(), None)
+ .await
+ .unwrap()
+ .disposition(),
+ DraftAppendDisposition::Replay
+ );
+ let second = first
+ .successor(
+ b"next".to_vec(),
+ AuthoredDraftStage::MediaPreparing,
+ None,
+ 11,
+ )
+ .unwrap();
+ store
+ .append_authored_draft(second.clone(), Some(first.revision()))
+ .await
+ .unwrap();
+ drop(store);
+ let reopened = open_store(&temp).await;
+ assert_eq!(
+ reopened
+ .authored_draft_head(first.draft_id())
+ .await
+ .unwrap(),
+ Some(second.clone())
+ );
+ assert_eq!(
+ reopened
+ .authored_draft_revision(first.draft_id(), first.revision())
+ .await
+ .unwrap(),
+ Some(first)
+ );
+ assert_eq!(
+ reopened.authored_draft_heads([7; 32], 10).await.unwrap(),
+ vec![second]
+ );
+ }
+
+ #[tokio::test]
+ async fn conflicts_and_query_bounds_fail_closed() {
+ let temp = TempDir::new().unwrap();
+ let store = open_store(&temp).await;
+ let first = draft(1, 10);
+ store
+ .append_authored_draft(first.clone(), None)
+ .await
+ .unwrap();
+ let conflicting = AuthoredDraft::initial(
+ first.draft_id(),
+ [7; 32],
+ first.payload_schema(),
+ b"conflict".to_vec(),
+ AuthoredDraftStage::Draft,
+ None,
+ 10,
+ )
+ .unwrap();
+ assert_eq!(
+ store.append_authored_draft(conflicting, None).await,
+ Err(Error::DraftRevisionConflict)
+ );
+ assert_eq!(
+ store.authored_draft_heads([7; 32], 0).await,
+ Err(Error::InvalidAuthoredDraft)
+ );
+ assert_eq!(
+ store.authored_draft_heads([0; 32], 1).await,
+ Err(Error::InvalidAuthoredDraft)
+ );
+ assert_eq!(
+ store
+ .authored_draft_heads([7; 32], AUTHORED_DRAFT_QUERY_LIMIT_MAX + 1)
+ .await,
+ Err(Error::InvalidAuthoredDraft)
+ );
+
+ let successor = first
+ .successor(
+ b"next".to_vec(),
+ AuthoredDraftStage::MediaPreparing,
+ None,
+ 11,
+ )
+ .unwrap();
+ assert_eq!(
+ store.append_authored_draft(successor.clone(), None).await,
+ Err(Error::DraftRevisionConflict)
+ );
+ assert_eq!(
+ store
+ .append_authored_draft(successor, Some(AuthoredDraftRevision::new(2).unwrap()),)
+ .await,
+ Err(Error::DraftRevisionConflict)
+ );
+ assert_eq!(
+ store
+ .append_authored_draft(draft(2, 20), Some(AuthoredDraftRevision::INITIAL))
+ .await,
+ Err(Error::DraftRevisionConflict)
+ );
+ let noninitial_first = AuthoredDraft::reconstruct(
+ AuthoredDraftId::new([3; 16]).unwrap(),
+ AuthoredDraftRevision::new(2).unwrap(),
+ [7; 32],
+ "radroots.phase1-draft.v1",
+ vec![3],
+ sha2::Sha256::digest([3]).into(),
+ AuthoredDraftStage::Draft,
+ None,
+ 30,
+ 30,
+ )
+ .unwrap();
+ assert_eq!(
+ store.append_authored_draft(noninitial_first, None).await,
+ Err(Error::DraftRevisionConflict)
+ );
+ }
+
+ #[tokio::test]
+ async fn simultaneous_first_append_has_one_insert_and_one_exact_replay() {
+ let temp = TempDir::new().unwrap();
+ let store = open_store(&temp).await;
+ let first = draft(2, 20);
+ let (left, right) = tokio::join!(
+ store.append_authored_draft(first.clone(), None),
+ store.append_authored_draft(first, None),
+ );
+ let dispositions = [left.unwrap().disposition(), right.unwrap().disposition()];
+ assert!(dispositions.contains(&DraftAppendDisposition::Inserted));
+ assert!(dispositions.contains(&DraftAppendDisposition::Replay));
+ }
+
+ #[tokio::test]
+ async fn read_only_append_and_u64_overflow_fail_before_mutation() {
+ let temp = TempDir::new().unwrap();
+ let store = open_store(&temp).await;
+ let paths = Paths::from_directory(temp.path()).unwrap();
+ drop(store);
+ let read_only = SqliteStorage::open(OpenOptions::new(paths, OpenMode::ReadOnly))
+ .await
+ .unwrap();
+ assert_eq!(
+ read_only.append_authored_draft(draft(3, 30), None).await,
+ Err(Error::BackendUnavailable)
+ );
+ drop(read_only);
+
+ let store = open_store(&temp).await;
+ let overflow = AuthoredDraft::initial(
+ AuthoredDraftId::new([4; 16]).unwrap(),
+ [7; 32],
+ "radroots.phase1-draft.v1",
+ vec![4],
+ AuthoredDraftStage::Draft,
+ None,
+ i64::MAX as u64 + 1,
+ )
+ .unwrap();
+ assert_eq!(
+ store.append_authored_draft(overflow, None).await,
+ Err(Error::InvalidAuthoredDraft)
+ );
+ }
+
+ #[allow(clippy::too_many_arguments)]
+ async fn insert_raw_and_decode(
+ store: &SqliteStorage,
+ draft: &AuthoredDraft,
+ draft_id: [u8; 16],
+ revision: i64,
+ author: [u8; 32],
+ stage: i64,
+ operation_id: Option<Vec<u8>>,
+ payload_sha256: [u8; 32],
+ created_at_unix_ms: i64,
+ updated_at_unix_ms: i64,
+ ) -> Result<AuthoredDraft, Error> {
+ sqlx::query(
+ "INSERT INTO radroots_runtime_authored_draft_revisions (
+ draft_id, revision, author, stage, operation_id, payload_sha256,
+ created_at_unix_ms, updated_at_unix_ms, snapshot
+ ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)",
+ )
+ .bind(draft_id.as_slice())
+ .bind(revision)
+ .bind(author.as_slice())
+ .bind(stage)
+ .bind(operation_id)
+ .bind(payload_sha256.as_slice())
+ .bind(created_at_unix_ms)
+ .bind(updated_at_unix_ms)
+ .bind(encode_snapshot(draft).unwrap())
+ .execute(store.pool())
+ .await
+ .unwrap();
+ let row = sqlx::query(
+ "SELECT * FROM radroots_runtime_authored_draft_revisions
+ WHERE draft_id = ? AND revision = ?",
+ )
+ .bind(draft_id.as_slice())
+ .bind(revision)
+ .fetch_one(store.pool())
+ .await
+ .unwrap();
+ decode_row(&row)
+ }
+
+ #[tokio::test]
+ async fn every_redundant_draft_column_is_verified_against_the_snapshot() {
+ let temp = TempDir::new().unwrap();
+ let store = open_store(&temp).await;
+
+ let value = draft(10, 100);
+ assert_eq!(
+ insert_raw_and_decode(
+ &store,
+ &value,
+ [99; 16],
+ 1,
+ *value.author(),
+ stage_code(value.stage()),
+ None,
+ *value.payload_sha256(),
+ 100,
+ 100,
+ )
+ .await,
+ Err(Error::CorruptAuthoredDraft)
+ );
+
+ for (id, revision, author, stage, operation_id, payload_sha256, created, updated) in [
+ (
+ 11,
+ 2,
+ [7; 32],
+ 0,
+ None,
+ *draft(11, 100).payload_sha256(),
+ 100,
+ 100,
+ ),
+ (
+ 12,
+ 1,
+ [8; 32],
+ 0,
+ None,
+ *draft(12, 100).payload_sha256(),
+ 100,
+ 100,
+ ),
+ (
+ 13,
+ 1,
+ [7; 32],
+ 1,
+ None,
+ *draft(13, 100).payload_sha256(),
+ 100,
+ 100,
+ ),
+ (
+ 14,
+ 1,
+ [7; 32],
+ 0,
+ Some(vec![1; 16]),
+ *draft(14, 100).payload_sha256(),
+ 100,
+ 100,
+ ),
+ (15, 1, [7; 32], 0, None, [9; 32], 100, 100),
+ (
+ 16,
+ 1,
+ [7; 32],
+ 0,
+ None,
+ *draft(16, 100).payload_sha256(),
+ 99,
+ 100,
+ ),
+ (
+ 17,
+ 1,
+ [7; 32],
+ 0,
+ None,
+ *draft(17, 100).payload_sha256(),
+ 100,
+ 101,
+ ),
+ ] {
+ let value = draft(id, 100);
+ assert_eq!(
+ insert_raw_and_decode(
+ &store,
+ &value,
+ [id; 16],
+ revision,
+ author,
+ stage,
+ operation_id,
+ payload_sha256,
+ created,
+ updated,
+ )
+ .await,
+ Err(Error::CorruptAuthoredDraft),
+ "redundant column case {id}"
+ );
+ }
+ }
+}
diff --git a/crates/storage_sqlite/src/lib.rs b/crates/storage_sqlite/src/lib.rs
@@ -13,6 +13,7 @@ pub mod status;
mod atomic;
mod authored;
+mod authored_draft;
mod event;
mod journal;
mod outbox;
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 <= 12
+ && plan.current_version <= 13
&& plan.steps.len() == usize::try_from(plan.current_version).unwrap_or(usize::MAX)
&& plan
.steps
@@ -488,6 +488,7 @@ const fn set_user_version_sql(version: u32) -> Option<&'static str> {
10 => Some("PRAGMA user_version = 10"),
11 => Some("PRAGMA user_version = 11"),
12 => Some("PRAGMA user_version = 12"),
+ 13 => Some("PRAGMA user_version = 13"),
_ => None,
}
}
@@ -714,7 +715,7 @@ mod tests {
.execute(&mut newer)
.await
.expect("application id");
- sqlx::raw_sql("PRAGMA user_version = 13")
+ sqlx::raw_sql("PRAGMA user_version = 14")
.execute(&mut newer)
.await
.expect("newer version");
@@ -723,10 +724,10 @@ mod tests {
Err(Error::SchemaTooNew {
database: RUNTIME_DATABASE,
supported: runtime::CURRENT_VERSION,
- actual: 13,
+ actual: 14,
})
));
- assert_eq!(pragma(&mut newer, "user_version").await, 13);
+ assert_eq!(pragma(&mut newer, "user_version").await, 14);
let mut wrong_identity = connection().await;
establish_runtime_version(&mut wrong_identity, 1).await;
diff --git a/crates/storage_sqlite/src/migration/runtime/0013_authored_draft_revisions.up.sql b/crates/storage_sqlite/src/migration/runtime/0013_authored_draft_revisions.up.sql
@@ -0,0 +1,27 @@
+CREATE TABLE radroots_runtime_authored_draft_revisions (
+ draft_id BLOB NOT NULL CHECK(length(draft_id) = 16),
+ revision INTEGER NOT NULL CHECK(revision > 0),
+ author BLOB NOT NULL CHECK(length(author) = 32),
+ stage INTEGER NOT NULL CHECK(stage BETWEEN 0 AND 5),
+ operation_id BLOB CHECK(operation_id IS NULL OR length(operation_id) = 16),
+ payload_sha256 BLOB NOT NULL CHECK(length(payload_sha256) = 32),
+ created_at_unix_ms INTEGER NOT NULL CHECK(created_at_unix_ms > 0),
+ updated_at_unix_ms INTEGER NOT NULL CHECK(updated_at_unix_ms >= created_at_unix_ms),
+ snapshot BLOB NOT NULL CHECK(length(snapshot) BETWEEN 1 AND 16777216),
+ PRIMARY KEY(draft_id, revision)
+) STRICT, WITHOUT ROWID;
+
+CREATE INDEX radroots_runtime_authored_draft_author_head_idx
+ON radroots_runtime_authored_draft_revisions(author, updated_at_unix_ms DESC, draft_id, revision DESC);
+
+CREATE TRIGGER radroots_runtime_authored_draft_revisions_update_guard
+BEFORE UPDATE ON radroots_runtime_authored_draft_revisions
+BEGIN
+ SELECT RAISE(ABORT, 'authored draft revisions are immutable');
+END;
+
+CREATE TRIGGER radroots_runtime_authored_draft_revisions_delete_guard
+BEFORE DELETE ON radroots_runtime_authored_draft_revisions
+BEGIN
+ SELECT RAISE(ABORT, 'authored draft revisions are immutable');
+END;
diff --git a/crates/storage_sqlite/src/migration/runtime/mod.rs b/crates/storage_sqlite/src/migration/runtime/mod.rs
@@ -6,7 +6,7 @@
/// Lowest runtime schema version this package can recognize.
pub const MINIMUM_VERSION: u32 = 1;
/// Current runtime schema version created by this package.
-pub const CURRENT_VERSION: u32 = 12;
+pub const CURRENT_VERSION: u32 = 13;
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");
@@ -22,6 +22,7 @@ const PROJECTION_REBUILD_SOURCE_BINDING_V10_SQL: &str =
const AUTHORED_OPERATIONS_V11_SQL: &str = include_str!("0011_authored_operations.up.sql");
const MATERIALIZED_PROJECTION_DOCUMENTS_V12_SQL: &str =
include_str!("0012_materialized_projection_documents.up.sql");
+const AUTHORED_DRAFT_REVISIONS_V13_SQL: &str = include_str!("0013_authored_draft_revisions.up.sql");
/// Stable, non-SQL description of one forward runtime migration.
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
@@ -477,6 +478,83 @@ const RUNTIME_V12_OBJECTS: &[&str] = &[
"radroots_runtime_source_generations_sequence_guard",
];
+const RUNTIME_V13_OBJECTS: &[&str] = &[
+ "radroots_runtime_atomic_commits",
+ "radroots_runtime_authored_artifacts",
+ "radroots_runtime_authored_artifacts_admission_ready_idx",
+ "radroots_runtime_authored_artifacts_signing_ready_idx",
+ "radroots_runtime_authored_atomic_commits",
+ "radroots_runtime_authored_atomic_commits_delete_guard",
+ "radroots_runtime_authored_atomic_commits_update_guard",
+ "radroots_runtime_authored_delivery_attempts",
+ "radroots_runtime_authored_delivery_plans",
+ "radroots_runtime_authored_delivery_ready_idx",
+ "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_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 {
@@ -551,6 +629,12 @@ pub const MIGRATIONS: &[MigrationDescriptor] = &[
up_sha256: "c5a083ed078d2d1d9623b8c7fd6521055a59c6f4c099318779c9fa32fa17b211",
owned_objects: RUNTIME_V12_OBJECTS,
},
+ MigrationDescriptor {
+ version: 13,
+ name: "authored_draft_revisions",
+ up_sha256: "7bb6466a14c7ca601b76852addf4a19996ee9ffd62436bdac8e9bf631dc46785",
+ owned_objects: RUNTIME_V13_OBJECTS,
+ },
];
pub(crate) const fn migration_sql(version: u32) -> Option<&'static str> {
@@ -567,6 +651,7 @@ pub(crate) const fn migration_sql(version: u32) -> Option<&'static str> {
10 => Some(PROJECTION_REBUILD_SOURCE_BINDING_V10_SQL),
11 => Some(AUTHORED_OPERATIONS_V11_SQL),
12 => Some(MATERIALIZED_PROJECTION_DOCUMENTS_V12_SQL),
+ 13 => Some(AUTHORED_DRAFT_REVISIONS_V13_SQL),
_ => None,
}
}
@@ -610,8 +695,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, 12);
- assert_eq!(MIGRATIONS.len(), 12);
+ assert_eq!(CURRENT_VERSION, 13);
+ assert_eq!(MIGRATIONS.len(), 13);
let migration = MIGRATIONS[8];
assert_eq!(snapshot.schema_version, 1);
assert_eq!(snapshot.database, "runtime.sqlite");
@@ -639,7 +724,7 @@ mod tests {
let sql = migration_sql(migration.version()).expect("registered SQL");
assert_eq!(format!("{:x}", Sha256::digest(sql)), migration.up_sha256());
}
- assert_eq!(migration_sql(13), None);
+ assert_eq!(migration_sql(14), None);
}
#[tokio::test]
diff --git a/crates/sync/src/lib.rs b/crates/sync/src/lib.rs
@@ -15,7 +15,7 @@ pub use engine::Engine;
pub use policy::Error;
pub use pull::{PullReceipt, PullRequest};
pub use push::{
- AdmissionRunReceipt, DeliveryExecutionReceipt, PushPreparation, PushRequest, PushStatus,
- SigningRunReceipt,
+ AdmissionRunReceipt, DeliveryExecutionReceipt, PushCancellationReceipt, PushPreparation,
+ PushRequest, PushStatus, SigningRunReceipt,
};
pub use status::SyncStatus;
diff --git a/crates/sync/src/push.rs b/crates/sync/src/push.rs
@@ -26,7 +26,7 @@ use radroots_storage::{
},
authored_delivery::{
AuthoredDeliveryIntent, AuthoredDeliveryPlan, AuthoredDeliveryPlanId,
- DELIVERY_PLAN_ATTEMPTS_MAX, DeliveryAttemptOutcome,
+ AuthoredDeliveryState, DELIVERY_PLAN_ATTEMPTS_MAX, DeliveryAttemptOutcome,
},
event::{AdmissionDisposition, EventAdmission, EventStore},
journal::{IdempotencyKey, OperationInstanceId},
@@ -216,6 +216,22 @@ pub struct DeliveryExecutionReceipt {
replay: bool,
}
+/// Durable result of cancelling all still-cancellable phases of one push.
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct PushCancellationReceipt {
+ status: PushStatus,
+ changed: bool,
+}
+
+impl PushCancellationReceipt {
+ pub const fn status(&self) -> &PushStatus {
+ &self.status
+ }
+ pub const fn changed(&self) -> bool {
+ self.changed
+ }
+}
+
impl DeliveryExecutionReceipt {
pub const fn plan(&self) -> &AuthoredDeliveryPlan {
&self.plan
@@ -346,6 +362,83 @@ impl Engine {
}))
}
+ /// Cancels every remaining local phase without discarding signed bytes or
+ /// delivery evidence. Re-entry after a phase-boundary crash finishes any
+ /// remaining cancellation and terminal replays are no-ops.
+ pub async fn cancel_push(
+ &self,
+ operation_id: SyncId,
+ ) -> Result<PushCancellationReceipt, Error> {
+ let mut status = self
+ .push_status(operation_id)
+ .await?
+ .ok_or(Error::StorageFailed)?;
+ let mut changed = false;
+ let now = self
+ .clock
+ .now_unix_ms()?
+ .max(status.artifact.updated_at_unix_ms())
+ .max(status.delivery_plan.updated_at_unix_ms());
+
+ let artifact_target = match status.artifact.signing_state() {
+ SigningState::Planned | SigningState::Retryable => Some(
+ CancelAuthoredTarget::ArtifactSigning(status.artifact.artifact_id()),
+ ),
+ SigningState::Signed
+ if matches!(
+ status.artifact.admission_state(),
+ AdmissionState::Pending | AdmissionState::Retryable
+ ) =>
+ {
+ Some(CancelAuthoredTarget::ArtifactAdmission(
+ status.artifact.artifact_id(),
+ ))
+ }
+ SigningState::Signed
+ | SigningState::Indeterminate
+ | SigningState::FailedTerminal
+ | SigningState::Cancelled => None,
+ };
+ if let Some(target) = artifact_target {
+ self.storage
+ .execute_authored(AuthoredAtomicCommand::Cancel(
+ CancelAuthoredWork::new(target, status.artifact.revision(), now)
+ .map_err(map_storage_error)?,
+ ))
+ .await
+ .map_err(map_storage_error)?;
+ changed = true;
+ status = self
+ .push_status(operation_id)
+ .await?
+ .ok_or(Error::StorageFailed)?;
+ }
+
+ if matches!(
+ status.delivery_plan.state(),
+ AuthoredDeliveryState::Pending | AuthoredDeliveryState::Retryable
+ ) {
+ self.storage
+ .execute_authored(AuthoredAtomicCommand::Cancel(
+ CancelAuthoredWork::new(
+ CancelAuthoredTarget::DeliveryPlan(status.delivery_plan.plan_id()),
+ status.delivery_plan.revision(),
+ now,
+ )
+ .map_err(map_storage_error)?,
+ ))
+ .await
+ .map_err(map_storage_error)?;
+ changed = true;
+ status = self
+ .push_status(operation_id)
+ .await?
+ .ok_or(Error::StorageFailed)?;
+ }
+
+ Ok(PushCancellationReceipt { status, changed })
+ }
+
/// Claims and executes one prepared signing artifact.
pub async fn sign_prepared(&self, request: PushRequest) -> Result<SigningRunReceipt, Error> {
let signer = self.signer.as_deref().ok_or(Error::MissingSigner)?;
diff --git a/crates/sync/tests/push_enqueue.rs b/crates/sync/tests/push_enqueue.rs
@@ -1010,6 +1010,88 @@ fn preparation_is_atomic_status_visible_and_replays_without_external_effects() {
}
#[test]
+fn phase_aware_cancellation_preserves_signed_artifacts_and_is_reentrant() {
+ let unsigned_signer = Arc::new(MockSigner::new(SignBehavior::Pending));
+ let (unsigned_engine, _) = setup_engine(unsigned_signer.clone());
+ let unsigned = request(61, "wss://relay.example");
+ block_on(unsigned_engine.prepare_push(unsigned.clone())).expect("prepare unsigned");
+ let cancelled = block_on(unsigned_engine.cancel_push(unsigned.operation_id()))
+ .expect("cancel unsigned operation");
+ assert!(cancelled.changed());
+ assert_eq!(
+ cancelled.status().artifact().signing_state(),
+ radroots_storage::authored::SigningState::Cancelled
+ );
+ assert_eq!(
+ cancelled.status().delivery_plan().state(),
+ AuthoredDeliveryState::Cancelled
+ );
+ assert!(cancelled.status().settlement().is_settled());
+ let replay = block_on(unsigned_engine.cancel_push(unsigned.operation_id()))
+ .expect("replay cancellation");
+ assert!(!replay.changed());
+ assert_eq!(replay.status(), cancelled.status());
+ assert_eq!(unsigned_signer.calls.load(Ordering::Relaxed), 0);
+
+ let signed_signer = Arc::new(MockSigner::new(SignBehavior::Success {
+ completed_at_unix_ms: 1_800_000_200_010,
+ }));
+ let (signed_engine, _) = setup_engine(signed_signer);
+ let signed = request(62, "wss://relay.example");
+ block_on(signed_engine.sign_prepared(signed.clone())).expect("sign operation");
+ let before = block_on(signed_engine.push_status(signed.operation_id()))
+ .unwrap()
+ .unwrap();
+ let raw = before
+ .artifact()
+ .signed()
+ .expect("exact signed artifact")
+ .event()
+ .raw_json()
+ .to_owned();
+ let cancelled = block_on(signed_engine.cancel_push(signed.operation_id()))
+ .expect("cancel signed operation");
+ assert_eq!(
+ cancelled.status().artifact().admission_state(),
+ radroots_storage::authored::AdmissionState::Cancelled
+ );
+ assert_eq!(
+ cancelled
+ .status()
+ .artifact()
+ .signed()
+ .expect("signed bytes retained")
+ .event()
+ .raw_json(),
+ raw
+ );
+ assert_eq!(
+ cancelled.status().delivery_plan().state(),
+ AuthoredDeliveryState::Cancelled
+ );
+
+ let admitted_signer = Arc::new(MockSigner::new(SignBehavior::Success {
+ completed_at_unix_ms: 1_800_000_200_010,
+ }));
+ let (admitted_engine, _) = setup_engine(admitted_signer);
+ let admitted = request(63, "wss://relay.example");
+ execute_to_admitted(&admitted_engine, &admitted);
+ let cancelled = block_on(admitted_engine.cancel_push(admitted.operation_id()))
+ .expect("cancel admitted delivery");
+ assert!(
+ cancelled
+ .status()
+ .artifact()
+ .admission_state()
+ .is_admitted()
+ );
+ assert_eq!(
+ cancelled.status().delivery_plan().state(),
+ AuthoredDeliveryState::Cancelled
+ );
+}
+
+#[test]
fn authored_storage_outcome_contracts_fail_closed_at_each_orchestration_phase() {
let storage = Arc::new(FaultStorage::new(84));
storage.fault_next_prepared();