lib

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

commit e6f296604363e5293d6abc6e18bf5eaea2177b9f
parent 77aa55c0894c0142dbae43e272a515fd093eac36
Author: triesap <tyson@radroots.org>
Date:   Mon, 21 Sep 2026 17:47:53 +0000

storage: atomically append paired authored revisions

- Bind two opaque drafts to one all-or-nothing commit
- Preserve exact historical replay and fail closed on conflicts
- Verify concurrency, cancellation and SQLite restart
- Qualify public APIs, owner coverage and workspace consumers

Diffstat:
Mcontracts/api_baselines/radroots_storage.txt | 9+++++++++
Mcontracts/api_baselines/radroots_storage_sqlite.txt | 1+
Acontracts/architecture/decisions/authored_draft_pair.v1.json | 23+++++++++++++++++++++++
Mcrates/storage/README.md | 13+++++++++++++
Mcrates/storage/src/authored_draft.rs | 12++++++++++++
Acrates/storage/src/authored_draft_pair.rs | 45+++++++++++++++++++++++++++++++++++++++++++++
Mcrates/storage/src/lib.rs | 1+
Mcrates/storage/src/memory.rs | 80++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----------------------
Acrates/storage/tests/authored_draft_pair.rs | 242+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/storage_sqlite/README.md | 7+++++++
Mcrates/storage_sqlite/src/authored_draft.rs | 119+++++++++++++++++++++++++++++++++++++++++++++++++++++--------------------------
Acrates/storage_sqlite/src/authored_draft_pair_tests.rs | 153+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acrates/storage_sqlite/tests/authored_draft_pair.rs | 209+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
13 files changed, 852 insertions(+), 62 deletions(-)

