lib

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

commit 6f274725752514996c9fbbbb534410c77e7dbc0e
parent 1acddf6c69521bbfd83936b4d34f73965efdb026
Author: triesap <tyson@radroots.org>
Date:   Mon,  3 Aug 2026 08:31:01 +0000

sync: implement sign and durable enqueue

- add replay-stable push requests with authorization and deadline policy
- verify exact signer output before canonical local admission
- atomically persist prepared, signed, committed, and outbox state
- verify rejection, challenge, timeout, conflict, replay, and cancellation paths

Diffstat:
MCargo.lock | 2++
Mcrates/sync/Cargo.toml | 2++
Mcrates/sync/src/lib.rs | 1+
Mcrates/sync/src/policy.rs | 10++++++++++
Mcrates/sync/src/push.rs | 449+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acrates/sync/tests/push_enqueue.rs | 300+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
6 files changed, 764 insertions(+), 0 deletions(-)

diff --git a/Cargo.lock b/Cargo.lock @@ -5155,6 +5155,7 @@ dependencies = [ name = "radroots_sync" version = "0.1.0-alpha" dependencies = [ + "futures", "futures-executor", "radroots_event", "radroots_event_codec", @@ -5163,6 +5164,7 @@ dependencies = [ "radroots_storage", "radroots_trade", "radroots_transport", + "secp256k1", "serde", "sha2", ] diff --git a/crates/sync/Cargo.toml b/crates/sync/Cargo.toml @@ -39,8 +39,10 @@ serde = { workspace = true, optional = true } sha2 = { workspace = true, default-features = false } [dev-dependencies] +futures = { workspace = true } futures-executor = { workspace = true } radroots_storage = { workspace = true, features = ["memory"] } +secp256k1 = { workspace = true } [lints] workspace = true diff --git a/crates/sync/src/lib.rs b/crates/sync/src/lib.rs @@ -12,3 +12,4 @@ pub mod status; pub use engine::Engine; pub use policy::Error; pub use pull::{PullReceipt, PullRequest}; +pub use push::{PushReceipt, PushRequest}; diff --git a/crates/sync/src/policy.rs b/crates/sync/src/policy.rs @@ -202,6 +202,11 @@ pub enum Error { InvalidProjectionRequest, ReducerFailed, InvalidReducerOutput, + InvalidPushRequest, + MissingSigner, + SignerFailed, + SignerDeadlineExceeded, + InvalidSignerOutput, } impl core::fmt::Display for Error { @@ -224,6 +229,11 @@ impl core::fmt::Display for Error { Self::InvalidProjectionRequest => "sync projection request is invalid", Self::ReducerFailed => "sync projection reducer failed", Self::InvalidReducerOutput => "sync projection reducer returned invalid progress", + Self::InvalidPushRequest => "sync push request is invalid", + Self::MissingSigner => "sync engine has no signer", + Self::SignerFailed => "sync signer did not produce an event", + Self::SignerDeadlineExceeded => "sync signer exceeded its deadline", + Self::InvalidSignerOutput => "sync signer output failed canonical verification", }) } } diff --git a/crates/sync/src/push.rs b/crates/sync/src/push.rs @@ -1 +1,450 @@ //! Signing, durable enqueue, delivery, and satisfaction orchestration. + +use radroots_event::{EventDraft, admission::RawEvent}; +use radroots_event_codec::verify::{self, Nip01SignatureVerifier}; +use radroots_protocol::runtime::v1::OperationId; +use radroots_signing::{ + Actor, + request::{CancellationPolicy, SignPolicy}, +}; +use radroots_storage::{ + Journal, Outbox, + atomic::{ + AtomicCommit, AtomicCommitDigest, AtomicCommitId, AtomicCommitOutcome, AtomicWorkflow, + CommitEnqueued, CommitSigned, + }, + event::EventAdmission, + journal::{ + IdempotencyDigest, IdempotencyKey, JournalStage, JournalState, OperationInstanceId, + PrepareOperation, + }, + outbox::{DeliveryPlanDigest, EnqueueOutboxItem, OutboxItemId, OutboxRecord}, +}; +use radroots_transport::{ + DeliveryRequest, Target, TransportId, + policy::SatisfactionPolicy, + sink::DeliveryPayload, + source::{EventProvenance, ObservedEvent}, + target::TargetSet, +}; +use sha2::{Digest, Sha256}; + +use crate::{ + Engine, + ingest::{AdmissionDecision, AdmissionPolicy, RegistryPolicy}, + policy::{Error, OperationKind, SyncId}, +}; + +/// Caller-owned, replay-stable inputs for one outbound operation. +#[derive(Clone)] +pub struct PushRequest { + operation_id: SyncId, + idempotency_key: IdempotencyKey, + actor: Actor, + draft: EventDraft, + targets: TargetSet, + satisfaction: SatisfactionPolicy, + cancellation: CancellationPolicy, +} + +impl PushRequest { + #[allow(clippy::too_many_arguments)] + pub fn new( + operation_id: SyncId, + idempotency_key: IdempotencyKey, + actor: Actor, + draft: EventDraft, + targets: TargetSet, + satisfaction: SatisfactionPolicy, + cancellation: CancellationPolicy, + ) -> Result<Self, Error> { + if !valid_satisfaction(&satisfaction, &targets) { + return Err(Error::InvalidPushRequest); + } + Ok(Self { + operation_id, + idempotency_key, + actor, + draft, + targets, + satisfaction, + cancellation, + }) + } + + pub const fn operation_id(&self) -> SyncId { + self.operation_id + } + pub const fn idempotency_key(&self) -> &IdempotencyKey { + &self.idempotency_key + } + pub const fn actor(&self) -> &Actor { + &self.actor + } + pub const fn draft(&self) -> &EventDraft { + &self.draft + } + pub const fn targets(&self) -> &TargetSet { + &self.targets + } + pub const fn satisfaction(&self) -> &SatisfactionPolicy { + &self.satisfaction + } + pub const fn cancellation(&self) -> CancellationPolicy { + self.cancellation + } +} + +impl core::fmt::Debug for PushRequest { + fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result { + formatter + .debug_struct("PushRequest") + .field("operation_id", &self.operation_id) + .field("idempotency_key", &self.idempotency_key) + .field("actor", &self.actor) + .field("draft", &"[redacted frozen event draft]") + .field("targets", &self.targets) + .field("satisfaction", &self.satisfaction) + .field("cancellation", &self.cancellation) + .finish() + } +} + +/// Durable result after signing and atomic outbox enqueue. +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct PushReceipt { + operation_id: SyncId, + outbox: OutboxRecord, + replay: bool, +} + +impl PushReceipt { + pub const fn operation_id(&self) -> SyncId { + self.operation_id + } + pub const fn outbox(&self) -> &OutboxRecord { + &self.outbox + } + pub const fn is_replay(&self) -> bool { + self.replay + } +} + +impl Engine { + /// Authorizes, signs, verifies, and durably enqueues one outbound event. + /// + /// Dropping the future before the final atomic enqueue leaves either a + /// prepared or signed recoverable journal record. Once that commit returns, + /// 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 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); + + if let Some(existing) = Journal::operation(self.storage.as_ref(), instance_id) + .await + .map_err(map_storage_error)? + { + if existing.operation_id() != OperationId::SyncPush + || existing.idempotency_key() != request.idempotency_key() + || existing.input_digest() != input_digest + { + return Err(Error::StorageConflict); + } + if existing.state().stage() == JournalStage::Committed { + let outbox = Outbox::item(self.storage.as_ref(), item_id) + .await + .map_err(map_storage_error)? + .ok_or(Error::StorageFailed)?; + return Ok(PushReceipt { + operation_id: request.operation_id, + outbox, + replay: true, + }); + } + } + + let prepared_at = self.clock.now_unix_ms()?; + let prepare = PrepareOperation::new( + instance_id, + OperationId::SyncPush, + request.idempotency_key.clone(), + input_digest, + prepared_at, + ) + .map_err(map_storage_error)?; + let prepare_commit = AtomicCommit::new( + next_commit_id(self, OperationKind::Sign)?, + AtomicCommitDigest::new(*input_digest.as_bytes()), + prepared_at, + AtomicWorkflow::Prepared(prepare), + ) + .map_err(map_storage_error)?; + let prepared_receipt = self + .storage + .commit(prepare_commit) + .await + .map_err(map_storage_error)?; + let AtomicCommitOutcome::Prepared { journal: prepared } = prepared_receipt.outcome() else { + return Err(Error::StorageFailed); + }; + + 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 sign_request = radroots_signing::SignRequest::new( + OperationId::SyncPush, + request.actor.clone(), + request.draft.clone(), + SignPolicy::new(sign_deadline_seconds, 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 { + return Err(Error::SignerDeadlineExceeded); + } + let event = signed.signed_event().clone(); + + let signed_record = if prepared.state().stage() == JournalStage::Signed { + if !matches!( + prepared.state(), + JournalState::Signed { event_id } if event_id == event.id() + ) { + return Err(Error::InvalidSignerOutput); + } + prepared.clone() + } else { + let signed_commit = AtomicCommit::new( + next_commit_id(self, OperationKind::Sign)?, + atomic_digest(b"radroots.sync.signed.v1", event.raw_json().as_bytes()), + self.clock.now_unix_ms()?, + AtomicWorkflow::Signed(Box::new(CommitSigned::new( + instance_id, + prepared.revision(), + event.clone(), + ))), + ) + .map_err(map_storage_error)?; + let receipt = self + .storage + .commit(signed_commit) + .await + .map_err(map_storage_error)?; + let AtomicCommitOutcome::Signed { journal, .. } = receipt.outcome() else { + return Err(Error::StorageFailed); + }; + journal.clone() + }; + + 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, + ) + .map_err(|_| Error::InvalidPushRequest)?; + let plan_digest = delivery_plan_digest(&delivery); + let committed_at = self.clock.now_unix_ms()?; + let outbox = + EnqueueOutboxItem::new(item_id, instance_id, plan_digest, delivery, committed_at) + .map_err(map_storage_error)?; + let enqueued = CommitEnqueued::new( + instance_id, + signed_record.revision(), + admission, + outbox, + committed_at, + ) + .map_err(map_storage_error)?; + let receipt = self + .storage + .commit( + AtomicCommit::new( + next_commit_id(self, OperationKind::Sign)?, + AtomicCommitDigest::new(*plan_digest.as_bytes()), + committed_at, + AtomicWorkflow::Enqueued(Box::new(enqueued)), + ) + .map_err(map_storage_error)?, + ) + .await + .map_err(map_storage_error)?; + let AtomicCommitOutcome::Enqueued { outbox, .. } = receipt.outcome() else { + return Err(Error::StorageFailed); + }; + Ok(PushReceipt { + operation_id: request.operation_id, + outbox: (**outbox).clone(), + replay: false, + }) + } +} + +fn outbound_admission( + event: &radroots_event::SignedEvent, + observed_at_unix_ms: u64, +) -> Result<EventAdmission, Error> { + let verified = verify::signature( + verify::id(RawEvent::new(event.envelope().clone())) + .map_err(|_| Error::InvalidSignerOutput)?, + &Nip01SignatureVerifier, + ) + .map_err(|_| Error::InvalidSignerOutput)?; + let validated = verify::contract(verified).map_err(|_| Error::InvalidSignerOutput)?; + let policy = RegistryPolicy::visible(); + if policy.decide(&validated) != AdmissionDecision::Visible { + return Err(Error::InvalidSignerOutput); + } + struct Evidence; + impl radroots_event::admission::AdmissionPolicy for Evidence { + type Error = core::convert::Infallible; + fn policy_id(&self) -> &'static str { + "radroots.registry_v7" + } + fn admit( + &self, + _: &radroots_event::admission::ContractValidatedEvent, + ) -> Result<(), Self::Error> { + Ok(()) + } + } + impl radroots_event::admission::VisibilityPolicy for Evidence { + type Error = core::convert::Infallible; + fn policy_id(&self) -> &'static str { + "radroots.registry_v7" + } + fn make_visible( + &self, + _: &radroots_event::admission::AdmittedEvent, + ) -> Result<(), Self::Error> { + Ok(()) + } + } + let visible = validated + .admit_with(&Evidence) + .and_then(|event| event.make_visible_with(&Evidence)) + .map_err(|never| match never {})?; + let local = Target::new(TransportId::LOCAL, "local:sync-authored") + .map_err(|_| Error::InvalidPushRequest)?; + let provenance = EventProvenance::new( + TransportId::LOCAL, + local.fingerprint().clone(), + observed_at_unix_ms, + ) + .map_err(|_| Error::InvalidPushRequest)?; + EventAdmission::visible(ObservedEvent::new(event.clone(), provenance), visible) + .map_err(map_storage_error) +} + +fn push_input_digest(request: &PushRequest) -> IdempotencyDigest { + 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()); + 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()) +} + +fn delivery_plan_digest(request: &DeliveryRequest) -> DeliveryPlanDigest { + let mut hasher = Sha256::new(); + hash_field(&mut hasher, b"radroots.sync.delivery.v1"); + hash_field(&mut hasher, request.payload().event().raw_json().as_bytes()); + for target in request.target_set().targets() { + hash_field(&mut hasher, target.fingerprint().as_str().as_bytes()); + } + hash_satisfaction(&mut hasher, request.satisfaction()); + DeliveryPlanDigest::new(hasher.finalize().into()) +} + +fn hash_satisfaction(hasher: &mut Sha256, policy: &SatisfactionPolicy) { + hasher.update([match policy.class() { + radroots_transport::policy::SatisfactionClass::Accepted => 0, + radroots_transport::policy::SatisfactionClass::Delivered => 1, + }]); + let targets = policy.targets(); + if targets.is_any() { + hasher.update([0]); + } else if targets.is_all() { + hasher.update([1]); + } else if let Some(threshold) = targets.quorum_threshold() { + hasher.update([2]); + hasher.update(threshold.to_be_bytes()); + } else if let Some(required) = targets.required_targets() { + hasher.update([3]); + for target in required { + hash_field(hasher, target.as_str().as_bytes()); + } + } +} + +fn valid_satisfaction(policy: &SatisfactionPolicy, targets: &TargetSet) -> bool { + let selection = policy.targets(); + if let Some(threshold) = selection.quorum_threshold() { + return usize::from(threshold) <= targets.len(); + } + if let Some(required) = selection.required_targets() { + return required.iter().all(|required| { + targets + .targets() + .iter() + .any(|target| target.fingerprint() == required) + }); + } + selection.is_any() || selection.is_all() +} + +fn atomic_digest(domain: &[u8], input: &[u8]) -> AtomicCommitDigest { + let mut hasher = Sha256::new(); + hash_field(&mut hasher, domain); + hash_field(&mut hasher, input); + AtomicCommitDigest::new(hasher.finalize().into()) +} + +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); +} + +fn next_commit_id(engine: &Engine, operation: OperationKind) -> Result<AtomicCommitId, Error> { + AtomicCommitId::new(*engine.ids.next_id(operation)?.as_bytes()).map_err(map_storage_error) +} + +fn delivery_request_id(id: SyncId) -> String { + const HEX: &[u8; 16] = b"0123456789abcdef"; + let mut value = String::from("push-"); + for byte in id.as_bytes() { + value.push(HEX[(byte >> 4) as usize] as char); + value.push(HEX[(byte & 0x0f) as usize] as char); + } + value +} + +fn map_storage_error(error: radroots_storage::Error) -> Error { + match error { + radroots_storage::Error::IdempotencyConflict + | radroots_storage::Error::OperationIdentityMismatch + | radroots_storage::Error::JournalRevisionConflict + | radroots_storage::Error::OutboxPlanConflict + | radroots_storage::Error::AtomicCommitConflict => Error::StorageConflict, + _ => Error::StorageFailed, + } +} diff --git a/crates/sync/tests/push_enqueue.rs b/crates/sync/tests/push_enqueue.rs @@ -0,0 +1,300 @@ +use std::sync::{ + Arc, + atomic::{AtomicU64, AtomicUsize, Ordering}, +}; + +use futures::{FutureExt, task::noop_waker_ref}; +use futures_executor::block_on; +use radroots_event::{EventDraft, SignedEvent, contract::AuthorRole, draft::SignedEventParts}; +use radroots_signing::{ + Actor, Error as SigningError, SignReceipt, SignRequest, Signer, SignerStatus, + actor::ActorSource, error::Kind as SigningErrorKind, request::CancellationPolicy, +}; +use radroots_storage::{ + EventStore, Journal, Outbox, Storage, + event::{EventQuery, EventQueryBounds, SourceGeneration}, + journal::{IdempotencyKey, JournalStage, OperationInstanceId}, + memory::MemoryStorage, + outbox::OutboxStage, +}; +use radroots_sync::{ + Engine, PushRequest, + policy::{Clock, DeadlinePolicy, Error, IdSource, OperationKind, SyncId}, +}; +use radroots_transport::{ + DeliveryReceipt, DeliveryRequest, Error as TransportError, EventSink, SinkStatus, Target, + TargetSet, TransportId, + policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy}, +}; +use secp256k1::{Keypair, Message, Secp256k1, SecretKey}; + +const CONTENT: &str = "frozen-content"; + +struct MockSink; + +impl EventSink for MockSink { + fn status(&self) -> radroots_transport::BoxFuture<'_, Result<SinkStatus, TransportError>> { + Box::pin(async { unreachable!("enqueue does not inspect sink") }) + } + fn deliver( + &self, + _request: DeliveryRequest, + ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, TransportError>> { + Box::pin(async { unreachable!("enqueue does not deliver") }) + } +} + +struct TestClock(AtomicU64); + +impl Clock for TestClock { + fn now_unix_ms(&self) -> Result<u64, Error> { + Ok(self.0.fetch_add(1, Ordering::Relaxed)) + } +} + +struct TestIds(AtomicU64); + +impl IdSource for TestIds { + fn next_id(&self, _operation: OperationKind) -> Result<SyncId, Error> { + let value = self.0.fetch_add(1, Ordering::Relaxed); + let byte = u8::try_from(value).map_err(|_| Error::InvalidSyncId)?; + SyncId::new([byte; 16]) + } +} + +#[derive(Clone, Copy)] +enum SignBehavior { + Success { completed_at_unix: u64 }, + Error(SigningErrorKind), + Pending, +} + +struct MockSigner { + behavior: SignBehavior, + calls: AtomicUsize, +} + +impl MockSigner { + fn new(behavior: SignBehavior) -> Self { + Self { + behavior, + calls: AtomicUsize::new(0), + } + } +} + +impl Signer for MockSigner { + fn status( + &self, + ) -> radroots_signing::signer::BoxFuture<'_, Result<SignerStatus, SigningError>> { + Box::pin(async { unreachable!("enqueue does not inspect signer status") }) + } + + fn sign( + &self, + request: SignRequest, + ) -> radroots_signing::signer::BoxFuture<'_, Result<SignReceipt, SigningError>> { + self.calls.fetch_add(1, Ordering::Relaxed); + Box::pin(async move { + match self.behavior { + SignBehavior::Success { completed_at_unix } => SignReceipt::from_signed_event( + &request, + signed_event(&request), + completed_at_unix, + ), + SignBehavior::Error(kind) => Err(SigningError::new(kind)), + SignBehavior::Pending => std::future::pending().await, + } + }) + } +} + +fn signing_keypair() -> Keypair { + let secret = SecretKey::from_slice(&[1; 32]).expect("secret key"); + Keypair::from_secret_key(&Secp256k1::new(), &secret) +} + +fn public_key_hex() -> String { + signing_keypair().x_only_public_key().0.to_string() +} + +fn signed_event(request: &SignRequest) -> SignedEvent { + let draft = request.draft(); + let id = draft.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()), + &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(), + ); + 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(), + sig: signature, + raw_json, + }) + .expect("signed event") +} + +fn request(operation_byte: u8, relay: &str) -> PushRequest { + let pubkey = public_key_hex(); + let draft = EventDraft::new( + "radroots.social.geochat.v1", + 20_000, + 1_800_000_100, + vec![], + CONTENT, + pubkey, + ) + .expect("draft"); + let actor = Actor::new( + *draft.expected_pubkey(), + ActorSource::ExplicitPublicKey, + [AuthorRole::Any], + ) + .expect("actor"); + PushRequest::new( + SyncId::new([operation_byte; 16]).expect("operation id"), + IdempotencyKey::parse(format!("push-{operation_byte}")).expect("idempotency key"), + actor, + draft, + TargetSet::new(vec![ + Target::new(TransportId::NOSTR, relay).expect("target"), + ]) + .expect("targets"), + SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::any()), + CancellationPolicy::PreservePublishedRequest, + ) + .expect("push request") +} + +fn setup_engine(signer: Arc<MockSigner>) -> (Engine, Arc<MemoryStorage>) { + let storage = Arc::new(MemoryStorage::new( + SourceGeneration::new([6; 32]).expect("generation"), + )); + let capability: Arc<dyn Storage> = storage.clone(); + let engine = Engine::builder( + capability, + Arc::new(TestClock(AtomicU64::new(1_800_000_200_000))), + Arc::new(TestIds(AtomicU64::new(10))), + DeadlinePolicy::new(10_000, 10_000, 10_000).expect("deadlines"), + ) + .sink(Arc::new(MockSink)) + .signer(signer) + .build() + .expect("engine"); + (engine, storage) +} + +#[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, + })); + let (engine, storage) = setup_engine(signer.clone()); + let request = request(1, "wss://relay.example"); + let receipt = block_on(engine.sign_and_enqueue(request.clone())).expect("enqueue"); + assert!(!receipt.is_replay()); + assert_eq!(receipt.outbox().stage(), OutboxStage::Pending); + assert_eq!(signer.calls.load(Ordering::Relaxed), 1); + let replay = block_on(engine.sign_and_enqueue(request)).expect("replay"); + assert!(replay.is_replay()); + assert_eq!(replay.outbox(), receipt.outbox()); + assert_eq!(signer.calls.load(Ordering::Relaxed), 1); + + let visible = block_on(storage.query_visible(EventQuery::all( + EventQueryBounds::first(10).expect("bounds"), + ))) + .expect("visible events"); + assert_eq!(visible.items().len(), 1); + assert!( + block_on(Outbox::item( + &*storage, + radroots_storage::outbox::OutboxItemId::new([1; 16]).expect("item id"), + )) + .expect("outbox item") + .is_some() + ); +} + +#[test] +fn signer_rejection_challenge_and_timeout_fail_before_enqueue() { + for kind in [ + SigningErrorKind::SignerRejected, + SigningErrorKind::SignerCapabilityMissing, + SigningErrorKind::SignerTimeout, + SigningErrorKind::SignerOutputInvalid, + ] { + let signer = Arc::new(MockSigner::new(SignBehavior::Error(kind))); + let (engine, storage) = setup_engine(signer); + assert_eq!( + block_on(engine.sign_and_enqueue(request(2, "wss://relay.example"))), + Err(Error::SignerFailed) + ); + assert!( + block_on(Outbox::item( + &*storage, + radroots_storage::outbox::OutboxItemId::new([2; 16]).expect("item id"), + )) + .expect("outbox lookup") + .is_none() + ); + } + + let late = Arc::new(MockSigner::new(SignBehavior::Success { + completed_at_unix: 1_800_000_212, + })); + let (engine, _) = setup_engine(late); + assert_eq!( + block_on(engine.sign_and_enqueue(request(3, "wss://relay.example"))), + Err(Error::SignerDeadlineExceeded) + ); +} + +#[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, + })); + let (engine, _) = setup_engine(signer.clone()); + block_on(engine.sign_and_enqueue(request(4, "wss://relay.example"))).expect("enqueue"); + assert_eq!( + block_on(engine.sign_and_enqueue(request(4, "wss://other.example"))), + Err(Error::StorageConflict) + ); + assert_eq!(signer.calls.load(Ordering::Relaxed), 1); + + let pending = Arc::new(MockSigner::new(SignBehavior::Pending)); + let (engine, storage) = setup_engine(pending); + let mut future = Box::pin(engine.sign_and_enqueue(request(5, "wss://relay.example"))).fuse(); + let mut context = std::task::Context::from_waker(noop_waker_ref()); + assert!(future.poll_unpin(&mut context).is_pending()); + drop(future); + let record = block_on(Journal::operation( + &*storage, + OperationInstanceId::new([5; 16]).expect("instance id"), + )) + .expect("journal lookup") + .expect("prepared record"); + assert_eq!(record.state().stage(), JournalStage::Prepared); + assert!( + block_on(Outbox::item( + &*storage, + radroots_storage::outbox::OutboxItemId::new([5; 16]).expect("item id"), + )) + .expect("outbox lookup") + .is_none() + ); +}