lib

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

commit f578069cf894b34695bdf0ed2298e0aa162b5c13
parent 640c5e9e12099c7a74e552f8b505aa6a7d51f240
Author: triesap <tyson@radroots.org>
Date:   Mon, 14 Sep 2026 09:21:33 +0000

sync: reconcile late signing evidence before scheduling

- Bind retained signatures to the original durable attempt
- Reload current state and preserve cancellation after late results
- Recover signed work without credentials and retain clock uncertainty
- Verify API stability, coverage and native generation consumers

Diffstat:
Acontracts/architecture/decisions/sync_signing_evidence.v1.json | 12++++++++++++
Mcrates/sync/README.md | 16++++++++++++++++
Mcrates/sync/src/push.rs | 107++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-------------------
Mcrates/sync/tests/push_enqueue.rs | 20++++++++++++++++++--
Acrates/sync/tests/push_enqueue/signing_evidence.rs | 400+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
5 files changed, 527 insertions(+), 28 deletions(-)

diff --git a/contracts/architecture/decisions/sync_signing_evidence.v1.json b/contracts/architecture/decisions/sync_signing_evidence.v1.json @@ -0,0 +1,12 @@ +{ + "schema": "radroots.sync-signing-evidence.v1", + "status": "approved", + "owner": "radroots_sync", + "producer_contracts": ["authored_signing_evidence.v1.json", "authored_signed_facts.v1.json"], + "binding": "Invoke sign_authored_evidence with the exact persisted plan and operation/artifact identity. Revalidate the evidence against the retained request and an injected observation time before RecordSignedArtifact with the complete original durable signing claim.", + "reconciliation": "Validate the storage receipt against its exact command, then reload current durable status even on receipt replay. Require the exact retained signed bytes. A stale error cannot overwrite a newer claim, stop or signed fact. Signed replay requires no signer or credential access.", + "wait_outcome": "Record valid evidence before reporting deadline or caller cancellation. Apply cancellation to the existing operation before scheduling further work. A cancelled delivery plan prevents new signing and admission, including late resolution of an Indeterminate artifact; already-completed admission remains historical evidence.", + "clock_failure": "If post-sign observation time is unavailable, return ClockUnavailable with the original durable attempt unresolved. Do not manufacture a terminal signer failure or observation time; subsequent recovery honors the declared replay capability.", + "lifetime": "Executor-neutral Sync creates no worker, timer or runtime. The host must retain and poll the future for late evidence delivery. Dropping it leaves durable claim/preimage recovery, not a promise of background completion.", + "boundaries": "No new key, event timestamp, transport, dependency or public type. Strict SignReceipt and expiring Blossom authorization remain unchanged. Global delivery-attempt stop reconciliation remains separate." +} diff --git a/crates/sync/README.md b/crates/sync/README.md @@ -29,3 +29,19 @@ Retain this value when composing a storage draft submission: composite replay compares captured timestamps exactly. Building it performs no storage, signing, clock or network operation. Ordinary preparation identity and replay semantics remain unchanged. + +Prepared signing consumes the authored-evidence signer hook, revalidates its +exact request binding, and records verified bytes against the original durable +claim even after that claim expires or is superseded. It reloads current state +after receipt replay before returning or permitting admission. Caller deadline +and cancellation outcomes are reported after retaining valid evidence; stopped +work cannot restart admission or signing. An already-signed replay does not +require a signer or credentials. The strict expiring authorization hook remains +unchanged. + +The host must continue polling an in-flight call to deliver its late evidence. +Dropping the future or losing the observation clock leaves the durable attempt +unresolved; recovery follows the declared replay capability and never invents +a new preimage, event timestamp, or key. Sync creates no worker or timer to +retain a discarded future. Delivery-wide stop reconciliation is a separate +contract from authored signing evidence. diff --git a/crates/sync/src/push.rs b/crates/sync/src/push.rs @@ -8,8 +8,7 @@ use radroots_event_codec::{ }; use radroots_protocol::runtime::v1::OperationId; use radroots_signing::{ - Actor, AuthoredArtifactId as SigningArtifactId, SignReceipt, SigningIntentId, - SigningOperationId, + Actor, AuthoredArtifactId as SigningArtifactId, SigningIntentId, SigningOperationId, recovery::{RecoveryDisposition, ReplayCapability, recovery_disposition}, request::{CancellationPolicy, SignPolicy}, }; @@ -20,9 +19,9 @@ use radroots_storage::{ OperationSettlement, RetrySchedule, SigningState, WorkClaim, WorkFailure, WorkPhase, }, authored_atomic::{ - ApplyAdmissionResult, ApplyDeliveryAttempt, ApplySignedArtifact, ApplyWorkFailure, - AuthoredAtomicCommand, AuthoredAtomicOutcome, AuthoredWorkTarget, CancelAuthoredTarget, - CancelAuthoredWork, ClaimAuthoredTarget, ClaimAuthoredWork, PrepareAuthoredOperation, + ApplyAdmissionResult, ApplyDeliveryAttempt, ApplyWorkFailure, AuthoredAtomicCommand, + AuthoredAtomicOutcome, AuthoredWorkTarget, CancelAuthoredTarget, CancelAuthoredWork, + ClaimAuthoredTarget, ClaimAuthoredWork, PrepareAuthoredOperation, RecordSignedArtifact, WorkFence, }, authored_delivery::{ @@ -457,12 +456,15 @@ impl Engine { /// Claims and executes one prepared signing artifact. pub async fn sign_prepared(&self, request: PushRequest) -> Result<SigningRunReceipt, Error> { - let signer = self.signer.as_deref().ok_or(Error::MissingSigner)?; self.prepare_push(request.clone()).await?; let status = self .push_status(request.operation_id) .await? .ok_or(Error::StorageFailed)?; + if status.delivery_plan.state() == AuthoredDeliveryState::Cancelled { + self.cancel_push(request.operation_id).await?; + return Err(Error::SigningCancelled); + } let artifact = status.artifact; match artifact.signing_state() { SigningState::Signed => { @@ -477,6 +479,7 @@ impl Engine { } SigningState::Planned | SigningState::Retryable => {} } + let signer = self.signer.as_deref().ok_or(Error::MissingSigner)?; let now = self.clock.now_unix_ms()?.max(artifact.updated_at_unix_ms()); if let Some(existing) = artifact.signing_claim() { @@ -525,6 +528,7 @@ impl Engine { if persisted_plan != request.plan { return Err(Error::StorageConflict); } + let signing_claim = claimed.signing_claim().ok_or(Error::StorageFailed)?.clone(); let deadline = self.deadlines.deadline_unix_ms(OperationKind::Sign, now)?; let signing_operation = SigningOperationId::new(*request.operation_id.as_bytes()) .map_err(|_| Error::InvalidPushRequest)?; @@ -541,46 +545,93 @@ impl Engine { .map_err(|_| Error::InvalidPushRequest)?; let expected_request = sign_request.clone(); - let signer_result = match signer.sign(sign_request).await { - Ok(receipt) => match self.clock.now_unix_ms() { - Ok(observed_at_unix_ms) => SignReceipt::from_signed_event( - &expected_request, - receipt.signed_event().clone(), - observed_at_unix_ms, - ), - Err(_) => Err(radroots_signing::Error::new( - radroots_signing::error::Kind::InternalError, - )), - }, + let signer_result = match signer.sign_authored_evidence(sign_request).await { + Ok(receipt) => { + // A missing observation time leaves the durable attempt unresolved; + // it is not evidence that signing failed without an effect. + let observed_at_unix_ms = self.clock.now_unix_ms()?; + receipt.revalidate(&expected_request, observed_at_unix_ms) + } Err(error) => Err(error), }; match signer_result { Ok(receipt) => { - let command = AuthoredAtomicCommand::ApplySigned( - ApplySignedArtifact::new( + let command = AuthoredAtomicCommand::RecordSigned( + RecordSignedArtifact::new( + claimed.operation_id(), claimed.artifact_id(), - fence, + signing_claim, receipt.signed_event().clone(), - receipt.completed_at_unix_ms(), + receipt.observed_at_unix_ms(), ) .map_err(map_storage_error)?, ); let applied = self .storage - .execute_authored(command) + .execute_authored(command.clone()) .await .map_err(map_storage_error)?; - let AuthoredAtomicOutcome::Artifact(artifact) = applied.outcome() else { + if !applied.matches_command(&command) { return Err(Error::StorageFailed); - }; + } + let status = self + .push_status(request.operation_id) + .await? + .ok_or(Error::StorageFailed)?; + if status + .artifact + .signed() + .is_none_or(|signed| signed.event() != receipt.signed_event()) + { + return Err(Error::StorageFailed); + } + if expected_request.cancellation_signal().is_cancelled() + || status.delivery_plan.state() == AuthoredDeliveryState::Cancelled + { + let stopped = self.cancel_push(request.operation_id).await?; + if stopped + .status + .artifact + .signed() + .is_none_or(|signed| signed.event() != receipt.signed_event()) + { + return Err(Error::StorageFailed); + } + return Err(Error::SigningCancelled); + } + match status.artifact.signing_state() { + SigningState::Cancelled => return Err(Error::SigningCancelled), + SigningState::FailedTerminal => return Err(Error::SignerFailed), + SigningState::Signed => {} + _ => return Err(Error::StorageFailed), + } + if receipt.observed_at_unix_ms() >= deadline { + return Err(Error::SignerDeadlineExceeded); + } Ok(SigningRunReceipt { - artifact: artifact.clone(), - replay: false, + artifact: status.artifact, + replay: applied.disposition() == AtomicCommitDisposition::Replay, }) } Err(error) => { let applied_at = self.clock.now_unix_ms()?; + let current = self + .push_status(request.operation_id) + .await? + .ok_or(Error::StorageFailed)?; + if current.artifact.signed().is_some() + || current.artifact.signing_claim() != claimed.signing_claim() + { + // Another worker may have committed evidence or a stop while this one waited. + return Err(match error.kind() { + radroots_signing::error::Kind::DeadlineExceeded => { + Error::SignerDeadlineExceeded + } + radroots_signing::error::Kind::SignerCancelled => Error::SigningCancelled, + _ => Error::SignerFailed, + }); + } if error.kind() == radroots_signing::error::Kind::DeadlineExceeded && claimed .signing_claim() @@ -674,6 +725,10 @@ impl Engine { replay: true, }); } + if status.delivery_plan.state() == AuthoredDeliveryState::Cancelled { + self.cancel_push(operation_id).await?; + return Err(Error::AdmissionFailed); + } if matches!( artifact.admission_state(), AdmissionState::Rejected | AdmissionState::Cancelled diff --git a/crates/sync/tests/push_enqueue.rs b/crates/sync/tests/push_enqueue.rs @@ -54,6 +54,9 @@ use secp256k1::{Keypair, Message, Secp256k1, SecretKey}; const CONTENT: &str = "frozen-content"; +#[path = "push_enqueue/signing_evidence.rs"] +mod signing_evidence; + struct MockSink; struct FaultStorage { @@ -1074,7 +1077,7 @@ fn typed_only_authored_contract_identity_survives_signing_and_local_admission() } #[test] -fn caller_revalidates_late_and_cancelled_signer_success_before_persistence() { +fn caller_retains_late_and_cancelled_signer_facts_before_reporting_wait_outcome() { for (byte, violation, expected) in [ ( 55, @@ -1112,7 +1115,20 @@ fn caller_revalidates_late_and_cancelled_signer_success_before_persistence() { let status = block_on(engine.push_status(push.operation_id())) .expect("status") .expect("prepared operation"); - assert!(status.artifact().signed().is_none()); + let signed = status + .artifact() + .signed() + .expect("retain verified signer facts"); + assert_eq!(signed.event().id(), push.plan().expected_event_id()); + assert_eq!(signed.event().created_at(), push.plan().created_at()); + assert_eq!( + status.artifact().admission_state(), + if matches!(violation, BoundaryViolation::CancelsBeforeReturn) { + radroots_storage::authored::AdmissionState::Cancelled + } else { + radroots_storage::authored::AdmissionState::Pending + } + ); assert!(status.delivery_plan().attempts().is_empty()); } } diff --git a/crates/sync/tests/push_enqueue/signing_evidence.rs b/crates/sync/tests/push_enqueue/signing_evidence.rs @@ -0,0 +1,400 @@ +use super::*; +use futures::channel::oneshot; +use radroots_signing::{AuthoredSignEvidence, SigningIntentId, SigningOperationId}; +use radroots_storage::{ + authored::{AdmissionState, FailureClass, SigningState, WorkFailure, WorkPhase}, + authored_atomic::{ApplyWorkFailure, AuthoredAtomicCommand, AuthoredWorkTarget, WorkFence}, +}; + +const NOW: u64 = 1_800_000_200_000; +type Pending = ( + SignRequest, + oneshot::Sender<Result<AuthoredSignEvidence, SigningError>>, +); + +#[derive(Default)] +struct HeldSigner { + pending: Mutex<VecDeque<Pending>>, + evidence_calls: AtomicUsize, + legacy_calls: AtomicUsize, +} + +impl HeldSigner { + fn take(&self) -> Pending { + self.pending + .lock() + .unwrap() + .pop_front() + .expect("started request") + } +} + +impl Signer for HeldSigner { + fn status( + &self, + ) -> radroots_signing::signer::BoxFuture<'_, Result<SignerStatus, SigningError>> { + Box::pin(async { + Ok(SignerStatus::new( + SignerAvailability::Ready, + vec![SignerCapability::new( + SignerKind::Remote, + ReplayCapability::ExactReplayByRequestId, + CancellationSupport::BeforeAndAfterPublication, + false, + false, + )], + None, + )) + }) + } + + fn sign( + &self, + _: SignRequest, + ) -> radroots_signing::signer::BoxFuture<'_, Result<SignReceipt, SigningError>> { + self.legacy_calls.fetch_add(1, Ordering::Relaxed); + Box::pin(async { Err(SigningError::new(SigningErrorKind::InternalError)) }) + } + + fn sign_authored_evidence( + &self, + request: SignRequest, + ) -> radroots_signing::signer::BoxFuture<'_, Result<AuthoredSignEvidence, SigningError>> { + Box::pin(async move { + self.evidence_calls.fetch_add(1, Ordering::Relaxed); + let (sender, receiver) = oneshot::channel(); + self.pending.lock().unwrap().push_back((request, sender)); + receiver + .await + .map_err(|_| SigningError::new(SigningErrorKind::SignerUnavailable))? + }) + } +} + +fn engine( + storage: Arc<dyn SyncStorage>, + clock: Arc<dyn Clock>, + signer: Option<Arc<dyn Signer>>, +) -> Engine { + let builder = Engine::builder( + storage, + clock, + Arc::new(TestIds(AtomicU64::new(100))), + DeadlinePolicy::new(10_000, 10_000, 10_000).unwrap(), + ) + .sink(Arc::new(MockSink)); + if let Some(signer) = signer { + builder.signer(signer).build().unwrap() + } else { + builder.build().unwrap() + } +} + +fn setup() -> (Engine, Arc<MemoryStorage>, Arc<TestClock>, Arc<HeldSigner>) { + let storage = Arc::new(MemoryStorage::new(SourceGeneration::new([17; 32]).unwrap())); + let clock = Arc::new(TestClock(AtomicU64::new(NOW))); + let signer = Arc::new(HeldSigner::default()); + ( + engine(storage.clone(), clock.clone(), Some(signer.clone())), + storage, + clock, + signer, + ) +} + +fn complete((request, sender): Pending, at: u64) -> SignedEvent { + let event = signed_event(&request); + let evidence = AuthoredSignEvidence::from_signed_event(&request, event.clone(), at).unwrap(); + sender.send(Ok(evidence)).unwrap(); + event +} + +#[test] +fn missing_observation_clock_preserves_uncertain_non_replayable_attempt() { + struct FaultClock(AtomicUsize); + impl Clock for FaultClock { + fn now_unix_ms(&self) -> Result<u64, Error> { + match self.0.fetch_add(1, Ordering::Relaxed) { + 2 => Err(Error::ClockUnavailable), + 0 | 1 => Ok(NOW), + _ => Ok(NOW + 11_000), + } + } + } + let storage = Arc::new(MemoryStorage::new(SourceGeneration::new([37; 32]).unwrap())); + let signer = Arc::new(MockSigner::with_replay( + SignBehavior::Success { + completed_at_unix_ms: NOW, + }, + ReplayCapability::NonReplayable, + )); + let engine = engine( + storage, + Arc::new(FaultClock(AtomicUsize::new(0))), + Some(signer.clone()), + ); + let push = request(37, "wss://relay.example"); + assert_eq!( + block_on(engine.sign_prepared(push.clone())), + Err(Error::ClockUnavailable) + ); + let status = block_on(engine.push_status(push.operation_id())) + .unwrap() + .unwrap(); + assert!(status.artifact().signing_claim().is_some()); + assert!(status.artifact().signed().is_none()); + assert!(status.artifact().last_failure().is_none()); + assert_eq!( + block_on(engine.sign_prepared(push)), + Err(Error::SigningIndeterminate) + ); + assert_eq!(signer.calls.load(Ordering::Relaxed), 1); +} + +#[test] +fn late_evidence_survives_durable_cancel_and_cannot_schedule_more_work() { + let (engine, storage, clock, signer) = setup(); + let push = request(31, "wss://relay.example"); + let mut future = Box::pin(engine.sign_prepared(push.clone())).fuse(); + assert!( + future + .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref())) + .is_pending() + ); + let pending = signer.take(); + block_on(engine.cancel_push(push.operation_id())).unwrap(); + clock.0.store(NOW + 11_000, Ordering::Release); + let event = complete(pending, NOW + 11_000); + assert_eq!(block_on(future), Err(Error::SigningCancelled)); + let status = block_on(engine.push_status(push.operation_id())) + .unwrap() + .unwrap(); + assert_eq!(status.artifact().signing_state(), SigningState::Cancelled); + assert_eq!(status.artifact().signed().unwrap().event(), &event); + assert_eq!( + status.delivery_plan().state(), + AuthoredDeliveryState::Cancelled + ); + assert!(block_on(engine.admit_signed(push.operation_id())).is_err()); + assert_eq!( + block_on(engine.sign_prepared(push)), + Err(Error::SigningCancelled) + ); + assert_eq!(signer.evidence_calls.load(Ordering::Relaxed), 1); + assert_eq!(signer.legacy_calls.load(Ordering::Relaxed), 0); + assert!( + block_on(storage.query_visible(EventQuery::all(EventQueryBounds::first(10).unwrap()))) + .unwrap() + .items() + .is_empty() + ); +} + +#[test] +fn superseded_attempts_retain_first_bytes_and_replay_without_a_signer() { + let (engine, storage, clock, signer) = setup(); + let push = request(32, "wss://relay.example"); + let mut first = Box::pin(engine.sign_prepared(push.clone())).fuse(); + let mut context = std::task::Context::from_waker(noop_waker_ref()); + assert!(first.poll_unpin(&mut context).is_pending()); + let first_result = signer.take(); + clock.0.store(NOW + 11_000, Ordering::Release); + let mut second = Box::pin(engine.sign_prepared(push.clone())).fuse(); + assert!(second.poll_unpin(&mut context).is_pending()); + let second_result = signer.take(); + assert_eq!( + first_result.0.signer_request_id(), + second_result.0.signer_request_id() + ); + let event = complete(first_result, NOW + 11_000); + assert_eq!(block_on(first), Err(Error::SignerDeadlineExceeded)); + let retained = block_on(engine.push_status(push.operation_id())) + .unwrap() + .unwrap(); + assert_eq!(complete(second_result, NOW + 11_000), event); + let result = block_on(second).unwrap(); + assert_eq!(result.artifact(), retained.artifact()); + assert_eq!(result.artifact().signed().unwrap().event(), &event); + let recovered = self::engine(storage, clock, None); + let replay = block_on(recovered.sign_prepared(push)).unwrap(); + assert!(replay.is_replay()); + assert_eq!(replay.artifact(), retained.artifact()); + assert_eq!(signer.evidence_calls.load(Ordering::Relaxed), 2); + assert_eq!(signer.legacy_calls.load(Ordering::Relaxed), 0); +} + +#[test] +fn stale_signer_failure_cannot_overwrite_another_workers_signed_fact() { + let (engine, _, clock, signer) = setup(); + let push = request(33, "wss://relay.example"); + let mut first = Box::pin(engine.sign_prepared(push.clone())).fuse(); + let mut context = std::task::Context::from_waker(noop_waker_ref()); + assert!(first.poll_unpin(&mut context).is_pending()); + let first_result = signer.take(); + clock.0.store(NOW + 11_000, Ordering::Release); + let mut second = Box::pin(engine.sign_prepared(push.clone())).fuse(); + assert!(second.poll_unpin(&mut context).is_pending()); + complete(signer.take(), NOW + 11_000); + let signed = block_on(second).unwrap(); + first_result + .1 + .send(Err(SigningError::new(SigningErrorKind::SignerRejected))) + .unwrap(); + assert_eq!(block_on(first), Err(Error::SignerFailed)); + let status = block_on(engine.push_status(push.operation_id())) + .unwrap() + .unwrap(); + assert_eq!(status.artifact(), signed.artifact()); + assert!(status.artifact().last_failure().is_none()); +} + +#[test] +fn valid_signature_under_another_operation_or_artifact_is_rejected() { + for change_operation in [false, true] { + let (engine, _, _, signer) = setup(); + let push = request(34, "wss://relay.example"); + let mut future = Box::pin(engine.sign_prepared(push.clone())).fuse(); + assert!( + future + .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref())) + .is_pending() + ); + let (original, sender) = signer.take(); + let intent = if change_operation { + SigningIntentId::new( + SigningOperationId::new([99; 16]).unwrap(), + original.intent_id().artifact_id(), + ) + } else { + SigningIntentId::new( + original.intent_id().operation_id(), + radroots_signing::AuthoredArtifactId::new([99; 16]).unwrap(), + ) + }; + let other = SignRequest::new( + original.operation_kind(), + intent, + original.actor().clone(), + original.authored_plan().unwrap().clone(), + original.policy(), + ) + .unwrap(); + let evidence = + AuthoredSignEvidence::from_signed_event(&other, signed_event(&other), NOW).unwrap(); + sender.send(Ok(evidence)).unwrap(); + assert_eq!(block_on(future), Err(Error::SignerFailed)); + let status = block_on(engine.push_status(push.operation_id())) + .unwrap() + .unwrap(); + assert!(status.artifact().signed().is_none()); + assert!(status.delivery_plan().request().is_none()); + assert!(status.delivery_plan().attempts().is_empty()); + } +} + +#[test] +fn cancelled_indeterminate_operation_keeps_late_facts_without_admission() { + let (engine, storage, clock, signer) = setup(); + let push = request(35, "wss://relay.example"); + let mut future = Box::pin(engine.sign_prepared(push.clone())).fuse(); + assert!( + future + .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref())) + .is_pending() + ); + let pending = signer.take(); + let status = block_on(engine.push_status(push.operation_id())) + .unwrap() + .unwrap(); + let claim = status.artifact().signing_claim().unwrap(); + block_on( + storage.execute_authored(AuthoredAtomicCommand::ApplyFailure( + ApplyWorkFailure::new( + AuthoredWorkTarget::Artifact(status.artifact().artifact_id()), + WorkFence::new(*claim.token(), claim.generation(), claim.row_revision()).unwrap(), + WorkFailure::new( + "signing_unknown", + WorkPhase::Signing, + FailureClass::Indeterminate, + None, + None, + ) + .unwrap(), + None, + NOW + 1, + ) + .unwrap(), + )), + ) + .unwrap(); + clock.0.store(NOW + 2, Ordering::Release); + let stopped = block_on(engine.cancel_push(push.operation_id())).unwrap(); + assert_eq!( + stopped.status().artifact().signing_state(), + SigningState::Indeterminate + ); + let event = complete(pending, NOW + 2); + assert_eq!(block_on(future), Err(Error::SigningCancelled)); + let status = block_on(engine.push_status(push.operation_id())) + .unwrap() + .unwrap(); + assert_eq!(status.artifact().signed().unwrap().event(), &event); + assert_eq!( + status.artifact().admission_state(), + AdmissionState::Cancelled + ); + assert_eq!( + block_on(engine.admit_signed(push.operation_id())), + Err(Error::AdmissionFailed) + ); + assert!(status.delivery_plan().attempts().is_empty()); +} + +#[tokio::test] +async fn late_sqlite_fact_reopens_without_requiring_missing_credentials() { + let directory = tempfile::tempdir().unwrap(); + let paths = Paths::from_directory(directory.path()).unwrap(); + let clock = Arc::new(TestClock(AtomicU64::new(NOW))); + let store = Arc::new( + SqliteStorage::open( + OpenOptions::new(paths.clone(), OpenMode::Create) + .with_source_generation(SourceGeneration::new([36; 32]).unwrap(), 1) + .unwrap(), + ) + .await + .unwrap(), + ); + let signer = Arc::new(BoundaryViolatingSigner { + violation: BoundaryViolation::CompletesAfterDeadline, + clock: clock.clone(), + }); + let original = engine(store.clone(), clock.clone(), Some(signer)); + let push = request(36, "wss://relay.example"); + assert_eq!( + original.sign_prepared(push.clone()).await, + Err(Error::SignerDeadlineExceeded) + ); + let before = original + .push_status(push.operation_id()) + .await + .unwrap() + .unwrap(); + assert!(before.artifact().signed().is_some()); + drop(original); + store.close().await.unwrap(); + let store = Arc::new( + SqliteStorage::open(OpenOptions::new(paths, OpenMode::ReadWriteExisting)) + .await + .unwrap(), + ); + let unavailable = Arc::new(MockSigner::new(SignBehavior::Error( + SigningErrorKind::SignerUnavailable, + ))); + let recovered = engine(store.clone(), clock, Some(unavailable.clone())); + let replay = recovered.sign_prepared(push).await.unwrap(); + assert!(replay.is_replay()); + assert_eq!(replay.artifact(), before.artifact()); + assert_eq!(unavailable.calls.load(Ordering::Relaxed), 0); + drop(recovered); + store.close().await.unwrap(); +}