diff --git a/contracts/api_baselines/radroots_storage.txt b/contracts/api_baselines/radroots_storage.txt @@ -492,16 +492,24 @@ pub const radroots_storage::authored_draft::AUTHORED_DRAFT_QUERY_LIMIT_MAX: u16 pub const radroots_storage::authored_draft::AUTHORED_DRAFT_SCHEMA_MAX_BYTES: usize pub trait radroots_storage::authored_draft::AuthoredDraftStore: core::marker::Send + core::marker::Sync pub fn radroots_storage::authored_draft::AuthoredDraftStore::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::authored_draft::AuthoredDraftStore::append_authored_draft_pair(&self, radroots_storage::authored_draft_pair::AuthoredDraftPair) -> radroots_transport::source::BoxFuture<'_, core::result::Result<[radroots_storage::authored_draft::DraftAppendReceipt; 2], radroots_storage::Error>> 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::append_authored_draft_pair(&self, radroots_storage::authored_draft_pair::AuthoredDraftPair) -> radroots_transport::source::BoxFuture<'_, core::result::Result<[radroots_storage::authored_draft::DraftAppendReceipt; 2], 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_pair +pub struct radroots_storage::authored_draft_pair::AuthoredDraftPair +impl radroots_storage::authored_draft_pair::AuthoredDraftPair +pub const fn radroots_storage::authored_draft_pair::AuthoredDraftPair::drafts(&self) -> &[radroots_storage::authored_draft::AuthoredDraft; 2] +pub const fn radroots_storage::authored_draft_pair::AuthoredDraftPair::expected_heads(&self) -> &[core::option::Option<radroots_storage::authored_draft::AuthoredDraftRevision>; 2] +pub fn radroots_storage::authored_draft_pair::AuthoredDraftPair::new(radroots_storage::authored_draft::AuthoredDraft, core::option::Option<radroots_storage::authored_draft::AuthoredDraftRevision>, radroots_storage::authored_draft::AuthoredDraft, core::option::Option<radroots_storage::authored_draft::AuthoredDraftRevision>) -> core::result::Result<Self, 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 @@ -961,6 +969,7 @@ pub fn radroots_storage::memory::MemoryStorage::authored_receipt(&self, radroots pub fn radroots_storage::memory::MemoryStorage::execute_authored(&self, radroots_storage::authored_atomic::AuthoredAtomicCommand) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_storage::authored_atomic::AuthoredAtomicReceipt, radroots_storage::Error>> impl radroots_storage::authored_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::append_authored_draft_pair(&self, radroots_storage::authored_draft_pair::AuthoredDraftPair) -> radroots_transport::source::BoxFuture<'_, core::result::Result<[radroots_storage::authored_draft::DraftAppendReceipt; 2], 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>> diff --git a/contracts/api_baselines/radroots_storage_sqlite.txt b/contracts/api_baselines/radroots_storage_sqlite.txt @@ -558,6 +558,7 @@ pub fn radroots_storage_sqlite::SqliteStorage::authored_receipt(&self, radroots_ pub fn radroots_storage_sqlite::SqliteStorage::execute_authored(&self, radroots_storage::authored_atomic::AuthoredAtomicCommand) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_storage::authored_atomic::AuthoredAtomicReceipt, radroots_storage::error::Error>> impl radroots_storage::authored_draft::AuthoredDraftStore for radroots_storage_sqlite::SqliteStorage pub fn radroots_storage_sqlite::SqliteStorage::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::Error>> +pub fn radroots_storage_sqlite::SqliteStorage::append_authored_draft_pair(&self, radroots_storage::authored_draft_pair::AuthoredDraftPair) -> radroots_transport::source::BoxFuture<'_, core::result::Result<[radroots_storage::authored_draft::DraftAppendReceipt; 2], radroots_storage::error::Error>> pub fn radroots_storage_sqlite::SqliteStorage::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::Error>> pub fn radroots_storage_sqlite::SqliteStorage::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::Error>> pub fn radroots_storage_sqlite::SqliteStorage::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::Error>> diff --git a/contracts/architecture/decisions/authored_draft_pair.v1.json b/contracts/architecture/decisions/authored_draft_pair.v1.json @@ -0,0 +1,23 @@ +{ + "schema": "radroots.authored-draft-pair.v1", + "status": "approved", + "owners": [ + "radroots_storage", + "radroots_storage_sqlite" + ], + "authority": "Existing opaque authored revision and backend-neutral atomic storage ownership.", + "request": "Exactly two validated distinct authored draft identifiers with the same author. Each retains its existing expected-head CAS and payload/successor constraints. No application payload interpretation.", + "atomicity": "Memory validates both under its one state lock before insertion. SQLite uses one BEGIN IMMEDIATE transaction and the existing authored-revision tables. Failure or cancellation before commit must retain neither new revision. Lost acknowledgment after commit has an unknown immediate outcome recoverable by exact replay.", + "replay": "Only both exact historical snapshots replay together. Mixed replay/insertion or any mismatch is a conflict and rolls back. Complete historical replay does not establish current head ownership; consumers must check current authority separately before effects.", + "compatibility": "Additive unpublished alpha API with fail-closed default BackendUnavailable. No sequential fallback. Existing single append, stored schema, migration identities, dependencies and limits remain unchanged.", + "non_goals": [ + "coordinate semantics", + "application reservation policy", + "signing", + "delivery", + "public SQL or transaction handles", + "arbitrary batch transactions", + "new database", + "release qualification" + ] +} diff --git a/crates/storage/README.md b/crates/storage/README.md @@ -210,6 +210,19 @@ Features are additive. `--no-default-features` exposes the backend-neutral SPI without an implementation. `memory` and `serde` are supported independently, and `--all-features` enables both. +## Paired authored revisions + +`AuthoredDraftStore::append_authored_draft_pair` installs exactly two distinct +opaque draft revisions for one author in one atomic commit. Both expected heads +and existing per-draft bounds apply. A partial existing pair or either conflict +leaves both heads unchanged. Unsupported backends return `BackendUnavailable`; +there is no sequential-write fallback. + +Exact replay returns both historical snapshots even after their heads advance. +It proves the requested snapshots exist, not current ownership or permission to +sign or deliver. Callers retain application policy and must recheck current +heads before effects. Losing the result after commit cannot establish rollback. + ## Intended consumers - `radroots_storage_sqlite` implements the contracts for native durable state. diff --git a/crates/storage/src/authored_draft.rs b/crates/storage/src/authored_draft.rs @@ -462,6 +462,18 @@ impl DraftAppendReceipt { } pub trait AuthoredDraftStore: Send + Sync { + /// Atomically appends two opaque revisions or replays both exact snapshots. + /// + /// Unsupported backends fail closed; sequential single-row calls are not an + /// implementation. Caller cancellation after commit may lose the receipt + /// without rolling back either row. Replay is historical, not a head lease. + fn append_authored_draft_pair( + &self, + _pair: crate::authored_draft_pair::AuthoredDraftPair, + ) -> BoxFuture<'_, Result<[DraftAppendReceipt; 2], Error>> { + Box::pin(async { Err(Error::BackendUnavailable) }) + } + fn query_authored_drafts( &self, query: AuthoredDraftQuery, diff --git a/crates/storage/src/authored_draft_pair.rs b/crates/storage/src/authored_draft_pair.rs @@ -0,0 +1,45 @@ +//! Atomic installation of two opaque authored revisions, without effects. + +use crate::{ + Error, + authored_draft::{AuthoredDraft, AuthoredDraftRevision}, +}; + +/// A fixed pair of distinct drafts belonging to one author. +/// +/// Both expected heads use the existing single-draft compare-and-swap contract. +/// A backend must install both revisions or neither. Only an exact replay of +/// both existing revisions succeeds as replay; a partially existing pair is a +/// conflict. Historical replay does not establish current head ownership. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct AuthoredDraftPair { + drafts: [AuthoredDraft; 2], + expected_heads: [Option<AuthoredDraftRevision>; 2], +} + +impl AuthoredDraftPair { + pub fn new( + first: AuthoredDraft, + first_expected: Option<AuthoredDraftRevision>, + second: AuthoredDraft, + second_expected: Option<AuthoredDraftRevision>, + ) -> Result<Self, Error> { + first.validate()?; + second.validate()?; + if first.draft_id() == second.draft_id() || first.author() != second.author() { + return Err(Error::InvalidAuthoredDraft); + } + Ok(Self { + drafts: [first, second], + expected_heads: [first_expected, second_expected], + }) + } + + pub const fn drafts(&self) -> &[AuthoredDraft; 2] { + &self.drafts + } + + pub const fn expected_heads(&self) -> &[Option<AuthoredDraftRevision>; 2] { + &self.expected_heads + } +} diff --git a/crates/storage/src/lib.rs b/crates/storage/src/lib.rs @@ -6,6 +6,7 @@ pub mod authored; pub mod authored_atomic; pub mod authored_delivery; pub mod authored_draft; +pub mod authored_draft_pair; pub mod authored_draft_query; pub mod authored_draft_submission; pub mod backup; diff --git a/crates/storage/src/memory.rs b/crates/storage/src/memory.rs @@ -1967,6 +1967,31 @@ impl AuthoredAtomicStorage for MemoryStorage { } impl AuthoredDraftStore for MemoryStorage { + fn append_authored_draft_pair( + &self, + pair: crate::authored_draft_pair::AuthoredDraftPair, + ) -> BoxFuture<'_, Result<[DraftAppendReceipt; 2], Error>> { + Box::pin(async move { + let mut state = self.state()?; + let [first, second] = pair.drafts(); + let [first_expected, second_expected] = *pair.expected_heads(); + let a = draft_append_disposition(&state.authored_drafts, first, first_expected)?; + let b = draft_append_disposition(&state.authored_drafts, second, second_expected)?; + if a != b { + return Err(Error::DraftRevisionConflict); + } + if a == DraftAppendDisposition::Inserted { + state + .authored_drafts + .extend([first.clone(), second.clone()]); + } + Ok([ + DraftAppendReceipt::new(first.clone(), a), + DraftAppendReceipt::new(second.clone(), b), + ]) + }) + } + fn query_authored_drafts( &self, query: crate::authored_draft_query::AuthoredDraftQuery, @@ -2031,29 +2056,10 @@ impl AuthoredDraftStore for MemoryStorage { Box::pin(async move { draft.validate()?; let mut state = self.state()?; - if let Some(existing) = state.authored_drafts.iter().find(|existing| { - existing.draft_id() == draft.draft_id() && existing.revision() == draft.revision() - }) { - return if existing == &draft { - Ok(DraftAppendReceipt::new( - existing.clone(), - DraftAppendDisposition::Replay, - )) - } else { - Err(Error::DraftRevisionConflict) - }; - } - let head = state - .authored_drafts - .iter() - .filter(|existing| existing.draft_id() == draft.draft_id()) - .max_by_key(|existing| existing.revision()); - match (head, expected_head) { - (None, None) if draft.revision() == AuthoredDraftRevision::INITIAL => {} - (Some(previous), Some(expected)) if previous.revision() == expected => { - draft.validate_successor_of(previous)?; - } - _ => return Err(Error::DraftRevisionConflict), + let disposition = + draft_append_disposition(&state.authored_drafts, &draft, expected_head)?; + if disposition == DraftAppendDisposition::Replay { + return Ok(DraftAppendReceipt::new(draft, disposition)); } state.authored_drafts.push(draft.clone()); Ok(DraftAppendReceipt::new( @@ -2150,3 +2156,31 @@ fn require_artifact_claim( } Ok(()) } + +fn draft_append_disposition( + drafts: &[AuthoredDraft], + draft: &AuthoredDraft, + expected_head: Option<AuthoredDraftRevision>, +) -> Result<DraftAppendDisposition, Error> { + if let Some(existing) = drafts.iter().find(|existing| { + existing.draft_id() == draft.draft_id() && existing.revision() == draft.revision() + }) { + return if existing == draft { + Ok(DraftAppendDisposition::Replay) + } else { + Err(Error::DraftRevisionConflict) + }; + } + let head = drafts + .iter() + .filter(|existing| existing.draft_id() == draft.draft_id()) + .max_by_key(|existing| existing.revision()); + match (head, expected_head) { + (None, None) if draft.revision() == AuthoredDraftRevision::INITIAL => {} + (Some(previous), Some(expected)) if previous.revision() == expected => { + draft.validate_successor_of(previous)?; + } + _ => return Err(Error::DraftRevisionConflict), + } + Ok(DraftAppendDisposition::Inserted) +} diff --git a/crates/storage/tests/authored_draft_pair.rs b/crates/storage/tests/authored_draft_pair.rs @@ -0,0 +1,242 @@ +use radroots_storage::{ + Error, + authored_draft::{ + AuthoredDraft, AuthoredDraftId, AuthoredDraftStage, AuthoredDraftStore, + DraftAppendDisposition, + }, + authored_draft_pair::AuthoredDraftPair, +}; +fn draft(id: u8, payload: &[u8]) -> AuthoredDraft { + AuthoredDraft::initial( + AuthoredDraftId::new([id; 16]).unwrap(), + [9; 32], + "fixture.pair.v1", + payload.to_vec(), + AuthoredDraftStage::Draft, + None, + 10, + ) + .unwrap() +} +fn pair(a: AuthoredDraft, b: AuthoredDraft) -> AuthoredDraftPair { + AuthoredDraftPair::new(a, None, b, None).unwrap() +} +async fn contract(store: &dyn AuthoredDraftStore) { + let a = draft(1, b"claim"); + let b = draft(2, b"captured"); + let request = pair(a.clone(), b.clone()); + let receipt = store + .append_authored_draft_pair(request.clone()) + .await + .unwrap(); + assert_eq!(receipt[0].draft(), &a); + assert_eq!(receipt[1].draft(), &b); + assert!( + receipt + .iter() + .all(|v| v.disposition() == DraftAppendDisposition::Inserted) + ); + let replay = store + .append_authored_draft_pair(request.clone()) + .await + .unwrap(); + assert!( + replay + .iter() + .all(|v| v.disposition() == DraftAppendDisposition::Replay) + ); + // First-member conflict and second-member conflict must retain no fresh row. + for candidate in [ + pair(draft(1, b"changed"), draft(3, b"new")), + pair(draft(3, b"new"), draft(2, b"changed")), + pair(a.clone(), draft(3, b"new")), + pair(draft(3, b"new"), b.clone()), + ] { + assert_eq!( + store.append_authored_draft_pair(candidate).await, + Err(Error::DraftRevisionConflict) + ); + assert!( + store + .authored_draft_head(draft(3, b"new").draft_id()) + .await + .unwrap() + .is_none() + ); + } + let a2 = a + .successor(b"next claim".to_vec(), AuthoredDraftStage::Draft, None, 11) + .unwrap(); + let b2 = b + .successor( + b"next capture".to_vec(), + AuthoredDraftStage::Draft, + None, + 11, + ) + .unwrap(); + let incorrect = + AuthoredDraftPair::new(a2.clone(), Some(a.revision()), b2.clone(), None).unwrap(); + assert_eq!( + store.append_authored_draft_pair(incorrect).await, + Err(Error::DraftRevisionConflict) + ); + assert_eq!( + store.authored_draft_head(a.draft_id()).await.unwrap(), + Some(a.clone()) + ); + let incorrect = + AuthoredDraftPair::new(a2.clone(), None, b2.clone(), Some(b.revision())).unwrap(); + assert_eq!( + store.append_authored_draft_pair(incorrect).await, + Err(Error::DraftRevisionConflict) + ); + let next = AuthoredDraftPair::new( + a2.clone(), + Some(a.revision()), + b2.clone(), + Some(b.revision()), + ) + .unwrap(); + store.append_authored_draft_pair(next).await.unwrap(); + // Replay returns immutable history and never rewinds the current head. + let replay = store.append_authored_draft_pair(request).await.unwrap(); + assert_eq!(replay[0].draft(), &a); + assert_eq!( + store.authored_draft_head(a.draft_id()).await.unwrap(), + Some(a2) + ); + assert_eq!( + store.authored_draft_head(b.draft_id()).await.unwrap(), + Some(b2) + ); + // The untouched single-row API retains its original behavior. + let c = draft(4, b"single"); + store.append_authored_draft(c.clone(), None).await.unwrap(); + assert_eq!( + store + .append_authored_draft(c, None) + .await + .unwrap() + .disposition(), + DraftAppendDisposition::Replay + ); +} + +#[test] +fn memory_atomic_pair_conflict_replay_and_successor_contract() { + futures_executor::block_on(contract(&radroots_storage::memory::MemoryStorage::default())); +} +#[test] +fn pair_rejects_duplicate_keys_and_foreign_author() { + let a = draft(1, b"a"); + assert_eq!( + AuthoredDraftPair::new(a.clone(), None, a.clone(), None), + Err(Error::InvalidAuthoredDraft) + ); + let foreign = AuthoredDraft::initial( + AuthoredDraftId::new([2; 16]).unwrap(), + [8; 32], + "fixture.pair.v1", + b"foreign".to_vec(), + AuthoredDraftStage::Draft, + None, + 10, + ) + .unwrap(); + assert_eq!( + AuthoredDraftPair::new(a, None, foreign, None), + Err(Error::InvalidAuthoredDraft) + ); +} +#[test] +fn concurrent_memory_reservations_commit_one_complete_pair() { + use std::sync::{Arc, Barrier}; + let store = Arc::new(radroots_storage::memory::MemoryStorage::default()); + let barrier = Arc::new(Barrier::new(2)); + let tasks: Vec<_> = + (2..=3) + .map(|id| { + let store = store.clone(); + let barrier = barrier.clone(); + std::thread::spawn(move || { + barrier.wait(); + ( + id, + futures_executor::block_on(store.append_authored_draft_pair(pair( + draft(1, &[id]), + draft(id, b"capture"), + ))), + ) + }) + }) + .collect(); + let results: Vec<_> = tasks.into_iter().map(|t| t.join().unwrap()).collect(); + assert_eq!(results.iter().filter(|(_, v)| v.is_ok()).count(), 1); + for (id, result) in results { + assert_eq!( + futures_executor::block_on(store.authored_draft_head(draft(id, b"lookup").draft_id())) + .unwrap() + .is_some(), + result.is_ok() + ); + } +} + +struct Unsupported(radroots_storage::memory::MemoryStorage); +impl AuthoredDraftStore for Unsupported { + fn query_authored_drafts( + &self, + q: radroots_storage::authored_draft_query::AuthoredDraftQuery, + ) -> radroots_transport::BoxFuture< + '_, + Result<radroots_storage::authored_draft_query::AuthoredDraftPage, Error>, + > { + self.0.query_authored_drafts(q) + } + fn append_authored_draft( + &self, + d: AuthoredDraft, + e: Option<radroots_storage::authored_draft::AuthoredDraftRevision>, + ) -> radroots_transport::BoxFuture< + '_, + Result<radroots_storage::authored_draft::DraftAppendReceipt, Error>, + > { + self.0.append_authored_draft(d, e) + } + fn authored_draft_head( + &self, + id: AuthoredDraftId, + ) -> radroots_transport::BoxFuture<'_, Result<Option<AuthoredDraft>, Error>> { + self.0.authored_draft_head(id) + } + fn authored_draft_revision( + &self, + id: AuthoredDraftId, + r: radroots_storage::authored_draft::AuthoredDraftRevision, + ) -> radroots_transport::BoxFuture<'_, Result<Option<AuthoredDraft>, Error>> { + self.0.authored_draft_revision(id, r) + } + fn authored_draft_heads( + &self, + a: [u8; 32], + n: u16, + ) -> radroots_transport::BoxFuture<'_, Result<Vec<AuthoredDraft>, Error>> { + self.0.authored_draft_heads(a, n) + } +} +#[test] +fn unsupported_backend_never_falls_back_to_separate_writes() { + let store = Unsupported(radroots_storage::memory::MemoryStorage::default()); + assert_eq!( + futures_executor::block_on( + store.append_authored_draft_pair(pair(draft(1, b"a"), draft(2, b"b"))) + ), + Err(Error::BackendUnavailable) + ); + assert!( + futures_executor::block_on(store.authored_draft_heads([9; 32], 10)) + .unwrap() + .is_empty() + ); +} diff --git a/crates/storage_sqlite/README.md b/crates/storage_sqlite/README.md @@ -168,3 +168,10 @@ let storage = SqliteStorage::open( # Ok(()) # } ``` + +Atomic paired authored revisions use the existing runtime revision tables in one +`BEGIN IMMEDIATE` transaction. Both insertions commit together; mismatched or +partial replay rolls back. Complete replay returns immutable historical rows and +does not establish current application authority. No migration or second storage +owner is introduced. Callers may recover an unknown commit result by replaying +the exact pair. diff --git a/crates/storage_sqlite/src/authored_draft.rs b/crates/storage_sqlite/src/authored_draft.rs @@ -18,6 +18,38 @@ mod query; mod bounded_row; impl AuthoredDraftStore for SqliteStorage { + fn append_authored_draft_pair( + &self, + pair: radroots_storage::authored_draft_pair::AuthoredDraftPair, + ) -> BoxFuture<'_, Result<[DraftAppendReceipt; 2], Error>> { + Box::pin(async move { + if self.event_mode() == radroots_storage::status::EventStoreMode::ReadOnly { + return Err(Error::BackendUnavailable); + } + let mut transaction = self + .pool() + .begin_with("BEGIN IMMEDIATE") + .await + .map_err(map_backend)?; + let [first, second] = pair.drafts(); + let [first_expected, second_expected] = *pair.expected_heads(); + let a = append_tx(&mut transaction, first, first_expected).await?; + let b = append_tx(&mut transaction, second, second_expected).await?; + if a != b { + return Err(Error::DraftRevisionConflict); + } + if a == DraftAppendDisposition::Replay { + transaction.rollback().await.map_err(map_backend)?; + } else { + transaction.commit().await.map_err(map_backend)?; + } + Ok([ + DraftAppendReceipt::new(first.clone(), a), + DraftAppendReceipt::new(second.clone(), b), + ]) + }) + } + fn query_authored_drafts( &self, query: radroots_storage::authored_draft_query::AuthoredDraftQuery, @@ -41,47 +73,13 @@ impl AuthoredDraftStore for SqliteStorage { .begin_with("BEGIN IMMEDIATE") .await .map_err(map_backend)?; - if let Some(row) = bounded_row::load( - &mut *transaction, - draft.draft_id().as_bytes(), - Some(draft.revision()), - ) - .await? - { - let existing = decode_row(&row)?; + let disposition = append_tx(&mut transaction, &draft, expected_head).await?; + if disposition == DraftAppendDisposition::Replay { transaction.rollback().await.map_err(map_backend)?; - return if existing == draft { - Ok(DraftAppendReceipt::new( - existing, - DraftAppendDisposition::Replay, - )) - } else { - Err(Error::DraftRevisionConflict) - }; - } - - let head = bounded_row::load(&mut *transaction, draft.draft_id().as_bytes(), None) - .await? - .as_ref() - .map(decode_row) - .transpose()?; - match (head.as_ref(), expected_head) { - (None, None) if draft.revision() == AuthoredDraftRevision::INITIAL => {} - (Some(previous), Some(expected)) if previous.revision() == expected => { - draft.validate_successor_of(previous)?; - } - _ => { - let _ = transaction.rollback().await; - return Err(Error::DraftRevisionConflict); - } + } else { + transaction.commit().await.map_err(map_backend)?; } - - insert_draft_tx(&mut transaction, &draft).await?; - transaction.commit().await.map_err(map_backend)?; - Ok(DraftAppendReceipt::new( - draft, - DraftAppendDisposition::Inserted, - )) + Ok(DraftAppendReceipt::new(draft, disposition)) }) } @@ -665,3 +663,46 @@ pub(crate) async fn insert_draft_tx( .map_err(map_backend)?; Ok(()) } + +async fn append_tx( + transaction: &mut sqlx::Transaction<'_, sqlx::Sqlite>, + draft: &AuthoredDraft, + expected_head: Option<AuthoredDraftRevision>, +) -> Result<DraftAppendDisposition, Error> { + if let Some(row) = bounded_row::load( + &mut **transaction, + draft.draft_id().as_bytes(), + Some(draft.revision()), + ) + .await? + { + let existing = decode_row(&row)?; + return if existing == *draft { + Ok(DraftAppendDisposition::Replay) + } else { + Err(Error::DraftRevisionConflict) + }; + } + + let head = bounded_row::load(&mut **transaction, draft.draft_id().as_bytes(), None) + .await? + .as_ref() + .map(decode_row) + .transpose()?; + match (head.as_ref(), expected_head) { + (None, None) if draft.revision() == AuthoredDraftRevision::INITIAL => {} + (Some(previous), Some(expected)) if previous.revision() == expected => { + draft.validate_successor_of(previous)?; + } + _ => { + return Err(Error::DraftRevisionConflict); + } + } + + insert_draft_tx(transaction, draft).await?; + Ok(DraftAppendDisposition::Inserted) +} + +#[cfg(test)] +#[path = "authored_draft_pair_tests.rs"] +mod pair_tests; diff --git a/crates/storage_sqlite/src/authored_draft_pair_tests.rs b/crates/storage_sqlite/src/authored_draft_pair_tests.rs @@ -0,0 +1,153 @@ +use super::*; +use crate::{OpenMode, OpenOptions, Paths}; +use radroots_storage::{authored_draft_pair::AuthoredDraftPair, event::SourceGeneration}; +use std::{future::poll_fn, task::Poll}; +fn draft(id: u8, schema: &str) -> AuthoredDraft { + AuthoredDraft::initial( + AuthoredDraftId::new([id; 16]).unwrap(), + [9; 32], + schema, + b"fixture".to_vec(), + AuthoredDraftStage::Draft, + None, + 10, + ) + .unwrap() +} +async fn open(path: &std::path::Path, mode: OpenMode) -> SqliteStorage { + let options = OpenOptions::new(Paths::from_directory(path).unwrap(), mode); + let options = if mode == OpenMode::Create { + options + .with_source_generation(SourceGeneration::new([7; 32]).unwrap(), 10) + .unwrap() + } else { + options + }; + SqliteStorage::open(options).await.unwrap() +} +#[tokio::test] +async fn second_sql_insert_failure_rolls_back_first_revision() { + let dir = tempfile::tempdir().unwrap(); + let store = open(dir.path(), OpenMode::Create).await; + sqlx::query("CREATE TRIGGER fixture_pair_abort BEFORE INSERT ON radroots_runtime_authored_draft_revisions WHEN NEW.payload_schema='fixture.reject.v1' BEGIN SELECT RAISE(ABORT,'fixture pair rejection'); END").execute(store.pool()).await.unwrap(); + let request = AuthoredDraftPair::new( + draft(1, "fixture.pair.v1"), + None, + draft(2, "fixture.reject.v1"), + None, + ) + .unwrap(); + assert!(store.append_authored_draft_pair(request).await.is_err()); + assert!( + store + .authored_draft_heads([9; 32], 10) + .await + .unwrap() + .is_empty() + ); + sqlx::query("DROP TRIGGER fixture_pair_abort") + .execute(store.pool()) + .await + .unwrap(); + store.close().await.unwrap(); + let store = open(dir.path(), OpenMode::ReadWriteExisting).await; + assert!( + store + .authored_draft_heads([9; 32], 10) + .await + .unwrap() + .is_empty() + ); + store.close().await.unwrap(); +} +#[tokio::test] +async fn cancelling_real_pair_futures_never_leaves_one_member_after_reopen() { + let dir = tempfile::tempdir().unwrap(); + let store = open(dir.path(), OpenMode::Create).await; + let mut cancelled = 0; + let mut completed = 0; + let mut durable = Vec::new(); + // Cancel at successive real Pending boundaries, including caller loss near + // commit. No test hook or surrogate transaction replaces the public method. + for boundary in 1..=64u8 { + let a = draft(boundary * 2, "fixture.pair.v1"); + let b = draft(boundary * 2 + 1, "fixture.pair.v1"); + let mut future = store.append_authored_draft_pair( + AuthoredDraftPair::new(a.clone(), None, b.clone(), None).unwrap(), + ); + let mut pending = 0; + let acknowledged = poll_fn(|cx| match future.as_mut().poll(cx) { + Poll::Ready(result) => { + result.unwrap(); + Poll::Ready(true) + } + Poll::Pending => { + pending += 1; + if pending == boundary { + Poll::Ready(false) + } else { + Poll::Pending + } + } + }) + .await; + drop(future); + if acknowledged { + completed += 1; + } else { + cancelled += 1; + } + // Drain any queued commit/rollback before observing both heads. + store + .pool() + .begin_with("BEGIN IMMEDIATE") + .await + .unwrap() + .rollback() + .await + .unwrap(); + let first = store.authored_draft_head(a.draft_id()).await.unwrap(); + let second = store.authored_draft_head(b.draft_id()).await.unwrap(); + assert_eq!(first.is_some(), second.is_some(), "boundary {boundary}"); + if acknowledged { + assert_eq!(first, Some(a.clone())); + } + durable.push((a, b, first.is_some())); + } + assert!( + cancelled > 0 && completed > 0, + "cancelled={cancelled}, completed={completed}" + ); + store.close().await.unwrap(); + let store = open(dir.path(), OpenMode::ReadWriteExisting).await; + for (a, b, present) in durable { + assert_eq!( + store + .authored_draft_head(a.draft_id()) + .await + .unwrap() + .is_some(), + present + ); + assert_eq!( + store + .authored_draft_head(b.draft_id()) + .await + .unwrap() + .is_some(), + present + ); + if present { + let receipts = store + .append_authored_draft_pair(AuthoredDraftPair::new(a, None, b, None).unwrap()) + .await + .unwrap(); + assert!( + receipts + .iter() + .all(|r| r.disposition() == DraftAppendDisposition::Replay) + ); + } + } + store.close().await.unwrap(); +} diff --git a/crates/storage_sqlite/tests/authored_draft_pair.rs b/crates/storage_sqlite/tests/authored_draft_pair.rs @@ -0,0 +1,209 @@ +use radroots_storage::{ + Error, + authored_draft::{ + AuthoredDraft, AuthoredDraftId, AuthoredDraftStage, AuthoredDraftStore, + DraftAppendDisposition, + }, + authored_draft_pair::AuthoredDraftPair, +}; +fn draft(id: u8, payload: &[u8]) -> AuthoredDraft { + AuthoredDraft::initial( + AuthoredDraftId::new([id; 16]).unwrap(), + [9; 32], + "fixture.pair.v1", + payload.to_vec(), + AuthoredDraftStage::Draft, + None, + 10, + ) + .unwrap() +} +fn pair(a: AuthoredDraft, b: AuthoredDraft) -> AuthoredDraftPair { + AuthoredDraftPair::new(a, None, b, None).unwrap() +} +async fn contract(store: &dyn AuthoredDraftStore) { + let a = draft(1, b"claim"); + let b = draft(2, b"captured"); + let request = pair(a.clone(), b.clone()); + let receipt = store + .append_authored_draft_pair(request.clone()) + .await + .unwrap(); + assert_eq!(receipt[0].draft(), &a); + assert_eq!(receipt[1].draft(), &b); + assert!( + receipt + .iter() + .all(|v| v.disposition() == DraftAppendDisposition::Inserted) + ); + let replay = store + .append_authored_draft_pair(request.clone()) + .await + .unwrap(); + assert!( + replay + .iter() + .all(|v| v.disposition() == DraftAppendDisposition::Replay) + ); + // First-member conflict and second-member conflict must retain no fresh row. + for candidate in [ + pair(draft(1, b"changed"), draft(3, b"new")), + pair(draft(3, b"new"), draft(2, b"changed")), + pair(a.clone(), draft(3, b"new")), + pair(draft(3, b"new"), b.clone()), + ] { + assert_eq!( + store.append_authored_draft_pair(candidate).await, + Err(Error::DraftRevisionConflict) + ); + assert!( + store + .authored_draft_head(draft(3, b"new").draft_id()) + .await + .unwrap() + .is_none() + ); + } + let a2 = a + .successor(b"next claim".to_vec(), AuthoredDraftStage::Draft, None, 11) + .unwrap(); + let b2 = b + .successor( + b"next capture".to_vec(), + AuthoredDraftStage::Draft, + None, + 11, + ) + .unwrap(); + let incorrect = + AuthoredDraftPair::new(a2.clone(), Some(a.revision()), b2.clone(), None).unwrap(); + assert_eq!( + store.append_authored_draft_pair(incorrect).await, + Err(Error::DraftRevisionConflict) + ); + assert_eq!( + store.authored_draft_head(a.draft_id()).await.unwrap(), + Some(a.clone()) + ); + let incorrect = + AuthoredDraftPair::new(a2.clone(), None, b2.clone(), Some(b.revision())).unwrap(); + assert_eq!( + store.append_authored_draft_pair(incorrect).await, + Err(Error::DraftRevisionConflict) + ); + let next = AuthoredDraftPair::new( + a2.clone(), + Some(a.revision()), + b2.clone(), + Some(b.revision()), + ) + .unwrap(); + store.append_authored_draft_pair(next).await.unwrap(); + // Replay returns immutable history and never rewinds the current head. + let replay = store.append_authored_draft_pair(request).await.unwrap(); + assert_eq!(replay[0].draft(), &a); + assert_eq!( + store.authored_draft_head(a.draft_id()).await.unwrap(), + Some(a2) + ); + assert_eq!( + store.authored_draft_head(b.draft_id()).await.unwrap(), + Some(b2) + ); + // The untouched single-row API retains its original behavior. + let c = draft(4, b"single"); + store.append_authored_draft(c.clone(), None).await.unwrap(); + assert_eq!( + store + .append_authored_draft(c, None) + .await + .unwrap() + .disposition(), + DraftAppendDisposition::Replay + ); +} + +use radroots_storage::event::SourceGeneration; +use radroots_storage_sqlite::{OpenMode, OpenOptions, Paths, SqliteStorage}; +async fn open(root: &std::path::Path, mode: OpenMode) -> SqliteStorage { + let options = OpenOptions::new(Paths::from_directory(root).unwrap(), mode); + let options = if mode == OpenMode::Create { + options + .with_source_generation(SourceGeneration::new([7; 32]).unwrap(), 10) + .unwrap() + } else { + options + }; + SqliteStorage::open(options).await.unwrap() +} +#[tokio::test] +async fn sqlite_pair_contract_persists_across_reopen_and_read_only_refuses() { + let dir = tempfile::tempdir().unwrap(); + let store = open(dir.path(), OpenMode::Create).await; + contract(&store).await; + store.close().await.unwrap(); + let store = open(dir.path(), OpenMode::ReadOnly).await; + for id in [1, 2] { + assert_eq!( + store + .authored_draft_head(draft(id, b"lookup").draft_id()) + .await + .unwrap() + .unwrap() + .revision() + .get(), + 2 + ); + } + assert_eq!( + store + .append_authored_draft_pair(pair(draft(5, b"a"), draft(6, b"b"))) + .await, + Err(Error::BackendUnavailable) + ); + store.close().await.unwrap(); + let store = open(dir.path(), OpenMode::ReadWriteExisting).await; + let replay = store + .append_authored_draft_pair(pair(draft(1, b"claim"), draft(2, b"captured"))) + .await + .unwrap(); + assert!( + replay + .iter() + .all(|v| v.disposition() == DraftAppendDisposition::Replay) + ); + store.close().await.unwrap(); +} +#[tokio::test] +async fn concurrent_sqlite_reservations_keep_exactly_one_complete_pair() { + let dir = tempfile::tempdir().unwrap(); + let store = open(dir.path(), OpenMode::Create).await; + let (a, b) = tokio::join!( + store.append_authored_draft_pair(pair(draft(1, b"a"), draft(2, b"capture a"))), + store.append_authored_draft_pair(pair(draft(1, b"b"), draft(3, b"capture b"))) + ); + assert_ne!(a.is_ok(), b.is_ok()); + assert_eq!( + store + .authored_draft_head(draft(2, b"lookup").draft_id()) + .await + .unwrap() + .is_some(), + a.is_ok() + ); + assert_eq!( + store + .authored_draft_head(draft(3, b"lookup").draft_id()) + .await + .unwrap() + .is_some(), + b.is_ok() + ); + store.close().await.unwrap(); + let store = open(dir.path(), OpenMode::ReadWriteExisting).await; + assert_eq!( + store.authored_draft_heads([9; 32], 10).await.unwrap().len(), + 2 + ); + store.close().await.unwrap(); +}