lib

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

commit 7afc4e41e3ad5b7125f09b43680ad2efeaf40f46
parent e783db193d63044c31d93d23148b8fa87cc939cc
Author: triesap <tyson@radroots.org>
Date:   Wed,  5 Aug 2026 03:19:00 +0000

refactor(sync): persist complete authored intent

- accept exact authored plans and caller-owned delivery deadlines
- atomically persist parent artifact and delivery intent before signing
- expose durable preparation status with exact idempotent replay
- prove memory and SQLite crash recovery without external effects


Diffstat:
Mcrates/storage/src/memory.rs | 5+++--
Mcrates/sync/src/lib.rs | 2+-
Mcrates/sync/src/policy.rs | 18+++++++++++++++---
Mcrates/sync/src/push.rs | 296++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---------
Mcrates/sync/tests/engine_composition.rs | 4++--
Mcrates/sync/tests/push_enqueue.rs | 241+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++----------------
Mcrates/sync/tests/reliability_scenarios.rs | 19++++++++++++++-----
7 files changed, 492 insertions(+), 93 deletions(-)

diff --git a/crates/storage/src/memory.rs b/crates/storage/src/memory.rs @@ -1323,8 +1323,9 @@ impl AuthoredAtomicStorage for MemoryStorage { if existing.digest() != command.digest() { return Err(Error::AtomicCommitConflict); } - return AuthoredAtomicReceipt::new( - &command, + return AuthoredAtomicReceipt::from_durable_parts( + existing.commit_id(), + existing.digest(), AtomicCommitDisposition::Replay, existing.committed_at_unix_ms(), existing.outcome().clone(), diff --git a/crates/sync/src/lib.rs b/crates/sync/src/lib.rs @@ -12,5 +12,5 @@ pub mod status; pub use engine::Engine; pub use policy::Error; pub use pull::{PullReceipt, PullRequest}; -pub use push::{PushReceipt, PushRequest}; +pub use push::{PushPreparation, PushReceipt, PushRequest, PushStatus}; pub use status::SyncStatus; diff --git a/crates/sync/src/policy.rs b/crates/sync/src/policy.rs @@ -5,7 +5,7 @@ use std::sync::Arc; use radroots_signing::Signer; use radroots_storage::{ EventStore, Journal, Outbox, ProjectionStore, atomic::AtomicStorage, - status::StorageStatusProvider, + authored_atomic::AuthoredAtomicStorage, status::StorageStatusProvider, }; use radroots_transport::{EventSink, EventSource}; @@ -15,12 +15,24 @@ const MAX_OPERATION_TIMEOUT_MS: u64 = 86_400_000; /// Exact backend-neutral storage capability required by sync orchestration. pub trait SyncStorage: - EventStore + Journal + Outbox + ProjectionStore + AtomicStorage + StorageStatusProvider + EventStore + + Journal + + Outbox + + ProjectionStore + + AtomicStorage + + AuthoredAtomicStorage + + StorageStatusProvider { } impl<T> SyncStorage for T where - T: EventStore + Journal + Outbox + ProjectionStore + AtomicStorage + StorageStatusProvider + T: EventStore + + Journal + + Outbox + + ProjectionStore + + AtomicStorage + + AuthoredAtomicStorage + + StorageStatusProvider { } diff --git a/crates/sync/src/push.rs b/crates/sync/src/push.rs @@ -1,18 +1,24 @@ //! Signing, durable enqueue, delivery, and satisfaction orchestration. -use radroots_event::{EventDraft, admission::RawEvent}; -use radroots_event_codec::verify::{self, Nip01SignatureVerifier}; +use radroots_event::admission::RawEvent; +use radroots_event_codec::{ + authoring::AuthoredEventPlan, + verify::{self, Nip01SignatureVerifier}, +}; use radroots_protocol::runtime::v1::OperationId; use radroots_signing::{ - Actor, + Actor, AuthoredArtifactId as SigningArtifactId, SigningIntentId, SigningOperationId, request::{CancellationPolicy, SignPolicy}, }; use radroots_storage::{ Journal, Outbox, atomic::{ - AtomicCommit, AtomicCommitDigest, AtomicCommitId, AtomicCommitOutcome, AtomicWorkflow, - CommitEnqueued, CommitSigned, + AtomicCommit, AtomicCommitDigest, AtomicCommitDisposition, AtomicCommitId, + AtomicCommitOutcome, AtomicWorkflow, CommitEnqueued, CommitSigned, }, + authored::{AuthoredArtifact, AuthoredArtifactId, AuthoredOperation}, + authored_atomic::{AuthoredAtomicCommand, AuthoredAtomicOutcome, PrepareAuthoredOperation}, + authored_delivery::{AuthoredDeliveryIntent, AuthoredDeliveryPlan, AuthoredDeliveryPlanId}, event::EventAdmission, journal::{ IdempotencyDigest, IdempotencyKey, JournalStage, JournalState, OperationInstanceId, @@ -47,9 +53,10 @@ pub struct PushRequest { operation_id: SyncId, idempotency_key: IdempotencyKey, actor: Actor, - draft: EventDraft, + plan: AuthoredEventPlan, targets: TargetSet, satisfaction: SatisfactionPolicy, + delivery_deadline_unix_ms: u64, cancellation: CancellationPolicy, } @@ -59,21 +66,26 @@ impl PushRequest { operation_id: SyncId, idempotency_key: IdempotencyKey, actor: Actor, - draft: EventDraft, + plan: AuthoredEventPlan, targets: TargetSet, satisfaction: SatisfactionPolicy, + delivery_deadline_unix_ms: u64, cancellation: CancellationPolicy, ) -> Result<Self, Error> { - if satisfaction.validate_for(&targets).is_err() { + if satisfaction.validate_for(&targets).is_err() + || delivery_deadline_unix_ms == 0 + || actor.public_key() != *plan.author() + { return Err(Error::InvalidPushRequest); } Ok(Self { operation_id, idempotency_key, actor, - draft, + plan, targets, satisfaction, + delivery_deadline_unix_ms, cancellation, }) } @@ -87,8 +99,8 @@ impl PushRequest { pub const fn actor(&self) -> &Actor { &self.actor } - pub const fn draft(&self) -> &EventDraft { - &self.draft + pub const fn plan(&self) -> &AuthoredEventPlan { + &self.plan } pub const fn targets(&self) -> &TargetSet { &self.targets @@ -96,6 +108,9 @@ impl PushRequest { pub const fn satisfaction(&self) -> &SatisfactionPolicy { &self.satisfaction } + pub const fn delivery_deadline_unix_ms(&self) -> u64 { + self.delivery_deadline_unix_ms + } pub const fn cancellation(&self) -> CancellationPolicy { self.cancellation } @@ -108,14 +123,59 @@ impl core::fmt::Debug for PushRequest { .field("operation_id", &self.operation_id) .field("idempotency_key", &self.idempotency_key) .field("actor", &self.actor) - .field("draft", &"[redacted frozen event draft]") + .field("plan", &"[redacted exact authored plan]") .field("targets", &self.targets) .field("satisfaction", &self.satisfaction) + .field("delivery_deadline_unix_ms", &self.delivery_deadline_unix_ms) .field("cancellation", &self.cancellation) .finish() } } +/// Complete durable intent created before any signer or transport effect. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct PushPreparation { + operation: AuthoredOperation, + artifact: AuthoredArtifact, + delivery_plan: AuthoredDeliveryPlan, + replay: bool, +} + +impl PushPreparation { + pub const fn operation(&self) -> &AuthoredOperation { + &self.operation + } + pub const fn artifact(&self) -> &AuthoredArtifact { + &self.artifact + } + pub const fn delivery_plan(&self) -> &AuthoredDeliveryPlan { + &self.delivery_plan + } + pub const fn is_replay(&self) -> bool { + self.replay + } +} + +/// Current durable authored state for one push operation. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct PushStatus { + operation: AuthoredOperation, + artifact: AuthoredArtifact, + delivery_plan: AuthoredDeliveryPlan, +} + +impl PushStatus { + pub const fn operation(&self) -> &AuthoredOperation { + &self.operation + } + pub const fn artifact(&self) -> &AuthoredArtifact { + &self.artifact + } + pub const fn delivery_plan(&self) -> &AuthoredDeliveryPlan { + &self.delivery_plan + } +} + /// Durable result after signing and atomic outbox enqueue. #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] #[derive(Clone, Debug, Eq, PartialEq)] @@ -191,6 +251,110 @@ impl DeliveryRunReceipt { } impl Engine { + /// Atomically persists the complete parent, exact plan, and delivery intent. + /// + /// This method invokes neither a signer nor a transport. Exact request + /// replay returns the original durable preparation; conflicting reuse of + /// the operation identity fails closed. + 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 receipt = self + .storage + .execute_authored(command) + .await + .map_err(map_storage_error)?; + let AuthoredAtomicOutcome::Prepared { + operation, + artifacts, + delivery_plans, + } = receipt.outcome() + else { + return Err(Error::StorageFailed); + }; + let [artifact] = artifacts.as_slice() else { + return Err(Error::StorageFailed); + }; + let [delivery_plan] = delivery_plans.as_slice() else { + return Err(Error::StorageFailed); + }; + if operation.operation_id() != operation_id + || artifact.artifact_id() != artifact_id + || delivery_plan.plan_id() != delivery_plan_id + { + return Err(Error::StorageFailed); + } + Ok(PushPreparation { + operation: operation.clone(), + artifact: artifact.clone(), + delivery_plan: delivery_plan.clone(), + replay: receipt.disposition() == AtomicCommitDisposition::Replay, + }) + } + + /// Loads the complete durable state for one prepared push operation. + pub async fn push_status(&self, operation_id: SyncId) -> Result<Option<PushStatus>, Error> { + let (operation_id, expected_artifact_id, expected_plan_id) = authored_ids(operation_id)?; + let Some(operation) = self + .storage + .authored_operation(operation_id) + .await + .map_err(map_storage_error)? + else { + return Ok(None); + }; + if operation.artifact_ids() != [expected_artifact_id] { + return Err(Error::StorageFailed); + } + let artifact = self + .storage + .authored_artifact(expected_artifact_id) + .await + .map_err(map_storage_error)? + .ok_or(Error::StorageFailed)?; + let delivery_plan = self + .storage + .authored_delivery_plan(expected_plan_id) + .await + .map_err(map_storage_error)? + .ok_or(Error::StorageFailed)?; + if artifact.operation_id() != operation.operation_id() + || delivery_plan.artifact_id() != artifact.artifact_id() + { + return Err(Error::StorageFailed); + } + Ok(Some(PushStatus { + operation, + artifact, + delivery_plan, + })) + } + /// Authorizes, signs, verifies, and durably enqueues one outbound event. /// /// Dropping the future before the final atomic enqueue leaves either a @@ -198,11 +362,12 @@ impl Engine { /// cancellation cannot claim rollback; replay returns the durable outbox. pub async fn sign_and_enqueue(&self, request: PushRequest) -> Result<PushReceipt, Error> { let signer = self.signer.as_deref().ok_or(Error::MissingSigner)?; + let preparation = self.prepare_push(request.clone()).await?; let instance_id = OperationInstanceId::new(*request.operation_id.as_bytes()) .map_err(map_storage_error)?; let item_id = OutboxItemId::new(*request.operation_id.as_bytes()).map_err(map_storage_error)?; - let input_digest = push_input_digest(&request); + let input_digest = push_input_digest(&request)?; if let Some(existing) = Journal::operation(self.storage.as_ref(), instance_id) .await @@ -258,23 +423,28 @@ impl Engine { let sign_deadline_ms = self .deadlines .deadline_unix_ms(OperationKind::Sign, self.clock.now_unix_ms()?)?; - let sign_deadline_seconds = sign_deadline_ms - .checked_add(999) - .ok_or(Error::DeadlineOverflow)? - / 1_000; + let signing_operation = SigningOperationId::new(*request.operation_id.as_bytes()) + .map_err(|_| Error::InvalidPushRequest)?; + let signing_artifact = + SigningArtifactId::new(*preparation.artifact().artifact_id().as_bytes()) + .map_err(|_| Error::InvalidPushRequest)?; let sign_request = radroots_signing::SignRequest::new( OperationId::SyncPush, + SigningIntentId::new(signing_operation, signing_artifact), request.actor.clone(), - request.draft.clone(), - SignPolicy::new(sign_deadline_seconds, request.cancellation) + request.plan.clone(), + SignPolicy::new(sign_deadline_ms, request.cancellation) .map_err(|_| Error::InvalidPushRequest)?, ) .map_err(|_| Error::InvalidPushRequest)?; - let signed = signer - .sign(sign_request) - .await - .map_err(|_| Error::SignerFailed)?; - if signed.completed_at_unix() > sign_deadline_seconds { + let signed = signer.sign(sign_request).await.map_err(|error| { + if error.kind() == radroots_signing::error::Kind::DeadlineExceeded { + Error::SignerDeadlineExceeded + } else { + Error::SignerFailed + } + })?; + if signed.completed_at_unix_ms() > sign_deadline_ms { return Err(Error::SignerDeadlineExceeded); } let event = signed.signed_event().clone(); @@ -311,15 +481,12 @@ impl Engine { }; let admission = outbound_admission(&event, self.clock.now_unix_ms()?)?; - let delivery_deadline = self - .deadlines - .deadline_unix_ms(OperationKind::Deliver, self.clock.now_unix_ms()?)?; let delivery = DeliveryRequest::new( delivery_request_id(request.operation_id), DeliveryPayload::new(event), request.targets, request.satisfaction, - delivery_deadline, + request.delivery_deadline_unix_ms, ) .map_err(|_| Error::InvalidPushRequest)?; let plan_digest = delivery_plan_digest(&delivery); @@ -566,15 +733,45 @@ fn outbound_admission( .map_err(map_storage_error) } -fn push_input_digest(request: &PushRequest) -> IdempotencyDigest { +fn push_input_digest(request: &PushRequest) -> Result<IdempotencyDigest, Error> { + Ok(IdempotencyDigest::new( + *authored_push_input_digest(request)?.as_bytes(), + )) +} + +fn authored_push_input_digest(request: &PushRequest) -> Result<AtomicCommitDigest, Error> { let mut hasher = Sha256::new(); - hash_field(&mut hasher, b"radroots.sync.push.v1"); - hash_field(&mut hasher, request.draft.expected_event_id().as_bytes()); + hash_field(&mut hasher, b"radroots.sync.authored-push.v2"); + hash_field(&mut hasher, request.operation_id.as_bytes()); + hash_field(&mut hasher, request.idempotency_key.as_str().as_bytes()); + hash_field(&mut hasher, request.plan.digest().as_bytes()); + hash_field(&mut hasher, request.actor.public_key().as_bytes()); + let source = match request.actor.source() { + radroots_signing::actor::ActorSource::LocalAccount(_) => 0, + radroots_signing::actor::ActorSource::ExplicitPublicKey => 1, + radroots_signing::actor::ActorSource::RemoteSigner(_) => 2, + radroots_signing::actor::ActorSource::Service(_) => 3, + _ => return Err(Error::InvalidPushRequest), + }; + hasher.update([source]); + if let Some(account_id) = request.actor.account_id() { + hash_field(&mut hasher, account_id.as_bytes()); + } + for role in request.actor.roles() { + hash_field(&mut hasher, role.as_str().as_bytes()); + } for target in request.targets.targets() { hash_field(&mut hasher, target.fingerprint().as_str().as_bytes()); } hash_satisfaction(&mut hasher, &request.satisfaction); - IdempotencyDigest::new(hasher.finalize().into()) + hasher.update(request.delivery_deadline_unix_ms.to_be_bytes()); + let cancellation = match request.cancellation { + CancellationPolicy::PreservePublishedRequest => 0, + CancellationPolicy::LocalCooperative => 1, + _ => return Err(Error::InvalidPushRequest), + }; + hasher.update([cancellation]); + Ok(AtomicCommitDigest::new(hasher.finalize().into())) } fn delivery_plan_digest(request: &DeliveryRequest) -> DeliveryPlanDigest { @@ -616,6 +813,41 @@ fn atomic_digest(domain: &[u8], input: &[u8]) -> AtomicCommitDigest { AtomicCommitDigest::new(hasher.finalize().into()) } +fn authored_ids( + operation_id: SyncId, +) -> Result< + ( + OperationInstanceId, + AuthoredArtifactId, + AuthoredDeliveryPlanId, + ), + Error, +> { + let operation = + OperationInstanceId::new(*operation_id.as_bytes()).map_err(map_storage_error)?; + let artifact = AuthoredArtifactId::new(derive_child_id( + b"radroots.sync.authored-artifact.v2", + operation_id.as_bytes(), + )) + .map_err(map_storage_error)?; + let delivery = AuthoredDeliveryPlanId::new(derive_child_id( + b"radroots.sync.authored-delivery.v2", + artifact.as_bytes(), + )) + .map_err(map_storage_error)?; + Ok((operation, artifact, delivery)) +} + +fn derive_child_id(domain: &[u8], parent: &[u8; 16]) -> [u8; 16] { + let mut hasher = Sha256::new(); + hash_field(&mut hasher, domain); + hash_field(&mut hasher, parent); + let digest: [u8; 32] = hasher.finalize().into(); + let mut child = [0_u8; 16]; + child.copy_from_slice(&digest[..16]); + child +} + fn hash_field(hasher: &mut Sha256, value: &[u8]) { hasher.update(u64::try_from(value.len()).unwrap_or(u64::MAX).to_be_bytes()); hasher.update(value); diff --git a/crates/sync/tests/engine_composition.rs b/crates/sync/tests/engine_composition.rs @@ -10,7 +10,7 @@ use radroots_sync::{ }; use radroots_transport::{ DeliveryReceipt, DeliveryRequest, Error as TransportError, EventSink, EventSource, FetchPage, - FetchRequest, SinkStatus, SourceStatus, + FetchRequest, SinkFailure, SinkStatus, SourceStatus, capability::{Availability, Maturity, SinkCapabilities, SourceCapabilities}, }; @@ -89,7 +89,7 @@ impl EventSink for MockSink { fn deliver( &self, _request: DeliveryRequest, - ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, TransportError>> { + ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> { Box::pin(async { unreachable!("composition does not deliver") }) } } diff --git a/crates/sync/tests/push_enqueue.rs b/crates/sync/tests/push_enqueue.rs @@ -8,7 +8,10 @@ use std::{ use futures::{FutureExt, task::noop_waker_ref}; use futures_executor::block_on; -use radroots_event::{EventDraft, SignedEvent, contract::AuthorRole, draft::SignedEventParts}; +use radroots_event::{ + GenericEventDraft, SignedEvent, contract::AuthorRole, draft::SignedEventParts, +}; +use radroots_event_codec::authoring::AuthoredEventPlan; use radroots_protocol::runtime::v1::SyncRetryDecision; use radroots_signing::{ Actor, Error as SigningError, SignReceipt, SignRequest, Signer, SignerStatus, @@ -21,15 +24,16 @@ use radroots_storage::{ memory::MemoryStorage, outbox::{ClaimOutboxItems, LeaseId, LeaseOwner, OutboxStage, SatisfactionResult}, }; +use radroots_storage_sqlite::{OpenMode, OpenOptions, Paths, SqliteStorage}; use radroots_sync::{ Engine, PushRequest, policy::{Clock, DeadlinePolicy, Error, IdSource, OperationKind, SyncId, SyncStorage}, push::DeliveryRunRequest, }; use radroots_transport::{ - DeliveryReceipt, DeliveryRequest, Error as TransportError, EventSink, SinkStatus, Target, - TargetSet, TransportId, - outcome::DeliveryOutcome, + DeliveryReceipt, DeliveryRequest, Error as TransportError, EventSink, SinkFailure, SinkStatus, + Target, TargetSet, TransportId, + outcome::{DeliveryOutcome, Retryability}, policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy}, sink::DeliveryTargetReceipt, }; @@ -46,7 +50,7 @@ impl EventSink for MockSink { fn deliver( &self, _request: DeliveryRequest, - ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, TransportError>> { + ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> { Box::pin(async { unreachable!("enqueue does not deliver") }) } } @@ -79,7 +83,7 @@ impl EventSink for ScriptedSink { fn deliver( &self, request: DeliveryRequest, - ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, TransportError>> { + ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> { self.requests .lock() .expect("scripted request lock") @@ -92,8 +96,18 @@ impl EventSink for ScriptedSink { .expect("scripted delivery behavior"); Box::pin(async move { match behavior { - DeliveryBehavior::Outcomes(outcomes) => receipt(&request, outcomes), - DeliveryBehavior::AdapterError => Err(TransportError::UnsupportedOperation), + DeliveryBehavior::Outcomes(outcomes) => { + Ok(receipt(&request, outcomes).expect("valid scripted receipt")) + } + DeliveryBehavior::AdapterError => Err(SinkFailure::for_request( + &request, + "adapter_unavailable", + Retryability::Retryable, + None, + None, + Vec::new(), + ) + .expect("valid scripted failure")), DeliveryBehavior::MismatchedRequest => { let mismatched = DeliveryRequest::new( "mismatched-request", @@ -101,11 +115,13 @@ impl EventSink for ScriptedSink { request.target_set().clone(), request.satisfaction().clone(), request.deadline_unix_ms(), - )?; - receipt( + ) + .expect("mismatched request"); + Ok(receipt( &mismatched, vec![DeliveryOutcome::accepted(); mismatched.target_set().len()], ) + .expect("mismatched receipt")) } } }) @@ -132,7 +148,7 @@ impl IdSource for TestIds { #[derive(Clone, Copy)] enum SignBehavior { - Success { completed_at_unix: u64 }, + Success { completed_at_unix_ms: u64 }, Error(SigningErrorKind), Pending, } @@ -165,10 +181,12 @@ impl Signer for MockSigner { self.calls.fetch_add(1, Ordering::Relaxed); Box::pin(async move { match self.behavior { - SignBehavior::Success { completed_at_unix } => SignReceipt::from_signed_event( + SignBehavior::Success { + completed_at_unix_ms, + } => SignReceipt::from_signed_event( &request, signed_event(&request), - completed_at_unix, + completed_at_unix_ms, ), SignBehavior::Error(kind) => Err(SigningError::new(kind)), SignBehavior::Pending => std::future::pending().await, @@ -187,29 +205,29 @@ fn public_key_hex() -> String { } fn signed_event(request: &SignRequest) -> SignedEvent { - let draft = request.draft(); - let id = draft.expected_event_id().to_hex(); + let plan = request.plan(); + let id = plan.expected_event_id().to_hex(); let pubkey = public_key_hex(); let signature = Secp256k1::new() .sign_schnorr_no_aux_rand( - &Message::from_digest(*draft.expected_event_id().as_bytes()), + &Message::from_digest(*plan.expected_event_id().as_bytes()), &signing_keypair(), ) .to_string(); let raw_json = format!( "{{\"id\":\"{id}\",\"pubkey\":\"{pubkey}\",\"created_at\":{},\"kind\":{},\"tags\":{:?},\"content\":{content:?},\"sig\":\"{signature}\"}}", - draft.created_at_u64(), - draft.kind_u32(), - draft.tags_as_vec(), - content = draft.content(), + plan.created_at(), + plan.body().kind(), + plan.body().tags(), + content = plan.body().content(), ); SignedEvent::new(SignedEventParts { id, pubkey, - created_at: draft.created_at_u64(), - kind: draft.kind_u32(), - tags: draft.tags_as_vec(), - content: draft.content().to_owned(), + created_at: plan.created_at(), + kind: plan.body().kind(), + tags: plan.body().tags().to_vec(), + content: plan.body().content().to_owned(), sig: signature, raw_json, }) @@ -232,17 +250,20 @@ fn request_with_policy( target_policy: TargetPolicy, ) -> PushRequest { let pubkey = public_key_hex(); - let draft = EventDraft::new( - "radroots.social.geochat.v1", - 20_000, - 1_800_000_100, - vec![], - CONTENT, - pubkey, + let plan = AuthoredEventPlan::from_generic( + GenericEventDraft::new( + "radroots.social.geochat.v1", + 20_000, + 1_800_000_100, + vec![], + CONTENT, + pubkey, + ) + .expect("draft"), ) - .expect("draft"); + .expect("authored plan"); let actor = Actor::new( - *draft.expected_pubkey(), + *plan.author(), ActorSource::ExplicitPublicKey, [AuthorRole::Any], ) @@ -251,7 +272,7 @@ fn request_with_policy( SyncId::new([operation_byte; 16]).expect("operation id"), IdempotencyKey::parse(format!("push-{operation_byte}")).expect("idempotency key"), actor, - draft, + plan, TargetSet::new( relays .iter() @@ -260,6 +281,7 @@ fn request_with_policy( ) .expect("targets"), SatisfactionPolicy::new(class, target_policy), + 1_800_000_300_000, CancellationPolicy::PreservePublishedRequest, ) .expect("push request") @@ -319,21 +341,21 @@ fn delivery_run(seed: u8, limit: u16) -> DeliveryRunRequest { #[test] fn authorized_signing_atomically_enqueues_and_replays_without_resigning() { let signer = Arc::new(MockSigner::new(SignBehavior::Success { - completed_at_unix: 1_800_000_200, + completed_at_unix_ms: 1_800_000_200_500, })); let (engine, storage) = setup_engine(signer.clone()); let request = request(1, "wss://relay.example"); assert_eq!(request.operation_id().as_bytes(), &[1; 16]); assert_eq!(request.idempotency_key().as_str(), "push-1"); assert_ne!(request.actor().public_key().as_bytes(), &[0; 32]); - assert!(!request.draft().content().is_empty()); + assert!(!request.plan().body().content().is_empty()); assert_eq!(request.targets().len(), 1); assert_eq!(request.satisfaction().class(), SatisfactionClass::Accepted); assert_eq!( request.cancellation(), CancellationPolicy::PreservePublishedRequest ); - assert!(format!("{request:?}").contains("redacted frozen event draft")); + assert!(format!("{request:?}").contains("redacted exact authored plan")); let receipt = block_on(engine.sign_and_enqueue(request.clone())).expect("enqueue"); assert_eq!(receipt.operation_id().as_bytes(), &[1; 16]); assert!(!receipt.is_replay()); @@ -360,6 +382,127 @@ fn authorized_signing_atomically_enqueues_and_replays_without_resigning() { } #[test] +fn preparation_is_atomic_status_visible_and_replays_without_external_effects() { + let signer = Arc::new(MockSigner::new(SignBehavior::Pending)); + let (engine, storage) = setup_engine(signer.clone()); + let push = request(6, "wss://relay.example"); + + let never_polled = engine.prepare_push(push.clone()); + drop(never_polled); + assert!( + block_on(engine.push_status(push.operation_id())) + .expect("status before prepare") + .is_none() + ); + + let prepared = block_on(engine.prepare_push(push.clone())).expect("prepare"); + assert!(!prepared.is_replay()); + assert_eq!( + prepared.operation().artifact_ids(), + &[prepared.artifact().artifact_id()] + ); + assert_eq!( + prepared.delivery_plan().artifact_id(), + prepared.artifact().artifact_id() + ); + assert!(prepared.delivery_plan().request().is_none()); + assert!( + block_on(Journal::operation( + &*storage, + OperationInstanceId::new(*push.operation_id().as_bytes()).expect("operation"), + )) + .expect("legacy journal lookup") + .is_none() + ); + + let status = block_on(engine.push_status(push.operation_id())) + .expect("status after prepare") + .expect("prepared status"); + assert_eq!(status.operation(), prepared.operation()); + assert_eq!(status.artifact(), prepared.artifact()); + assert_eq!(status.delivery_plan(), prepared.delivery_plan()); + + let replay = block_on(engine.prepare_push(push)).expect("exact replay"); + assert!(replay.is_replay()); + assert_eq!(replay.operation(), prepared.operation()); + assert_eq!(replay.artifact(), prepared.artifact()); + assert_eq!(replay.delivery_plan(), prepared.delivery_plan()); + assert_eq!(signer.calls.load(Ordering::Relaxed), 0); + + assert_eq!( + block_on(engine.prepare_push(request(6, "wss://other.example"))), + Err(Error::StorageConflict) + ); + assert_eq!(signer.calls.load(Ordering::Relaxed), 0); +} + +#[tokio::test] +async fn sqlite_preparation_reopens_and_replays_the_original_complete_intent() { + let directory = tempfile::tempdir().expect("database directory"); + let paths = Paths::from_directory(directory.path()).expect("paths"); + let push = request(7, "wss://relay.example"); + let original = { + let store = Arc::new( + SqliteStorage::open( + OpenOptions::new(paths.clone(), OpenMode::Create) + .with_source_generation(SourceGeneration::new([7; 32]).expect("generation"), 1) + .expect("source generation"), + ) + .await + .expect("open SQLite"), + ); + let capability: Arc<dyn SyncStorage> = store; + let signer = Arc::new(MockSigner::new(SignBehavior::Pending)); + let engine = Engine::builder( + capability, + Arc::new(TestClock(AtomicU64::new(1_800_000_200_000))), + Arc::new(TestIds(AtomicU64::new(80))), + DeadlinePolicy::new(10_000, 10_000, 10_000).expect("deadlines"), + ) + .sink(Arc::new(MockSink)) + .signer(signer.clone()) + .build() + .expect("engine"); + let prepared = engine.prepare_push(push.clone()).await.expect("prepare"); + assert!(!prepared.is_replay()); + assert_eq!(signer.calls.load(Ordering::Relaxed), 0); + prepared + }; + + let store = Arc::new( + SqliteStorage::open(OpenOptions::new(paths, OpenMode::ReadWriteExisting)) + .await + .expect("reopen SQLite"), + ); + let capability: Arc<dyn SyncStorage> = store; + let signer = Arc::new(MockSigner::new(SignBehavior::Pending)); + let engine = Engine::builder( + capability, + Arc::new(TestClock(AtomicU64::new(1_800_000_210_000))), + Arc::new(TestIds(AtomicU64::new(90))), + DeadlinePolicy::new(10_000, 10_000, 10_000).expect("deadlines"), + ) + .sink(Arc::new(MockSink)) + .signer(signer.clone()) + .build() + .expect("engine"); + let status = engine + .push_status(push.operation_id()) + .await + .expect("status") + .expect("durable preparation"); + assert_eq!(status.operation(), original.operation()); + assert_eq!(status.artifact(), original.artifact()); + assert_eq!(status.delivery_plan(), original.delivery_plan()); + let replay = engine.prepare_push(push).await.expect("replay"); + assert!(replay.is_replay()); + assert_eq!(replay.operation(), original.operation()); + assert_eq!(replay.artifact(), original.artifact()); + assert_eq!(replay.delivery_plan(), original.delivery_plan()); + assert_eq!(signer.calls.load(Ordering::Relaxed), 0); +} + +#[test] fn signer_rejection_challenge_and_timeout_fail_before_enqueue() { for kind in [ SigningErrorKind::SignerRejected, @@ -384,7 +527,7 @@ fn signer_rejection_challenge_and_timeout_fail_before_enqueue() { } let late = Arc::new(MockSigner::new(SignBehavior::Success { - completed_at_unix: 1_800_000_212, + completed_at_unix_ms: 1_800_000_212_000, })); let (engine, _) = setup_engine(late); assert_eq!( @@ -396,7 +539,7 @@ fn signer_rejection_challenge_and_timeout_fail_before_enqueue() { #[test] fn idempotency_conflict_and_cancellation_before_commit_fail_closed() { let signer = Arc::new(MockSigner::new(SignBehavior::Success { - completed_at_unix: 1_800_000_200, + completed_at_unix_ms: 1_800_000_200_500, })); let (engine, _) = setup_engine(signer.clone()); block_on(engine.sign_and_enqueue(request(4, "wss://relay.example"))).expect("enqueue"); @@ -467,7 +610,7 @@ fn delivery_evaluates_any_all_quorum_required_and_partial_outcomes() { ]), ])); let signer = Arc::new(MockSigner::new(SignBehavior::Success { - completed_at_unix: 1_800_000_200, + completed_at_unix_ms: 1_800_000_200_500, })); let ((engine, _), _) = setup_engine_with_sink(signer, sink.clone()); let two = ["wss://one.example", "wss://two.example"]; @@ -565,7 +708,7 @@ fn transport_failure_is_durable_and_retry_preserves_the_exact_plan() { ]), ])); let signer = Arc::new(MockSigner::new(SignBehavior::Success { - completed_at_unix: 1_800_000_200, + completed_at_unix_ms: 1_800_000_200_500, })); let ((engine, _), _) = setup_engine_with_sink(signer, sink.clone()); block_on(engine.sign_and_enqueue(request_with_policy( @@ -616,7 +759,7 @@ fn transport_failure_is_durable_and_retry_preserves_the_exact_plan() { fn malformed_receipts_release_work_and_expired_plans_terminalize() { let malformed_sink = Arc::new(ScriptedSink::new([DeliveryBehavior::MismatchedRequest])); let signer = Arc::new(MockSigner::new(SignBehavior::Success { - completed_at_unix: 1_800_000_200, + completed_at_unix_ms: 1_800_000_200_500, })); let ((engine, storage), _) = setup_engine_with_sink(signer, malformed_sink); let enqueued = block_on(engine.sign_and_enqueue(request(31, "wss://one.example"))) @@ -633,7 +776,7 @@ fn malformed_receipts_release_work_and_expired_plans_terminalize() { let expired_sink = Arc::new(ScriptedSink::new([])); let signer = Arc::new(MockSigner::new(SignBehavior::Success { - completed_at_unix: 1_800_000_200, + completed_at_unix_ms: 1_800_000_200_500, })); let ((engine, _), clock) = setup_engine_with_sink(signer, expired_sink.clone()); let enqueued = block_on(engine.sign_and_enqueue(request(32, "wss://one.example"))) @@ -686,7 +829,7 @@ fn delivery_run_rejects_unbounded_claims() { Err(Error::InvalidDeliveryRequest) ); let signer = Arc::new(MockSigner::new(SignBehavior::Success { - completed_at_unix: 1_800_000_200, + completed_at_unix_ms: 1_800_000_200_500, })); let (engine, _) = setup_engine(signer); let over_engine_budget = DeliveryRunRequest::new( @@ -706,12 +849,13 @@ fn delivery_run_rejects_unbounded_claims() { valid.operation_id(), valid.idempotency_key().clone(), valid.actor().clone(), - valid.draft().clone(), + valid.plan().clone(), valid.targets().clone(), SatisfactionPolicy::new( SatisfactionClass::Accepted, TargetPolicy::quorum(2).expect("quorum"), ), + valid.delivery_deadline_unix_ms(), valid.cancellation(), ); assert!(matches!(invalid_quorum, Err(Error::InvalidPushRequest))); @@ -723,12 +867,13 @@ fn delivery_run_rejects_unbounded_claims() { valid.operation_id(), valid.idempotency_key().clone(), valid.actor().clone(), - valid.draft().clone(), + valid.plan().clone(), valid.targets().clone(), SatisfactionPolicy::new( SatisfactionClass::Accepted, TargetPolicy::required(vec![absent]).expect("required"), ), + valid.delivery_deadline_unix_ms(), valid.cancellation(), ); assert!(matches!(invalid_required, Err(Error::InvalidPushRequest))); @@ -737,7 +882,7 @@ fn delivery_run_rejects_unbounded_claims() { #[test] fn retry_decisions_are_passive_typed_and_deadline_aware() { let signer = Arc::new(MockSigner::new(SignBehavior::Success { - completed_at_unix: 1_800_000_200, + completed_at_unix_ms: 1_800_000_200_500, })); let (engine, storage) = setup_engine(signer); let enqueued = block_on(engine.sign_and_enqueue(request(61, "wss://one.example"))) @@ -800,7 +945,7 @@ fn retry_decisions_are_passive_typed_and_deadline_aware() { DeliveryBehavior::Outcomes(vec![DeliveryOutcome::rejected()]), ])); let signer = Arc::new(MockSigner::new(SignBehavior::Success { - completed_at_unix: 1_800_000_200, + completed_at_unix_ms: 1_800_000_200_500, })); let ((engine, _), _) = setup_engine_with_sink(signer, sink); block_on(engine.sign_and_enqueue(request(62, "wss://one.example"))) diff --git a/crates/sync/tests/reliability_scenarios.rs b/crates/sync/tests/reliability_scenarios.rs @@ -20,8 +20,8 @@ use radroots_sync::{ push::DeliveryRunRequest, }; use radroots_transport::{ - DeliveryReceipt, DeliveryRequest, Error as TransportError, EventSink, SinkStatus, Target, - TargetSet, TransportId, + DeliveryReceipt, DeliveryRequest, Error as TransportError, EventSink, SinkFailure, SinkStatus, + Target, TargetSet, TransportId, capability::{Availability, Maturity, SinkCapabilities}, outcome::DeliveryOutcome, policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy}, @@ -65,13 +65,21 @@ impl EventSink for RecoverySink { fn deliver( &self, request: DeliveryRequest, - ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, TransportError>> { + ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> { let call = self.0.fetch_add(1, Ordering::Relaxed); Box::pin(async move { if call == 0 { - return Err(TransportError::UnsupportedOperation); + return Err(SinkFailure::for_request( + &request, + "recovery_unavailable", + radroots_transport::outcome::Retryability::Retryable, + None, + None, + Vec::new(), + ) + .expect("valid recovery failure")); } - DeliveryReceipt::for_request( + Ok(DeliveryReceipt::for_request( &request, request .target_set() @@ -83,6 +91,7 @@ impl EventSink for RecoverySink { }) .collect(), ) + .expect("valid recovery receipt")) }) } }