lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

commit 363b9be18831be2fd6ef505d6154fb878ed8efd7
parent 13ac1f2135199702a07ba6aa14ab22101f17fcb0
Author: triesap <tyson@radroots.org>
Date:   Thu, 10 Sep 2026 12:12:04 +0000

storage: Support governed draft submission and queries

- Commit captured draft intent and authored preparation atomically
- Bind exact replay before source CAS and preserve ordinary Prepare receipts
- Bound scoped recovery pages and isolate corrupt revision evidence
- Qualify migration compatibility API coverage and complete workspace gates

Diffstat:
Mcontracts/api_baselines/radroots_storage.txt | 73+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcontracts/api_baselines/radroots_sync.txt | 2++
Acontracts/architecture/decisions/authored_draft_submission.v1.json | 17+++++++++++++++++
Mcontracts/architecture/deviations.toml | 21+++++++++++++++++++++
Mcrates/storage/README.md | 22++++++++++++++++++++++
Mcrates/storage/src/authored_atomic.rs | 97++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcrates/storage/src/authored_draft.rs | 64+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---
Acrates/storage/src/authored_draft_query.rs | 265+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acrates/storage/src/authored_draft_submission.rs | 279+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/storage/src/lib.rs | 2++
Mcrates/storage/src/memory.rs | 159+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--------------
Mcrates/storage/tests/authored_atomic.rs | 3+++
Acrates/storage/tests/authored_atomic/draft_submission.rs | 532+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acrates/storage/tests/authored_draft_query.rs | 260+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/storage_sqlite/README.md | 14++++++++++++++
Mcrates/storage_sqlite/src/authored.rs | 95++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---------------------
Mcrates/storage_sqlite/src/authored_draft.rs | 91++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-------------------
Acrates/storage_sqlite/src/authored_draft_query.rs | 101+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acrates/storage_sqlite/src/authored_draft_query_tests.rs | 234+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acrates/storage_sqlite/src/authored_draft_submission_tests.rs | 549+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/storage_sqlite/src/migration.rs | 20+++++++++++++-------
Acrates/storage_sqlite/src/migration/runtime/0014_authored_draft_query_metadata.up.sql | 19+++++++++++++++++++
Mcrates/storage_sqlite/src/migration/runtime/mod.rs | 96+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++----
Acrates/storage_sqlite/src/migration_draft_query_tests.rs | 143+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/sync/README.md | 7+++++++
Mcrates/sync/src/push.rs | 65++++++++++++++++++++++++++++++++++++++++-------------------------
Mcrates/sync/tests/push_enqueue.rs | 41+++++++++++++++++++++++++++++++++++++++++
27 files changed, 3157 insertions(+), 114 deletions(-)

