lib

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

commit 934f377e5dc22771d044a99699f7d3a6078b517a
parent 7afc4e41e3ad5b7125f09b43680ad2efeaf40f46
Author: triesap <tyson@radroots.org>
Date:   Wed,  5 Aug 2026 03:39:10 +0000

refactor(sync): recover signing and admission work

- execute signing and admission behind durable fenced claims
- preserve exact replay and classify uncertain remote effects
- persist retry terminal indeterminate and cancellation outcomes
- prove stale-claim crash and SQLite reopen recovery paths


Diffstat:
Mcrates/storage/src/authored.rs | 26++++++++++++++++++++++++--
Mcrates/storage/src/authored_delivery.rs | 6+++++-
Mcrates/storage/tests/authored.rs | 2+-
Mcrates/sync/src/lib.rs | 4+++-
Mcrates/sync/src/policy.rs | 14++++++++++++++
Mcrates/sync/src/push.rs | 508+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++------
Mcrates/sync/tests/push_enqueue.rs | 409++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-------
7 files changed, 897 insertions(+), 72 deletions(-)

diff --git a/crates/storage/src/authored.rs b/crates/storage/src/authored.rs @@ -737,6 +737,8 @@ impl AuthoredArtifact { self.signing_state, SigningState::Planned | SigningState::Retryable ) + || claim.row_revision().get().checked_add(1) != Some(self.revision.get()) + || claim.acquired_at_unix_ms() != self.updated_at_unix_ms }) || self.admission_claim.as_ref().is_some_and(|claim| { claim.validate().is_err() || self.signed.is_none() @@ -744,6 +746,8 @@ impl AuthoredArtifact { self.admission_state, AdmissionState::Pending | AdmissionState::Retryable ) + || claim.row_revision().get().checked_add(1) != Some(self.revision.get()) + || claim.acquired_at_unix_ms() != self.updated_at_unix_ms }) { return Err(Error::InvalidAuthoredArtifact); } @@ -852,13 +856,22 @@ impl AuthoredArtifact { } pub fn set_signing_claim(&mut self, claim: WorkClaim, at_unix_ms: u64) -> Result<(), Error> { + let existing_blocks = self.signing_claim.as_ref().is_some_and(|existing| { + at_unix_ms < existing.expires_at_unix_ms() + || claim.generation() <= existing.generation() + }); if !self.origin.is_resignable() || !matches!( self.signing_state, SigningState::Planned | SigningState::Retryable ) - || self.signing_claim.is_some() + || existing_blocks || claim.row_revision() != self.revision + || claim.acquired_at_unix_ms() != at_unix_ms + || self + .signing_retry + .as_ref() + .is_some_and(|retry| at_unix_ms < retry.not_before_unix_ms()) { return Err(Error::InvalidAuthoredTransition); } @@ -939,13 +952,22 @@ impl AuthoredArtifact { } pub fn set_admission_claim(&mut self, claim: WorkClaim, at_unix_ms: u64) -> Result<(), Error> { + let existing_blocks = self.admission_claim.as_ref().is_some_and(|existing| { + at_unix_ms < existing.expires_at_unix_ms() + || claim.generation() <= existing.generation() + }); if self.signing_state != SigningState::Signed || !matches!( self.admission_state, AdmissionState::Pending | AdmissionState::Retryable ) - || self.admission_claim.is_some() + || existing_blocks || claim.row_revision() != self.revision + || claim.acquired_at_unix_ms() != at_unix_ms + || self + .admission_retry + .as_ref() + .is_some_and(|retry| at_unix_ms < retry.not_before_unix_ms()) { return Err(Error::InvalidAuthoredTransition); } diff --git a/crates/storage/src/authored_delivery.rs b/crates/storage/src/authored_delivery.rs @@ -360,9 +360,13 @@ impl AuthoredDeliveryPlan { } pub fn claim(&mut self, claim: WorkClaim, now_unix_ms: u64) -> Result<(), Error> { + let existing_blocks = self.claim.as_ref().is_some_and(|existing| { + now_unix_ms < existing.expires_at_unix_ms() + || claim.generation() <= existing.generation() + }); if self.state.is_terminal() || self.request.is_none() - || self.claim.is_some() + || existing_blocks || claim.row_revision() != self.revision || claim.acquired_at_unix_ms() != now_unix_ms || self diff --git a/crates/storage/tests/authored.rs b/crates/storage/tests/authored.rs @@ -104,7 +104,7 @@ fn planned_artifacts_enforce_exact_signing_and_admission_transitions() { [3; 16], "signer", NonZeroU64::MIN, - 10, + 11, 20, artifact.revision(), ) diff --git a/crates/sync/src/lib.rs b/crates/sync/src/lib.rs @@ -12,5 +12,7 @@ pub mod status; pub use engine::Engine; pub use policy::Error; pub use pull::{PullReceipt, PullRequest}; -pub use push::{PushPreparation, PushReceipt, PushRequest, PushStatus}; +pub use push::{ + AdmissionRunReceipt, PushPreparation, PushReceipt, PushRequest, PushStatus, SigningRunReceipt, +}; pub use status::SyncStatus; diff --git a/crates/sync/src/policy.rs b/crates/sync/src/policy.rs @@ -230,9 +230,14 @@ pub enum Error { InvalidReducerOutput, InvalidPushRequest, MissingSigner, + SignerCapabilityUnavailable, SignerFailed, SignerDeadlineExceeded, + SigningCancelled, + SigningIndeterminate, + WorkClaimConflict, InvalidSignerOutput, + AdmissionFailed, InvalidDeliveryRequest, MissingSink, InvalidStatusRequest, @@ -260,9 +265,18 @@ impl core::fmt::Display for Error { Self::InvalidReducerOutput => "sync projection reducer returned invalid progress", Self::InvalidPushRequest => "sync push request is invalid", Self::MissingSigner => "sync engine has no signer", + Self::SignerCapabilityUnavailable => { + "sync signer did not declare one usable replay capability" + } Self::SignerFailed => "sync signer did not produce an event", Self::SignerDeadlineExceeded => "sync signer exceeded its deadline", + Self::SigningCancelled => "sync signing was durably cancelled", + Self::SigningIndeterminate => { + "sync signing may have produced a non-replayable remote effect" + } + Self::WorkClaimConflict => "sync authored work is claimed by another execution", Self::InvalidSignerOutput => "sync signer output failed canonical verification", + Self::AdmissionFailed => "sync local admission did not complete", Self::InvalidDeliveryRequest => "sync delivery request is invalid", Self::MissingSink => "sync engine has no event sink", Self::InvalidStatusRequest => "sync status request is invalid", diff --git a/crates/sync/src/push.rs b/crates/sync/src/push.rs @@ -1,5 +1,6 @@ //! Signing, durable enqueue, delivery, and satisfaction orchestration. +use core::num::{NonZeroU32, NonZeroU64}; use radroots_event::admission::RawEvent; use radroots_event_codec::{ authoring::AuthoredEventPlan, @@ -8,6 +9,7 @@ use radroots_event_codec::{ use radroots_protocol::runtime::v1::OperationId; use radroots_signing::{ Actor, AuthoredArtifactId as SigningArtifactId, SigningIntentId, SigningOperationId, + recovery::{RecoveryDisposition, ReplayCapability, recovery_disposition}, request::{CancellationPolicy, SignPolicy}, }; use radroots_storage::{ @@ -16,10 +18,17 @@ use radroots_storage::{ AtomicCommit, AtomicCommitDigest, AtomicCommitDisposition, AtomicCommitId, AtomicCommitOutcome, AtomicWorkflow, CommitEnqueued, CommitSigned, }, - authored::{AuthoredArtifact, AuthoredArtifactId, AuthoredOperation}, - authored_atomic::{AuthoredAtomicCommand, AuthoredAtomicOutcome, PrepareAuthoredOperation}, + authored::{ + AdmissionState, AuthoredArtifact, AuthoredArtifactId, AuthoredOperation, FailureClass, + RetrySchedule, SigningState, WorkClaim, WorkFailure, WorkPhase, + }, + authored_atomic::{ + ApplyAdmissionResult, ApplySignedArtifact, ApplyWorkFailure, AuthoredAtomicCommand, + AuthoredAtomicOutcome, AuthoredWorkTarget, CancelAuthoredTarget, CancelAuthoredWork, + ClaimAuthoredTarget, ClaimAuthoredWork, PrepareAuthoredOperation, WorkFence, + }, authored_delivery::{AuthoredDeliveryIntent, AuthoredDeliveryPlan, AuthoredDeliveryPlanId}, - event::EventAdmission, + event::{AdmissionDisposition, EventAdmission, EventStore}, journal::{ IdempotencyDigest, IdempotencyKey, JournalStage, JournalState, OperationInstanceId, PrepareOperation, @@ -46,6 +55,11 @@ use crate::{ }; const MAX_DELIVERY_LEASE_MS: u64 = 86_400_000; +const SIGNING_CLAIM_OWNER_EXACT: &str = "radroots-sync-signing-exact"; +const SIGNING_CLAIM_OWNER_LOCAL: &str = "radroots-sync-signing-local"; +const SIGNING_CLAIM_OWNER_NON_REPLAYABLE: &str = "radroots-sync-signing-non-replayable"; +const ADMISSION_CLAIM_OWNER: &str = "radroots-sync-admission"; +const WORK_RETRY_DELAY_MS: u64 = 1_000; /// Caller-owned, replay-stable inputs for one outbound operation. #[derive(Clone)] @@ -176,6 +190,38 @@ impl PushStatus { } } +/// Durable result of one bounded signing execution. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct SigningRunReceipt { + artifact: AuthoredArtifact, + replay: bool, +} + +impl SigningRunReceipt { + pub const fn artifact(&self) -> &AuthoredArtifact { + &self.artifact + } + pub const fn is_replay(&self) -> bool { + self.replay + } +} + +/// Durable result of one bounded local-admission execution. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct AdmissionRunReceipt { + artifact: AuthoredArtifact, + replay: bool, +} + +impl AdmissionRunReceipt { + pub const fn artifact(&self) -> &AuthoredArtifact { + &self.artifact + } + pub const fn is_replay(&self) -> bool { + self.replay + } +} + /// Durable result after signing and atomic outbox enqueue. #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] #[derive(Clone, Debug, Eq, PartialEq)] @@ -355,14 +401,378 @@ 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)?; + let artifact = status.artifact; + match artifact.signing_state() { + SigningState::Signed => { + return Ok(SigningRunReceipt { + artifact, + replay: true, + }); + } + SigningState::Indeterminate => return Err(Error::SigningIndeterminate), + SigningState::FailedTerminal | SigningState::Cancelled => { + return Err(Error::SignerFailed); + } + SigningState::Planned | SigningState::Retryable => {} + } + + let now = self.clock.now_unix_ms()?.max(artifact.updated_at_unix_ms()); + if let Some(existing) = artifact.signing_claim() { + if now < existing.expires_at_unix_ms() { + return Err(Error::WorkClaimConflict); + } + if existing.owner() == SIGNING_CLAIM_OWNER_NON_REPLAYABLE { + let (claimed, fence) = self + .claim_artifact( + artifact, + ClaimAuthoredTarget::ArtifactSigning, + SIGNING_CLAIM_OWNER_NON_REPLAYABLE, + now, + ) + .await?; + let failure = WorkFailure::new( + "signing_effect_unknown_after_restart", + WorkPhase::Signing, + FailureClass::Indeterminate, + None, + None, + ) + .map_err(map_storage_error)?; + self.apply_artifact_failure(claimed.artifact_id(), fence, failure, None) + .await?; + return Err(Error::SigningIndeterminate); + } + } + + let signer_status = signer.status().await.map_err(|_| Error::SignerFailed)?; + let replay_capability = signer_replay_capability(&signer_status)?; + let (claimed, fence) = self + .claim_artifact( + artifact, + ClaimAuthoredTarget::ArtifactSigning, + signing_claim_owner(replay_capability), + now, + ) + .await?; + let persisted_plan = claimed + .plan() + .ok_or(Error::StorageFailed)? + .decode() + .map_err(map_storage_error)? + .into_plan(); + if persisted_plan != request.plan { + return Err(Error::StorageConflict); + } + let deadline = self.deadlines.deadline_unix_ms(OperationKind::Sign, now)?; + let signing_operation = SigningOperationId::new(*request.operation_id.as_bytes()) + .map_err(|_| Error::InvalidPushRequest)?; + let signing_artifact = SigningArtifactId::new(*claimed.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, + persisted_plan, + SignPolicy::new(deadline, request.cancellation) + .map_err(|_| Error::InvalidPushRequest)?, + ) + .map_err(|_| Error::InvalidPushRequest)?; + + match signer.sign(sign_request).await { + Ok(receipt) => { + let command = AuthoredAtomicCommand::ApplySigned( + ApplySignedArtifact::new( + claimed.artifact_id(), + fence, + receipt.signed_event().clone(), + receipt.completed_at_unix_ms(), + ) + .map_err(map_storage_error)?, + ); + let applied = self + .storage + .execute_authored(command) + .await + .map_err(map_storage_error)?; + let AuthoredAtomicOutcome::Artifact(artifact) = applied.outcome() else { + return Err(Error::StorageFailed); + }; + Ok(SigningRunReceipt { + artifact: artifact.clone(), + replay: false, + }) + } + Err(error) => { + let applied_at = self.clock.now_unix_ms()?; + if error.kind() == radroots_signing::error::Kind::SignerCancelled { + let command = AuthoredAtomicCommand::Cancel( + CancelAuthoredWork::new( + CancelAuthoredTarget::ArtifactSigning(claimed.artifact_id()), + claimed.revision(), + applied_at, + ) + .map_err(map_storage_error)?, + ); + self.storage + .execute_authored(command) + .await + .map_err(map_storage_error)?; + return Err(Error::SigningCancelled); + } + let disposition = recovery_disposition( + replay_capability, + error.remote_effect(), + error.retryable(), + ); + let class = match disposition { + RecoveryDisposition::RetryExactRequest | RecoveryDisposition::RetryLocal => { + FailureClass::Retryable + } + RecoveryDisposition::Indeterminate => FailureClass::Indeterminate, + RecoveryDisposition::Failed => FailureClass::Terminal, + _ => FailureClass::Indeterminate, + }; + let retry_at = if class == FailureClass::Retryable { + Some( + applied_at + .checked_add(WORK_RETRY_DELAY_MS) + .ok_or(Error::DeadlineOverflow)?, + ) + } else { + None + }; + let failure = + WorkFailure::new(error.code(), WorkPhase::Signing, class, retry_at, None) + .map_err(map_storage_error)?; + let retry = retry_schedule(claimed.signing_retry(), &failure, retry_at)?; + self.apply_artifact_failure(claimed.artifact_id(), fence, failure, retry) + .await?; + match disposition { + RecoveryDisposition::Indeterminate => Err(Error::SigningIndeterminate), + RecoveryDisposition::RetryExactRequest + | RecoveryDisposition::RetryLocal + | RecoveryDisposition::Failed => { + if error.kind() == radroots_signing::error::Kind::DeadlineExceeded { + Err(Error::SignerDeadlineExceeded) + } else { + Err(Error::SignerFailed) + } + } + _ => Err(Error::SigningIndeterminate), + } + } + } + } + + /// Claims and executes local admission for one durably signed artifact. + pub async fn admit_signed(&self, operation_id: SyncId) -> Result<AdmissionRunReceipt, Error> { + let status = self + .push_status(operation_id) + .await? + .ok_or(Error::StorageFailed)?; + let artifact = status.artifact; + if artifact.signing_state() != SigningState::Signed { + return Err(Error::InvalidSignerOutput); + } + if artifact.admission_state().is_admitted() { + return Ok(AdmissionRunReceipt { + artifact, + replay: true, + }); + } + if matches!( + artifact.admission_state(), + AdmissionState::Rejected | AdmissionState::Cancelled + ) { + return Err(Error::AdmissionFailed); + } + let now = self.clock.now_unix_ms()?.max(artifact.updated_at_unix_ms()); + if artifact + .admission_claim() + .is_some_and(|claim| now < claim.expires_at_unix_ms()) + { + return Err(Error::WorkClaimConflict); + } + let (claimed, fence) = self + .claim_artifact( + artifact, + ClaimAuthoredTarget::ArtifactAdmission, + ADMISSION_CLAIM_OWNER, + now, + ) + .await?; + let event = claimed + .signed() + .ok_or(Error::InvalidSignerOutput)? + .event() + .clone(); + let admission = match outbound_admission(&event, now) { + Ok(admission) => admission, + Err(error) => { + let failure = WorkFailure::new( + "invalid_signed_artifact", + WorkPhase::Admission, + FailureClass::Terminal, + None, + None, + ) + .map_err(map_storage_error)?; + self.apply_artifact_failure(claimed.artifact_id(), fence, failure, None) + .await?; + return Err(error); + } + }; + let admission_receipt = match EventStore::admit(self.storage.as_ref(), admission).await { + Ok(receipt) => receipt, + Err(error) => { + let applied_at = self.clock.now_unix_ms()?; + let terminal = matches!( + error, + radroots_storage::Error::EventConflict + | radroots_storage::Error::AdmissionRegression + ); + let retry_at = if terminal { + None + } else { + Some( + applied_at + .checked_add(WORK_RETRY_DELAY_MS) + .ok_or(Error::DeadlineOverflow)?, + ) + }; + let failure = WorkFailure::new( + if terminal { + "admission_conflict" + } else { + "admission_storage_unavailable" + }, + WorkPhase::Admission, + if terminal { + FailureClass::Terminal + } else { + FailureClass::Retryable + }, + retry_at, + None, + ) + .map_err(map_storage_error)?; + let retry = retry_schedule(claimed.admission_retry(), &failure, retry_at)?; + self.apply_artifact_failure(claimed.artifact_id(), fence, failure, retry) + .await?; + return Err(Error::AdmissionFailed); + } + }; + let state = match admission_receipt.disposition() { + AdmissionDisposition::Duplicate => AdmissionState::Duplicate, + AdmissionDisposition::Inserted | AdmissionDisposition::Advanced => { + AdmissionState::Inserted + } + }; + let applied_at = self.clock.now_unix_ms()?.max(claimed.updated_at_unix_ms()); + let command = AuthoredAtomicCommand::ApplyAdmission( + ApplyAdmissionResult::new(claimed.artifact_id(), fence, state, None, None, applied_at) + .map_err(map_storage_error)?, + ); + let applied = self + .storage + .execute_authored(command) + .await + .map_err(map_storage_error)?; + let AuthoredAtomicOutcome::Artifact(artifact) = applied.outcome() else { + return Err(Error::StorageFailed); + }; + Ok(AdmissionRunReceipt { + artifact: artifact.clone(), + replay: false, + }) + } + + async fn claim_artifact( + &self, + artifact: AuthoredArtifact, + target: fn(AuthoredArtifactId) -> ClaimAuthoredTarget, + owner: &'static str, + acquired_at: u64, + ) -> Result<(AuthoredArtifact, WorkFence), Error> { + let existing = match target(artifact.artifact_id()) { + ClaimAuthoredTarget::ArtifactSigning(_) => artifact.signing_claim(), + ClaimAuthoredTarget::ArtifactAdmission(_) => artifact.admission_claim(), + ClaimAuthoredTarget::DeliveryPlan(_) => return Err(Error::StorageFailed), + }; + let generation = existing.map_or(1, |claim| claim.generation().get().saturating_add(1)); + let generation = NonZeroU64::new(generation).ok_or(Error::StorageFailed)?; + let expires_at = acquired_at + .checked_add(self.deadlines.timeout_ms(OperationKind::Sign)) + .ok_or(Error::DeadlineOverflow)?; + let claim = WorkClaim::new( + *self.ids.next_id(OperationKind::Sign)?.as_bytes(), + owner, + generation, + acquired_at, + expires_at, + artifact.revision(), + ) + .map_err(map_storage_error)?; + let fence = WorkFence::new(*claim.token(), claim.generation(), claim.row_revision()) + .map_err(map_storage_error)?; + let command = AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( + target(artifact.artifact_id()), + claim, + )); + let receipt = self + .storage + .execute_authored(command) + .await + .map_err(map_claim_error)?; + let AuthoredAtomicOutcome::Artifact(claimed) = receipt.outcome() else { + return Err(Error::StorageFailed); + }; + Ok((claimed.clone(), fence)) + } + + async fn apply_artifact_failure( + &self, + artifact_id: AuthoredArtifactId, + fence: WorkFence, + failure: WorkFailure, + retry: Option<RetrySchedule>, + ) -> Result<AuthoredArtifact, Error> { + let applied_at = self.clock.now_unix_ms()?; + let command = AuthoredAtomicCommand::ApplyFailure( + ApplyWorkFailure::new( + AuthoredWorkTarget::Artifact(artifact_id), + fence, + failure, + retry, + applied_at, + ) + .map_err(map_storage_error)?, + ); + let receipt = self + .storage + .execute_authored(command) + .await + .map_err(map_storage_error)?; + let AuthoredAtomicOutcome::Artifact(artifact) = receipt.outcome() else { + return Err(Error::StorageFailed); + }; + Ok(artifact.clone()) + } + /// 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 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 = @@ -395,6 +805,15 @@ impl Engine { } } + let signed = self.sign_prepared(request.clone()).await?; + self.admit_signed(request.operation_id).await?; + let event = signed + .artifact() + .signed() + .ok_or(Error::InvalidSignerOutput)? + .event() + .clone(); + let prepared_at = self.clock.now_unix_ms()?; let prepare = PrepareOperation::new( instance_id, @@ -420,35 +839,6 @@ impl Engine { return Err(Error::StorageFailed); }; - let sign_deadline_ms = self - .deadlines - .deadline_unix_ms(OperationKind::Sign, self.clock.now_unix_ms()?)?; - 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.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| { - 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(); - let signed_record = if prepared.state().stage() == JournalStage::Signed { if !matches!( prepared.state(), @@ -635,6 +1025,46 @@ fn synthetic_receipt(request: &DeliveryRequest, retryable: bool) -> Result<Deliv DeliveryReceipt::for_request(request, targets).map_err(|_| Error::InvalidDeliveryRequest) } +fn signer_replay_capability( + status: &radroots_signing::SignerStatus, +) -> Result<ReplayCapability, Error> { + let mut capabilities = status.capabilities().iter(); + let first = capabilities + .next() + .ok_or(Error::SignerCapabilityUnavailable)? + .replay(); + if capabilities.all(|capability| capability.replay() == first) { + Ok(first) + } else { + Ok(ReplayCapability::NonReplayable) + } +} + +fn signing_claim_owner(replay: ReplayCapability) -> &'static str { + match replay { + ReplayCapability::ExactReplayByRequestId => SIGNING_CLAIM_OWNER_EXACT, + ReplayCapability::LocalReplaySafe => SIGNING_CLAIM_OWNER_LOCAL, + ReplayCapability::NonReplayable => SIGNING_CLAIM_OWNER_NON_REPLAYABLE, + _ => SIGNING_CLAIM_OWNER_NON_REPLAYABLE, + } +} + +fn retry_schedule( + previous: Option<&RetrySchedule>, + failure: &WorkFailure, + retry_at: Option<u64>, +) -> Result<Option<RetrySchedule>, Error> { + let Some(retry_at) = retry_at else { + return Ok(None); + }; + let schedule = match previous { + Some(previous) => previous.next_attempt(retry_at, failure.clone()), + None => RetrySchedule::new(NonZeroU32::MIN, retry_at, failure.clone()), + } + .map_err(map_storage_error)?; + Ok(Some(schedule)) +} + fn delivery_evidence_digest(evidence: &DeliveryAttemptEvidence) -> AtomicCommitDigest { let mut hasher = Sha256::new(); hash_field(&mut hasher, b"radroots.sync.delivery-evidence.v1"); @@ -877,3 +1307,13 @@ fn map_storage_error(error: radroots_storage::Error) -> Error { _ => Error::StorageFailed, } } + +fn map_claim_error(error: radroots_storage::Error) -> Error { + match error { + radroots_storage::Error::AtomicCommitConflict + | radroots_storage::Error::InvalidAuthoredTransition + | radroots_storage::Error::InvalidWorkClaim + | radroots_storage::Error::DeliveryPlanClaimConflict => Error::WorkClaimConflict, + other => map_storage_error(other), + } +} diff --git a/crates/sync/tests/push_enqueue.rs b/crates/sync/tests/push_enqueue.rs @@ -15,12 +15,17 @@ use radroots_event_codec::authoring::AuthoredEventPlan; use radroots_protocol::runtime::v1::SyncRetryDecision; use radroots_signing::{ Actor, Error as SigningError, SignReceipt, SignRequest, Signer, SignerStatus, - actor::ActorSource, error::Kind as SigningErrorKind, request::CancellationPolicy, + actor::ActorSource, + capability::{CancellationSupport, SignerCapability, SignerKind}, + error::Kind as SigningErrorKind, + recovery::ReplayCapability, + request::CancellationPolicy, + status::SignerAvailability, }; use radroots_storage::{ EventStore, Journal, Outbox, event::{EventQuery, EventQueryBounds, SourceGeneration}, - journal::{IdempotencyKey, JournalStage, OperationInstanceId}, + journal::{IdempotencyKey, OperationInstanceId}, memory::MemoryStorage, outbox::{ClaimOutboxItems, LeaseId, LeaseOwner, OutboxStage, SatisfactionResult}, }; @@ -150,11 +155,13 @@ impl IdSource for TestIds { enum SignBehavior { Success { completed_at_unix_ms: u64 }, Error(SigningErrorKind), + Uncertain(SigningErrorKind), Pending, } struct MockSigner { behavior: SignBehavior, + replay: ReplayCapability, calls: AtomicUsize, } @@ -162,6 +169,15 @@ impl MockSigner { fn new(behavior: SignBehavior) -> Self { Self { behavior, + replay: ReplayCapability::LocalReplaySafe, + calls: AtomicUsize::new(0), + } + } + + fn with_replay(behavior: SignBehavior, replay: ReplayCapability) -> Self { + Self { + behavior, + replay, calls: AtomicUsize::new(0), } } @@ -171,7 +187,24 @@ impl Signer for MockSigner { fn status( &self, ) -> radroots_signing::signer::BoxFuture<'_, Result<SignerStatus, SigningError>> { - Box::pin(async { unreachable!("enqueue does not inspect signer status") }) + let replay = self.replay; + Box::pin(async move { + Ok(SignerStatus::new( + SignerAvailability::Ready, + vec![SignerCapability::new( + if replay == ReplayCapability::LocalReplaySafe { + SignerKind::Local + } else { + SignerKind::Remote + }, + replay, + CancellationSupport::BeforePublication, + false, + false, + )], + None, + )) + }) } fn sign( @@ -189,6 +222,9 @@ impl Signer for MockSigner { completed_at_unix_ms, ), SignBehavior::Error(kind) => Err(SigningError::new(kind)), + SignBehavior::Uncertain(kind) => { + Err(SigningError::new(kind).with_possible_remote_effect()) + } SignBehavior::Pending => std::future::pending().await, } }) @@ -436,8 +472,222 @@ fn preparation_is_atomic_status_visible_and_replays_without_external_effects() { assert_eq!(signer.calls.load(Ordering::Relaxed), 0); } +#[test] +fn signing_claim_recovery_respects_exact_and_non_replayable_capabilities() { + let exact_storage = Arc::new(MemoryStorage::new( + SourceGeneration::new([8; 32]).expect("generation"), + )); + let exact_pending = Arc::new(MockSigner::with_replay( + SignBehavior::Pending, + ReplayCapability::ExactReplayByRequestId, + )); + let exact_engine = Engine::builder( + exact_storage.clone(), + Arc::new(TestClock(AtomicU64::new(1_800_000_200_000))), + Arc::new(TestIds(AtomicU64::new(100))), + DeadlinePolicy::new(10_000, 10_000, 10_000).expect("deadlines"), + ) + .sink(Arc::new(MockSink)) + .signer(exact_pending.clone()) + .build() + .expect("exact engine"); + let exact_push = request(8, "wss://relay.example"); + let mut exact_future = Box::pin(exact_engine.sign_prepared(exact_push.clone())).fuse(); + let mut context = std::task::Context::from_waker(noop_waker_ref()); + assert!(exact_future.poll_unpin(&mut context).is_pending()); + drop(exact_future); + let claimed = block_on(exact_engine.push_status(exact_push.operation_id())) + .expect("status") + .expect("claimed status"); + assert!(claimed.artifact().signing_claim().is_some()); + + let exact_success = Arc::new(MockSigner::with_replay( + SignBehavior::Success { + completed_at_unix_ms: 1_800_000_220_000, + }, + ReplayCapability::ExactReplayByRequestId, + )); + let exact_recovery = Engine::builder( + exact_storage, + Arc::new(TestClock(AtomicU64::new(1_800_000_211_000))), + Arc::new(TestIds(AtomicU64::new(110))), + DeadlinePolicy::new(10_000, 10_000, 10_000).expect("deadlines"), + ) + .sink(Arc::new(MockSink)) + .signer(exact_success.clone()) + .build() + .expect("recovery engine"); + let recovered = + block_on(exact_recovery.sign_prepared(exact_push.clone())).expect("exact replay recovery"); + assert_eq!( + recovered.artifact().signing_state(), + radroots_storage::authored::SigningState::Signed + ); + assert_eq!(exact_success.calls.load(Ordering::Relaxed), 1); + let replay = block_on(exact_recovery.sign_prepared(exact_push)).expect("signed replay"); + assert!(replay.is_replay()); + assert_eq!(exact_success.calls.load(Ordering::Relaxed), 1); + + let unsafe_storage = Arc::new(MemoryStorage::new( + SourceGeneration::new([9; 32]).expect("generation"), + )); + let unsafe_pending = Arc::new(MockSigner::with_replay( + SignBehavior::Pending, + ReplayCapability::NonReplayable, + )); + let unsafe_engine = Engine::builder( + unsafe_storage.clone(), + Arc::new(TestClock(AtomicU64::new(1_800_000_200_000))), + Arc::new(TestIds(AtomicU64::new(120))), + DeadlinePolicy::new(10_000, 10_000, 10_000).expect("deadlines"), + ) + .sink(Arc::new(MockSink)) + .signer(unsafe_pending) + .build() + .expect("unsafe engine"); + let unsafe_push = request(9, "wss://relay.example"); + let mut unsafe_future = Box::pin(unsafe_engine.sign_prepared(unsafe_push.clone())).fuse(); + let mut context = std::task::Context::from_waker(noop_waker_ref()); + assert!(unsafe_future.poll_unpin(&mut context).is_pending()); + drop(unsafe_future); + + let unsafe_success = Arc::new(MockSigner::with_replay( + SignBehavior::Success { + completed_at_unix_ms: 1_800_000_220_000, + }, + ReplayCapability::NonReplayable, + )); + let unsafe_recovery = Engine::builder( + unsafe_storage, + Arc::new(TestClock(AtomicU64::new(1_800_000_211_000))), + Arc::new(TestIds(AtomicU64::new(130))), + DeadlinePolicy::new(10_000, 10_000, 10_000).expect("deadlines"), + ) + .sink(Arc::new(MockSink)) + .signer(unsafe_success.clone()) + .build() + .expect("unsafe recovery engine"); + assert_eq!( + block_on(unsafe_recovery.sign_prepared(unsafe_push.clone())), + Err(Error::SigningIndeterminate) + ); + assert_eq!(unsafe_success.calls.load(Ordering::Relaxed), 0); + let unsafe_status = block_on(unsafe_recovery.push_status(unsafe_push.operation_id())) + .expect("unsafe status") + .expect("unsafe operation"); + assert_eq!( + unsafe_status.artifact().signing_state(), + radroots_storage::authored::SigningState::Indeterminate + ); +} + +#[test] +fn signing_failures_persist_retry_indeterminate_terminal_and_cancelled_states() { + let exact = Arc::new(MockSigner::with_replay( + SignBehavior::Uncertain(SigningErrorKind::SignerTimeout), + ReplayCapability::ExactReplayByRequestId, + )); + let (engine, storage) = setup_engine(exact); + let retryable = request(10, "wss://relay.example"); + assert_eq!( + block_on(engine.sign_prepared(retryable.clone())), + Err(Error::SignerFailed) + ); + let retryable_status = block_on(engine.push_status(retryable.operation_id())) + .expect("retryable status") + .expect("retryable operation"); + assert_eq!( + retryable_status.artifact().signing_state(), + radroots_storage::authored::SigningState::Retryable + ); + let retry_at = retryable_status + .artifact() + .signing_retry() + .expect("retry schedule") + .not_before_unix_ms(); + let success = Arc::new(MockSigner::with_replay( + SignBehavior::Success { + completed_at_unix_ms: retry_at + 2, + }, + ReplayCapability::ExactReplayByRequestId, + )); + let recovery = Engine::builder( + storage, + Arc::new(TestClock(AtomicU64::new(retry_at + 1))), + Arc::new(TestIds(AtomicU64::new(140))), + DeadlinePolicy::new(10_000, 10_000, 10_000).expect("deadlines"), + ) + .sink(Arc::new(MockSink)) + .signer(success.clone()) + .build() + .expect("recovery engine"); + assert_eq!( + block_on(recovery.sign_prepared(retryable)) + .expect("retry exact request") + .artifact() + .signing_state(), + radroots_storage::authored::SigningState::Signed + ); + assert_eq!(success.calls.load(Ordering::Relaxed), 1); + + let non_replayable = Arc::new(MockSigner::with_replay( + SignBehavior::Uncertain(SigningErrorKind::SignerTimeout), + ReplayCapability::NonReplayable, + )); + let (engine, _) = setup_engine(non_replayable); + let uncertain = request(11, "wss://relay.example"); + assert_eq!( + block_on(engine.sign_prepared(uncertain.clone())), + Err(Error::SigningIndeterminate) + ); + assert_eq!( + block_on(engine.push_status(uncertain.operation_id())) + .expect("indeterminate status") + .expect("indeterminate operation") + .artifact() + .signing_state(), + radroots_storage::authored::SigningState::Indeterminate + ); + + let cancelled = Arc::new(MockSigner::new(SignBehavior::Error( + SigningErrorKind::SignerCancelled, + ))); + let (engine, _) = setup_engine(cancelled); + let cancelled_push = request(12, "wss://relay.example"); + assert_eq!( + block_on(engine.sign_prepared(cancelled_push.clone())), + Err(Error::SigningCancelled) + ); + assert_eq!( + block_on(engine.push_status(cancelled_push.operation_id())) + .expect("cancelled status") + .expect("cancelled operation") + .artifact() + .signing_state(), + radroots_storage::authored::SigningState::Cancelled + ); + + let terminal = Arc::new(MockSigner::new(SignBehavior::Error( + SigningErrorKind::SignerRejected, + ))); + let (engine, _) = setup_engine(terminal); + let terminal_push = request(13, "wss://relay.example"); + assert_eq!( + block_on(engine.sign_prepared(terminal_push.clone())), + Err(Error::SignerFailed) + ); + assert_eq!( + block_on(engine.push_status(terminal_push.operation_id())) + .expect("terminal status") + .expect("terminal operation") + .artifact() + .signing_state(), + radroots_storage::authored::SigningState::FailedTerminal + ); +} + #[tokio::test] -async fn sqlite_preparation_reopens_and_replays_the_original_complete_intent() { +async fn sqlite_signing_and_admission_recover_across_every_reopen_boundary() { let directory = tempfile::tempdir().expect("database directory"); let paths = Paths::from_directory(directory.path()).expect("paths"); let push = request(7, "wss://relay.example"); @@ -469,37 +719,120 @@ async fn sqlite_preparation_reopens_and_replays_the_original_complete_intent() { prepared }; + { + let store = Arc::new( + SqliteStorage::open(OpenOptions::new(paths.clone(), OpenMode::ReadWriteExisting)) + .await + .expect("reopen for signing"), + ); + let capability: Arc<dyn SyncStorage> = store; + let signer = Arc::new(MockSigner::new(SignBehavior::Success { + completed_at_unix_ms: 1_800_000_210_500, + })); + 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.clone()) + .await + .expect("prepare replay"); + assert!(replay.is_replay()); + assert_eq!(signer.calls.load(Ordering::Relaxed), 0); + let signed = engine + .sign_prepared(push.clone()) + .await + .expect("sign after reopen"); + assert_eq!( + signed.artifact().signing_state(), + radroots_storage::authored::SigningState::Signed + ); + assert_eq!( + signed.artifact().admission_state(), + radroots_storage::authored::AdmissionState::Pending + ); + assert!(signed.artifact().signed().is_some()); + assert_eq!(signer.calls.load(Ordering::Relaxed), 1); + } + + { + let store = Arc::new( + SqliteStorage::open(OpenOptions::new(paths.clone(), OpenMode::ReadWriteExisting)) + .await + .expect("reopen for admission"), + ); + let capability: Arc<dyn SyncStorage> = store.clone(); + let signer = Arc::new(MockSigner::new(SignBehavior::Pending)); + let engine = Engine::builder( + capability, + Arc::new(TestClock(AtomicU64::new(1_800_000_211_000))), + Arc::new(TestIds(AtomicU64::new(100))), + DeadlinePolicy::new(10_000, 10_000, 10_000).expect("deadlines"), + ) + .sink(Arc::new(MockSink)) + .signer(signer.clone()) + .build() + .expect("engine"); + let admitted = engine + .admit_signed(push.operation_id()) + .await + .expect("admit after reopen"); + assert!(!admitted.is_replay()); + assert!(admitted.artifact().admission_state().is_admitted()); + let replay = engine + .admit_signed(push.operation_id()) + .await + .expect("admission replay"); + assert!(replay.is_replay()); + assert_eq!(signer.calls.load(Ordering::Relaxed), 0); + let visible = store + .query_visible(EventQuery::all( + EventQueryBounds::first(10).expect("bounds"), + )) + .await + .expect("visible events"); + assert_eq!(visible.items().len(), 1); + } + let store = Arc::new( - SqliteStorage::open(OpenOptions::new(paths, OpenMode::ReadWriteExisting)) + SqliteStorage::open(OpenOptions::new(paths, OpenMode::ReadOnly)) .await - .expect("reopen SQLite"), + .expect("final read-only reopen"), ); - 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))), + store, + Arc::new(TestClock(AtomicU64::new(1_800_000_212_000))), + Arc::new(TestIds(AtomicU64::new(110))), DeadlinePolicy::new(10_000, 10_000, 10_000).expect("deadlines"), ) .sink(Arc::new(MockSink)) - .signer(signer.clone()) .build() - .expect("engine"); - let status = engine + .expect("read-only engine"); + let final_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); + .expect("final status") + .expect("durable operation"); + assert_eq!( + final_status.artifact().signing_state(), + radroots_storage::authored::SigningState::Signed + ); + assert!(final_status.artifact().admission_state().is_admitted()); + assert!(final_status.delivery_plan().request().is_some()); } #[test] @@ -551,17 +884,27 @@ fn idempotency_conflict_and_cancellation_before_commit_fail_closed() { 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 pending_push = request(5, "wss://relay.example"); + let mut future = Box::pin(engine.sign_and_enqueue(pending_push.clone())).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); + let status = block_on(engine.push_status(SyncId::new([5; 16]).expect("operation"))) + .expect("authored status") + .expect("prepared operation"); + assert!(status.artifact().signing_claim().is_some()); + assert_eq!( + block_on(engine.sign_prepared(pending_push)), + Err(Error::WorkClaimConflict) + ); + assert!( + block_on(Journal::operation( + &*storage, + OperationInstanceId::new([5; 16]).expect("instance id"), + )) + .expect("legacy journal lookup") + .is_none() + ); assert!( block_on(Outbox::item( &*storage,