diff --git a/contracts/api_baselines/radroots_storage.txt b/contracts/api_baselines/radroots_storage.txt @@ -200,6 +200,7 @@ pub radroots_storage::authored_atomic::AuthoredAtomicCommand::ApplySigned(radroo pub radroots_storage::authored_atomic::AuthoredAtomicCommand::Cancel(radroots_storage::authored_atomic::CancelAuthoredWork) pub radroots_storage::authored_atomic::AuthoredAtomicCommand::Claim(radroots_storage::authored_atomic::ClaimAuthoredWork) pub radroots_storage::authored_atomic::AuthoredAtomicCommand::Prepare(radroots_storage::authored_atomic::PrepareAuthoredOperation) +pub radroots_storage::authored_atomic::AuthoredAtomicCommand::PrepareFromDraft(alloc::boxed::Box<radroots_storage::authored_draft_submission::PrepareFromDraft>) impl radroots_storage::authored_atomic::AuthoredAtomicCommand pub fn radroots_storage::authored_atomic::AuthoredAtomicCommand::commit_id(&self) -> radroots_storage::atomic::AtomicCommitId pub fn radroots_storage::authored_atomic::AuthoredAtomicCommand::digest(&self) -> radroots_storage::atomic::AtomicCommitDigest @@ -211,6 +212,7 @@ pub radroots_storage::authored_atomic::AuthoredAtomicOutcome::Prepared pub radroots_storage::authored_atomic::AuthoredAtomicOutcome::Prepared::artifacts: alloc::vec::Vec<radroots_storage::authored::AuthoredArtifact> pub radroots_storage::authored_atomic::AuthoredAtomicOutcome::Prepared::delivery_plans: alloc::vec::Vec<radroots_storage::authored_delivery::AuthoredDeliveryPlan> pub radroots_storage::authored_atomic::AuthoredAtomicOutcome::Prepared::operation: radroots_storage::authored::AuthoredOperation +pub radroots_storage::authored_atomic::AuthoredAtomicOutcome::Submitted(alloc::boxed::Box<radroots_storage::authored_draft_submission::PrepareFromDraft>) pub enum radroots_storage::authored_atomic::AuthoredWorkTarget pub radroots_storage::authored_atomic::AuthoredWorkTarget::Artifact(radroots_storage::authored::AuthoredArtifactId) pub radroots_storage::authored_atomic::AuthoredWorkTarget::DeliveryPlan(radroots_storage::authored_delivery::AuthoredDeliveryPlanId) @@ -261,6 +263,7 @@ pub const fn radroots_storage::authored_atomic::AuthoredAtomicReceipt::committed pub const fn radroots_storage::authored_atomic::AuthoredAtomicReceipt::digest(&self) -> radroots_storage::atomic::AtomicCommitDigest pub const fn radroots_storage::authored_atomic::AuthoredAtomicReceipt::disposition(&self) -> radroots_storage::atomic::AtomicCommitDisposition pub fn radroots_storage::authored_atomic::AuthoredAtomicReceipt::from_durable_parts(radroots_storage::atomic::AtomicCommitId, radroots_storage::atomic::AtomicCommitDigest, radroots_storage::atomic::AtomicCommitDisposition, u64, radroots_storage::authored_atomic::AuthoredAtomicOutcome) -> core::result::Result<Self, radroots_storage::Error> +pub fn radroots_storage::authored_atomic::AuthoredAtomicReceipt::matches_command(&self, &radroots_storage::authored_atomic::AuthoredAtomicCommand) -> bool pub fn radroots_storage::authored_atomic::AuthoredAtomicReceipt::new(&radroots_storage::authored_atomic::AuthoredAtomicCommand, radroots_storage::atomic::AtomicCommitDisposition, u64, radroots_storage::authored_atomic::AuthoredAtomicOutcome) -> core::result::Result<Self, radroots_storage::Error> pub const fn radroots_storage::authored_atomic::AuthoredAtomicReceipt::outcome(&self) -> &radroots_storage::authored_atomic::AuthoredAtomicOutcome pub struct radroots_storage::authored_atomic::CancelAuthoredWork @@ -390,11 +393,13 @@ pub fn radroots_storage::authored_draft::AuthoredDraft::payload_schema(&self) -> pub const fn radroots_storage::authored_draft::AuthoredDraft::payload_sha256(&self) -> &[u8; 32] pub fn radroots_storage::authored_draft::AuthoredDraft::reconstruct(radroots_storage::authored_draft::AuthoredDraftId, radroots_storage::authored_draft::AuthoredDraftRevision, [u8; 32], impl core::convert::Into<alloc::string::String>, alloc::vec::Vec<u8>, [u8; 32], radroots_storage::authored_draft::AuthoredDraftStage, core::option::Option<radroots_storage::journal::OperationInstanceId>, u64, u64) -> core::result::Result<Self, radroots_storage::Error> pub const fn radroots_storage::authored_draft::AuthoredDraft::revision(&self) -> radroots_storage::authored_draft::AuthoredDraftRevision +pub const fn radroots_storage::authored_draft::AuthoredDraft::scope(&self) -> core::option::Option<radroots_storage::authored_draft_query::AuthoredDraftScope> pub const fn radroots_storage::authored_draft::AuthoredDraft::stage(&self) -> radroots_storage::authored_draft::AuthoredDraftStage pub fn radroots_storage::authored_draft::AuthoredDraft::successor(&self, alloc::vec::Vec<u8>, radroots_storage::authored_draft::AuthoredDraftStage, core::option::Option<radroots_storage::journal::OperationInstanceId>, u64) -> core::result::Result<Self, radroots_storage::Error> pub const fn radroots_storage::authored_draft::AuthoredDraft::updated_at_unix_ms(&self) -> u64 pub fn radroots_storage::authored_draft::AuthoredDraft::validate(&self) -> core::result::Result<(), radroots_storage::Error> pub fn radroots_storage::authored_draft::AuthoredDraft::validate_successor_of(&self, &Self) -> core::result::Result<(), radroots_storage::Error> +pub fn radroots_storage::authored_draft::AuthoredDraft::with_scope(self, radroots_storage::authored_draft_query::AuthoredDraftScope) -> core::result::Result<Self, radroots_storage::Error> pub struct radroots_storage::authored_draft::AuthoredDraftId(_) impl radroots_storage::authored_draft::AuthoredDraftId pub const fn radroots_storage::authored_draft::AuthoredDraftId::as_bytes(&self) -> &[u8; 16] @@ -428,11 +433,78 @@ pub fn radroots_storage::authored_draft::AuthoredDraftStore::append_authored_dra pub fn radroots_storage::authored_draft::AuthoredDraftStore::authored_draft_head(&self, radroots_storage::authored_draft::AuthoredDraftId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::authored_draft::AuthoredDraft>, radroots_storage::Error>> pub fn radroots_storage::authored_draft::AuthoredDraftStore::authored_draft_heads(&self, [u8; 32], u16) -> radroots_transport::source::BoxFuture<'_, core::result::Result<alloc::vec::Vec<radroots_storage::authored_draft::AuthoredDraft>, radroots_storage::Error>> pub fn radroots_storage::authored_draft::AuthoredDraftStore::authored_draft_revision(&self, radroots_storage::authored_draft::AuthoredDraftId, radroots_storage::authored_draft::AuthoredDraftRevision) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::authored_draft::AuthoredDraft>, radroots_storage::Error>> +pub fn radroots_storage::authored_draft::AuthoredDraftStore::query_authored_drafts(&self, radroots_storage::authored_draft_query::AuthoredDraftQuery) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_storage::authored_draft_query::AuthoredDraftPage, radroots_storage::Error>> impl radroots_storage::authored_draft::AuthoredDraftStore for radroots_storage::memory::MemoryStorage pub fn radroots_storage::memory::MemoryStorage::append_authored_draft(&self, radroots_storage::authored_draft::AuthoredDraft, core::option::Option<radroots_storage::authored_draft::AuthoredDraftRevision>) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_storage::authored_draft::DraftAppendReceipt, radroots_storage::Error>> pub fn radroots_storage::memory::MemoryStorage::authored_draft_head(&self, radroots_storage::authored_draft::AuthoredDraftId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::authored_draft::AuthoredDraft>, radroots_storage::Error>> pub fn radroots_storage::memory::MemoryStorage::authored_draft_heads(&self, [u8; 32], u16) -> radroots_transport::source::BoxFuture<'_, core::result::Result<alloc::vec::Vec<radroots_storage::authored_draft::AuthoredDraft>, radroots_storage::Error>> pub fn radroots_storage::memory::MemoryStorage::authored_draft_revision(&self, radroots_storage::authored_draft::AuthoredDraftId, radroots_storage::authored_draft::AuthoredDraftRevision) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::authored_draft::AuthoredDraft>, radroots_storage::Error>> +pub fn radroots_storage::memory::MemoryStorage::query_authored_drafts(&self, radroots_storage::authored_draft_query::AuthoredDraftQuery) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_storage::authored_draft_query::AuthoredDraftPage, radroots_storage::Error>> +pub mod radroots_storage::authored_draft_query +pub enum radroots_storage::authored_draft_query::AuthoredDraftQueryRecord +pub radroots_storage::authored_draft_query::AuthoredDraftQueryRecord::Corrupt +pub radroots_storage::authored_draft_query::AuthoredDraftQueryRecord::Corrupt::draft_key: [u8; 16] +pub radroots_storage::authored_draft_query::AuthoredDraftQueryRecord::Corrupt::revision: radroots_storage::authored_draft::AuthoredDraftRevision +pub radroots_storage::authored_draft_query::AuthoredDraftQueryRecord::Draft(radroots_storage::authored_draft::AuthoredDraft) +impl radroots_storage::authored_draft_query::AuthoredDraftQueryRecord +pub fn radroots_storage::authored_draft_query::AuthoredDraftQueryRecord::draft_id(&self) -> core::result::Result<radroots_storage::authored_draft::AuthoredDraftId, radroots_storage::Error> +pub fn radroots_storage::authored_draft_query::AuthoredDraftQueryRecord::draft_key(&self) -> [u8; 16] +pub const fn radroots_storage::authored_draft_query::AuthoredDraftQueryRecord::revision(&self) -> radroots_storage::authored_draft::AuthoredDraftRevision +impl core::fmt::Debug for radroots_storage::authored_draft_query::AuthoredDraftQueryRecord +pub fn radroots_storage::authored_draft_query::AuthoredDraftQueryRecord::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct radroots_storage::authored_draft_query::AuthoredDraftCursor +pub struct radroots_storage::authored_draft_query::AuthoredDraftPage +impl radroots_storage::authored_draft_query::AuthoredDraftPage +pub fn radroots_storage::authored_draft_query::AuthoredDraftPage::into_records(self) -> alloc::vec::Vec<radroots_storage::authored_draft_query::AuthoredDraftQueryRecord> +pub fn radroots_storage::authored_draft_query::AuthoredDraftPage::new(&radroots_storage::authored_draft_query::AuthoredDraftQuery, alloc::vec::Vec<radroots_storage::authored_draft_query::AuthoredDraftQueryRecord>, bool) -> core::result::Result<Self, radroots_storage::Error> +pub const fn radroots_storage::authored_draft_query::AuthoredDraftPage::next_cursor(&self) -> core::option::Option<&radroots_storage::authored_draft_query::AuthoredDraftCursor> +pub fn radroots_storage::authored_draft_query::AuthoredDraftPage::records(&self) -> &[radroots_storage::authored_draft_query::AuthoredDraftQueryRecord] +pub struct radroots_storage::authored_draft_query::AuthoredDraftQuery +impl radroots_storage::authored_draft_query::AuthoredDraftQuery +pub const fn radroots_storage::authored_draft_query::AuthoredDraftQuery::after(&self) -> core::option::Option<[u8; 16]> +pub const fn radroots_storage::authored_draft_query::AuthoredDraftQuery::author(&self) -> &[u8; 32] +pub fn radroots_storage::authored_draft_query::AuthoredDraftQuery::cursor_after(&self, [u8; 16]) -> radroots_storage::authored_draft_query::AuthoredDraftCursor +pub const fn radroots_storage::authored_draft_query::AuthoredDraftQuery::limit(&self) -> u16 +pub fn radroots_storage::authored_draft_query::AuthoredDraftQuery::matches(&self, &radroots_storage::authored_draft::AuthoredDraft) -> bool +pub fn radroots_storage::authored_draft_query::AuthoredDraftQuery::new([u8; 32], impl core::convert::AsRef<str>, core::option::Option<radroots_storage::authored_draft_query::AuthoredDraftScope>, u16) -> core::result::Result<Self, radroots_storage::Error> +pub fn radroots_storage::authored_draft_query::AuthoredDraftQuery::payload_schema(&self) -> &str +pub const fn radroots_storage::authored_draft_query::AuthoredDraftQuery::scope(&self) -> core::option::Option<radroots_storage::authored_draft_query::AuthoredDraftScope> +pub fn radroots_storage::authored_draft_query::AuthoredDraftQuery::with_cursor(self, &radroots_storage::authored_draft_query::AuthoredDraftCursor) -> core::result::Result<Self, radroots_storage::Error> +pub struct radroots_storage::authored_draft_query::AuthoredDraftScope(_) +impl radroots_storage::authored_draft_query::AuthoredDraftScope +pub const fn radroots_storage::authored_draft_query::AuthoredDraftScope::as_bytes(&self) -> &[u8; 32] +pub fn radroots_storage::authored_draft_query::AuthoredDraftScope::new([u8; 32]) -> core::result::Result<Self, radroots_storage::Error> +impl core::convert::From<radroots_storage::authored_draft_query::AuthoredDraftScope> for [u8; 32] +pub fn [u8; 32]::from(radroots_storage::authored_draft_query::AuthoredDraftScope) -> Self +impl core::convert::TryFrom<[u8; 32]> for radroots_storage::authored_draft_query::AuthoredDraftScope +pub type radroots_storage::authored_draft_query::AuthoredDraftScope::Error = radroots_storage::Error +pub fn radroots_storage::authored_draft_query::AuthoredDraftScope::try_from([u8; 32]) -> core::result::Result<Self, radroots_storage::Error> +pub const radroots_storage::authored_draft_query::AUTHORED_DRAFT_CURSOR_SCHEMA_VERSION: u16 +pub const radroots_storage::authored_draft_query::AUTHORED_DRAFT_PAGE_PAYLOAD_MAX_BYTES: usize +pub const radroots_storage::authored_draft_query::AUTHORED_DRAFT_PAGE_SNAPSHOT_MAX_BYTES: usize +pub mod radroots_storage::authored_draft_submission +pub struct radroots_storage::authored_draft_submission::AuthoredDraftSource +impl radroots_storage::authored_draft_submission::AuthoredDraftSource +pub const fn radroots_storage::authored_draft_submission::AuthoredDraftSource::author(&self) -> &[u8; 32] +pub fn radroots_storage::authored_draft_submission::AuthoredDraftSource::capture(&radroots_storage::authored_draft::AuthoredDraft) -> core::result::Result<Self, radroots_storage::Error> +pub const fn radroots_storage::authored_draft_submission::AuthoredDraftSource::draft_id(&self) -> radroots_storage::authored_draft::AuthoredDraftId +pub fn radroots_storage::authored_draft_submission::AuthoredDraftSource::matches(&self, &radroots_storage::authored_draft::AuthoredDraft) -> bool +pub fn radroots_storage::authored_draft_submission::AuthoredDraftSource::payload_schema(&self) -> &str +pub const fn radroots_storage::authored_draft_submission::AuthoredDraftSource::payload_sha256(&self) -> &[u8; 32] +pub const fn radroots_storage::authored_draft_submission::AuthoredDraftSource::revision(&self) -> radroots_storage::authored_draft::AuthoredDraftRevision +pub const fn radroots_storage::authored_draft_submission::AuthoredDraftSource::scope(&self) -> core::option::Option<radroots_storage::authored_draft_query::AuthoredDraftScope> +pub struct radroots_storage::authored_draft_submission::PrepareFromDraft +impl radroots_storage::authored_draft_submission::PrepareFromDraft +pub const fn radroots_storage::authored_draft_submission::PrepareFromDraft::command_id(&self) -> radroots_storage::atomic::AtomicCommitId +pub fn radroots_storage::authored_draft_submission::PrepareFromDraft::commit_id(&self) -> radroots_storage::atomic::AtomicCommitId +pub fn radroots_storage::authored_draft_submission::PrepareFromDraft::commit_id_for(&[u8; 32], radroots_storage::atomic::AtomicCommitId) -> radroots_storage::atomic::AtomicCommitId +pub const fn radroots_storage::authored_draft_submission::PrepareFromDraft::intent(&self) -> &radroots_storage::authored_draft::AuthoredDraft +pub fn radroots_storage::authored_draft_submission::PrepareFromDraft::new(radroots_storage::atomic::AtomicCommitId, radroots_storage::authored_draft_submission::AuthoredDraftSource, radroots_storage::authored_draft::AuthoredDraft, radroots_storage::authored_atomic::PrepareAuthoredOperation) -> core::result::Result<Self, radroots_storage::Error> +pub const fn radroots_storage::authored_draft_submission::PrepareFromDraft::preparation(&self) -> &radroots_storage::authored_atomic::PrepareAuthoredOperation +pub const fn radroots_storage::authored_draft_submission::PrepareFromDraft::source(&self) -> &radroots_storage::authored_draft_submission::AuthoredDraftSource +pub fn radroots_storage::authored_draft_submission::PrepareFromDraft::validate(&self) -> core::result::Result<(), radroots_storage::Error> +impl core::fmt::Debug for radroots_storage::authored_draft_submission::PrepareFromDraft +pub fn radroots_storage::authored_draft_submission::PrepareFromDraft::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result pub mod radroots_storage::backup pub enum radroots_storage::backup::BackupMemberKind pub radroots_storage::backup::BackupMemberKind::Metadata @@ -823,6 +895,7 @@ pub fn radroots_storage::memory::MemoryStorage::append_authored_draft(&self, rad pub fn radroots_storage::memory::MemoryStorage::authored_draft_head(&self, radroots_storage::authored_draft::AuthoredDraftId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::authored_draft::AuthoredDraft>, radroots_storage::Error>> pub fn radroots_storage::memory::MemoryStorage::authored_draft_heads(&self, [u8; 32], u16) -> radroots_transport::source::BoxFuture<'_, core::result::Result<alloc::vec::Vec<radroots_storage::authored_draft::AuthoredDraft>, radroots_storage::Error>> pub fn radroots_storage::memory::MemoryStorage::authored_draft_revision(&self, radroots_storage::authored_draft::AuthoredDraftId, radroots_storage::authored_draft::AuthoredDraftRevision) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::authored_draft::AuthoredDraft>, radroots_storage::Error>> +pub fn radroots_storage::memory::MemoryStorage::query_authored_drafts(&self, radroots_storage::authored_draft_query::AuthoredDraftQuery) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_storage::authored_draft_query::AuthoredDraftPage, radroots_storage::Error>> impl radroots_storage::backup::StorageReliability for radroots_storage::memory::MemoryStorage pub fn radroots_storage::memory::MemoryStorage::begin_backup(&self, radroots_storage::backup::BackupPlan) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_storage::backup::BackupOperation, radroots_storage::Error>> pub fn radroots_storage::memory::MemoryStorage::begin_restore(&self, radroots_storage::backup::RestorePlan) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_storage::backup::RestoreOperation, radroots_storage::Error>> diff --git a/contracts/api_baselines/radroots_sync.txt b/contracts/api_baselines/radroots_sync.txt @@ -201,6 +201,7 @@ pub const fn radroots_sync::push::PushPreparation::operation(&self) -> &radroots pub struct radroots_sync::push::PushRequest impl radroots_sync::push::PushRequest pub const fn radroots_sync::push::PushRequest::actor(&self) -> &radroots_signing::actor::Actor +pub fn radroots_sync::push::PushRequest::authored_preparation(&self, u64) -> core::result::Result<radroots_storage::authored_atomic::PrepareAuthoredOperation, radroots_sync::policy::Error> pub const fn radroots_sync::push::PushRequest::cancellation(&self) -> radroots_signing::request::CancellationPolicy pub const fn radroots_sync::push::PushRequest::delivery_deadline_unix_ms(&self) -> u64 pub const fn radroots_sync::push::PushRequest::idempotency_key(&self) -> &radroots_storage::journal::IdempotencyKey @@ -351,6 +352,7 @@ pub const fn radroots_sync::push::PushPreparation::operation(&self) -> &radroots pub struct radroots_sync::PushRequest impl radroots_sync::push::PushRequest pub const fn radroots_sync::push::PushRequest::actor(&self) -> &radroots_signing::actor::Actor +pub fn radroots_sync::push::PushRequest::authored_preparation(&self, u64) -> core::result::Result<radroots_storage::authored_atomic::PrepareAuthoredOperation, radroots_sync::policy::Error> pub const fn radroots_sync::push::PushRequest::cancellation(&self) -> radroots_signing::request::CancellationPolicy pub const fn radroots_sync::push::PushRequest::delivery_deadline_unix_ms(&self) -> u64 pub const fn radroots_sync::push::PushRequest::idempotency_key(&self) -> &radroots_storage::journal::IdempotencyKey diff --git a/contracts/architecture/decisions/authored_draft_submission.v1.json b/contracts/architecture/decisions/authored_draft_submission.v1.json @@ -0,0 +1,17 @@ +{ + "schema": "radroots.authored-draft-submission.v1", + "status": "implemented", + "owners": ["radroots_storage", "radroots_storage_sqlite", "radroots_sync"], + "source": "Capture the existing immutable draft ID, revision, author, payload schema, optional scope, payload digest, editable stage and timestamps. Initial submission requires the current head to match that capture. The source remains editable and is never rewritten by submission.", + "identity": "The stable submission lookup key hashes the fixed v1 domain, the 32-byte author account and the nonzero 16-byte caller command ID. Context/scope and semantic input are compared within that key; they never change its identity. Caller generation and request digest do not determine identity.", + "intent": "The caller owns a distinct initial ReadyToSign or Queued draft containing the complete immutable semantic request, including app actor/account context, media and target policy. Storage requires matching source author and scope, operation association and captured time. Every prepared artifact must be an initial unsent plan by that author. Storage does not interpret an application payload schema or authorize signing.", + "transaction": "One existing BEGIN IMMEDIATE transaction installs operation, artifacts, delivery plans, immutable intent, ordinary Prepare receipt and the composite submission receipt. The latter is the source-to-intent association. No second connection, generic SQL API, journal, signer or network callback is introduced. All intermediate failure and abandoned precommit work rolls back; success is returned only after commit.", + "replay": "Lookup the composite receipt before source-head CAS. Compare its validated full captured source, full intent bytes and all preparation fields independently of the caller input digest. Exact replay returns the original operation even after later edits; differing input under the same author/command conflicts. Ordinary Prepare identity/digest/replay behavior stays unchanged so Sync resumes signing using its existing path.", + "query": "Require exact independent author, payload schema and optional scope. Scan current heads by stable ascending draft ID with a version-1 context-bound continuation and at most 256 results. Each page uses and releases one read snapshot. A changed revision cannot move its ID; new IDs behind a cursor require a later sweep. This is not a frozen inventory across pages.", + "corruption": "Return typed per-record corruption locators, retaining source evidence and continuing after each ID. Known foreign schemas/scopes are excluded before snapshot loading. Historical records with unknown schema metadata expose only an author-bound opaque ID/revision locator across that author's contexts, never payload. Protected-data/backend unavailability is an operation error, not row corruption.", + "bounds": {"page_records": 256, "metadata_lookahead": 1, "decoded_page_payload_bytes": 4194304, "serialized_page_snapshot_bytes": 16777216, "existing_atomic_receipt_snapshot_bytes": 4194304}, + "migration": "Forward runtime schema 14 atomically adds generic schema/scope metadata and a query index. Safely backfill known legacy schema strings; malformed snapshots retain unknown metadata. Original snapshots and v1-v13 migration checksums remain unchanged. Failure restores prior schema/guards/user_version; a prior schema plan rejects the newer database. WAL/FULL and native SQLx SQLite linkage stay unchanged.", + "compatibility": "The pre-release storage API gains query/submission types, one required draft query method and command/outcome enum variants. Downstream exhaustive matches and draft-store adapters require the ordered producer adoption. Legacy unscoped draft JSON omits the optional scope and remains byte-identical. Existing Prepare IDs and wire bytes remain unchanged. Public package identity, dependency versions and verification thresholds stay unchanged.", + "verification": ["validated serde and redacted diagnostics", "Memory/SQLite exact and conflicting replay with unchanged supplied digest", "later-save replay and fresh stale-CAS rejection", "simultaneous save/submit and duplicate submissions", "every intermediate insert/receipt fault rolls back", "lost callback and reopened durable receipt", "migration forward and rollback with old-plan refusal", "full current affected profiles, API review, coverage and workspace/release preflight"], + "non_goals": ["application schema policy", "signing authorization", "new service host", "physical power-loss or device qualification", "signed release", "deferred platforms"] +} diff --git a/contracts/architecture/deviations.toml b/contracts/architecture/deviations.toml @@ -2,6 +2,27 @@ schema_version = 1 architecture_id = "radroots.crates.release.v1" [[deviation]] +id = "RCRV1-DEV-018" +date = "2026-09-10" +status = "closed" +approval = "Standing user authorization covers all necessary shared-owner prerequisites, verified commits and non-force publication through RCLD-TERA-100." +affected_steps = ["158", "163", "173", "178", "207"] +spec_anchors = ["contracts/crates/release_v1/radroots_crates_release_v1.toml#package.radroots_storage", "contracts/crates/release_v1/radroots_crates_release_v1.toml#package.radroots_storage_sqlite", "contracts/crates/release_v1/radroots_crates_release_v1.toml#package.radroots_sync"] +source_evidence = ["The current shared host has immutable draft append CAS and authored Prepare, but no typed transaction joining them.", "The unfiltered draft-head Vec can fail wholesale on a foreign or corrupt snapshot and offers no continuation."] +replacement_action = "Implement contracts/architecture/decisions/authored_draft_submission.v1.json within existing storage, SQLite and Sync ownership for C054/P054, before ordered consumer adoption." +verification = ["Qualify typed exact replay, atomic rollback/reopen, concurrent save/submit, bounded isolated pages and forward migration compatibility.", "Review public enum/trait additions and pass affected profiles, unchanged coverage, workspace and release preflight."] +unresolved_risk = "Current owner, SDK, API, unchanged coverage and complete workspace/release-preflight qualification pass. Application schema policy and ordered consumer adoption remain with the application; deferred platform, physical-device and signed-release qualification are not claimed." +normative_architecture_change = false +adr_required = false +closure_evidence = [ + "Memory and real SQLite qualify stable author/command identity, complete semantic comparison independent of caller digest, replay before CAS after later edits, distinct intentional submissions and ordinary Prepare resumption.", + "SQLite qualifies sixteen before/after record fault windows, abandoned precommit, actual COMMIT failure, capacity exhaustion, concurrent save/submit, lost callback, reopen, read-only and explicit close.", + "Bounded scoped pages isolate corrupt metadata/payloads, preserve legacy unscoped bytes and enforce count and byte budgets. The reference backend uses at most limit+1 scratch heads and scans1000 reversed records without losing revisions or continuation.", + "Runtime14 migrates atomically without rewriting original snapshots or historical checksums. The actual prior13ac binary rejects14 in read-only and writable modes with both database files unchanged.", + "Reviewed pre-release trait/enum additions and the pure Sync preparation helper pass current owner, SDK, full workspace, Rustdoc, contract, dependency, portable and release-preflight gates. All45 required coverage reports pass unchanged90% thresholds with explicit retained-measurement provenance.", +] + +[[deviation]] id = "RCRV1-DEV-017" date = "2026-09-10" status = "closed" diff --git a/crates/storage/README.md b/crates/storage/README.md @@ -111,6 +111,28 @@ Idempotency-key construction validates borrowed input before allocating its bounded owned representation, so rejected oversized input cannot force a second attacker-sized allocation at this public boundary. +## Draft queries and atomic submission + +`AuthoredDraftStore::query_authored_drafts` returns bounded current-head pages +under an independently selected author, payload schema and optional immutable +scope. Continue using the returned scoped cursor. Corrupt rows have individual +repair locators; applications retain their evidence and continue other work. +Pages scan stable draft IDs, so a later sweep must revisit new IDs inserted +behind the cursor. A page holds at most 256 records and 4 MiB of decoded payload. + +`AuthoredAtomicCommand::PrepareFromDraft` joins a captured source revision, +a distinct initial intent and the existing authored preparation in one commit. +The application puts the complete frozen semantic request in the intent payload +and owns strict payload validation. Reserve a stable command ID before effects. +The author and command identify the receipt; context, source, intent and every +preparation field are compared for exact replay. Replay precedes fresh source +CAS and survives later editing. The source remains editable. The stored +submission receipt retains the immutable source-to-intent association, and the +ordinary Prepare receipt allows existing signing orchestration to resume. + +No signing or transport effect occurs in this storage transaction. Public +backends must implement the new query method and submission command explicitly. + ## Protected metadata and security `private_artifact` stores bounded metadata and opaque durable secret references. diff --git a/crates/storage/src/authored_atomic.rs b/crates/storage/src/authored_atomic.rs @@ -14,6 +14,7 @@ use crate::{ WorkFailure, WorkPhase, }, authored_delivery::{AuthoredDeliveryPlan, AuthoredDeliveryPlanId, DeliveryAttemptOutcome}, + authored_draft_submission::PrepareFromDraft, journal::OperationInstanceId, }; @@ -50,6 +51,11 @@ impl WorkFence { } } +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[cfg_attr( + feature = "serde", + serde(try_from = "PrepareWire", into = "PrepareWire") +)] #[derive(Clone, Debug, Eq, PartialEq)] pub struct PrepareAuthoredOperation { operation: AuthoredOperation, @@ -59,6 +65,42 @@ pub struct PrepareAuthoredOperation { requested_at_unix_ms: u64, } +#[cfg(feature = "serde")] +#[derive(serde::Serialize, serde::Deserialize)] +#[serde(deny_unknown_fields)] +struct PrepareWire { + operation: AuthoredOperation, + artifacts: Vec<AuthoredArtifact>, + delivery_plans: Vec<AuthoredDeliveryPlan>, + input_digest: AtomicCommitDigest, + requested_at_unix_ms: u64, +} +#[cfg(feature = "serde")] +impl TryFrom<PrepareWire> for PrepareAuthoredOperation { + type Error = Error; + fn try_from(v: PrepareWire) -> Result<Self, Error> { + Self::new( + v.operation, + v.artifacts, + v.delivery_plans, + v.input_digest, + v.requested_at_unix_ms, + ) + } +} +#[cfg(feature = "serde")] +impl From<PrepareAuthoredOperation> for PrepareWire { + fn from(v: PrepareAuthoredOperation) -> Self { + Self { + operation: v.operation, + artifacts: v.artifacts, + delivery_plans: v.delivery_plans, + input_digest: v.input_digest, + requested_at_unix_ms: v.requested_at_unix_ms, + } + } +} + impl PrepareAuthoredOperation { pub fn new( operation: AuthoredOperation, @@ -372,6 +414,7 @@ impl CancelAuthoredWork { #[derive(Clone, Debug, Eq, PartialEq)] pub enum AuthoredAtomicCommand { Prepare(PrepareAuthoredOperation), + PrepareFromDraft(Box<PrepareFromDraft>), Claim(ClaimAuthoredWork), ApplySigned(ApplySignedArtifact), ApplyAdmission(ApplyAdmissionResult), @@ -382,6 +425,9 @@ pub enum AuthoredAtomicCommand { impl AuthoredAtomicCommand { pub fn commit_id(&self) -> AtomicCommitId { + if let Self::PrepareFromDraft(value) = self { + return value.commit_id(); + } let digest = self.digest(); let mut hasher = Sha256::new(); hash_field(&mut hasher, b"radroots.authored.atomic.id.v2"); @@ -403,6 +449,12 @@ impl AuthoredAtomicCommand { hash_field(&mut hasher, self.phase_bytes()); hash_field(&mut hasher, &self.target_bytes()); match self { + Self::PrepareFromDraft(value) => { + hash_field(&mut hasher, value.commit_id().as_bytes()); + hash_field(&mut hasher, value.source().payload_sha256()); + hash_field(&mut hasher, value.intent().payload_sha256()); + hash_field(&mut hasher, value.preparation().input_digest().as_bytes()); + } Self::Prepare(value) => hash_field(&mut hasher, value.input_digest.as_bytes()), Self::Claim(value) => { hash_field(&mut hasher, value.claim.token()); @@ -423,6 +475,7 @@ impl AuthoredAtomicCommand { pub const fn requested_at_unix_ms(&self) -> u64 { match self { + Self::PrepareFromDraft(value) => value.preparation().requested_at_unix_ms(), Self::Prepare(value) => value.requested_at_unix_ms, Self::Claim(value) => value.claim.acquired_at_unix_ms(), Self::ApplySigned(value) => value.applied_at_unix_ms, @@ -435,6 +488,7 @@ impl AuthoredAtomicCommand { fn phase_bytes(&self) -> &'static [u8] { match self { + Self::PrepareFromDraft(_) => b"draft_submission_v1", Self::Prepare(_) => b"prepare", Self::Claim(_) => b"claim", Self::ApplySigned(_) => b"signing", @@ -451,6 +505,9 @@ impl AuthoredAtomicCommand { fn target_bytes(&self) -> [u8; 16] { match self { + Self::PrepareFromDraft(value) => { + *value.preparation().operation().operation_id().as_bytes() + } Self::Prepare(value) => *value.operation.operation_id().as_bytes(), Self::Claim(value) => match &value.target { ClaimAuthoredTarget::ArtifactSigning(id) @@ -479,7 +536,7 @@ impl AuthoredAtomicCommand { Self::ApplyAdmission(value) => Some(value.fence.generation), Self::ApplyDelivery(value) => Some(value.fence.generation), Self::ApplyFailure(value) => Some(value.fence.generation), - Self::Prepare(_) | Self::Cancel(_) => None, + Self::Prepare(_) | Self::PrepareFromDraft(_) | Self::Cancel(_) => None, } } } @@ -488,6 +545,7 @@ impl AuthoredAtomicCommand { #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] #[derive(Clone, Debug, Eq, PartialEq)] pub enum AuthoredAtomicOutcome { + Submitted(Box<PrepareFromDraft>), Prepared { operation: AuthoredOperation, artifacts: Vec<AuthoredArtifact>, @@ -507,6 +565,22 @@ pub struct AuthoredAtomicReceipt { } impl AuthoredAtomicReceipt { + /// Exact submission replay compares the entire validated captured request. + pub fn matches_command(&self, command: &AuthoredAtomicCommand) -> bool { + if self.commit_id != command.commit_id() || self.digest != command.digest() { + return false; + } + match (command, &self.outcome) { + ( + AuthoredAtomicCommand::PrepareFromDraft(request), + AuthoredAtomicOutcome::Submitted(committed), + ) => request == committed, + (AuthoredAtomicCommand::PrepareFromDraft(_), _) + | (_, AuthoredAtomicOutcome::Submitted(_)) => false, + _ => true, + } + } + pub fn new( command: &AuthoredAtomicCommand, disposition: AtomicCommitDisposition, @@ -516,6 +590,17 @@ impl AuthoredAtomicReceipt { if committed_at_unix_ms < command.requested_at_unix_ms() { return Err(Error::AtomicWorkflowMismatch); } + match (command, &outcome) { + ( + AuthoredAtomicCommand::PrepareFromDraft(request), + AuthoredAtomicOutcome::Submitted(value), + ) if request == value => value.validate()?, + (AuthoredAtomicCommand::PrepareFromDraft(_), _) + | (_, AuthoredAtomicOutcome::Submitted(_)) => { + return Err(Error::AtomicWorkflowMismatch); + } + _ => {} + } Ok(Self { commit_id: command.commit_id(), digest: command.digest(), @@ -535,6 +620,15 @@ impl AuthoredAtomicReceipt { if committed_at_unix_ms == 0 || !outcome.is_valid() { return Err(Error::AtomicWorkflowMismatch); } + if let AuthoredAtomicOutcome::Submitted(value) = &outcome { + let command = AuthoredAtomicCommand::PrepareFromDraft(value.clone()); + if commit_id != command.commit_id() + || digest != command.digest() + || committed_at_unix_ms < command.requested_at_unix_ms() + { + return Err(Error::AtomicWorkflowMismatch); + } + } Ok(Self { commit_id, digest, @@ -582,6 +676,7 @@ impl AuthoredAtomicOutcome { && plan.validate().is_ok() }) } + Self::Submitted(value) => value.validate().is_ok(), Self::Artifact(artifact) => artifact.validate().is_ok(), Self::DeliveryPlan(plan) => plan.validate().is_ok(), } diff --git a/crates/storage/src/authored_draft.rs b/crates/storage/src/authored_draft.rs @@ -5,7 +5,11 @@ use radroots_transport::BoxFuture; use sha2::{Digest, Sha256}; use std::{string::String, vec::Vec}; -use crate::{Error, journal::OperationInstanceId}; +use crate::{ + Error, + authored_draft_query::{AuthoredDraftPage, AuthoredDraftQuery, AuthoredDraftScope}, + journal::OperationInstanceId, +}; pub const AUTHORED_DRAFT_PAYLOAD_MAX_BYTES: usize = 4 * 1024 * 1024; pub const AUTHORED_DRAFT_SCHEMA_MAX_BYTES: usize = 128; @@ -119,6 +123,7 @@ pub struct AuthoredDraft { revision: AuthoredDraftRevision, author: [u8; 32], payload_schema: String, + scope: Option<AuthoredDraftScope>, payload: Vec<u8>, payload_sha256: [u8; 32], stage: AuthoredDraftStage, @@ -134,6 +139,8 @@ struct AuthoredDraftRevisionWire { revision: AuthoredDraftRevision, author: [u8; 32], payload_schema: String, + #[serde(default, skip_serializing_if = "Option::is_none")] + scope: Option<AuthoredDraftScope>, payload: Vec<u8>, payload_sha256: [u8; 32], stage: AuthoredDraftStage, @@ -147,11 +154,12 @@ impl TryFrom<AuthoredDraftRevisionWire> for AuthoredDraft { type Error = Error; fn try_from(value: AuthoredDraftRevisionWire) -> Result<Self, Self::Error> { - Self::reconstruct( + Self::reconstruct_scoped( value.draft_id, value.revision, value.author, value.payload_schema, + value.scope, value.payload, value.payload_sha256, value.stage, @@ -170,6 +178,7 @@ impl From<AuthoredDraft> for AuthoredDraftRevisionWire { revision: value.revision, author: value.author, payload_schema: value.payload_schema, + scope: value.scope, payload: value.payload, payload_sha256: value.payload_sha256, stage: value.stage, @@ -222,11 +231,12 @@ impl AuthoredDraft { updated_at_unix_ms: u64, ) -> Result<Self, Error> { let payload_sha256 = Sha256::digest(payload.as_slice()).into(); - let next = Self::reconstruct( + let next = Self::reconstruct_scoped( self.draft_id, self.revision.next()?, self.author, self.payload_schema.clone(), + self.scope, payload, payload_sha256, stage, @@ -251,11 +261,41 @@ impl AuthoredDraft { created_at_unix_ms: u64, updated_at_unix_ms: u64, ) -> Result<Self, Error> { + Self::reconstruct_scoped( + draft_id, + revision, + author, + payload_schema, + None, + payload, + payload_sha256, + stage, + operation_id, + created_at_unix_ms, + updated_at_unix_ms, + ) + } + + #[allow(clippy::too_many_arguments)] + fn reconstruct_scoped( + draft_id: AuthoredDraftId, + revision: AuthoredDraftRevision, + author: [u8; 32], + payload_schema: impl Into<String>, + scope: Option<AuthoredDraftScope>, + 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(), + scope, payload, payload_sha256, stage, @@ -297,6 +337,7 @@ impl AuthoredDraft { let identity_matches = self.draft_id == previous.draft_id && self.author == previous.author && self.payload_schema == previous.payload_schema + && self.scope == previous.scope && 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; @@ -364,6 +405,18 @@ impl AuthoredDraft { pub fn payload_schema(&self) -> &str { self.payload_schema.as_str() } + /// Selects an immutable scope while constructing an initial local record. + pub fn with_scope(mut self, scope: AuthoredDraftScope) -> Result<Self, Error> { + if self.revision != AuthoredDraftRevision::INITIAL || self.scope.is_some() { + return Err(Error::InvalidAuthoredDraft); + } + self.scope = Some(scope); + Ok(self) + } + pub const fn scope(&self) -> Option<AuthoredDraftScope> { + self.scope + } + pub fn payload(&self) -> &[u8] { self.payload.as_slice() } @@ -409,6 +462,11 @@ impl DraftAppendReceipt { } pub trait AuthoredDraftStore: Send + Sync { + fn query_authored_drafts( + &self, + query: AuthoredDraftQuery, + ) -> BoxFuture<'_, Result<AuthoredDraftPage, Error>>; + fn append_authored_draft( &self, draft: AuthoredDraft, diff --git a/crates/storage/src/authored_draft_query.rs b/crates/storage/src/authored_draft_query.rs @@ -0,0 +1,265 @@ +//! Bounded, independently scoped pages of current opaque draft revisions. +use crate::{ + Error, + authored_draft::{ + AUTHORED_DRAFT_QUERY_LIMIT_MAX, AUTHORED_DRAFT_SCHEMA_MAX_BYTES, AuthoredDraft, + AuthoredDraftId, AuthoredDraftRevision, + }, +}; + +/// Maximum decoded payload bytes retained by one query page. +pub const AUTHORED_DRAFT_PAGE_PAYLOAD_MAX_BYTES: usize = 4 * 1024 * 1024; +/// Maximum serialized snapshot bytes read by one native query page. +pub const AUTHORED_DRAFT_PAGE_SNAPSHOT_MAX_BYTES: usize = 16 * 1024 * 1024; +pub const AUTHORED_DRAFT_CURSOR_SCHEMA_VERSION: u16 = 1; + +/// An opaque, stable application-selected scope digest; never a credential. +#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[cfg_attr(feature = "serde", serde(try_from = "[u8; 32]", into = "[u8; 32]"))] +pub struct AuthoredDraftScope([u8; 32]); + +impl AuthoredDraftScope { + pub fn new(bytes: [u8; 32]) -> Result<Self, Error> { + if bytes.iter().all(|byte| *byte == 0) { + return Err(Error::InvalidAuthoredDraft); + } + Ok(Self(bytes)) + } + pub const fn as_bytes(&self) -> &[u8; 32] { + &self.0 + } +} +impl TryFrom<[u8; 32]> for AuthoredDraftScope { + type Error = Error; + fn try_from(value: [u8; 32]) -> Result<Self, Error> { + Self::new(value) + } +} +impl From<AuthoredDraftScope> for [u8; 32] { + fn from(value: AuthoredDraftScope) -> Self { + value.0 + } +} + +/// Scope is selected independently on every call, including continuations. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct AuthoredDraftQuery { + author: [u8; 32], + payload_schema: String, + scope: Option<AuthoredDraftScope>, + limit: u16, + after: Option<[u8; 16]>, +} +impl AuthoredDraftQuery { + pub fn new( + author: [u8; 32], + payload_schema: impl AsRef<str>, + scope: Option<AuthoredDraftScope>, + limit: u16, + ) -> Result<Self, Error> { + let schema = payload_schema.as_ref(); + if author.iter().all(|byte| *byte == 0) + || schema.is_empty() + || schema.len() > AUTHORED_DRAFT_SCHEMA_MAX_BYTES + || schema != schema.trim() + || schema.chars().any(char::is_control) + || limit == 0 + || limit > AUTHORED_DRAFT_QUERY_LIMIT_MAX + { + return Err(Error::InvalidAuthoredDraft); + } + Ok(Self { + author, + payload_schema: schema.to_owned(), + scope, + limit, + after: None, + }) + } + pub fn with_cursor(mut self, cursor: &AuthoredDraftCursor) -> Result<Self, Error> { + if cursor.schema_version != AUTHORED_DRAFT_CURSOR_SCHEMA_VERSION + || cursor.author != self.author + || cursor.payload_schema != self.payload_schema + || cursor.scope != self.scope + { + return Err(Error::InvalidAuthoredDraft); + } + self.after = Some(cursor.after_id); + Ok(self) + } + pub const fn author(&self) -> &[u8; 32] { + &self.author + } + pub fn payload_schema(&self) -> &str { + &self.payload_schema + } + pub const fn scope(&self) -> Option<AuthoredDraftScope> { + self.scope + } + pub const fn limit(&self) -> u16 { + self.limit + } + pub const fn after(&self) -> Option<[u8; 16]> { + self.after + } + pub fn matches(&self, draft: &AuthoredDraft) -> bool { + draft.author() == &self.author + && draft.payload_schema() == self.payload_schema + && draft.scope() == self.scope + } + pub fn cursor_after(&self, after_id: [u8; 16]) -> AuthoredDraftCursor { + AuthoredDraftCursor { + schema_version: AUTHORED_DRAFT_CURSOR_SCHEMA_VERSION, + author: self.author, + payload_schema: self.payload_schema.clone(), + scope: self.scope, + after_id, + } + } +} + +/// Stable ID ordering does not move when a draft receives another revision. +/// This is a bounded scan, not an immutable cross-page database snapshot. +#[derive(Clone, Debug, Eq, PartialEq)] +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[cfg_attr(feature = "serde", serde(try_from = "CursorWire", into = "CursorWire"))] +pub struct AuthoredDraftCursor { + schema_version: u16, + author: [u8; 32], + payload_schema: String, + scope: Option<AuthoredDraftScope>, + after_id: [u8; 16], +} +#[cfg(feature = "serde")] +#[derive(serde::Serialize, serde::Deserialize)] +#[serde(deny_unknown_fields)] +struct CursorWire { + schema_version: u16, + author: [u8; 32], + payload_schema: String, + scope: Option<AuthoredDraftScope>, + after_id: [u8; 16], +} +#[cfg(feature = "serde")] +impl TryFrom<CursorWire> for AuthoredDraftCursor { + type Error = Error; + fn try_from(value: CursorWire) -> Result<Self, Error> { + if value.schema_version != AUTHORED_DRAFT_CURSOR_SCHEMA_VERSION { + return Err(Error::InvalidAuthoredDraft); + } + Ok( + AuthoredDraftQuery::new(value.author, value.payload_schema, value.scope, 1)? + .cursor_after(value.after_id), + ) + } +} +#[cfg(feature = "serde")] +impl From<AuthoredDraftCursor> for CursorWire { + fn from(value: AuthoredDraftCursor) -> Self { + Self { + schema_version: value.schema_version, + author: value.author, + payload_schema: value.payload_schema, + scope: value.scope, + after_id: value.after_id, + } + } +} + +/// Corrupt rows expose only their local position, never untrusted payloads. +/// A raw position can identify a malformed all-zero draft ID for repair. +#[derive(Clone, Eq, PartialEq)] +pub enum AuthoredDraftQueryRecord { + Draft(AuthoredDraft), + Corrupt { + draft_key: [u8; 16], + revision: AuthoredDraftRevision, + }, +} +impl core::fmt::Debug for AuthoredDraftQueryRecord { + fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result { + f.debug_struct(match self { + Self::Draft(_) => "Draft", + Self::Corrupt { .. } => "Corrupt", + }) + .field("draft_key", &self.draft_key()) + .field("revision", &self.revision()) + .finish_non_exhaustive() + } +} + +impl AuthoredDraftQueryRecord { + pub fn draft_key(&self) -> [u8; 16] { + match self { + Self::Draft(draft) => *draft.draft_id().as_bytes(), + Self::Corrupt { draft_key, .. } => *draft_key, + } + } + pub const fn revision(&self) -> AuthoredDraftRevision { + match self { + Self::Draft(draft) => draft.revision(), + Self::Corrupt { revision, .. } => *revision, + } + } + pub fn draft_id(&self) -> Result<AuthoredDraftId, Error> { + AuthoredDraftId::new(self.draft_key()) + } +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct AuthoredDraftPage { + records: Vec<AuthoredDraftQueryRecord>, + next_cursor: Option<AuthoredDraftCursor>, +} +impl AuthoredDraftPage { + /// Backends return an ordered, bounded page after applying independent scope. + pub fn new( + query: &AuthoredDraftQuery, + records: Vec<AuthoredDraftQueryRecord>, + has_more: bool, + ) -> Result<Self, Error> { + if records.len() > usize::from(query.limit) || (has_more && records.is_empty()) { + return Err(Error::InvalidAuthoredDraft); + } + let mut previous = query.after; + let mut payload_bytes = 0usize; + for record in &records { + let key = record.draft_key(); + if previous.is_some_and(|previous| key <= previous) { + return Err(Error::InvalidAuthoredDraft); + } + previous = Some(key); + if let AuthoredDraftQueryRecord::Draft(draft) = record { + if !query.matches(draft) { + return Err(Error::InvalidAuthoredDraft); + } + draft.validate()?; + payload_bytes = payload_bytes + .checked_add(draft.payload().len()) + .ok_or(Error::InvalidAuthoredDraft)?; + if payload_bytes > AUTHORED_DRAFT_PAGE_PAYLOAD_MAX_BYTES { + return Err(Error::InvalidAuthoredDraft); + } + } + } + let next_cursor = if has_more { + previous.map(|key| query.cursor_after(key)) + } else { + None + }; + Ok(Self { + records, + next_cursor, + }) + } + pub fn records(&self) -> &[AuthoredDraftQueryRecord] { + &self.records + } + pub const fn next_cursor(&self) -> Option<&AuthoredDraftCursor> { + self.next_cursor.as_ref() + } + pub fn into_records(self) -> Vec<AuthoredDraftQueryRecord> { + self.records + } +} diff --git a/crates/storage/src/authored_draft_submission.rs b/crates/storage/src/authored_draft_submission.rs @@ -0,0 +1,279 @@ +//! Atomic association of an immutable captured draft with an authored preparation. +//! +//! The application owns the complete semantic request in the intent payload. +//! Storage compares that payload and every preparation field on replay. A command +//! key is scoped by the stable author, never by a process generation or a digest. +use crate::{ + Error, + atomic::AtomicCommitId, + authored::{ArtifactOrigin, SigningState}, + authored_atomic::PrepareAuthoredOperation, + authored_draft::{AuthoredDraft, AuthoredDraftId, AuthoredDraftRevision, AuthoredDraftStage}, + authored_draft_query::{AuthoredDraftQuery, AuthoredDraftScope}, +}; +use sha2::{Digest, Sha256}; + +/// Captured source metadata; contains no application payload or credentials. +#[derive(Clone, Debug, Eq, PartialEq)] +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[cfg_attr(feature = "serde", serde(try_from = "SourceWire", into = "SourceWire"))] +pub struct AuthoredDraftSource { + draft_id: AuthoredDraftId, + revision: AuthoredDraftRevision, + author: [u8; 32], + payload_schema: String, + scope: Option<AuthoredDraftScope>, + payload_sha256: [u8; 32], + stage: AuthoredDraftStage, + created_at_unix_ms: u64, + updated_at_unix_ms: u64, +} +impl AuthoredDraftSource { + pub fn capture(draft: &AuthoredDraft) -> Result<Self, Error> { + draft.validate()?; + let value = Self { + draft_id: draft.draft_id(), + revision: draft.revision(), + author: *draft.author(), + payload_schema: draft.payload_schema().to_owned(), + scope: draft.scope(), + payload_sha256: *draft.payload_sha256(), + stage: draft.stage(), + created_at_unix_ms: draft.created_at_unix_ms(), + updated_at_unix_ms: draft.updated_at_unix_ms(), + }; + value.validate()?; + Ok(value) + } + fn validate(&self) -> Result<(), Error> { + AuthoredDraftQuery::new(self.author, &self.payload_schema, self.scope, 1)?; + if !matches!( + self.stage, + AuthoredDraftStage::Draft + | AuthoredDraftStage::MediaPreparing + | AuthoredDraftStage::MediaUploading + ) || self.created_at_unix_ms == 0 + || self.updated_at_unix_ms < self.created_at_unix_ms + { + return Err(Error::InvalidAuthoredDraft); + } + Ok(()) + } + pub fn matches(&self, draft: &AuthoredDraft) -> bool { + Self::capture(draft).is_ok_and(|captured| captured == *self) + } + 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 + } + pub const fn scope(&self) -> Option<AuthoredDraftScope> { + self.scope + } + pub const fn payload_sha256(&self) -> &[u8; 32] { + &self.payload_sha256 + } +} + +/// A captured request and its distinct immutable intent, installed without effects. +#[derive(Clone, Eq, PartialEq)] +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[cfg_attr( + feature = "serde", + serde(try_from = "SubmissionWire", into = "SubmissionWire") +)] +pub struct PrepareFromDraft { + command_id: AtomicCommitId, + source: AuthoredDraftSource, + intent: AuthoredDraft, + preparation: PrepareAuthoredOperation, +} +impl PrepareFromDraft { + pub fn new( + command_id: AtomicCommitId, + source: AuthoredDraftSource, + intent: AuthoredDraft, + preparation: PrepareAuthoredOperation, + ) -> Result<Self, Error> { + let value = Self { + command_id, + source, + intent, + preparation, + }; + value.validate()?; + Ok(value) + } + pub fn validate(&self) -> Result<(), Error> { + AtomicCommitId::new(*self.command_id.as_bytes())?; + self.source.validate()?; + self.intent.validate()?; + let operation = self.preparation.operation(); + let at = self.preparation.requested_at_unix_ms(); + if self.intent.draft_id() == self.source.draft_id + || self.intent.author() != &self.source.author + || self.intent.scope() != self.source.scope + || self.intent.revision() != AuthoredDraftRevision::INITIAL + || !matches!( + self.intent.stage(), + AuthoredDraftStage::ReadyToSign | AuthoredDraftStage::Queued + ) + || self.intent.operation_id() != Some(operation.operation_id()) + || self.intent.created_at_unix_ms() != at + || self.intent.updated_at_unix_ms() != at + || at < self.source.updated_at_unix_ms + || operation.created_at_unix_ms() != at + || operation.updated_at_unix_ms() != at + || operation.revision().get() != 1 + { + return Err(Error::AtomicWorkflowMismatch); + } + for artifact in self.preparation.artifacts() { + if artifact.origin() != ArtifactOrigin::Planned + || artifact.signing_state() != SigningState::Planned + || artifact.revision().get() != 1 + || artifact.created_at_unix_ms() != at + || artifact.updated_at_unix_ms() != at + || artifact + .plan() + .ok_or(Error::AtomicWorkflowMismatch)? + .decode()? + .plan() + .author() + .as_bytes() + != &self.source.author + { + return Err(Error::AtomicWorkflowMismatch); + } + } + if self.preparation.delivery_plans().iter().any(|plan| { + plan.revision().get() != 1 + || plan.created_at_unix_ms() != at + || plan.updated_at_unix_ms() != at + }) { + return Err(Error::AtomicWorkflowMismatch); + } + Ok(()) + } + /// Stable lookup key, available before persistence and after restart. + pub fn commit_id_for(author: &[u8; 32], command_id: AtomicCommitId) -> AtomicCommitId { + let mut hash = Sha256::new(); + hash.update(b"radroots.authored.draft.submission.id.v1\0"); + hash.update(author); + hash.update(command_id.as_bytes()); + let digest = hash.finalize(); + let mut id = [0; 16]; + id.copy_from_slice(&digest[..16]); + AtomicCommitId::new(id).expect("SHA-256 derived submission identity is nonzero") + } + pub fn commit_id(&self) -> AtomicCommitId { + Self::commit_id_for(&self.source.author, self.command_id) + } + pub const fn command_id(&self) -> AtomicCommitId { + self.command_id + } + pub const fn source(&self) -> &AuthoredDraftSource { + &self.source + } + pub const fn intent(&self) -> &AuthoredDraft { + &self.intent + } + pub const fn preparation(&self) -> &PrepareAuthoredOperation { + &self.preparation + } +} +impl core::fmt::Debug for PrepareFromDraft { + fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result { + formatter + .debug_struct("PrepareFromDraft") + .field("command_id", &self.command_id) + .field("source", &self.source) + .field("intent_id", &self.intent.draft_id()) + .field("operation_id", &self.preparation.operation().operation_id()) + .finish_non_exhaustive() + } +} + +#[cfg(feature = "serde")] +#[derive(serde::Serialize, serde::Deserialize)] +#[serde(deny_unknown_fields)] +struct SourceWire { + draft_id: AuthoredDraftId, + revision: AuthoredDraftRevision, + author: [u8; 32], + payload_schema: String, + scope: Option<AuthoredDraftScope>, + payload_sha256: [u8; 32], + stage: AuthoredDraftStage, + created_at_unix_ms: u64, + updated_at_unix_ms: u64, +} +#[cfg(feature = "serde")] +impl TryFrom<SourceWire> for AuthoredDraftSource { + type Error = Error; + fn try_from(v: SourceWire) -> Result<Self, Error> { + let value = Self { + draft_id: v.draft_id, + revision: v.revision, + author: v.author, + payload_schema: v.payload_schema, + scope: v.scope, + payload_sha256: v.payload_sha256, + stage: v.stage, + created_at_unix_ms: v.created_at_unix_ms, + updated_at_unix_ms: v.updated_at_unix_ms, + }; + value.validate()?; + Ok(value) + } +} +#[cfg(feature = "serde")] +impl From<AuthoredDraftSource> for SourceWire { + fn from(v: AuthoredDraftSource) -> Self { + Self { + draft_id: v.draft_id, + revision: v.revision, + author: v.author, + payload_schema: v.payload_schema, + scope: v.scope, + payload_sha256: v.payload_sha256, + stage: v.stage, + created_at_unix_ms: v.created_at_unix_ms, + updated_at_unix_ms: v.updated_at_unix_ms, + } + } +} +#[cfg(feature = "serde")] +#[derive(serde::Serialize, serde::Deserialize)] +#[serde(deny_unknown_fields)] +struct SubmissionWire { + command_id: AtomicCommitId, + source: AuthoredDraftSource, + intent: AuthoredDraft, + preparation: PrepareAuthoredOperation, +} +#[cfg(feature = "serde")] +impl TryFrom<SubmissionWire> for PrepareFromDraft { + type Error = Error; + fn try_from(v: SubmissionWire) -> Result<Self, Error> { + Self::new(v.command_id, v.source, v.intent, v.preparation) + } +} +#[cfg(feature = "serde")] +impl From<PrepareFromDraft> for SubmissionWire { + fn from(v: PrepareFromDraft) -> Self { + Self { + command_id: v.command_id, + source: v.source, + intent: v.intent, + preparation: v.preparation, + } + } +} diff --git a/crates/storage/src/lib.rs b/crates/storage/src/lib.rs @@ -6,6 +6,8 @@ pub mod authored; pub mod authored_atomic; pub mod authored_delivery; pub mod authored_draft; +pub mod authored_draft_query; +pub mod authored_draft_submission; pub mod backup; mod error; pub mod event; diff --git a/crates/storage/src/memory.rs b/crates/storage/src/memory.rs @@ -1417,6 +1417,45 @@ impl AtomicStorage for MemoryStorage { } } +fn prepare_authored_memory( + candidate: &mut State, + value: crate::authored_atomic::PrepareAuthoredOperation, +) -> Result<AuthoredAtomicOutcome, Error> { + if candidate + .authored_operations + .iter() + .any(|operation| operation.operation_id() == value.operation().operation_id()) + || value.artifacts().iter().any(|artifact| { + candidate + .authored_artifacts + .iter() + .any(|existing| existing.artifact_id() == artifact.artifact_id()) + }) + || value.delivery_plans().iter().any(|plan| { + candidate + .authored_delivery_plans + .iter() + .any(|existing| existing.plan_id() == plan.plan_id()) + }) + { + return Err(Error::AtomicCommitConflict); + } + candidate + .authored_operations + .push(value.operation().clone()); + candidate + .authored_artifacts + .extend(value.artifacts().iter().cloned()); + candidate + .authored_delivery_plans + .extend(value.delivery_plans().iter().cloned()); + Ok(AuthoredAtomicOutcome::Prepared { + operation: value.operation().clone(), + artifacts: value.artifacts().to_vec(), + delivery_plans: value.delivery_plans().to_vec(), + }) +} + impl AuthoredAtomicStorage for MemoryStorage { fn execute_authored( &self, @@ -1429,7 +1468,7 @@ impl AuthoredAtomicStorage for MemoryStorage { .iter() .find(|receipt| receipt.commit_id() == command.commit_id()) { - if existing.digest() != command.digest() { + if !existing.matches_command(&command) { return Err(Error::AtomicCommitConflict); } return AuthoredAtomicReceipt::from_durable_parts( @@ -1444,35 +1483,45 @@ impl AuthoredAtomicStorage for MemoryStorage { let mut candidate = state.clone(); let outcome = match command.clone() { AuthoredAtomicCommand::Prepare(value) => { - if candidate.authored_operations.iter().any(|operation| { - operation.operation_id() == value.operation().operation_id() - }) || value.artifacts().iter().any(|artifact| { - candidate - .authored_artifacts - .iter() - .any(|existing| existing.artifact_id() == artifact.artifact_id()) - }) || value.delivery_plans().iter().any(|plan| { - candidate - .authored_delivery_plans - .iter() - .any(|existing| existing.plan_id() == plan.plan_id()) - }) { + prepare_authored_memory(&mut candidate, value)? + } + AuthoredAtomicCommand::PrepareFromDraft(value) => { + value.validate()?; + let head = candidate + .authored_drafts + .iter() + .filter(|draft| draft.draft_id() == value.source().draft_id()) + .max_by_key(|draft| draft.revision()); + if !head.is_some_and(|head| value.source().matches(head)) { + return Err(Error::DraftRevisionConflict); + } + if candidate + .authored_drafts + .iter() + .any(|draft| draft.draft_id() == value.intent().draft_id()) + { + return Err(Error::DraftRevisionConflict); + } + let ordinary = AuthoredAtomicCommand::Prepare(value.preparation().clone()); + if candidate + .authored_atomic_receipts + .iter() + .any(|receipt| receipt.commit_id() == ordinary.commit_id()) + { return Err(Error::AtomicCommitConflict); } + let prepared = + prepare_authored_memory(&mut candidate, value.preparation().clone())?; + candidate.authored_drafts.push(value.intent().clone()); candidate - .authored_operations - .push(value.operation().clone()); - candidate - .authored_artifacts - .extend(value.artifacts().iter().cloned()); - candidate - .authored_delivery_plans - .extend(value.delivery_plans().iter().cloned()); - AuthoredAtomicOutcome::Prepared { - operation: value.operation().clone(), - artifacts: value.artifacts().to_vec(), - delivery_plans: value.delivery_plans().to_vec(), - } + .authored_atomic_receipts + .push(AuthoredAtomicReceipt::new( + &ordinary, + AtomicCommitDisposition::Committed, + ordinary.requested_at_unix_ms(), + prepared, + )?); + AuthoredAtomicOutcome::Submitted(value) } AuthoredAtomicCommand::Claim(value) => match value.target() { ClaimAuthoredTarget::ArtifactSigning(artifact_id) => { @@ -1783,6 +1832,62 @@ impl AuthoredAtomicStorage for MemoryStorage { } impl AuthoredDraftStore for MemoryStorage { + fn query_authored_drafts( + &self, + query: crate::authored_draft_query::AuthoredDraftQuery, + ) -> BoxFuture<'_, Result<crate::authored_draft_query::AuthoredDraftPage, Error>> { + Box::pin(async move { + use crate::authored_draft_query::{ + AUTHORED_DRAFT_PAGE_PAYLOAD_MAX_BYTES, AuthoredDraftPage, AuthoredDraftQueryRecord, + }; + let state = self.state()?; + let mut heads: std::collections::BTreeMap<AuthoredDraftId, &AuthoredDraft> = + std::collections::BTreeMap::new(); + let capacity = usize::from(query.limit()) + 1; + for draft in &state.authored_drafts { + if !query.matches(draft) + || query + .after() + .is_some_and(|after| *draft.draft_id().as_bytes() <= after) + { + continue; + } + if let Some(head) = heads.get_mut(&draft.draft_id()) { + if draft.revision() > head.revision() { + *head = draft; + } + continue; + } + // Retain only the smallest requested IDs and one lookahead. + // Draft identity metadata is immutable across revisions. + if heads.len() == capacity { + if heads + .last_key_value() + .is_some_and(|(last, _)| draft.draft_id() >= *last) + { + continue; + } + heads.pop_last(); + } + heads.insert(draft.draft_id(), draft); + } + let mut records = Vec::new(); + let mut bytes = 0usize; + let mut has_more = false; + for draft in heads.values() { + if records.len() == usize::from(query.limit()) + || bytes + draft.payload().len() > AUTHORED_DRAFT_PAGE_PAYLOAD_MAX_BYTES + { + has_more = true; + break; + } + bytes += draft.payload().len(); + records.push(AuthoredDraftQueryRecord::Draft((*draft).clone())); + } + AuthoredDraftPage::new(&query, records, has_more) + }) + } + fn append_authored_draft( &self, draft: AuthoredDraft, diff --git a/crates/storage/tests/authored_atomic.rs b/crates/storage/tests/authored_atomic.rs @@ -1016,3 +1016,6 @@ fn memory_executes_all_cancellation_targets_and_revision_fences() { AuthoredDeliveryState::Cancelled ); } + +#[path = "authored_atomic/draft_submission.rs"] +mod draft_submission; diff --git a/crates/storage/tests/authored_atomic/draft_submission.rs b/crates/storage/tests/authored_atomic/draft_submission.rs @@ -0,0 +1,532 @@ +use super::*; +use radroots_storage::{ + atomic::AtomicCommitId, + authored_draft::{ + AuthoredDraft, AuthoredDraftId, AuthoredDraftRevision, AuthoredDraftStage, + AuthoredDraftStore, + }, + authored_draft_query::AuthoredDraftScope, + authored_draft_submission::{AuthoredDraftSource, PrepareFromDraft}, +}; +use sha2::{Digest, Sha256}; + +fn source() -> AuthoredDraft { + AuthoredDraft::initial( + AuthoredDraftId::new([8; 16]).unwrap(), + *authored_plan().author().as_bytes(), + "fixture.partial.v1", + b"unfinished 0.".to_vec(), + AuthoredDraftStage::Draft, + None, + 9, + ) + .unwrap() + .with_scope(AuthoredDraftScope::new([5; 32]).unwrap()) + .unwrap() +} +fn request(source: &AuthoredDraft, key: u8, id: u8) -> PrepareFromDraft { + let AuthoredAtomicCommand::Prepare(base) = prepare(7).0 else { + unreachable!() + }; + let operation_id = OperationInstanceId::new([id; 16]).unwrap(); + let artifact_id = AuthoredArtifactId::new([id; 16]).unwrap(); + let delivery_id = AuthoredDeliveryPlanId::new([id; 16]).unwrap(); + let preparation = PrepareAuthoredOperation::new( + AuthoredOperation::new(operation_id, vec![artifact_id], 10).unwrap(), + vec![ + AuthoredArtifact::planned(artifact_id, operation_id, 0, &authored_plan(), 10).unwrap(), + ], + vec![ + AuthoredDeliveryPlan::new( + delivery_id, + artifact_id, + base.delivery_plans()[0].intent().clone(), + 10, + ) + .unwrap(), + ], + base.input_digest(), + 10, + ) + .unwrap(); + let payload = b"complete immutable semantic intent".to_vec(); + let mut intent = AuthoredDraft::reconstruct( + AuthoredDraftId::new([id; 16]).unwrap(), + AuthoredDraftRevision::INITIAL, + *source.author(), + "fixture.intent.v1", + payload.clone(), + Sha256::digest(&payload).into(), + AuthoredDraftStage::Queued, + Some(operation_id), + 10, + 10, + ) + .unwrap(); + if let Some(scope) = source.scope() { + intent = intent.with_scope(scope).unwrap(); + } + PrepareFromDraft::new( + AtomicCommitId::new([key; 16]).unwrap(), + AuthoredDraftSource::capture(source).unwrap(), + intent, + preparation, + ) + .unwrap() +} +fn command(request: PrepareFromDraft) -> AuthoredAtomicCommand { + AuthoredAtomicCommand::PrepareFromDraft(Box::new(request)) +} + +#[test] +fn submitted_request_replays_before_cas_and_compares_actual_fields_without_trusting_digest() { + block_on(async { + let store = MemoryStorage::default(); + let source = source(); + store + .append_authored_draft(source.clone(), None) + .await + .unwrap(); + let request = request(&source, 4, 1); + let cmd = command(request.clone()); + let receipt = store.execute_authored(cmd.clone()).await.unwrap(); + assert_eq!(receipt.disposition(), AtomicCommitDisposition::Committed); + assert_eq!( + receipt.outcome(), + &AuthoredAtomicOutcome::Submitted(Box::new(request.clone())) + ); + assert_eq!( + store.authored_draft_head(source.draft_id()).await.unwrap(), + Some(source.clone()) + ); + assert_eq!( + store + .authored_draft_head(request.intent().draft_id()) + .await + .unwrap(), + Some(request.intent().clone()) + ); + let newer = source + .successor( + b"later editing".to_vec(), + AuthoredDraftStage::Draft, + None, + 11, + ) + .unwrap(); + store + .append_authored_draft(newer, Some(source.revision())) + .await + .unwrap(); + let replay = store.execute_authored(cmd.clone()).await.unwrap(); + assert_eq!(replay.disposition(), AtomicCommitDisposition::Replay); + assert_eq!(replay.outcome(), receipt.outcome()); + // This is the exact ordinary command used by later Sync signing. + let ordinary = AuthoredAtomicCommand::Prepare(request.preparation().clone()); + let ordinary_receipt = store.execute_authored(ordinary.clone()).await.unwrap(); + assert_eq!( + ordinary_receipt.disposition(), + AtomicCommitDisposition::Replay + ); + assert!(matches!( + ordinary_receipt.outcome(), + AuthoredAtomicOutcome::Prepared { .. } + )); + assert!(!receipt.matches_command(&ordinary)); + assert!(!ordinary_receipt.matches_command(&cmd)); + let base = request.preparation(); + let old_plan = &base.delivery_plans()[0]; + let changed_plan = AuthoredDeliveryPlan::new( + old_plan.plan_id(), + old_plan.artifact_id(), + AuthoredDeliveryIntent::new( + "changed-delivery-request", + old_plan.intent().target_set().clone(), + SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::all()), + 101, + ) + .unwrap(), + 10, + ) + .unwrap(); + let changed = PrepareFromDraft::new( + request.command_id(), + request.source().clone(), + request.intent().clone(), + PrepareAuthoredOperation::new( + base.operation().clone(), + base.artifacts().to_vec(), + vec![changed_plan], + base.input_digest(), + 10, + ) + .unwrap(), + ) + .unwrap(); + let changed = command(changed); + assert_eq!(cmd.commit_id(), changed.commit_id()); + assert_eq!( + cmd.digest(), + changed.digest(), + "caller digest is intentionally unchanged" + ); + assert_eq!( + store.execute_authored(changed).await, + Err(Error::AtomicCommitConflict) + ); + let fresh = command(self::request(&source, 6, 2)); + assert_eq!( + store.execute_authored(fresh.clone()).await, + Err(Error::DraftRevisionConflict) + ); + assert!( + store + .authored_receipt(fresh.commit_id()) + .await + .unwrap() + .is_none() + ); + assert!( + store + .authored_operation(OperationInstanceId::new([2; 16]).unwrap()) + .await + .unwrap() + .is_none() + ); + assert_eq!( + store + .authored_receipt(cmd.commit_id()) + .await + .unwrap() + .unwrap() + .outcome(), + receipt.outcome() + ); + }); +} + +#[test] +fn intentional_identical_submissions_use_distinct_commands_and_operations() { + block_on(async { + let store = MemoryStorage::default(); + let source = source(); + store + .append_authored_draft(source.clone(), None) + .await + .unwrap(); + for (key, id) in [(1, 1), (2, 2)] { + let value = request(&source, key, id); + let receipt = store + .execute_authored(command(value.clone())) + .await + .unwrap(); + assert_eq!(receipt.disposition(), AtomicCommitDisposition::Committed); + assert_eq!( + receipt.commit_id(), + PrepareFromDraft::commit_id_for(source.author(), value.command_id()) + ); + } + assert_eq!( + store + .authored_draft_heads(*source.author(), 10) + .await + .unwrap() + .len(), + 3 + ); + }); +} + +#[test] +fn submission_wire_and_receipts_validate_bindings_and_redact_payloads() { + let source = source(); + let request = request(&source, 1, 2); + let wire = serde_json::to_value(&request).unwrap(); + assert_eq!( + serde_json::from_value::<PrepareFromDraft>(wire.clone()).unwrap(), + request + ); + assert!(!format!("{:?}", command(request.clone())).contains("semantic intent")); + assert!(!format!("{request:?}").contains("atomic authored plan")); + assert_eq!(request.source().draft_id(), source.draft_id()); + assert_eq!(request.source().revision(), source.revision()); + assert_eq!(request.source().author(), source.author()); + assert_eq!(request.source().payload_schema(), source.payload_schema()); + assert_eq!(request.source().scope(), source.scope()); + assert_eq!(request.source().payload_sha256(), source.payload_sha256()); + let invalid = [ + ("/command_id", serde_json::json!([0; 16].to_vec())), + ("/source/author", serde_json::json!([0; 32].to_vec())), + ("/source/payload_schema", serde_json::json!(" invalid")), + ("/source/stage", serde_json::json!("queued")), + ("/source/created_at_unix_ms", serde_json::json!(0)), + ("/source/updated_at_unix_ms", serde_json::json!(8)), + ("/source/updated_at_unix_ms", serde_json::json!(11)), + ("/source/draft_id", serde_json::json!([2; 16].to_vec())), + ("/source/author", serde_json::json!([1; 32].to_vec())), + ("/source/scope", serde_json::Value::Null), + ("/intent/revision", serde_json::json!(2)), + ("/intent/operation_id", serde_json::json!([3; 16].to_vec())), + ("/intent/updated_at_unix_ms", serde_json::json!(11)), + ("/preparation/requested_at_unix_ms", serde_json::json!(0)), + ("/preparation/artifacts", serde_json::json!([])), + ]; + for (pointer, replacement) in invalid { + let mut invalid = wire.clone(); + *invalid.pointer_mut(pointer).unwrap() = replacement; + assert!( + serde_json::from_value::<PrepareFromDraft>(invalid).is_err(), + "{pointer}" + ); + } + let cmd = command(request.clone()); + let outcome = AuthoredAtomicOutcome::Submitted(Box::new(request)); + for (id, digest, at) in [ + (AtomicCommitId::new([7; 16]).unwrap(), cmd.digest(), 10), + (cmd.commit_id(), AtomicCommitDigest::new([7; 32]), 10), + (cmd.commit_id(), cmd.digest(), 9), + ] { + assert!( + AuthoredAtomicReceipt::from_durable_parts( + id, + digest, + AtomicCommitDisposition::Committed, + at, + outcome.clone() + ) + .is_err() + ); + } +} + +#[test] +fn submission_rejects_preexisting_work_state_and_mismatched_captured_times() { + let source = source(); + let request = request(&source, 1, 2); + let base = request.preparation(); + let reject = |operation: AuthoredOperation, + artifacts: Vec<AuthoredArtifact>, + plans: Vec<AuthoredDeliveryPlan>| { + let preparation = + PrepareAuthoredOperation::new(operation, artifacts, plans, base.input_digest(), 10) + .unwrap(); + assert_eq!( + PrepareFromDraft::new( + request.command_id(), + request.source().clone(), + request.intent().clone(), + preparation + ), + Err(Error::AtomicWorkflowMismatch) + ); + }; + for (created, updated, revision) in [(9, 10, 1), (10, 11, 1), (10, 10, 2)] { + reject( + AuthoredOperation::reconstruct( + base.operation().operation_id(), + base.operation().artifact_ids().to_vec(), + created, + updated, + NonZeroU64::new(revision).unwrap(), + ) + .unwrap(), + base.artifacts().to_vec(), + base.delivery_plans().to_vec(), + ); + } + let artifact = &base.artifacts()[0]; + let mut already_signed = artifact.clone(); + already_signed + .record_signed(signed(&authored_plan()), 10) + .unwrap(); + let mut claimed = artifact.clone(); + claimed + .set_signing_claim( + WorkClaim::new( + [1; 16], + "existing worker", + NonZeroU64::MIN, + 10, + 20, + NonZeroU64::MIN, + ) + .unwrap(), + 10, + ) + .unwrap(); + let mut later_wire = serde_json::to_value(artifact).unwrap(); + later_wire["updated_at_unix_ms"] = serde_json::json!(11); + let other_author = AuthoredEventPlan::from_generic( + GenericEventDraft::new( + "radroots.social.geochat.v1", + 20_000, + 1_800_000_100, + vec![], + "another author", + "79be667ef9dcbbac55a06295ce870b07029bfcdb2dce28d959f2815b16f81798", + ) + .unwrap(), + ) + .unwrap(); + for changed in [ + AuthoredArtifact::imported_signed( + artifact.artifact_id(), + artifact.operation_id(), + 0, + signed(&authored_plan()), + 10, + ) + .unwrap(), + already_signed, + claimed, + AuthoredArtifact::planned( + artifact.artifact_id(), + artifact.operation_id(), + 0, + &authored_plan(), + 9, + ) + .unwrap(), + serde_json::from_value(later_wire).unwrap(), + AuthoredArtifact::planned( + artifact.artifact_id(), + artifact.operation_id(), + 0, + &other_author, + 10, + ) + .unwrap(), + ] { + reject( + base.operation().clone(), + vec![changed], + base.delivery_plans().to_vec(), + ); + } + let plan = &base.delivery_plans()[0]; + let mut cancelled = plan.clone(); + cancelled.cancel(10).unwrap(); + let mut later_wire = serde_json::to_value(plan).unwrap(); + later_wire["updated_at_unix_ms"] = serde_json::json!(11); + for changed in [ + cancelled, + AuthoredDeliveryPlan::new(plan.plan_id(), plan.artifact_id(), plan.intent().clone(), 9) + .unwrap(), + serde_json::from_value(later_wire).unwrap(), + ] { + reject( + base.operation().clone(), + base.artifacts().to_vec(), + vec![changed], + ); + } + let original = request.intent(); + for (stage, operation, created) in [ + (AuthoredDraftStage::Draft, None, 10), + (AuthoredDraftStage::Queued, original.operation_id(), 9), + ] { + let intent = AuthoredDraft::reconstruct( + original.draft_id(), + original.revision(), + *original.author(), + original.payload_schema(), + original.payload().to_vec(), + *original.payload_sha256(), + stage, + operation, + created, + 10, + ) + .unwrap() + .with_scope(original.scope().unwrap()) + .unwrap(); + assert_eq!( + PrepareFromDraft::new( + request.command_id(), + request.source().clone(), + intent, + base.clone() + ), + Err(Error::AtomicWorkflowMismatch) + ); + } + let mut ready = serde_json::to_value(&request).unwrap(); + ready["intent"]["stage"] = serde_json::json!("ready_to_sign"); + let ready: PrepareFromDraft = serde_json::from_value(ready).unwrap(); + assert_eq!(ready.intent().stage(), AuthoredDraftStage::ReadyToSign); + let cmd = command(request.clone()); + let ordinary = AuthoredAtomicCommand::Prepare(base.clone()); + let submitted = AuthoredAtomicOutcome::Submitted(Box::new(request.clone())); + assert!( + AuthoredAtomicReceipt::new(&ordinary, AtomicCommitDisposition::Committed, 10, submitted) + .is_err() + ); + let prepared = AuthoredAtomicOutcome::Prepared { + operation: base.operation().clone(), + artifacts: base.artifacts().to_vec(), + delivery_plans: base.delivery_plans().to_vec(), + }; + assert!( + AuthoredAtomicReceipt::new(&cmd, AtomicCommitDisposition::Committed, 10, prepared).is_err() + ); +} + +#[test] +fn submission_cannot_adopt_unassociated_existing_intent_or_operation() { + block_on(async { + let source = source(); + let request = request(&source, 1, 2); + let empty = MemoryStorage::default(); + assert_eq!( + empty.execute_authored(command(request.clone())).await, + Err(Error::DraftRevisionConflict) + ); + for existing_intent in [true, false] { + let store = MemoryStorage::default(); + store + .append_authored_draft(source.clone(), None) + .await + .unwrap(); + if existing_intent { + store + .append_authored_draft(request.intent().clone(), None) + .await + .unwrap(); + } else { + store + .execute_authored(AuthoredAtomicCommand::Prepare( + request.preparation().clone(), + )) + .await + .unwrap(); + } + assert_eq!( + store.execute_authored(command(request.clone())).await, + Err(if existing_intent { + Error::DraftRevisionConflict + } else { + Error::AtomicCommitConflict + }) + ); + assert!( + store + .authored_receipt(request.commit_id()) + .await + .unwrap() + .is_none() + ); + assert_eq!( + store.authored_draft_head(source.draft_id()).await.unwrap(), + Some(source.clone()) + ); + assert_eq!( + store + .authored_operation(request.preparation().operation().operation_id()) + .await + .unwrap() + .is_some(), + !existing_intent + ); + } + }); +} diff --git a/crates/storage/tests/authored_draft_query.rs b/crates/storage/tests/authored_draft_query.rs @@ -0,0 +1,260 @@ +use futures_executor::block_on; +use radroots_storage::{ + Error, + authored_draft::{AuthoredDraft, AuthoredDraftId, AuthoredDraftStage, AuthoredDraftStore}, + authored_draft_query::{ + AUTHORED_DRAFT_PAGE_PAYLOAD_MAX_BYTES, AuthoredDraftCursor, AuthoredDraftPage, + AuthoredDraftQuery, AuthoredDraftQueryRecord, AuthoredDraftScope, + }, + memory::MemoryStorage, +}; + +fn draft( + id: u8, + schema: &str, + scope: Option<AuthoredDraftScope>, + payload: Vec<u8>, +) -> AuthoredDraft { + let draft = AuthoredDraft::initial( + AuthoredDraftId::new([id; 16]).unwrap(), + [9; 32], + schema, + payload, + AuthoredDraftStage::Draft, + None, + 10, + ) + .unwrap(); + scope.map_or_else( + || draft.clone(), + |scope| draft.clone().with_scope(scope).unwrap(), + ) +} +fn query(scope: Option<AuthoredDraftScope>, limit: u16) -> AuthoredDraftQuery { + AuthoredDraftQuery::new([9; 32], "fixture.composer.v1", scope, limit).unwrap() +} + +#[test] +fn scoped_pages_preserve_revisions_and_do_not_mix_other_schemas_or_scopes() { + let store = MemoryStorage::default(); + let scope = AuthoredDraftScope::new([7; 32]).unwrap(); + let first = draft( + 2, + "fixture.composer.v1", + Some(scope), + b"unfinished 1.".to_vec(), + ); + let second = draft( + 4, + "fixture.composer.v1", + Some(scope), + b"partial date 2026-".to_vec(), + ); + for value in [ + first.clone(), + second.clone(), + draft(1, "fixture.profile.v1", Some(scope), b"profile".to_vec()), + draft(3, "fixture.composer.v1", None, b"other context".to_vec()), + ] { + block_on(store.append_authored_draft(value, None)).unwrap(); + } + let page = block_on(store.query_authored_drafts(query(Some(scope), 1))).unwrap(); + assert_eq!( + page.records(), + [AuthoredDraftQueryRecord::Draft(first.clone())] + ); + assert_eq!(page.records()[0].draft_id().unwrap(), first.draft_id()); + assert_eq!(page.records()[0].revision(), first.revision()); + let cursor = page.next_cursor().unwrap(); + let changed = first + .successor(b"later edit".to_vec(), AuthoredDraftStage::Draft, None, 11) + .unwrap(); + assert_eq!(changed.scope(), Some(scope)); + block_on(store.append_authored_draft(changed.clone(), Some(first.revision()))).unwrap(); + let next = + block_on(store.query_authored_drafts(query(Some(scope), 2).with_cursor(cursor).unwrap())) + .unwrap(); + assert_eq!(next.records(), [AuthoredDraftQueryRecord::Draft(second)]); + assert!(next.next_cursor().is_none()); + let fresh = block_on(store.query_authored_drafts(query(Some(scope), 1))).unwrap(); + assert_eq!( + fresh.records(), + [AuthoredDraftQueryRecord::Draft(changed.clone())] + ); + assert!(changed.with_scope(scope).is_err()); + assert!(first.clone().with_scope(scope).is_err()); + let mut forged = serde_json::to_value( + first + .successor(b"next".to_vec(), AuthoredDraftStage::Draft, None, 11) + .unwrap(), + ) + .unwrap(); + forged["scope"] = serde_json::json!([8; 32].to_vec()); + let forged: AuthoredDraft = serde_json::from_value(forged).unwrap(); + assert_eq!( + forged.validate_successor_of(&first), + Err(Error::DraftRevisionConflict) + ); +} + +#[test] +fn query_and_cursor_reject_invalid_or_changed_authority() { + assert!(AuthoredDraftScope::new([0; 32]).is_err()); + for schema in ["", " x", "x\n", &"x".repeat(129)] { + assert!(AuthoredDraftQuery::new([9; 32], schema, None, 1).is_err()); + } + for (author, limit) in [([0; 32], 1), ([9; 32], 0), ([9; 32], 257)] { + assert!(AuthoredDraftQuery::new(author, "fixture.composer.v1", None, limit).is_err()); + } + let q = query(None, 256); + let cursor = q.cursor_after([0; 16]); + for changed in [ + AuthoredDraftQuery::new([8; 32], q.payload_schema(), None, 1).unwrap(), + AuthoredDraftQuery::new([9; 32], "fixture.profile.v1", None, 1).unwrap(), + query(Some(AuthoredDraftScope::new([7; 32]).unwrap()), 1), + ] { + assert!(changed.with_cursor(&cursor).is_err()); + } + let bytes = serde_json::to_vec(&cursor).unwrap(); + let decoded: AuthoredDraftCursor = serde_json::from_slice(&bytes).unwrap(); + assert_eq!(decoded, cursor); + assert_eq!( + q.clone().with_cursor(&decoded).unwrap().after(), + Some([0; 16]) + ); + for (field, value) in [ + ("schema_version", serde_json::json!(2)), + ("payload_schema", serde_json::json!("")), + ("scope", serde_json::json!([0; 32].to_vec())), + ("unknown", serde_json::json!(true)), + ] { + let mut wire = serde_json::to_value(&cursor).unwrap(); + wire[field] = value; + assert!(serde_json::from_value::<AuthoredDraftCursor>(wire).is_err()); + } +} + +#[test] +fn legacy_unscoped_wire_and_scoped_serde_keep_exact_payload_bytes() { + let legacy = draft(1, "fixture.composer.v1", None, b"0.\n2026-".to_vec()); + let bytes = serde_json::to_vec(&legacy).unwrap(); + assert!( + !String::from_utf8(bytes.clone()) + .unwrap() + .contains("\"scope\"") + ); + let restored: AuthoredDraft = serde_json::from_slice(&bytes).unwrap(); + assert_eq!(restored, legacy); + assert_eq!(serde_json::to_vec(&restored).unwrap(), bytes); + let scope = AuthoredDraftScope::new([7; 32]).unwrap(); + assert_eq!( + AuthoredDraftScope::try_from(<[u8; 32]>::from(scope)).unwrap(), + scope + ); + let scoped = legacy.with_scope(scope).unwrap(); + let scoped_bytes = serde_json::to_vec(&scoped).unwrap(); + assert_eq!( + serde_json::from_slice::<AuthoredDraft>(&scoped_bytes).unwrap(), + scoped + ); + let mut malformed = serde_json::to_value(&scoped).unwrap(); + malformed["scope"] = serde_json::json!([0; 32].to_vec()); + assert!(serde_json::from_value::<AuthoredDraft>(malformed).is_err()); +} + +#[test] +fn page_payload_budget_advances_without_losing_the_next_large_record() { + let store = MemoryStorage::default(); + for id in [1, 2] { + block_on(store.append_authored_draft( + draft(id, "fixture.composer.v1", None, vec![id; 3 * 1024 * 1024]), + None, + )) + .unwrap(); + } + let q = query(None, 256); + let first = block_on(store.query_authored_drafts(q.clone())).unwrap(); + assert_eq!(first.records().len(), 1); + let next = block_on( + store.query_authored_drafts(q.clone().with_cursor(first.next_cursor().unwrap()).unwrap()), + ) + .unwrap(); + assert_eq!(next.records()[0].draft_key(), [2; 16]); + assert!(next.next_cursor().is_none()); + let mut records = first.into_records(); + records.extend(next.into_records()); + let total_payload: usize = records + .iter() + .map(|record| match record { + AuthoredDraftQueryRecord::Draft(draft) => draft.payload().len(), + AuthoredDraftQueryRecord::Corrupt { .. } => 0, + }) + .sum(); + assert!(total_payload > AUTHORED_DRAFT_PAGE_PAYLOAD_MAX_BYTES); + assert!(AuthoredDraftPage::new(&q, records, false).is_err()); + assert!(AuthoredDraftPage::new(&q, vec![], true).is_err()); + let record = AuthoredDraftQueryRecord::Draft(draft(1, "fixture.profile.v1", None, vec![1])); + assert!(AuthoredDraftPage::new(&q, vec![record], false).is_err()); + let record = AuthoredDraftQueryRecord::Draft(draft(1, "fixture.composer.v1", None, vec![1])); + assert!(AuthoredDraftPage::new(&q, vec![record.clone(), record.clone()], false).is_err()); + assert!(AuthoredDraftPage::new(&query(None, 1), vec![record.clone(), record], false).is_err()); +} + +#[test] +fn bounded_reference_pages_scan_a_thousand_reversed_heads_and_preserve_newer_revisions() { + let store = MemoryStorage::default(); + for number in (1u16..=1000).rev() { + let mut id = [0; 16]; + id[14..].copy_from_slice(&number.to_be_bytes()); + let value = AuthoredDraft::initial( + radroots_storage::authored_draft::AuthoredDraftId::new(id).unwrap(), + [9; 32], + "fixture.composer.v1", + number.to_be_bytes().to_vec(), + AuthoredDraftStage::Draft, + None, + 10, + ) + .unwrap(); + block_on(store.append_authored_draft(value, None)).unwrap(); + } + let q = query(None, 127); + let mut cursor = None; + let mut found = Vec::new(); + loop { + let current = cursor.as_ref().map_or_else( + || q.clone(), + |cursor| q.clone().with_cursor(cursor).unwrap(), + ); + let page = block_on(store.query_authored_drafts(current)).unwrap(); + assert!(page.records().len() <= usize::from(q.limit())); + for record in page.records() { + let key = record.draft_key(); + found.push(u16::from_be_bytes([key[14], key[15]])); + } + cursor = page.next_cursor().cloned(); + if found.len() == 127 { + let AuthoredDraftQueryRecord::Draft(first) = &page.records()[0] else { + panic!("valid fixture") + }; + let edited = first + .successor( + b"newer source".to_vec(), + AuthoredDraftStage::Draft, + None, + 11, + ) + .unwrap(); + block_on(store.append_authored_draft(edited, Some(first.revision()))).unwrap(); + } + if cursor.is_none() { + break; + } + } + assert_eq!(found, (1u16..=1000).collect::<Vec<_>>()); + let fresh = block_on(store.query_authored_drafts(query(None, 1))).unwrap(); + let AuthoredDraftQueryRecord::Draft(first) = &fresh.records()[0] else { + panic!("valid fixture") + }; + assert_eq!(first.payload(), b"newer source"); +} diff --git a/crates/storage_sqlite/README.md b/crates/storage_sqlite/README.md @@ -8,6 +8,20 @@ apply only the governed forward migrations. Fresh stores require a host-supplied `SourceGeneration` and creation timestamp; the crate never reads hidden entropy or a wall clock. +Runtime schema v14 adds generic schema/scope metadata and an index for bounded +draft-head queries. Migration preserves every original revision snapshot, +including corrupt historical evidence. Known foreign schemas and scopes are +filtered before decoding; unknown historical schemas expose only an opaque +author-bound corruption locator. Each page releases its read snapshot before +returning a continuation and reads at most 16 MiB of serialized snapshots. + +Draft submission uses the existing WAL/FULL writer and authored tables. One +transaction writes the immutable intent, operation, artifacts, delivery plans, +ordinary Prepare receipt and source-association receipt. Exact replay precedes +source-head CAS, including after reopen or later editing. Existing atomic +receipt snapshots retain their 4 MiB limit; oversized submissions fail without +partial writes. No additional public database or connection API is introduced. + Backend status and the last integrity result are passive. Hosts invoke `check_integrity` explicitly with their own positive timestamp when they want full SQLite and foreign-key validation across both owned files. `close` drains diff --git a/crates/storage_sqlite/src/authored.rs b/crates/storage_sqlite/src/authored.rs @@ -183,7 +183,7 @@ async fn execute_transaction( .map_err(map_backend)? { let committed = decode_receipt_row(&row)?; - if committed.digest() != command.digest() { + if !committed.matches_command(command) { return Err(Error::AtomicCommitConflict); } return AuthoredAtomicReceipt::from_durable_parts( @@ -196,6 +196,14 @@ async fn execute_transaction( } let outcome = execute_command(transaction, command.clone()).await?; + commit_outcome(transaction, command, outcome).await +} + +async fn commit_outcome( + transaction: &mut sqlx::Transaction<'_, Sqlite>, + command: &AuthoredAtomicCommand, + outcome: AuthoredAtomicOutcome, +) -> Result<AuthoredAtomicReceipt, Error> { let receipt = AuthoredAtomicReceipt::new( command, AtomicCommitDisposition::Committed, @@ -223,35 +231,62 @@ async fn execute_transaction( Ok(receipt) } +async fn prepare_operation( + transaction: &mut sqlx::Transaction<'_, Sqlite>, + value: radroots_storage::authored_atomic::PrepareAuthoredOperation, +) -> Result<AuthoredAtomicOutcome, Error> { + if row_exists( + transaction, + "SELECT 1 FROM radroots_runtime_authored_operations WHERE operation_id = ?", + value.operation().operation_id().as_bytes(), + ) + .await? + || any_artifact_exists(transaction, value.artifacts()).await? + || any_plan_exists(transaction, value.delivery_plans()).await? + { + return Err(Error::AtomicCommitConflict); + } + persist_operation(transaction, value.operation()).await?; + for artifact in value.artifacts() { + persist_artifact(transaction, artifact).await?; + } + for plan in value.delivery_plans() { + persist_plan(transaction, plan).await?; + } + Ok(AuthoredAtomicOutcome::Prepared { + operation: value.operation().clone(), + artifacts: value.artifacts().to_vec(), + delivery_plans: value.delivery_plans().to_vec(), + }) +} + async fn execute_command( transaction: &mut sqlx::Transaction<'_, Sqlite>, command: AuthoredAtomicCommand, ) -> Result<AuthoredAtomicOutcome, Error> { match command { - AuthoredAtomicCommand::Prepare(value) => { - if row_exists( - transaction, - "SELECT 1 FROM radroots_runtime_authored_operations WHERE operation_id = ?", - value.operation().operation_id().as_bytes(), - ) - .await? - || any_artifact_exists(transaction, value.artifacts()).await? - || any_plan_exists(transaction, value.delivery_plans()).await? + AuthoredAtomicCommand::Prepare(value) => prepare_operation(transaction, value).await, + AuthoredAtomicCommand::PrepareFromDraft(value) => { + value.validate()?; + let head = + crate::authored_draft::load_head_tx(transaction, value.source().draft_id()).await?; + if !head + .as_ref() + .is_some_and(|head| value.source().matches(head)) { - return Err(Error::AtomicCommitConflict); - } - persist_operation(transaction, value.operation()).await?; - for artifact in value.artifacts() { - persist_artifact(transaction, artifact).await?; + return Err(Error::DraftRevisionConflict); } - for plan in value.delivery_plans() { - persist_plan(transaction, plan).await?; + if crate::authored_draft::load_head_tx(transaction, value.intent().draft_id()) + .await? + .is_some() + { + return Err(Error::DraftRevisionConflict); } - Ok(AuthoredAtomicOutcome::Prepared { - operation: value.operation().clone(), - artifacts: value.artifacts().to_vec(), - delivery_plans: value.delivery_plans().to_vec(), - }) + let ordinary = AuthoredAtomicCommand::Prepare(value.preparation().clone()); + let prepared = prepare_operation(transaction, value.preparation().clone()).await?; + crate::authored_draft::insert_draft_tx(transaction, value.intent()).await?; + commit_outcome(transaction, &ordinary, prepared).await?; + Ok(AuthoredAtomicOutcome::Submitted(value)) } AuthoredAtomicCommand::Claim(value) => match value.target() { ClaimAuthoredTarget::ArtifactSigning(id) => { @@ -942,7 +977,9 @@ fn decode_receipt_row(row: &SqliteRow) -> Result<AuthoredAtomicReceipt, Error> { let requested = u64_from_i64(column(row, "requested_at_unix_ms")?)?; let committed = u64_from_i64(column(row, "committed_at_unix_ms")?)?; let snapshot = decode_snapshot::<ReceiptSnapshot>(column(row, "receipt")?)?; - if committed < requested { + if committed < requested + || matches!(&snapshot.outcome, AuthoredAtomicOutcome::Submitted(value) if value.preparation().requested_at_unix_ms() != requested) + { return Err(Error::AtomicCommitFailed); } AuthoredAtomicReceipt::from_durable_parts( @@ -1046,6 +1083,9 @@ fn decode_snapshot<T: DeserializeOwned>(bytes: Vec<u8>) -> Result<T, Error> { fn command_target(command: &AuthoredAtomicCommand) -> [u8; 16] { match command { + AuthoredAtomicCommand::PrepareFromDraft(value) => { + *value.preparation().operation().operation_id().as_bytes() + } AuthoredAtomicCommand::Prepare(value) => *value.operation().operation_id().as_bytes(), AuthoredAtomicCommand::Claim(value) => match value.target() { ClaimAuthoredTarget::ArtifactSigning(id) @@ -1069,7 +1109,7 @@ fn command_target(command: &AuthoredAtomicCommand) -> [u8; 16] { fn command_phase(command: &AuthoredAtomicCommand) -> &'static str { match command { - AuthoredAtomicCommand::Prepare(_) => "prepare", + AuthoredAtomicCommand::Prepare(_) | AuthoredAtomicCommand::PrepareFromDraft(_) => "prepare", AuthoredAtomicCommand::Claim(_) => "claim", AuthoredAtomicCommand::ApplySigned(_) => "signing", AuthoredAtomicCommand::ApplyAdmission(_) => "admission", @@ -1257,7 +1297,7 @@ mod tests { ) } - fn prepare() -> (AuthoredAtomicCommand, AuthoredEventPlan) { + pub(super) fn prepare() -> (AuthoredAtomicCommand, AuthoredEventPlan) { let (operation_id, artifact_id, plan_id) = ids(); let event_plan = plan(); let artifact = AuthoredArtifact::planned(artifact_id, operation_id, 0, &event_plan, 10) @@ -2462,3 +2502,8 @@ mod tests { assert!(plan.contains("radroots_runtime_authored_delivery_ready_idx")); } } + +#[cfg(test)] +#[cfg_attr(coverage_nightly, coverage(off))] +#[path = "authored_draft_submission_tests.rs"] +mod draft_submission_tests; diff --git a/crates/storage_sqlite/src/authored_draft.rs b/crates/storage_sqlite/src/authored_draft.rs @@ -11,7 +11,18 @@ use sqlx::{Row, sqlite::SqliteRow}; const SNAPSHOT_MAX_BYTES: usize = 16 * 1024 * 1024; +#[path = "authored_draft_query.rs"] +mod query; + impl AuthoredDraftStore for SqliteStorage { + fn query_authored_drafts( + &self, + query: radroots_storage::authored_draft_query::AuthoredDraftQuery, + ) -> BoxFuture<'_, Result<radroots_storage::authored_draft_query::AuthoredDraftPage, Error>> + { + Box::pin(async move { query::page(self, query).await }) + } + fn append_authored_draft( &self, draft: AuthoredDraft, @@ -71,25 +82,7 @@ impl AuthoredDraftStore for SqliteStorage { } } - 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)?; + insert_draft_tx(&mut transaction, &draft).await?; transaction.commit().await.map_err(map_backend)?; Ok(DraftAppendReceipt::new( draft, @@ -217,6 +210,17 @@ fn decode_row(row: &SqliteRow) -> Result<AuthoredDraft, Error> { (Some(raw), Some(expected)) => raw.as_slice() == expected.as_bytes(), _ => false, }; + let schema = row + .try_get::<String, _>("payload_schema") + .map_err(|_| Error::CorruptAuthoredDraft)?; + let scope = row + .try_get::<Option<Vec<u8>>, _>("payload_scope") + .map_err(|_| Error::CorruptAuthoredDraft)?; + if schema != draft.payload_schema() + || scope != draft.scope().map(|scope| scope.as_bytes().to_vec()) + { + return Err(Error::CorruptAuthoredDraft); + } if draft.draft_id().as_bytes() != &draft_id || draft.revision().get() != revision || draft.author() != &author @@ -282,7 +286,7 @@ mod tests { .unwrap() } - async fn open_store(temp: &TempDir) -> SqliteStorage { + pub(super) async fn open_store(temp: &TempDir) -> SqliteStorage { let paths = Paths::from_directory(temp.path()).unwrap(); SqliteStorage::open( OpenOptions::new(paths, OpenMode::Create) @@ -491,8 +495,8 @@ mod tests { 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 (?, ?, ?, ?, ?, ?, ?, ?, ?)", + created_at_unix_ms, updated_at_unix_ms, snapshot, payload_schema, payload_scope + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", ) .bind(draft_id.as_slice()) .bind(revision) @@ -503,6 +507,8 @@ mod tests { .bind(created_at_unix_ms) .bind(updated_at_unix_ms) .bind(encode_snapshot(draft).unwrap()) + .bind(draft.payload_schema()) + .bind(draft.scope().map(|scope| scope.as_bytes().to_vec())) .execute(store.pool()) .await .unwrap(); @@ -625,3 +631,44 @@ mod tests { } } } + +#[cfg(test)] +#[cfg_attr(coverage_nightly, coverage(off))] +#[path = "authored_draft_query_tests.rs"] +mod query_tests; + +pub(crate) async fn load_head_tx( + transaction: &mut sqlx::Transaction<'_, sqlx::Sqlite>, + id: AuthoredDraftId, +) -> Result<Option<AuthoredDraft>, Error> { + sqlx::query("SELECT * FROM radroots_runtime_authored_draft_revisions WHERE draft_id = ? ORDER BY revision DESC LIMIT 1") + .bind(id.as_bytes().as_slice()).fetch_optional(&mut **transaction).await.map_err(map_backend)? + .as_ref().map(decode_row).transpose() +} +pub(crate) async fn insert_draft_tx( + transaction: &mut sqlx::Transaction<'_, sqlx::Sqlite>, + draft: &AuthoredDraft, +) -> Result<(), Error> { + 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, payload_schema, payload_scope + ) 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) + .bind(draft.payload_schema()) + .bind(draft.scope().map(|scope| scope.as_bytes().to_vec())) + .execute(&mut **transaction) + .await + .map_err(map_backend)?; + Ok(()) +} diff --git a/crates/storage_sqlite/src/authored_draft_query.rs b/crates/storage_sqlite/src/authored_draft_query.rs @@ -0,0 +1,101 @@ +use super::{SqliteStorage, decode_row, map_backend}; +use radroots_storage::{ + Error, + authored_draft::AuthoredDraftRevision, + authored_draft_query::{ + AUTHORED_DRAFT_PAGE_PAYLOAD_MAX_BYTES, AUTHORED_DRAFT_PAGE_SNAPSHOT_MAX_BYTES, + AuthoredDraftPage, AuthoredDraftQuery, AuthoredDraftQueryRecord, + }, +}; +use sqlx::Row; + +pub(super) async fn page( + store: &SqliteStorage, + query: AuthoredDraftQuery, +) -> Result<AuthoredDraftPage, Error> { + let mut transaction = store.pool().begin().await.map_err(map_backend)?; + let result = read_page(&mut transaction, &query).await; + // The read snapshot is always released before returning any continuation. + let rollback = transaction.rollback().await.map_err(map_backend); + match result { + Ok(page) => { + rollback?; + Ok(page) + } + Err(error) => Err(error), + } +} + +async fn read_page( + transaction: &mut sqlx::Transaction<'_, sqlx::Sqlite>, + query: &AuthoredDraftQuery, +) -> Result<AuthoredDraftPage, Error> { + let after = query.after().map(|value| value.to_vec()); + let rows = sqlx::query( + "SELECT revisions.draft_id, revisions.revision, length(revisions.snapshot) AS snapshot_bytes, + revisions.payload_schema = '' AS unknown_schema + FROM radroots_runtime_authored_draft_revisions AS revisions + WHERE revisions.author = ? + AND ((revisions.payload_schema = ? AND revisions.payload_scope IS ?) + OR revisions.payload_schema = '') + AND (? IS NULL OR revisions.draft_id > ?) + 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.draft_id LIMIT ?" + ).bind(query.author().as_slice()).bind(query.payload_schema()) + .bind(query.scope().map(|value| value.as_bytes().to_vec())) + .bind(&after).bind(after).bind(i64::from(query.limit()) + 1) + .fetch_all(&mut **transaction).await.map_err(map_backend)?; + let mut records = Vec::new(); + let mut snapshot_bytes = 0usize; + let mut payload_bytes = 0usize; + let mut has_more = rows.len() > usize::from(query.limit()); + for row in rows.iter().take(usize::from(query.limit())) { + // Keys and positive revisions are enforced by the STRICT table. A zero + // key is still a valid scan/repair position for a corrupt domain ID. + let key: [u8; 16] = row + .try_get::<Vec<u8>, _>("draft_id") + .map_err(map_backend)? + .try_into() + .map_err(|_| Error::CorruptAuthoredDraft)?; + let revision = row.try_get::<i64, _>("revision").map_err(map_backend)?; + let revision = AuthoredDraftRevision::new( + u64::try_from(revision).map_err(|_| Error::CorruptAuthoredDraft)?, + )?; + let size = row + .try_get::<i64, _>("snapshot_bytes") + .map_err(map_backend)?; + let unknown = row + .try_get::<bool, _>("unknown_schema") + .map_err(map_backend)?; + let corrupt = AuthoredDraftQueryRecord::Corrupt { + draft_key: key, + revision, + }; + if unknown || size <= 0 || size as u64 > AUTHORED_DRAFT_PAGE_SNAPSHOT_MAX_BYTES as u64 { + records.push(corrupt); + continue; + } + let size = size as usize; + if snapshot_bytes + size > AUTHORED_DRAFT_PAGE_SNAPSHOT_MAX_BYTES { + has_more = true; + break; + } + snapshot_bytes += size; + let row = sqlx::query("SELECT * FROM radroots_runtime_authored_draft_revisions WHERE draft_id = ? AND revision = ?") + .bind(key.as_slice()).bind(revision.get() as i64) + .fetch_one(&mut **transaction).await.map_err(map_backend)?; + match decode_row(&row) { + Ok(draft) if query.matches(&draft) => { + if payload_bytes + draft.payload().len() > AUTHORED_DRAFT_PAGE_PAYLOAD_MAX_BYTES { + has_more = true; + break; + } + payload_bytes += draft.payload().len(); + records.push(AuthoredDraftQueryRecord::Draft(draft)); + } + Ok(_) | Err(_) => records.push(corrupt), + } + } + AuthoredDraftPage::new(query, records, has_more) +} diff --git a/crates/storage_sqlite/src/authored_draft_query_tests.rs b/crates/storage_sqlite/src/authored_draft_query_tests.rs @@ -0,0 +1,234 @@ +use super::{SqliteStorage, tests::open_store}; +use radroots_storage::{ + Error, + authored_draft::{AuthoredDraft, AuthoredDraftId, AuthoredDraftStage, AuthoredDraftStore}, + authored_draft_query::{AuthoredDraftQuery, AuthoredDraftQueryRecord, AuthoredDraftScope}, +}; +use tempfile::TempDir; + +fn draft( + id: u8, + schema: &str, + scope: Option<AuthoredDraftScope>, + payload: Vec<u8>, +) -> AuthoredDraft { + let draft = AuthoredDraft::initial( + AuthoredDraftId::new([id; 16]).unwrap(), + [7; 32], + schema, + payload, + AuthoredDraftStage::Draft, + None, + 10, + ) + .unwrap(); + match scope { + Some(scope) => draft.with_scope(scope).unwrap(), + None => draft, + } +} +fn query(scope: Option<AuthoredDraftScope>, limit: u16) -> AuthoredDraftQuery { + AuthoredDraftQuery::new([7; 32], "fixture.composer.v1", scope, limit).unwrap() +} +async fn corrupt( + store: &SqliteStorage, + id: u8, + author: u8, + schema: &str, + scope: Option<AuthoredDraftScope>, +) { + sqlx::query("INSERT INTO radroots_runtime_authored_draft_revisions + (draft_id, revision, author, stage, payload_sha256, created_at_unix_ms, updated_at_unix_ms, snapshot, payload_schema, payload_scope) + VALUES (?, 1, ?, 0, ?, 10, 10, ?, ?, ?)") + .bind([id; 16].as_slice()).bind([author; 32].as_slice()).bind([3; 32].as_slice()) + .bind(b"{malformed private payload".as_slice()).bind(schema).bind(scope.map(|scope| scope.as_bytes().to_vec())) + .execute(store.pool()).await.unwrap(); +} + +#[tokio::test] +async fn scoped_sqlite_pages_isolate_corruption_foreign_schemas_and_contexts() { + let temp = TempDir::new().unwrap(); + let store = open_store(&temp).await; + let scope = AuthoredDraftScope::new([9; 32]).unwrap(); + let first = draft( + 2, + "fixture.composer.v1", + Some(scope), + b"incomplete 0.".to_vec(), + ); + let second = draft(6, "fixture.composer.v1", Some(scope), b"2026-".to_vec()); + for value in [ + first.clone(), + second.clone(), + draft(1, "fixture.profile.v1", Some(scope), b"profile".to_vec()), + draft(3, "fixture.composer.v1", None, b"unscoped".to_vec()), + ] { + store.append_authored_draft(value, None).await.unwrap(); + } + corrupt(&store, 4, 7, "fixture.composer.v1", Some(scope)).await; + // Unknown historical schema yields only an author-bound corruption locator. + corrupt(&store, 5, 7, "", None).await; + corrupt(&store, 7, 8, "", None).await; + corrupt(&store, 8, 7, "fixture.profile.v1", Some(scope)).await; + let q = query(Some(scope), 1); + let first_page = store.query_authored_drafts(q.clone()).await.unwrap(); + assert_eq!( + first_page.records(), + [AuthoredDraftQueryRecord::Draft(first.clone())] + ); + let changed = first + .successor( + b"later saved edit".to_vec(), + AuthoredDraftStage::Draft, + None, + 11, + ) + .unwrap(); + store + .append_authored_draft(changed, Some(first.revision())) + .await + .unwrap(); + let mut cursor = first_page.next_cursor().cloned(); + let mut records = first_page.into_records(); + while let Some(next) = cursor { + let page = store + .query_authored_drafts(q.clone().with_cursor(&next).unwrap()) + .await + .unwrap(); + cursor = page.next_cursor().cloned(); + records.extend(page.into_records()); + } + assert_eq!( + records + .iter() + .map(AuthoredDraftQueryRecord::draft_key) + .collect::<Vec<_>>(), + [[2; 16], [4; 16], [5; 16], [6; 16]] + ); + assert!(matches!( + records[1], + AuthoredDraftQueryRecord::Corrupt { .. } + )); + assert!(matches!( + records[2], + AuthoredDraftQueryRecord::Corrupt { .. } + )); + assert_eq!(records[3], AuthoredDraftQueryRecord::Draft(second)); + assert!(!format!("{records:?}").contains("private payload")); + assert_eq!( + store.authored_draft_heads([7; 32], 10).await, + Err(Error::CorruptAuthoredDraft) + ); + store.close().await.unwrap(); + assert!(store.query_authored_drafts(q.clone()).await.is_err()); + let reopened = open_store(&temp).await; + let page = reopened + .query_authored_drafts(query(Some(scope), 10)) + .await + .unwrap(); + assert_eq!(page.records().len(), 4); + assert!(page.next_cursor().is_none()); + reopened.close().await.unwrap(); +} + +#[tokio::test] +async fn sqlite_snapshot_budget_keeps_large_valid_records_on_later_pages() { + // Small numeric bytes exhaust decoded payload first; larger ones exhaust serialized snapshots first. + for payload_byte in [0, 99] { + let temp = TempDir::new().unwrap(); + let store = open_store(&temp).await; + for id in [1, 2] { + store + .append_authored_draft( + draft( + id, + "fixture.composer.v1", + None, + vec![payload_byte; 3 * 1024 * 1024], + ), + None, + ) + .await + .unwrap(); + } + let q = query(None, 256); + let first = store.query_authored_drafts(q.clone()).await.unwrap(); + assert_eq!(first.records().len(), 1); + let next = store + .query_authored_drafts(q.with_cursor(first.next_cursor().unwrap()).unwrap()) + .await + .unwrap(); + assert_eq!(next.records()[0].draft_key(), [2; 16]); + assert!(matches!( + next.records()[0], + AuthoredDraftQueryRecord::Draft(_) + )); + assert!(next.next_cursor().is_none()); + store.close().await.unwrap(); + } +} + +#[tokio::test] +async fn zero_domain_identity_is_an_isolated_repair_position() { + let temp = TempDir::new().unwrap(); + let store = open_store(&temp).await; + corrupt(&store, 0, 7, "", None).await; + let good = draft(1, "fixture.composer.v1", None, b"partial".to_vec()); + store + .append_authored_draft(good.clone(), None) + .await + .unwrap(); + let q = query(None, 1); + let page = store.query_authored_drafts(q.clone()).await.unwrap(); + assert!(page.records()[0].draft_id().is_err()); + let next = store + .query_authored_drafts(q.with_cursor(page.next_cursor().unwrap()).unwrap()) + .await + .unwrap(); + assert_eq!(next.records(), [AuthoredDraftQueryRecord::Draft(good)]); + store.close().await.unwrap(); +} + +#[tokio::test] +async fn inconsistent_query_metadata_is_isolated_and_never_returns_foreign_payload() { + for corrupt_scope in [false, true] { + let temp = TempDir::new().unwrap(); + let store = open_store(&temp).await; + let value = draft(1, "fixture.composer.v1", None, b"private original".to_vec()); + store + .append_authored_draft(value.clone(), None) + .await + .unwrap(); + sqlx::query("DROP TRIGGER radroots_runtime_authored_draft_revisions_update_guard") + .execute(store.pool()) + .await + .unwrap(); + let (schema, scope) = if corrupt_scope { + ( + "fixture.composer.v1", + Some(AuthoredDraftScope::new([4; 32]).unwrap()), + ) + } else { + ("fixture.other.v1", None) + }; + sqlx::query("UPDATE radroots_runtime_authored_draft_revisions SET payload_schema = ?, payload_scope = ?") + .bind(schema).bind(scope.map(|v|v.as_bytes().to_vec())).execute(store.pool()).await.unwrap(); + assert_eq!( + store.authored_draft_head(value.draft_id()).await, + Err(Error::CorruptAuthoredDraft) + ); + let page = store + .query_authored_drafts(AuthoredDraftQuery::new([7; 32], schema, scope, 1).unwrap()) + .await + .unwrap(); + assert_eq!( + page.records(), + [AuthoredDraftQueryRecord::Corrupt { + draft_key: *value.draft_id().as_bytes(), + revision: value.revision() + }] + ); + assert!(!format!("{page:?}").contains("private original")); + store.close().await.unwrap(); + } +} diff --git a/crates/storage_sqlite/src/authored_draft_submission_tests.rs b/crates/storage_sqlite/src/authored_draft_submission_tests.rs @@ -0,0 +1,549 @@ +use super::*; +use crate::{OpenMode, OpenOptions, Paths}; +use radroots_storage::{ + authored_draft::{ + AuthoredDraft, AuthoredDraftId, AuthoredDraftRevision, AuthoredDraftStage, + AuthoredDraftStore, + }, + authored_draft_query::AuthoredDraftScope, + authored_draft_submission::{AuthoredDraftSource, PrepareFromDraft}, + event::SourceGeneration, + memory::MemoryStorage, +}; +use sha2::{Digest, Sha256}; +use tempfile::TempDir; + +async fn open(temp: &TempDir, mode: OpenMode) -> SqliteStorage { + let options = OpenOptions::new(Paths::from_directory(temp.path()).unwrap(), mode); + let options = if mode == OpenMode::Create { + options + .with_source_generation(SourceGeneration::new([9; 32]).unwrap(), 9) + .unwrap() + } else { + options + }; + SqliteStorage::open(options).await.unwrap() +} +fn fixture() -> (AuthoredDraft, PrepareFromDraft) { + let (ordinary, plan) = super::tests::prepare(); + let AuthoredAtomicCommand::Prepare(preparation) = ordinary else { + unreachable!() + }; + let scope = AuthoredDraftScope::new([3; 32]).unwrap(); + let source = AuthoredDraft::initial( + AuthoredDraftId::new([8; 16]).unwrap(), + *plan.author().as_bytes(), + "fixture.partial.v1", + b"partial 0.".to_vec(), + AuthoredDraftStage::Draft, + None, + 9, + ) + .unwrap() + .with_scope(scope) + .unwrap(); + let payload = b"complete captured request".to_vec(); + let intent = AuthoredDraft::reconstruct( + AuthoredDraftId::new([7; 16]).unwrap(), + AuthoredDraftRevision::INITIAL, + *source.author(), + "fixture.intent.v1", + payload.clone(), + Sha256::digest(&payload).into(), + AuthoredDraftStage::Queued, + Some(preparation.operation().operation_id()), + 10, + 10, + ) + .unwrap() + .with_scope(scope) + .unwrap(); + let request = PrepareFromDraft::new( + AtomicCommitId::new([4; 16]).unwrap(), + AuthoredDraftSource::capture(&source).unwrap(), + intent, + preparation, + ) + .unwrap(); + (source, request) +} +fn command(request: &PrepareFromDraft) -> AuthoredAtomicCommand { + AuthoredAtomicCommand::PrepareFromDraft(Box::new(request.clone())) +} +async fn empty_submission( + store: &SqliteStorage, + source: &AuthoredDraft, + request: &PrepareFromDraft, +) { + for table in [ + "operations", + "artifacts", + "delivery_plans", + "delivery_targets", + "atomic_commits", + ] { + // Audited test-only identifiers come exclusively from the literal list above. + let count: i64 = sqlx::query_scalar(sqlx::AssertSqlSafe(format!( + "SELECT COUNT(*) FROM radroots_runtime_authored_{table}" + ))) + .fetch_one(store.pool()) + .await + .unwrap(); + assert_eq!(count, 0, "{table}"); + } + assert_eq!( + store + .authored_draft_heads(*source.author(), 10) + .await + .unwrap(), + std::slice::from_ref(source) + ); + assert!( + store + .authored_receipt(request.commit_id()) + .await + .unwrap() + .is_none() + ); +} + +#[tokio::test] +async fn submission_survives_lost_callback_reopen_and_matches_memory_after_later_save() { + let temp = TempDir::new().unwrap(); + let store = open(&temp, OpenMode::Create).await; + let memory = MemoryStorage::default(); + let (source, request) = fixture(); + store + .append_authored_draft(source.clone(), None) + .await + .unwrap(); + memory + .append_authored_draft(source.clone(), None) + .await + .unwrap(); + let expected = memory.execute_authored(command(&request)).await.unwrap(); + assert_eq!( + store.execute_authored(command(&request)).await.unwrap(), + expected + ); + let newer = source + .successor(b"later edit".to_vec(), AuthoredDraftStage::Draft, None, 11) + .unwrap(); + for target in [&store as &dyn AuthoredDraftStore, &memory] { + target + .append_authored_draft(newer.clone(), Some(source.revision())) + .await + .unwrap(); + } + // Discard the callback value and release the actual writer before recovery. + store.close().await.unwrap(); + assert!(store.execute_authored(command(&request)).await.is_err()); + let store = open(&temp, OpenMode::ReadWriteExisting).await; + let replay = store.execute_authored(command(&request)).await.unwrap(); + assert_eq!( + replay, + memory.execute_authored(command(&request)).await.unwrap() + ); + assert_eq!(replay.disposition(), AtomicCommitDisposition::Replay); + assert_eq!(replay.outcome(), expected.outcome()); + let ordinary = AuthoredAtomicCommand::Prepare(request.preparation().clone()); + assert_eq!( + store.execute_authored(ordinary.clone()).await.unwrap(), + memory.execute_authored(ordinary).await.unwrap() + ); + // A valid changed context under the same author/command conflicts before CAS. + let mut changed = serde_json::to_value(&request).unwrap(); + changed["source"]["scope"] = serde_json::json!([6; 32].to_vec()); + changed["intent"]["scope"] = serde_json::json!([6; 32].to_vec()); + let changed = serde_json::from_value::<PrepareFromDraft>(changed).unwrap(); + assert_eq!(command(&changed).digest(), command(&request).digest()); + assert_eq!( + store.execute_authored(command(&changed)).await, + Err(Error::AtomicCommitConflict) + ); + let reader = open(&temp, OpenMode::ReadOnly).await; + assert!( + reader + .authored_receipt(request.commit_id()) + .await + .unwrap() + .is_some() + ); + assert!(reader.execute_authored(command(&request)).await.is_err()); + reader.close().await.unwrap(); + store.close().await.unwrap(); +} + +#[tokio::test] +async fn submission_rolls_back_before_and_after_every_record_and_receipt_insert() { + let (source, request) = fixture(); + let ordinary = AuthoredAtomicCommand::Prepare(request.preparation().clone()); + let hex = |id: AtomicCommitId| { + id.as_bytes() + .iter() + .map(|b| format!("{b:02x}")) + .collect::<String>() + }; + let points = [ + ("operations", String::new()), + ("artifacts", String::new()), + ("delivery_plans", String::new()), + ("delivery_targets", "WHEN NEW.ordinal = 0".into()), + ("delivery_targets", "WHEN NEW.ordinal = 1".into()), + ( + "draft_revisions", + "WHEN NEW.draft_id != x'08080808080808080808080808080808'".into(), + ), + ( + "atomic_commits", + format!("WHEN NEW.commit_id = x'{}'", hex(ordinary.commit_id())), + ), + ( + "atomic_commits", + format!("WHEN NEW.commit_id = x'{}'", hex(request.commit_id())), + ), + ]; + for timing in ["BEFORE", "AFTER"] { + for (table, condition) in &points { + let temp = TempDir::new().unwrap(); + let store = open(&temp, OpenMode::Create).await; + store + .append_authored_draft(source.clone(), None) + .await + .unwrap(); + // Timing/table/condition are fixed test literals plus hex-encoded typed IDs. + sqlx::query(sqlx::AssertSqlSafe(format!("CREATE TRIGGER submission_fault {timing} INSERT ON radroots_runtime_authored_{table} {condition} BEGIN SELECT RAISE(ABORT, 'injected submission failure'); END"))) + .execute(store.pool()).await.unwrap(); + assert!( + store.execute_authored(command(&request)).await.is_err(), + "{timing} {table} {condition}" + ); + empty_submission(&store, &source, &request).await; + sqlx::query("DROP TRIGGER submission_fault") + .execute(store.pool()) + .await + .unwrap(); + store.close().await.unwrap(); + let store = open(&temp, OpenMode::ReadWriteExisting).await; + empty_submission(&store, &source, &request).await; + assert_eq!( + store + .execute_authored(command(&request)) + .await + .unwrap() + .disposition(), + AtomicCommitDisposition::Committed + ); + store.close().await.unwrap(); + } + } +} + +#[tokio::test] +async fn abandoned_precommit_and_sqlite_full_preserve_source_without_success_association() { + let temp = TempDir::new().unwrap(); + let store = open(&temp, OpenMode::Create).await; + let (source, request) = fixture(); + store + .append_authored_draft(source.clone(), None) + .await + .unwrap(); + let mut transaction = store.pool().begin_with("BEGIN IMMEDIATE").await.unwrap(); + let receipt = execute_transaction(&mut transaction, &command(&request)) + .await + .unwrap(); + assert_eq!(receipt.disposition(), AtomicCommitDisposition::Committed); + // This internal result has not crossed the public commit boundary. + drop(transaction); + empty_submission(&store, &source, &request).await; + let mut transaction = store.pool().begin_with("BEGIN IMMEDIATE").await.unwrap(); + let page_count: i64 = sqlx::query_scalar("PRAGMA page_count") + .fetch_one(&mut *transaction) + .await + .unwrap(); + // The PRAGMA accepts an integer, not a bind parameter; this is a typed i64. + sqlx::query(sqlx::AssertSqlSafe(format!( + "PRAGMA max_page_count = {page_count}" + ))) + .execute(&mut *transaction) + .await + .unwrap(); + let mut large = serde_json::to_value(&request).unwrap(); + let payload = vec![97u8; 100_000]; + large["intent"]["payload"] = serde_json::json!(payload); + large["intent"]["payload_sha256"] = serde_json::json!(Sha256::digest(&payload).to_vec()); + let large: PrepareFromDraft = serde_json::from_value(large).unwrap(); + assert!( + execute_transaction(&mut transaction, &command(&large)) + .await + .is_err() + ); + // SQLITE_FULL may already roll back the transaction at the engine boundary. + let _ = transaction.rollback().await; + empty_submission(&store, &source, &request).await; + store.close().await.unwrap(); + let store = open(&temp, OpenMode::ReadWriteExisting).await; + assert_eq!( + store + .execute_authored(command(&request)) + .await + .unwrap() + .disposition(), + AtomicCommitDisposition::Committed + ); + store.close().await.unwrap(); +} + +#[tokio::test] +async fn concurrent_duplicates_and_save_have_one_atomic_order() { + for _ in 0..4 { + let temp = TempDir::new().unwrap(); + let store = open(&temp, OpenMode::Create).await; + let (source, request) = fixture(); + store + .append_authored_draft(source.clone(), None) + .await + .unwrap(); + let newer = source + .successor( + b"concurrent edit".to_vec(), + AuthoredDraftStage::Draft, + None, + 11, + ) + .unwrap(); + let (first, second, save) = tokio::join!( + store.execute_authored(command(&request)), + store.execute_authored(command(&request)), + store.append_authored_draft(newer.clone(), Some(source.revision())) + ); + save.unwrap(); + match (first, second) { + (Ok(a), Ok(b)) => { + assert_eq!(a.outcome(), b.outcome()); + assert_ne!(a.disposition(), b.disposition()); + assert_eq!( + store + .execute_authored(command(&request)) + .await + .unwrap() + .disposition(), + AtomicCommitDisposition::Replay + ); + let operations: i64 = + sqlx::query_scalar("SELECT COUNT(*) FROM radroots_runtime_authored_operations") + .fetch_one(store.pool()) + .await + .unwrap(); + assert_eq!(operations, 1); + } + (Err(Error::DraftRevisionConflict), Err(Error::DraftRevisionConflict)) => { + assert!( + store + .authored_receipt(request.commit_id()) + .await + .unwrap() + .is_none() + ); + assert!( + store + .authored_draft_head(request.intent().draft_id()) + .await + .unwrap() + .is_none() + ); + } + other => panic!("non-atomic submit/save result: {other:?}"), + } + assert_eq!( + store.authored_draft_head(source.draft_id()).await.unwrap(), + Some(newer) + ); + store.close().await.unwrap(); + } +} + +#[tokio::test] +async fn commit_failure_returns_no_success_and_rolls_back_every_authored_write() { + let temp = TempDir::new().unwrap(); + let store = open(&temp, OpenMode::Create).await; + let (source, request) = fixture(); + store + .append_authored_draft(source.clone(), None) + .await + .unwrap(); + sqlx::query("CREATE TABLE submission_commit_fault (parent BLOB REFERENCES radroots_runtime_authored_operations(operation_id) DEFERRABLE INITIALLY DEFERRED)") + .execute(store.pool()).await.unwrap(); + // The last receipt succeeds; only the actual transaction COMMIT fails. + sqlx::query("CREATE TRIGGER submission_commit_fault_trigger AFTER INSERT ON radroots_runtime_authored_atomic_commits BEGIN INSERT INTO submission_commit_fault VALUES (x'99999999999999999999999999999999'); END") + .execute(store.pool()).await.unwrap(); + assert!(store.execute_authored(command(&request)).await.is_err()); + empty_submission(&store, &source, &request).await; + sqlx::query("DROP TRIGGER submission_commit_fault_trigger") + .execute(store.pool()) + .await + .unwrap(); + sqlx::query("DROP TABLE submission_commit_fault") + .execute(store.pool()) + .await + .unwrap(); + assert_eq!( + store + .execute_authored(command(&request)) + .await + .unwrap() + .disposition(), + AtomicCommitDisposition::Committed + ); + store.close().await.unwrap(); +} + +#[test] +fn storage_submission_machine_policy_matches_native_bounds() { + let policy: serde_json::Value = serde_json::from_str(include_str!( + "../../../contracts/architecture/decisions/authored_draft_submission.v1.json" + )) + .unwrap(); + let bounds = &policy["bounds"]; + assert_eq!( + bounds["page_records"], + radroots_storage::authored_draft::AUTHORED_DRAFT_QUERY_LIMIT_MAX + ); + assert_eq!( + bounds["decoded_page_payload_bytes"], + radroots_storage::authored_draft_query::AUTHORED_DRAFT_PAGE_PAYLOAD_MAX_BYTES + ); + assert_eq!( + bounds["serialized_page_snapshot_bytes"], + radroots_storage::authored_draft_query::AUTHORED_DRAFT_PAGE_SNAPSHOT_MAX_BYTES + ); + assert_eq!( + bounds["existing_atomic_receipt_snapshot_bytes"], + super::SNAPSHOT_MAX_BYTES + ); + assert_eq!( + policy["owners"], + serde_json::json!([ + "radroots_storage", + "radroots_storage_sqlite", + "radroots_sync" + ]) + ); +} + +#[tokio::test] +async fn existing_intent_and_partial_graph_id_collisions_cannot_create_associations() { + let (source, request) = fixture(); + let temp = TempDir::new().unwrap(); + let store = open(&temp, OpenMode::Create).await; + store + .append_authored_draft(source.clone(), None) + .await + .unwrap(); + store + .append_authored_draft(request.intent().clone(), None) + .await + .unwrap(); + assert_eq!( + store.execute_authored(command(&request)).await, + Err(Error::DraftRevisionConflict) + ); + assert!( + store + .authored_receipt(request.commit_id()) + .await + .unwrap() + .is_none() + ); + assert!( + store + .authored_operation(request.preparation().operation().operation_id()) + .await + .unwrap() + .is_none() + ); + store.close().await.unwrap(); + let temp = TempDir::new().unwrap(); + let store = open(&temp, OpenMode::Create).await; + store + .execute_authored(AuthoredAtomicCommand::Prepare( + request.preparation().clone(), + )) + .await + .unwrap(); + for artifact_collision in [true, false] { + let base = request.preparation(); + let operation_id = OperationInstanceId::new([20; 16]).unwrap(); + let artifact_id = if artifact_collision { + base.artifacts()[0].artifact_id() + } else { + AuthoredArtifactId::new([21; 16]).unwrap() + }; + let plan_id = if artifact_collision { + AuthoredDeliveryPlanId::new([22; 16]).unwrap() + } else { + base.delivery_plans()[0].plan_id() + }; + let artifact = AuthoredArtifact::planned( + artifact_id, + operation_id, + 0, + base.artifacts()[0].plan().unwrap().decode().unwrap().plan(), + 10, + ) + .unwrap(); + let preparation = radroots_storage::authored_atomic::PrepareAuthoredOperation::new( + AuthoredOperation::new(operation_id, vec![artifact_id], 10).unwrap(), + vec![artifact], + vec![ + AuthoredDeliveryPlan::new( + plan_id, + artifact_id, + base.delivery_plans()[0].intent().clone(), + 10, + ) + .unwrap(), + ], + base.input_digest(), + 10, + ) + .unwrap(); + assert_eq!( + store + .execute_authored(AuthoredAtomicCommand::Prepare(preparation)) + .await, + Err(Error::AtomicCommitConflict) + ); + assert!( + store + .authored_operation(operation_id) + .await + .unwrap() + .is_none() + ); + } + store.close().await.unwrap(); +} + +#[tokio::test] +async fn submission_receipt_shadow_times_must_match_the_captured_request() { + let temp = TempDir::new().unwrap(); + let store = open(&temp, OpenMode::Create).await; + let (source, request) = fixture(); + store.append_authored_draft(source, None).await.unwrap(); + store.execute_authored(command(&request)).await.unwrap(); + for requested in [9i64, 11] { + let row = sqlx::query("SELECT commit_id, commit_digest, ? AS requested_at_unix_ms, committed_at_unix_ms, receipt FROM radroots_runtime_authored_atomic_commits WHERE commit_id = ?") + .bind(requested).bind(request.commit_id().as_bytes().as_slice()).fetch_one(store.pool()).await.unwrap(); + assert_eq!(decode_receipt_row(&row), Err(Error::AtomicCommitFailed)); + } + assert_eq!( + store + .execute_authored(command(&request)) + .await + .unwrap() + .disposition(), + AtomicCommitDisposition::Replay + ); + store.close().await.unwrap(); +} diff --git a/crates/storage_sqlite/src/migration.rs b/crates/storage_sqlite/src/migration.rs @@ -371,7 +371,7 @@ async fn metadata( fn validate_plan(plan: &MigrationPlan) -> Result<(), Error> { let valid = plan.minimum_version > 0 && plan.minimum_version <= plan.current_version - && plan.current_version <= 13 + && plan.current_version <= 14 && plan.steps.len() == usize::try_from(plan.current_version).unwrap_or(usize::MAX) && plan .steps @@ -489,6 +489,7 @@ const fn set_user_version_sql(version: u32) -> Option<&'static str> { 11 => Some("PRAGMA user_version = 11"), 12 => Some("PRAGMA user_version = 12"), 13 => Some("PRAGMA user_version = 13"), + 14 => Some("PRAGMA user_version = 14"), _ => None, } } @@ -503,13 +504,13 @@ mod tests { const TEST_V1_OBJECTS: &[&str] = &["radroots_test_one"]; const TEST_V2_OBJECTS: &[&str] = &["radroots_test_one", "radroots_test_two"]; - async fn connection() -> SqliteConnection { + pub(super) async fn connection() -> SqliteConnection { SqliteConnection::connect("sqlite::memory:") .await .expect("memory SQLite") } - async fn pragma(connection: &mut SqliteConnection, name: &str) -> i64 { + pub(super) async fn pragma(connection: &mut SqliteConnection, name: &str) -> i64 { let sql = match name { "application_id" => "PRAGMA application_id", "user_version" => "PRAGMA user_version", @@ -521,7 +522,7 @@ mod tests { .expect("pragma") } - async fn establish_runtime_version(connection: &mut SqliteConnection, version: u32) { + pub(super) async fn establish_runtime_version(connection: &mut SqliteConnection, version: u32) { for migration_version in 1..=version { sqlx::raw_sql( runtime::migration_sql(migration_version).expect("registered runtime SQL"), @@ -715,7 +716,7 @@ mod tests { .execute(&mut newer) .await .expect("application id"); - sqlx::raw_sql("PRAGMA user_version = 14") + sqlx::raw_sql("PRAGMA user_version = 15") .execute(&mut newer) .await .expect("newer version"); @@ -724,10 +725,10 @@ mod tests { Err(Error::SchemaTooNew { database: RUNTIME_DATABASE, supported: runtime::CURRENT_VERSION, - actual: 14, + actual: 15, }) )); - assert_eq!(pragma(&mut newer, "user_version").await, 14); + assert_eq!(pragma(&mut newer, "user_version").await, 15); let mut wrong_identity = connection().await; establish_runtime_version(&mut wrong_identity, 1).await; @@ -856,3 +857,8 @@ mod tests { } } } + +#[cfg(test)] +#[cfg_attr(coverage_nightly, coverage(off))] +#[path = "migration_draft_query_tests.rs"] +mod draft_query_tests; diff --git a/crates/storage_sqlite/src/migration/runtime/0014_authored_draft_query_metadata.up.sql b/crates/storage_sqlite/src/migration/runtime/0014_authored_draft_query_metadata.up.sql @@ -0,0 +1,19 @@ +-- Backfill only generic query metadata; immutable revision snapshots are retained. +DROP TRIGGER radroots_runtime_authored_draft_revisions_update_guard; +ALTER TABLE radroots_runtime_authored_draft_revisions + ADD COLUMN payload_schema TEXT NOT NULL DEFAULT '' CHECK(length(CAST(payload_schema AS BLOB)) <= 128); +ALTER TABLE radroots_runtime_authored_draft_revisions + ADD COLUMN payload_scope BLOB CHECK(payload_scope IS NULL OR length(payload_scope) = 32); +UPDATE radroots_runtime_authored_draft_revisions SET payload_schema = + CASE WHEN json_valid(CAST(snapshot AS TEXT)) THEN + CASE WHEN json_type(CAST(snapshot AS TEXT), '$.payload_schema') = 'text' + AND length(CAST(json_extract(CAST(snapshot AS TEXT), '$.payload_schema') AS BLOB)) BETWEEN 1 AND 128 + THEN json_extract(CAST(snapshot AS TEXT), '$.payload_schema') ELSE '' END + ELSE '' END; +CREATE INDEX radroots_runtime_authored_draft_scope_head_idx + ON radroots_runtime_authored_draft_revisions(author, payload_schema, payload_scope, 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; 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 = 13; +pub const CURRENT_VERSION: u32 = 14; 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"); @@ -24,6 +24,9 @@ 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"); +const AUTHORED_DRAFT_QUERY_V14_SQL: &str = + include_str!("0014_authored_draft_query_metadata.up.sql"); + /// Stable, non-SQL description of one forward runtime migration. #[derive(Clone, Copy, Debug, Eq, PartialEq)] pub struct MigrationDescriptor { @@ -555,6 +558,84 @@ const RUNTIME_V13_OBJECTS: &[&str] = &[ "radroots_runtime_source_generations_sequence_guard", ]; +const RUNTIME_V14_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_draft_scope_head_idx", + "radroots_runtime_authored_migration_evidence", + "radroots_runtime_authored_migration_evidence_delete_guard", + "radroots_runtime_authored_migration_evidence_update_guard", + "radroots_runtime_authored_operations", + "radroots_runtime_delivery_evidence", + "radroots_runtime_delivery_evidence_item_idx", + "radroots_runtime_event_index_checkpoints", + "radroots_runtime_event_index_manifests", + "radroots_runtime_event_index_shards", + "radroots_runtime_event_provenance", + "radroots_runtime_event_provenance_observed_idx", + "radroots_runtime_events", + "radroots_runtime_events_admission_idx", + "radroots_runtime_events_contract_metadata_guard", + "radroots_runtime_events_contract_metadata_insert_guard", + "radroots_runtime_events_delete_guard", + "radroots_runtime_events_event_id_idx", + "radroots_runtime_events_raw_update_guard", + "radroots_runtime_journal_idempotency_idx", + "radroots_runtime_journal_operations", + "radroots_runtime_journal_recovery_idx", + "radroots_runtime_legacy_event_staging", + "radroots_runtime_legacy_event_staging_delete_guard", + "radroots_runtime_legacy_event_staging_insert_guard", + "radroots_runtime_legacy_event_staging_update_guard", + "radroots_runtime_legacy_import_commit_delete_guard", + "radroots_runtime_legacy_import_commit_update_guard", + "radroots_runtime_legacy_import_commits", + "radroots_runtime_legacy_import_delete_guard", + "radroots_runtime_legacy_import_identity_guard", + "radroots_runtime_legacy_import_member_delete_guard", + "radroots_runtime_legacy_import_member_identity_guard", + "radroots_runtime_legacy_import_member_state_guard", + "radroots_runtime_legacy_import_members", + "radroots_runtime_legacy_import_state_guard", + "radroots_runtime_legacy_import_state_idx", + "radroots_runtime_legacy_imports", + "radroots_runtime_legacy_outbox_staging", + "radroots_runtime_legacy_outbox_staging_delete_guard", + "radroots_runtime_legacy_outbox_staging_insert_guard", + "radroots_runtime_legacy_outbox_staging_parent_idx", + "radroots_runtime_legacy_outbox_staging_update_guard", + "radroots_runtime_outbox_items", + "radroots_runtime_outbox_operation_idx", + "radroots_runtime_outbox_ready_idx", + "radroots_runtime_outbox_targets", + "radroots_runtime_projection_checkpoints", + "radroots_runtime_projection_documents", + "radroots_runtime_projection_invalidations", + "radroots_runtime_projection_rebuilds", + "radroots_runtime_projection_rebuilds_stage_idx", + "radroots_runtime_projection_snapshots", + "radroots_runtime_projection_snapshots_created_idx", + "radroots_runtime_projection_snapshots_update_guard", + "radroots_runtime_source_generations", + "radroots_runtime_source_generations_active_idx", + "radroots_runtime_source_generations_delete_guard", + "radroots_runtime_source_generations_identity_guard", + "radroots_runtime_source_generations_sequence_guard", +]; + /// Ordered, immutable runtime migration plan. pub const MIGRATIONS: &[MigrationDescriptor] = &[ MigrationDescriptor { @@ -635,6 +716,12 @@ pub const MIGRATIONS: &[MigrationDescriptor] = &[ up_sha256: "7bb6466a14c7ca601b76852addf4a19996ee9ffd62436bdac8e9bf631dc46785", owned_objects: RUNTIME_V13_OBJECTS, }, + MigrationDescriptor { + version: 14, + name: "authored_draft_query_metadata", + up_sha256: "3e108ddf9fbbc559b5e36bbdcb1032343b771f91a074aebc62ae8ad6c1521b55", + owned_objects: RUNTIME_V14_OBJECTS, + }, ]; pub(crate) const fn migration_sql(version: u32) -> Option<&'static str> { @@ -652,6 +739,7 @@ pub(crate) const fn migration_sql(version: u32) -> Option<&'static str> { 11 => Some(AUTHORED_OPERATIONS_V11_SQL), 12 => Some(MATERIALIZED_PROJECTION_DOCUMENTS_V12_SQL), 13 => Some(AUTHORED_DRAFT_REVISIONS_V13_SQL), + 14 => Some(AUTHORED_DRAFT_QUERY_V14_SQL), _ => None, } } @@ -695,8 +783,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, 13); - assert_eq!(MIGRATIONS.len(), 13); + assert_eq!(CURRENT_VERSION, 14); + assert_eq!(MIGRATIONS.len(), 14); let migration = MIGRATIONS[8]; assert_eq!(snapshot.schema_version, 1); assert_eq!(snapshot.database, "runtime.sqlite"); @@ -724,7 +812,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(14), None); + assert_eq!(migration_sql(CURRENT_VERSION + 1), None); } #[tokio::test] diff --git a/crates/storage_sqlite/src/migration_draft_query_tests.rs b/crates/storage_sqlite/src/migration_draft_query_tests.rs @@ -0,0 +1,143 @@ +use super::tests::{connection, establish_runtime_version, pragma}; +use super::*; + +const FAILING_V14: &str = concat!( + include_str!("migration/runtime/0014_authored_draft_query_metadata.up.sql"), + "\nINSERT INTO missing_fixture_table VALUES (1);" +); + +async fn seed(connection: &mut SqliteConnection) -> Vec<Vec<u8>> { + let snapshots = vec![ + br#"{"payload_schema":"fixture.composer.v1","payload":[49,46]}"#.to_vec(), + b"{broken historical row".to_vec(), + ]; + for (index, snapshot) in snapshots.iter().enumerate() { + sqlx::query("INSERT INTO radroots_runtime_authored_draft_revisions + (draft_id, revision, author, stage, payload_sha256, created_at_unix_ms, updated_at_unix_ms, snapshot) + VALUES (?, 1, ?, 0, ?, 10, 10, ?)") + .bind([index as u8 + 1; 16].as_slice()).bind([7; 32].as_slice()).bind([3; 32].as_slice()).bind(snapshot) + .execute(&mut *connection).await.unwrap(); + } + snapshots +} +fn plan(current: u32, fail: bool) -> MigrationPlan { + let steps = runtime::MIGRATIONS + .iter() + .take(current as usize) + .map(|step| MigrationStep { + version: step.version(), + sql: if fail && step.version() == 14 { + FAILING_V14 + } else { + runtime::migration_sql(step.version()).unwrap() + }, + owned_objects: step.owned_objects(), + }) + .collect(); + MigrationPlan { + database: RUNTIME_DATABASE, + application_id: RUNTIME_APPLICATION_ID, + set_application_id_sql: SET_RUNTIME_APPLICATION_ID, + minimum_version: runtime::MINIMUM_VERSION, + current_version: current, + steps, + } +} + +#[tokio::test] +async fn draft_metadata_upgrade_preserves_source_and_rejects_prior_schema_policy() { + let mut connection = connection().await; + establish_runtime_version(&mut connection, 13).await; + let snapshots = seed(&mut connection).await; + assert!(matches!( + migrate_runtime(&mut connection, OpenMode::ReadOnly).await, + Err(Error::SchemaMigrationRequired { + current: 14, + actual: 13, + .. + }) + )); + let report = migrate_runtime(&mut connection, OpenMode::ReadWriteExisting) + .await + .unwrap(); + assert_eq!(report.applied(), 1); + assert_eq!(pragma(&mut connection, "user_version").await, 14); + let actual: Vec<Vec<u8>> = sqlx::query_scalar( + "SELECT snapshot FROM radroots_runtime_authored_draft_revisions ORDER BY draft_id", + ) + .fetch_all(&mut connection) + .await + .unwrap(); + assert_eq!(actual, snapshots); + let schemas: Vec<String> = sqlx::query_scalar( + "SELECT payload_schema FROM radroots_runtime_authored_draft_revisions ORDER BY draft_id", + ) + .fetch_all(&mut connection) + .await + .unwrap(); + assert_eq!(schemas, ["fixture.composer.v1", ""]); + assert!(matches!( + migrate(&mut connection, OpenMode::ReadOnly, &plan(13, false)).await, + Err(Error::SchemaTooNew { + supported: 13, + actual: 14, + .. + }) + )); + assert!( + sqlx::query( + "UPDATE radroots_runtime_authored_draft_revisions SET payload_schema = 'changed'" + ) + .execute(&mut connection) + .await + .is_err() + ); + assert_eq!( + migrate_runtime(&mut connection, OpenMode::ReadOnly) + .await + .unwrap() + .applied(), + 0 + ); + connection.close().await.unwrap(); +} + +#[tokio::test] +async fn metadata_migration_failure_rolls_back_columns_index_guards_and_version() { + let mut connection = connection().await; + establish_runtime_version(&mut connection, 13).await; + let snapshots = seed(&mut connection).await; + assert!( + migrate( + &mut connection, + OpenMode::ReadWriteExisting, + &plan(14, true) + ) + .await + .is_err() + ); + assert_eq!(pragma(&mut connection, "user_version").await, 13); + assert!( + sqlx::query("SELECT payload_schema FROM radroots_runtime_authored_draft_revisions") + .fetch_all(&mut connection) + .await + .is_err() + ); + let actual: Vec<Vec<u8>> = sqlx::query_scalar( + "SELECT snapshot FROM radroots_runtime_authored_draft_revisions ORDER BY draft_id", + ) + .fetch_all(&mut connection) + .await + .unwrap(); + assert_eq!(actual, snapshots); + assert!( + sqlx::query("UPDATE radroots_runtime_authored_draft_revisions SET stage = 1") + .execute(&mut connection) + .await + .is_err() + ); + migrate_runtime(&mut connection, OpenMode::ReadWriteExisting) + .await + .unwrap(); + connection.close().await.unwrap(); +} diff --git a/crates/sync/README.md b/crates/sync/README.md @@ -22,3 +22,10 @@ scope; a receipt never proves complete global history. Publication remains disabled while behavior is implemented and qualified in the subsequent Release V1 sync checkpoints. + +Advanced hosts can call `PushRequest::authored_preparation` with an explicitly +captured timestamp to build the same pure preparation used by `prepare_push`. +Retain this value when composing a storage draft submission: composite replay +compares captured timestamps exactly. Building it performs no storage, signing, +clock or network operation. Ordinary preparation identity and replay semantics +remain unchanged. diff --git a/crates/sync/src/push.rs b/crates/sync/src/push.rs @@ -97,6 +97,45 @@ impl PushRequest { }) } + /// Builds the existing pure preparation at a caller-captured stable time. + /// No storage, signer, clock or network is invoked. Retain this value for + /// atomic submission replay, which compares captured timestamps exactly. + pub fn authored_preparation( + &self, + captured_at_unix_ms: u64, + ) -> Result<PrepareAuthoredOperation, Error> { + let (operation_id, artifact_id, delivery_plan_id) = authored_ids(self.operation_id)?; + let operation = + AuthoredOperation::new(operation_id, vec![artifact_id], captured_at_unix_ms) + .map_err(map_storage_error)?; + let artifact = AuthoredArtifact::planned( + artifact_id, + operation_id, + 0, + self.plan(), + captured_at_unix_ms, + ) + .map_err(map_storage_error)?; + let intent = AuthoredDeliveryIntent::new( + delivery_request_id(self.operation_id), + self.targets.clone(), + self.satisfaction.clone(), + self.delivery_deadline_unix_ms, + ) + .map_err(map_storage_error)?; + let delivery_plan = + AuthoredDeliveryPlan::new(delivery_plan_id, artifact_id, intent, captured_at_unix_ms) + .map_err(map_storage_error)?; + PrepareAuthoredOperation::new( + operation, + vec![artifact], + vec![delivery_plan], + authored_push_input_digest(self)?, + captured_at_unix_ms, + ) + .map_err(map_storage_error) + } + pub const fn operation_id(&self) -> SyncId { self.operation_id } @@ -260,31 +299,7 @@ impl Engine { pub async fn prepare_push(&self, request: PushRequest) -> Result<PushPreparation, Error> { let prepared_at = self.clock.now_unix_ms()?; let (operation_id, artifact_id, delivery_plan_id) = authored_ids(request.operation_id)?; - let operation = AuthoredOperation::new(operation_id, vec![artifact_id], prepared_at) - .map_err(map_storage_error)?; - let artifact = - AuthoredArtifact::planned(artifact_id, operation_id, 0, request.plan(), prepared_at) - .map_err(map_storage_error)?; - let intent = AuthoredDeliveryIntent::new( - delivery_request_id(request.operation_id), - request.targets.clone(), - request.satisfaction.clone(), - request.delivery_deadline_unix_ms, - ) - .map_err(map_storage_error)?; - let delivery_plan = - AuthoredDeliveryPlan::new(delivery_plan_id, artifact_id, intent, prepared_at) - .map_err(map_storage_error)?; - let command = AuthoredAtomicCommand::Prepare( - PrepareAuthoredOperation::new( - operation, - vec![artifact], - vec![delivery_plan], - authored_push_input_digest(&request)?, - prepared_at, - ) - .map_err(map_storage_error)?, - ); + let command = AuthoredAtomicCommand::Prepare(request.authored_preparation(prepared_at)?); let receipt = self .storage .execute_authored(command) diff --git a/crates/sync/tests/push_enqueue.rs b/crates/sync/tests/push_enqueue.rs @@ -487,6 +487,7 @@ impl AuthoredAtomicStorage for FaultStorage { radroots_storage::authored_atomic::AuthoredAtomicOutcome::Prepared { .. } => 1, radroots_storage::authored_atomic::AuthoredAtomicOutcome::Artifact(_) => 2, radroots_storage::authored_atomic::AuthoredAtomicOutcome::DeliveryPlan(_) => 3, + radroots_storage::authored_atomic::AuthoredAtomicOutcome::Submitted(_) => 5, }; let fault = self.fault_kind.load(Ordering::Relaxed); if (fault != kind && !(fault == 4 && kind == 1)) @@ -2264,3 +2265,43 @@ async fn sqlite_authored_delivery_retry_survives_reopen() { assert_eq!(completed.plan().state(), AuthoredDeliveryState::Satisfied); assert_eq!(sink.requests.lock().expect("request log").len(), 1); } + +#[test] +fn captured_preparation_matches_engine_without_changing_ordinary_replay_identity() { + use radroots_storage::authored_atomic::AuthoredAtomicCommand; + let signer = Arc::new(MockSigner::new(SignBehavior::Pending)); + let (engine, _) = setup_engine(signer.clone()); + let push = request(91, "wss://captured.example"); + let captured = push.authored_preparation(1_800_000_200_000).unwrap(); + assert!(push.authored_preparation(0).is_err()); + assert!( + block_on(engine.push_status(push.operation_id())) + .unwrap() + .is_none() + ); + assert_eq!(signer.calls.load(Ordering::Relaxed), 0); + let later = push.authored_preparation(1_800_000_200_001).unwrap(); + assert_ne!( + captured, later, + "the composite command compares captured time exactly" + ); + let original_command = AuthoredAtomicCommand::Prepare(captured.clone()); + let later_command = AuthoredAtomicCommand::Prepare(later); + assert_eq!(original_command.commit_id(), later_command.commit_id()); + assert_eq!(original_command.digest(), later_command.digest()); + let prepared = block_on(engine.prepare_push(push)).unwrap(); + assert_eq!(prepared.operation(), captured.operation()); + assert_eq!( + [prepared.artifact()], + captured.artifacts().iter().collect::<Vec<_>>().as_slice() + ); + assert_eq!( + [prepared.delivery_plan()], + captured + .delivery_plans() + .iter() + .collect::<Vec<_>>() + .as_slice() + ); + assert_eq!(signer.calls.load(Ordering::Relaxed), 0); +}