lib

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

push.rs (60862B)


      1 //! Signing, durable enqueue, delivery, and satisfaction orchestration.
      2 
      3 use core::num::{NonZeroU32, NonZeroU64};
      4 use radroots_event::admission::RawEvent;
      5 use radroots_event_codec::{
      6     authoring::AuthoredEventPlan,
      7     verify::{self, Nip01SignatureVerifier},
      8 };
      9 use radroots_protocol::runtime::v1::OperationId;
     10 use radroots_signing::{
     11     Actor, AuthoredArtifactId as SigningArtifactId, SigningIntentId, SigningOperationId,
     12     recovery::{RecoveryDisposition, ReplayCapability, recovery_disposition},
     13     request::{CancellationPolicy, SignPolicy},
     14 };
     15 use radroots_storage::{
     16     atomic::{AtomicCommitDigest, AtomicCommitDisposition},
     17     authored::{
     18         AdmissionState, AuthoredArtifact, AuthoredArtifactId, AuthoredOperation, FailureClass,
     19         OperationSettlement, RetrySchedule, SigningState, WorkClaim, WorkFailure, WorkPhase,
     20     },
     21     authored_atomic::{
     22         ApplyAdmissionResult, ApplyWorkFailure, AuthoredAtomicCommand, AuthoredAtomicOutcome,
     23         AuthoredWorkTarget, CancelAuthoredTarget, CancelAuthoredWork, ClaimAuthoredTarget,
     24         ClaimAuthoredWork, PrepareAuthoredOperation, RecordSignedArtifact, WorkFence,
     25     },
     26     authored_delivery::{
     27         AuthoredDeliveryHistory, AuthoredDeliveryIntent, AuthoredDeliveryPlan,
     28         AuthoredDeliveryPlanId, AuthoredDeliveryState, DELIVERY_PLAN_ATTEMPTS_MAX,
     29         DeliveryAttemptOutcome,
     30     },
     31     event::{AdmissionDisposition, EventAdmission, EventStore},
     32     journal::{IdempotencyKey, OperationInstanceId},
     33 };
     34 use radroots_transport::{
     35     DeliveryRequest, SinkFailure, Target, TransportId,
     36     outcome::Retryability,
     37     policy::{SatisfactionPolicy, SatisfactionState},
     38     source::{EventProvenance, ObservedEvent},
     39     target::TargetSet,
     40 };
     41 use sha2::{Digest, Sha256};
     42 
     43 use crate::{
     44     Engine,
     45     ingest::{AdmissionDecision, AdmissionPolicy, RegistryPolicy},
     46     policy::{Error, OperationKind, SyncId},
     47 };
     48 
     49 mod delivery;
     50 
     51 const SIGNING_CLAIM_OWNER_EXACT: &str = "radroots-sync-signing-exact";
     52 const SIGNING_CLAIM_OWNER_LOCAL: &str = "radroots-sync-signing-local";
     53 const SIGNING_CLAIM_OWNER_NON_REPLAYABLE: &str = "radroots-sync-signing-non-replayable";
     54 const ADMISSION_CLAIM_OWNER: &str = "radroots-sync-admission";
     55 const DELIVERY_CLAIM_OWNER: &str = "radroots-sync-delivery";
     56 const WORK_RETRY_DELAY_MS: u64 = 1_000;
     57 
     58 /// Caller-owned, replay-stable inputs for one outbound operation.
     59 #[derive(Clone)]
     60 pub struct PushRequest {
     61     operation_id: SyncId,
     62     idempotency_key: IdempotencyKey,
     63     actor: Actor,
     64     plan: AuthoredEventPlan,
     65     targets: TargetSet,
     66     satisfaction: SatisfactionPolicy,
     67     delivery_deadline_unix_ms: u64,
     68     cancellation: CancellationPolicy,
     69 }
     70 
     71 impl PushRequest {
     72     #[allow(clippy::too_many_arguments)]
     73     pub fn new(
     74         operation_id: SyncId,
     75         idempotency_key: IdempotencyKey,
     76         actor: Actor,
     77         plan: AuthoredEventPlan,
     78         targets: TargetSet,
     79         satisfaction: SatisfactionPolicy,
     80         delivery_deadline_unix_ms: u64,
     81         cancellation: CancellationPolicy,
     82     ) -> Result<Self, Error> {
     83         if satisfaction.validate_for(&targets).is_err()
     84             || delivery_deadline_unix_ms == 0
     85             || actor.public_key() != *plan.author()
     86         {
     87             return Err(Error::InvalidPushRequest);
     88         }
     89         Ok(Self {
     90             operation_id,
     91             idempotency_key,
     92             actor,
     93             plan,
     94             targets,
     95             satisfaction,
     96             delivery_deadline_unix_ms,
     97             cancellation,
     98         })
     99     }
    100 
    101     /// Builds the existing pure preparation at a caller-captured stable time.
    102     /// No storage, signer, clock or network is invoked. Retain this value for
    103     /// atomic submission replay, which compares captured timestamps exactly.
    104     pub fn authored_preparation(
    105         &self,
    106         captured_at_unix_ms: u64,
    107     ) -> Result<PrepareAuthoredOperation, Error> {
    108         let (operation_id, artifact_id, delivery_plan_id) = authored_ids(self.operation_id)?;
    109         let operation =
    110             AuthoredOperation::new(operation_id, vec![artifact_id], captured_at_unix_ms)
    111                 .map_err(map_storage_error)?;
    112         let artifact = AuthoredArtifact::planned(
    113             artifact_id,
    114             operation_id,
    115             0,
    116             self.plan(),
    117             captured_at_unix_ms,
    118         )
    119         .map_err(map_storage_error)?;
    120         let intent = AuthoredDeliveryIntent::new(
    121             delivery_request_id(self.operation_id),
    122             self.targets.clone(),
    123             self.satisfaction.clone(),
    124             self.delivery_deadline_unix_ms,
    125         )
    126         .map_err(map_storage_error)?;
    127         let delivery_plan =
    128             AuthoredDeliveryPlan::new(delivery_plan_id, artifact_id, intent, captured_at_unix_ms)
    129                 .map_err(map_storage_error)?;
    130         PrepareAuthoredOperation::new(
    131             operation,
    132             vec![artifact],
    133             vec![delivery_plan],
    134             authored_push_input_digest(self)?,
    135             captured_at_unix_ms,
    136         )
    137         .map_err(map_storage_error)
    138     }
    139 
    140     pub const fn operation_id(&self) -> SyncId {
    141         self.operation_id
    142     }
    143     pub const fn idempotency_key(&self) -> &IdempotencyKey {
    144         &self.idempotency_key
    145     }
    146     pub const fn actor(&self) -> &Actor {
    147         &self.actor
    148     }
    149     pub const fn plan(&self) -> &AuthoredEventPlan {
    150         &self.plan
    151     }
    152     pub const fn targets(&self) -> &TargetSet {
    153         &self.targets
    154     }
    155     pub const fn satisfaction(&self) -> &SatisfactionPolicy {
    156         &self.satisfaction
    157     }
    158     pub const fn delivery_deadline_unix_ms(&self) -> u64 {
    159         self.delivery_deadline_unix_ms
    160     }
    161     pub const fn cancellation(&self) -> CancellationPolicy {
    162         self.cancellation
    163     }
    164 }
    165 
    166 impl core::fmt::Debug for PushRequest {
    167     fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
    168         formatter
    169             .debug_struct("PushRequest")
    170             .field("operation_id", &self.operation_id)
    171             .field("idempotency_key", &self.idempotency_key)
    172             .field("actor", &self.actor)
    173             .field("plan", &"[redacted exact authored plan]")
    174             .field("targets", &self.targets)
    175             .field("satisfaction", &self.satisfaction)
    176             .field("delivery_deadline_unix_ms", &self.delivery_deadline_unix_ms)
    177             .field("cancellation", &self.cancellation)
    178             .finish()
    179     }
    180 }
    181 
    182 /// Complete durable intent created before any signer or transport effect.
    183 #[derive(Clone, Debug, Eq, PartialEq)]
    184 pub struct PushPreparation {
    185     operation: AuthoredOperation,
    186     artifact: AuthoredArtifact,
    187     delivery_plan: AuthoredDeliveryPlan,
    188     replay: bool,
    189 }
    190 
    191 impl PushPreparation {
    192     pub const fn operation(&self) -> &AuthoredOperation {
    193         &self.operation
    194     }
    195     pub const fn artifact(&self) -> &AuthoredArtifact {
    196         &self.artifact
    197     }
    198     pub const fn delivery_plan(&self) -> &AuthoredDeliveryPlan {
    199         &self.delivery_plan
    200     }
    201     pub const fn is_replay(&self) -> bool {
    202         self.replay
    203     }
    204 }
    205 
    206 /// Current durable authored state for one push operation.
    207 #[derive(Clone, Debug, Eq, PartialEq)]
    208 pub struct PushStatus {
    209     operation: AuthoredOperation,
    210     artifact: AuthoredArtifact,
    211     delivery_history: AuthoredDeliveryHistory,
    212     settlement: OperationSettlement,
    213 }
    214 
    215 impl PushStatus {
    216     pub const fn operation(&self) -> &AuthoredOperation {
    217         &self.operation
    218     }
    219     pub const fn artifact(&self) -> &AuthoredArtifact {
    220         &self.artifact
    221     }
    222     pub const fn delivery_plan(&self) -> &AuthoredDeliveryPlan {
    223         self.delivery_history.plan()
    224     }
    225     /// Consistent original-claim history; missing provenance remains unknown.
    226     pub const fn delivery_history(&self) -> &AuthoredDeliveryHistory {
    227         &self.delivery_history
    228     }
    229     pub const fn settlement(&self) -> OperationSettlement {
    230         self.settlement
    231     }
    232 }
    233 
    234 /// Durable result of one bounded signing execution.
    235 #[derive(Clone, Debug, Eq, PartialEq)]
    236 pub struct SigningRunReceipt {
    237     artifact: AuthoredArtifact,
    238     replay: bool,
    239 }
    240 
    241 impl SigningRunReceipt {
    242     pub const fn artifact(&self) -> &AuthoredArtifact {
    243         &self.artifact
    244     }
    245     pub const fn is_replay(&self) -> bool {
    246         self.replay
    247     }
    248 }
    249 
    250 /// Durable result of one bounded local-admission execution.
    251 #[derive(Clone, Debug, Eq, PartialEq)]
    252 pub struct AdmissionRunReceipt {
    253     artifact: AuthoredArtifact,
    254     replay: bool,
    255 }
    256 
    257 /// Durable result of one bounded authored-delivery execution.
    258 #[derive(Clone, Debug, Eq, PartialEq)]
    259 pub struct DeliveryExecutionReceipt {
    260     plan: AuthoredDeliveryPlan,
    261     replay: bool,
    262 }
    263 
    264 /// Durable result of cancelling all still-cancellable phases of one push.
    265 #[derive(Clone, Debug, Eq, PartialEq)]
    266 pub struct PushCancellationReceipt {
    267     status: PushStatus,
    268     changed: bool,
    269 }
    270 
    271 impl PushCancellationReceipt {
    272     pub const fn status(&self) -> &PushStatus {
    273         &self.status
    274     }
    275     pub const fn changed(&self) -> bool {
    276         self.changed
    277     }
    278 }
    279 
    280 impl DeliveryExecutionReceipt {
    281     pub const fn plan(&self) -> &AuthoredDeliveryPlan {
    282         &self.plan
    283     }
    284     pub const fn is_replay(&self) -> bool {
    285         self.replay
    286     }
    287 }
    288 
    289 impl AdmissionRunReceipt {
    290     pub const fn artifact(&self) -> &AuthoredArtifact {
    291         &self.artifact
    292     }
    293     pub const fn is_replay(&self) -> bool {
    294         self.replay
    295     }
    296 }
    297 
    298 impl Engine {
    299     /// Atomically persists the complete parent, exact plan, and delivery intent.
    300     ///
    301     /// This method invokes neither a signer nor a transport. Exact request
    302     /// replay returns the original durable preparation; conflicting reuse of
    303     /// the operation identity fails closed.
    304     pub async fn prepare_push(&self, request: PushRequest) -> Result<PushPreparation, Error> {
    305         let prepared_at = self.clock.now_unix_ms()?;
    306         let (operation_id, artifact_id, delivery_plan_id) = authored_ids(request.operation_id)?;
    307         let command = AuthoredAtomicCommand::Prepare(request.authored_preparation(prepared_at)?);
    308         let receipt = self
    309             .storage
    310             .execute_authored(command)
    311             .await
    312             .map_err(map_storage_error)?;
    313         let AuthoredAtomicOutcome::Prepared {
    314             operation,
    315             artifacts,
    316             delivery_plans,
    317         } = receipt.outcome()
    318         else {
    319             return Err(Error::StorageFailed);
    320         };
    321         let [artifact] = artifacts.as_slice() else {
    322             return Err(Error::StorageFailed);
    323         };
    324         let [delivery_plan] = delivery_plans.as_slice() else {
    325             return Err(Error::StorageFailed);
    326         };
    327         if operation.operation_id() != operation_id
    328             || artifact.artifact_id() != artifact_id
    329             || delivery_plan.plan_id() != delivery_plan_id
    330         {
    331             return Err(Error::StorageFailed);
    332         }
    333         Ok(PushPreparation {
    334             operation: operation.clone(),
    335             artifact: artifact.clone(),
    336             delivery_plan: delivery_plan.clone(),
    337             replay: receipt.disposition() == AtomicCommitDisposition::Replay,
    338         })
    339     }
    340 
    341     /// Loads the complete durable state for one prepared push operation.
    342     pub async fn push_status(&self, operation_id: SyncId) -> Result<Option<PushStatus>, Error> {
    343         let (operation_id, expected_artifact_id, expected_plan_id) = authored_ids(operation_id)?;
    344         let Some(operation) = self
    345             .storage
    346             .authored_operation(operation_id)
    347             .await
    348             .map_err(map_storage_error)?
    349         else {
    350             return Ok(None);
    351         };
    352         if operation.artifact_ids() != [expected_artifact_id] {
    353             return Err(Error::StorageFailed);
    354         }
    355         let artifact = self
    356             .storage
    357             .authored_artifact(expected_artifact_id)
    358             .await
    359             .map_err(map_storage_error)?
    360             .ok_or(Error::StorageFailed)?;
    361         let delivery_history = self
    362             .storage
    363             .authored_delivery_history(expected_plan_id)
    364             .await
    365             .map_err(map_storage_error)?
    366             .ok_or(Error::StorageFailed)?;
    367         delivery_history.validate().map_err(map_storage_error)?;
    368         let delivery_plan = delivery_history.plan();
    369         if artifact.artifact_id() != expected_artifact_id
    370             || delivery_plan.plan_id() != expected_plan_id
    371             || artifact.operation_id() != operation.operation_id()
    372             || delivery_plan.artifact_id() != artifact.artifact_id()
    373         {
    374             return Err(Error::StorageFailed);
    375         }
    376         let settlement = OperationSettlement::evaluate_complete(
    377             &operation,
    378             core::slice::from_ref(&artifact),
    379             core::slice::from_ref(delivery_plan),
    380         )
    381         .map_err(map_storage_error)?;
    382         Ok(Some(PushStatus {
    383             operation,
    384             artifact,
    385             delivery_history,
    386             settlement,
    387         }))
    388     }
    389 
    390     /// Cancels every remaining local phase without discarding signed bytes or
    391     /// delivery evidence. Re-entry after a phase-boundary crash finishes any
    392     /// remaining cancellation and terminal replays are no-ops.
    393     pub async fn cancel_push(
    394         &self,
    395         operation_id: SyncId,
    396     ) -> Result<PushCancellationReceipt, Error> {
    397         let mut status = self
    398             .push_status(operation_id)
    399             .await?
    400             .ok_or(Error::StorageFailed)?;
    401         let mut changed = false;
    402         let now = self
    403             .clock
    404             .now_unix_ms()?
    405             .max(status.artifact.updated_at_unix_ms())
    406             .max(status.delivery_plan().updated_at_unix_ms());
    407 
    408         let artifact_target = match status.artifact.signing_state() {
    409             SigningState::Planned | SigningState::Retryable => Some(
    410                 CancelAuthoredTarget::ArtifactSigning(status.artifact.artifact_id()),
    411             ),
    412             SigningState::Signed
    413                 if matches!(
    414                     status.artifact.admission_state(),
    415                     AdmissionState::Pending | AdmissionState::Retryable
    416                 ) =>
    417             {
    418                 Some(CancelAuthoredTarget::ArtifactAdmission(
    419                     status.artifact.artifact_id(),
    420                 ))
    421             }
    422             SigningState::Signed
    423             | SigningState::Indeterminate
    424             | SigningState::FailedTerminal
    425             | SigningState::Cancelled => None,
    426         };
    427         if let Some(target) = artifact_target {
    428             self.storage
    429                 .execute_authored(AuthoredAtomicCommand::Cancel(
    430                     CancelAuthoredWork::new(target, status.artifact.revision(), now)
    431                         .map_err(map_storage_error)?,
    432                 ))
    433                 .await
    434                 .map_err(map_storage_error)?;
    435             changed = true;
    436             status = self
    437                 .push_status(operation_id)
    438                 .await?
    439                 .ok_or(Error::StorageFailed)?;
    440         }
    441 
    442         if status.delivery_plan().stop_requested_at_unix_ms().is_none() {
    443             let command = AuthoredAtomicCommand::Cancel(
    444                 CancelAuthoredWork::new(
    445                     CancelAuthoredTarget::DeliveryPlan(status.delivery_plan().plan_id()),
    446                     status.delivery_plan().revision(),
    447                     now,
    448                 )
    449                 .map_err(map_storage_error)?,
    450             );
    451             let receipt = self
    452                 .storage
    453                 .execute_authored(command.clone())
    454                 .await
    455                 .map_err(map_storage_error)?;
    456             if !receipt.matches_command(&command) {
    457                 return Err(Error::StorageFailed);
    458             }
    459             let AuthoredAtomicOutcome::DeliveryPlan(stopped) = receipt.outcome() else {
    460                 return Err(Error::StorageFailed);
    461             };
    462             if stopped.validate().is_err()
    463                 || stopped.plan_id() != status.delivery_plan().plan_id()
    464                 || stopped.artifact_id() != status.artifact.artifact_id()
    465                 || stopped.request() != status.delivery_plan().request()
    466                 || stopped.stop_requested_at_unix_ms() != Some(now)
    467                 || status.delivery_plan().revision().get().checked_add(1)
    468                     != Some(stopped.revision().get())
    469             {
    470                 return Err(Error::StorageFailed);
    471             }
    472             changed = true;
    473             status = self
    474                 .push_status(operation_id)
    475                 .await?
    476                 .ok_or(Error::StorageFailed)?;
    477             if status.delivery_plan().stop_requested_at_unix_ms().is_none() {
    478                 return Err(Error::StorageFailed);
    479             }
    480         }
    481 
    482         Ok(PushCancellationReceipt { status, changed })
    483     }
    484 
    485     /// Claims and executes one prepared signing artifact.
    486     pub async fn sign_prepared(&self, request: PushRequest) -> Result<SigningRunReceipt, Error> {
    487         self.prepare_push(request.clone()).await?;
    488         let status = self
    489             .push_status(request.operation_id)
    490             .await?
    491             .ok_or(Error::StorageFailed)?;
    492         if status.delivery_plan().state() == AuthoredDeliveryState::Cancelled {
    493             self.cancel_push(request.operation_id).await?;
    494             return Err(Error::SigningCancelled);
    495         }
    496         let artifact = status.artifact;
    497         match artifact.signing_state() {
    498             SigningState::Signed => {
    499                 return Ok(SigningRunReceipt {
    500                     artifact,
    501                     replay: true,
    502                 });
    503             }
    504             SigningState::Indeterminate => return Err(Error::SigningIndeterminate),
    505             SigningState::FailedTerminal | SigningState::Cancelled => {
    506                 return Err(Error::SignerFailed);
    507             }
    508             SigningState::Planned | SigningState::Retryable => {}
    509         }
    510         let signer = self.signer.as_deref().ok_or(Error::MissingSigner)?;
    511 
    512         let now = self.clock.now_unix_ms()?.max(artifact.updated_at_unix_ms());
    513         if let Some(existing) = artifact.signing_claim() {
    514             if now < existing.expires_at_unix_ms() {
    515                 return Err(Error::WorkClaimConflict);
    516             }
    517             if existing.owner() == SIGNING_CLAIM_OWNER_NON_REPLAYABLE {
    518                 let (claimed, fence) = self
    519                     .claim_artifact(
    520                         artifact,
    521                         ClaimAuthoredTarget::ArtifactSigning,
    522                         SIGNING_CLAIM_OWNER_NON_REPLAYABLE,
    523                         now,
    524                     )
    525                     .await?;
    526                 let failure = WorkFailure::new(
    527                     "signing_effect_unknown_after_restart",
    528                     WorkPhase::Signing,
    529                     FailureClass::Indeterminate,
    530                     None,
    531                     None,
    532                 )
    533                 .map_err(map_storage_error)?;
    534                 self.apply_artifact_failure(claimed.artifact_id(), fence, failure, None, now)
    535                     .await?;
    536                 return Err(Error::SigningIndeterminate);
    537             }
    538         }
    539 
    540         let signer_status = signer.status().await.map_err(|_| Error::SignerFailed)?;
    541         let replay_capability = signer_replay_capability(&signer_status)?;
    542         let (claimed, fence) = self
    543             .claim_artifact(
    544                 artifact,
    545                 ClaimAuthoredTarget::ArtifactSigning,
    546                 signing_claim_owner(replay_capability),
    547                 now,
    548             )
    549             .await?;
    550         let persisted_plan = claimed
    551             .plan()
    552             .ok_or(Error::StorageFailed)?
    553             .decode()
    554             .map_err(map_storage_error)?
    555             .into_plan();
    556         if persisted_plan != request.plan {
    557             return Err(Error::StorageConflict);
    558         }
    559         let signing_claim = claimed.signing_claim().ok_or(Error::StorageFailed)?.clone();
    560         let deadline = self.deadlines.deadline_unix_ms(OperationKind::Sign, now)?;
    561         let signing_operation = SigningOperationId::new(*request.operation_id.as_bytes())
    562             .map_err(|_| Error::InvalidPushRequest)?;
    563         let signing_artifact = SigningArtifactId::new(*claimed.artifact_id().as_bytes())
    564             .map_err(|_| Error::InvalidPushRequest)?;
    565         let sign_request = radroots_signing::SignRequest::new(
    566             OperationId::SyncPush,
    567             SigningIntentId::new(signing_operation, signing_artifact),
    568             request.actor,
    569             persisted_plan,
    570             SignPolicy::new(deadline, request.cancellation)
    571                 .map_err(|_| Error::InvalidPushRequest)?,
    572         )
    573         .map_err(|_| Error::InvalidPushRequest)?;
    574 
    575         let expected_request = sign_request.clone();
    576         let signer_result = match signer.sign_authored_evidence(sign_request).await {
    577             Ok(receipt) => {
    578                 // A missing observation time leaves the durable attempt unresolved;
    579                 // it is not evidence that signing failed without an effect.
    580                 let observed_at_unix_ms = self.clock.now_unix_ms()?;
    581                 receipt.revalidate(&expected_request, observed_at_unix_ms)
    582             }
    583             Err(error) => Err(error),
    584         };
    585 
    586         match signer_result {
    587             Ok(receipt) => {
    588                 let command = AuthoredAtomicCommand::RecordSigned(
    589                     RecordSignedArtifact::new(
    590                         claimed.operation_id(),
    591                         claimed.artifact_id(),
    592                         signing_claim,
    593                         receipt.signed_event().clone(),
    594                         receipt.observed_at_unix_ms(),
    595                     )
    596                     .map_err(map_storage_error)?,
    597                 );
    598                 let applied = self
    599                     .storage
    600                     .execute_authored(command.clone())
    601                     .await
    602                     .map_err(map_storage_error)?;
    603                 if !applied.matches_command(&command) {
    604                     return Err(Error::StorageFailed);
    605                 }
    606                 let status = self
    607                     .push_status(request.operation_id)
    608                     .await?
    609                     .ok_or(Error::StorageFailed)?;
    610                 if status
    611                     .artifact
    612                     .signed()
    613                     .is_none_or(|signed| signed.event() != receipt.signed_event())
    614                 {
    615                     return Err(Error::StorageFailed);
    616                 }
    617                 if expected_request.cancellation_signal().is_cancelled()
    618                     || status.delivery_plan().state() == AuthoredDeliveryState::Cancelled
    619                 {
    620                     let stopped = self.cancel_push(request.operation_id).await?;
    621                     if stopped
    622                         .status
    623                         .artifact
    624                         .signed()
    625                         .is_none_or(|signed| signed.event() != receipt.signed_event())
    626                     {
    627                         return Err(Error::StorageFailed);
    628                     }
    629                     return Err(Error::SigningCancelled);
    630                 }
    631                 match status.artifact.signing_state() {
    632                     SigningState::Cancelled => return Err(Error::SigningCancelled),
    633                     SigningState::FailedTerminal => return Err(Error::SignerFailed),
    634                     SigningState::Signed => {}
    635                     _ => return Err(Error::StorageFailed),
    636                 }
    637                 if receipt.observed_at_unix_ms() >= deadline {
    638                     return Err(Error::SignerDeadlineExceeded);
    639                 }
    640                 Ok(SigningRunReceipt {
    641                     artifact: status.artifact,
    642                     replay: applied.disposition() == AtomicCommitDisposition::Replay,
    643                 })
    644             }
    645             Err(error) => {
    646                 let applied_at = self.clock.now_unix_ms()?;
    647                 let current = self
    648                     .push_status(request.operation_id)
    649                     .await?
    650                     .ok_or(Error::StorageFailed)?;
    651                 if current.artifact.signed().is_some()
    652                     || current.artifact.signing_claim() != claimed.signing_claim()
    653                 {
    654                     // Another worker may have committed evidence or a stop while this one waited.
    655                     return Err(match error.kind() {
    656                         radroots_signing::error::Kind::DeadlineExceeded => {
    657                             Error::SignerDeadlineExceeded
    658                         }
    659                         radroots_signing::error::Kind::SignerCancelled => Error::SigningCancelled,
    660                         _ => Error::SignerFailed,
    661                     });
    662                 }
    663                 if error.kind() == radroots_signing::error::Kind::DeadlineExceeded
    664                     && claimed
    665                         .signing_claim()
    666                         .is_some_and(|claim| applied_at >= claim.expires_at_unix_ms())
    667                 {
    668                     // The signer result arrived after this worker's fence
    669                     // expired. It is rejected, but this stale worker must not
    670                     // mutate durable state; recovery will reclaim the exact
    671                     // request under a fresh fence.
    672                     return Err(Error::SignerDeadlineExceeded);
    673                 }
    674                 if error.kind() == radroots_signing::error::Kind::SignerCancelled {
    675                     let command = AuthoredAtomicCommand::Cancel(
    676                         CancelAuthoredWork::new(
    677                             CancelAuthoredTarget::ArtifactSigning(claimed.artifact_id()),
    678                             claimed.revision(),
    679                             applied_at,
    680                         )
    681                         .map_err(map_storage_error)?,
    682                     );
    683                     self.storage
    684                         .execute_authored(command)
    685                         .await
    686                         .map_err(map_storage_error)?;
    687                     return Err(Error::SigningCancelled);
    688                 }
    689                 let disposition = recovery_disposition(
    690                     replay_capability,
    691                     error.remote_effect(),
    692                     error.retryable(),
    693                 );
    694                 let class = match disposition {
    695                     RecoveryDisposition::RetryExactRequest | RecoveryDisposition::RetryLocal => {
    696                         FailureClass::Retryable
    697                     }
    698                     RecoveryDisposition::Indeterminate => FailureClass::Indeterminate,
    699                     RecoveryDisposition::Failed => FailureClass::Terminal,
    700                     _ => FailureClass::Indeterminate,
    701                 };
    702                 let retry_at = if class == FailureClass::Retryable {
    703                     Some(
    704                         applied_at
    705                             .checked_add(WORK_RETRY_DELAY_MS)
    706                             .ok_or(Error::DeadlineOverflow)?,
    707                     )
    708                 } else {
    709                     None
    710                 };
    711                 let failure =
    712                     WorkFailure::new(error.code(), WorkPhase::Signing, class, retry_at, None)
    713                         .map_err(map_storage_error)?;
    714                 let retry = retry_schedule(claimed.signing_retry(), &failure, retry_at)?;
    715                 self.apply_artifact_failure(
    716                     claimed.artifact_id(),
    717                     fence,
    718                     failure,
    719                     retry,
    720                     applied_at,
    721                 )
    722                 .await?;
    723                 match disposition {
    724                     RecoveryDisposition::Indeterminate => Err(Error::SigningIndeterminate),
    725                     RecoveryDisposition::RetryExactRequest
    726                     | RecoveryDisposition::RetryLocal
    727                     | RecoveryDisposition::Failed => {
    728                         if error.kind() == radroots_signing::error::Kind::DeadlineExceeded {
    729                             Err(Error::SignerDeadlineExceeded)
    730                         } else {
    731                             Err(Error::SignerFailed)
    732                         }
    733                     }
    734                     _ => Err(Error::SigningIndeterminate),
    735                 }
    736             }
    737         }
    738     }
    739 
    740     /// Claims and executes local admission for one durably signed artifact.
    741     pub async fn admit_signed(&self, operation_id: SyncId) -> Result<AdmissionRunReceipt, Error> {
    742         let status = self
    743             .push_status(operation_id)
    744             .await?
    745             .ok_or(Error::StorageFailed)?;
    746         let delivery_stopped = status.delivery_plan().stop_requested_at_unix_ms().is_some();
    747         let artifact = status.artifact;
    748         if artifact.signing_state() != SigningState::Signed {
    749             return Err(Error::InvalidSignerOutput);
    750         }
    751         if artifact.admission_state().is_admitted() {
    752             return Ok(AdmissionRunReceipt {
    753                 artifact,
    754                 replay: true,
    755             });
    756         }
    757         if delivery_stopped {
    758             self.cancel_push(operation_id).await?;
    759             return Err(Error::AdmissionFailed);
    760         }
    761         if matches!(
    762             artifact.admission_state(),
    763             AdmissionState::Rejected | AdmissionState::Cancelled
    764         ) {
    765             return Err(Error::AdmissionFailed);
    766         }
    767         let now = self.clock.now_unix_ms()?.max(artifact.updated_at_unix_ms());
    768         if artifact
    769             .admission_claim()
    770             .is_some_and(|claim| now < claim.expires_at_unix_ms())
    771         {
    772             return Err(Error::WorkClaimConflict);
    773         }
    774         let (claimed, fence) = self
    775             .claim_artifact(
    776                 artifact,
    777                 ClaimAuthoredTarget::ArtifactAdmission,
    778                 ADMISSION_CLAIM_OWNER,
    779                 now,
    780             )
    781             .await?;
    782         let event = claimed
    783             .signed()
    784             .ok_or(Error::InvalidSignerOutput)?
    785             .event()
    786             .clone();
    787         let persisted_plan = claimed
    788             .plan()
    789             .ok_or(Error::StorageFailed)?
    790             .decode()
    791             .map_err(map_storage_error)?
    792             .into_plan();
    793         let contract_id = persisted_plan.body().contract().contract_id().as_str();
    794         let admission = match outbound_admission(&event, contract_id, now) {
    795             Ok(admission) => admission,
    796             Err(error) => {
    797                 let failure = WorkFailure::new(
    798                     "invalid_signed_artifact",
    799                     WorkPhase::Admission,
    800                     FailureClass::Terminal,
    801                     None,
    802                     None,
    803                 )
    804                 .map_err(map_storage_error)?;
    805                 self.apply_artifact_failure(claimed.artifact_id(), fence, failure, None, now)
    806                     .await?;
    807                 return Err(error);
    808             }
    809         };
    810         let admission_receipt = match EventStore::admit(self.storage.as_ref(), admission).await {
    811             Ok(receipt) => receipt,
    812             Err(error) => {
    813                 let applied_at = self.clock.now_unix_ms()?.max(claimed.updated_at_unix_ms());
    814                 let terminal = matches!(
    815                     error,
    816                     radroots_storage::Error::EventConflict
    817                         | radroots_storage::Error::AdmissionRegression
    818                 );
    819                 let retry_at = if terminal {
    820                     None
    821                 } else {
    822                     Some(
    823                         applied_at
    824                             .checked_add(WORK_RETRY_DELAY_MS)
    825                             .ok_or(Error::DeadlineOverflow)?,
    826                     )
    827                 };
    828                 let failure = WorkFailure::new(
    829                     if terminal {
    830                         "admission_conflict"
    831                     } else {
    832                         "admission_storage_unavailable"
    833                     },
    834                     WorkPhase::Admission,
    835                     if terminal {
    836                         FailureClass::Terminal
    837                     } else {
    838                         FailureClass::Retryable
    839                     },
    840                     retry_at,
    841                     None,
    842                 )
    843                 .map_err(map_storage_error)?;
    844                 let retry = retry_schedule(claimed.admission_retry(), &failure, retry_at)?;
    845                 self.apply_artifact_failure(
    846                     claimed.artifact_id(),
    847                     fence,
    848                     failure,
    849                     retry,
    850                     applied_at,
    851                 )
    852                 .await?;
    853                 return Err(match error {
    854                     radroots_storage::Error::SpaceInsufficient => Error::StorageSpaceInsufficient,
    855                     _ => Error::AdmissionFailed,
    856                 });
    857             }
    858         };
    859         let state = match admission_receipt.disposition() {
    860             AdmissionDisposition::Duplicate => AdmissionState::Duplicate,
    861             AdmissionDisposition::Inserted | AdmissionDisposition::Advanced => {
    862                 AdmissionState::Inserted
    863             }
    864         };
    865         let applied_at = self.clock.now_unix_ms()?.max(claimed.updated_at_unix_ms());
    866         let command = AuthoredAtomicCommand::ApplyAdmission(
    867             ApplyAdmissionResult::new(claimed.artifact_id(), fence, state, None, None, applied_at)
    868                 .map_err(map_storage_error)?,
    869         );
    870         let applied = self
    871             .storage
    872             .execute_authored(command)
    873             .await
    874             .map_err(map_storage_error)?;
    875         let AuthoredAtomicOutcome::Artifact(artifact) = applied.outcome() else {
    876             return Err(Error::StorageFailed);
    877         };
    878         Ok(AdmissionRunReceipt {
    879             artifact: artifact.clone(),
    880             replay: false,
    881         })
    882     }
    883 
    884     async fn claim_artifact(
    885         &self,
    886         artifact: AuthoredArtifact,
    887         target: fn(AuthoredArtifactId) -> ClaimAuthoredTarget,
    888         owner: &'static str,
    889         acquired_at: u64,
    890     ) -> Result<(AuthoredArtifact, WorkFence), Error> {
    891         let existing = match target(artifact.artifact_id()) {
    892             ClaimAuthoredTarget::ArtifactSigning(_) => artifact.signing_claim(),
    893             ClaimAuthoredTarget::ArtifactAdmission(_) => artifact.admission_claim(),
    894             ClaimAuthoredTarget::DeliveryPlan(_) => return Err(Error::StorageFailed),
    895         };
    896         let generation = existing.map_or(1, |claim| claim.generation().get().saturating_add(1));
    897         let generation = NonZeroU64::new(generation).ok_or(Error::StorageFailed)?;
    898         let expires_at = acquired_at
    899             .checked_add(self.deadlines.timeout_ms(OperationKind::Sign))
    900             .ok_or(Error::DeadlineOverflow)?;
    901         let claim = WorkClaim::new(
    902             *self.ids.next_id(OperationKind::Sign)?.as_bytes(),
    903             owner,
    904             generation,
    905             acquired_at,
    906             expires_at,
    907             artifact.revision(),
    908         )
    909         .map_err(map_storage_error)?;
    910         let fence = WorkFence::new(*claim.token(), claim.generation(), claim.row_revision())
    911             .map_err(map_storage_error)?;
    912         let command = AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new(
    913             target(artifact.artifact_id()),
    914             claim,
    915         ));
    916         let receipt = self
    917             .storage
    918             .execute_authored(command)
    919             .await
    920             .map_err(map_claim_error)?;
    921         let AuthoredAtomicOutcome::Artifact(claimed) = receipt.outcome() else {
    922             return Err(Error::StorageFailed);
    923         };
    924         Ok((claimed.clone(), fence))
    925     }
    926 
    927     async fn apply_artifact_failure(
    928         &self,
    929         artifact_id: AuthoredArtifactId,
    930         fence: WorkFence,
    931         failure: WorkFailure,
    932         retry: Option<RetrySchedule>,
    933         applied_at: u64,
    934     ) -> Result<AuthoredArtifact, Error> {
    935         let command = AuthoredAtomicCommand::ApplyFailure(
    936             ApplyWorkFailure::new(
    937                 AuthoredWorkTarget::Artifact(artifact_id),
    938                 fence,
    939                 failure,
    940                 retry,
    941                 applied_at,
    942             )
    943             .map_err(map_storage_error)?,
    944         );
    945         let receipt = self
    946             .storage
    947             .execute_authored(command)
    948             .await
    949             .map_err(map_storage_error)?;
    950         let AuthoredAtomicOutcome::Artifact(artifact) = receipt.outcome() else {
    951             return Err(Error::StorageFailed);
    952         };
    953         Ok(artifact.clone())
    954     }
    955 }
    956 
    957 fn signer_replay_capability(
    958     status: &radroots_signing::SignerStatus,
    959 ) -> Result<ReplayCapability, Error> {
    960     let mut capabilities = status.capabilities().iter();
    961     let first = capabilities
    962         .next()
    963         .ok_or(Error::SignerCapabilityUnavailable)?
    964         .replay();
    965     if capabilities.all(|capability| capability.replay() == first) {
    966         Ok(first)
    967     } else {
    968         Ok(ReplayCapability::NonReplayable)
    969     }
    970 }
    971 
    972 fn signing_claim_owner(replay: ReplayCapability) -> &'static str {
    973     match replay {
    974         ReplayCapability::ExactReplayByRequestId => SIGNING_CLAIM_OWNER_EXACT,
    975         ReplayCapability::LocalReplaySafe => SIGNING_CLAIM_OWNER_LOCAL,
    976         ReplayCapability::NonReplayable => SIGNING_CLAIM_OWNER_NON_REPLAYABLE,
    977         _ => SIGNING_CLAIM_OWNER_NON_REPLAYABLE,
    978     }
    979 }
    980 
    981 fn retry_schedule(
    982     previous: Option<&RetrySchedule>,
    983     failure: &WorkFailure,
    984     retry_at: Option<u64>,
    985 ) -> Result<Option<RetrySchedule>, Error> {
    986     let Some(retry_at) = retry_at else {
    987         return Ok(None);
    988     };
    989     let schedule = match previous {
    990         Some(previous) => previous.next_attempt(retry_at, failure.clone()),
    991         None => RetrySchedule::new(NonZeroU32::MIN, retry_at, failure.clone()),
    992     }
    993     .map_err(map_storage_error)?;
    994     Ok(Some(schedule))
    995 }
    996 
    997 fn normalize_sink_failure(
    998     request: &DeliveryRequest,
    999     failure: SinkFailure,
   1000     attempted_at_unix_ms: u64,
   1001 ) -> Result<SinkFailure, Error> {
   1002     if failure.retryability() != Retryability::Retryable {
   1003         return Ok(failure);
   1004     }
   1005     let default_retry_at = attempted_at_unix_ms
   1006         .checked_add(WORK_RETRY_DELAY_MS)
   1007         .ok_or(Error::DeadlineOverflow)?;
   1008     let retry_at = failure
   1009         .retry_after_unix_ms()
   1010         .unwrap_or(default_retry_at)
   1011         .max(
   1012             attempted_at_unix_ms
   1013                 .checked_add(1)
   1014                 .ok_or(Error::DeadlineOverflow)?,
   1015         );
   1016     if failure.retry_after_unix_ms() == Some(retry_at) {
   1017         return Ok(failure);
   1018     }
   1019     SinkFailure::for_request(
   1020         request,
   1021         failure.code(),
   1022         failure.retryability(),
   1023         Some(retry_at),
   1024         failure.message().map(str::to_owned),
   1025         failure.partial_evidence().to_vec(),
   1026     )
   1027     .map_err(|_| Error::InvalidDeliveryRequest)
   1028 }
   1029 
   1030 fn delivery_retry_schedule(
   1031     attempt: NonZeroU32,
   1032     outcome: &DeliveryAttemptOutcome,
   1033     satisfaction: SatisfactionState,
   1034     attempted_at_unix_ms: u64,
   1035 ) -> Result<Option<RetrySchedule>, Error> {
   1036     if satisfaction != SatisfactionState::Pending || attempt.get() >= DELIVERY_PLAN_ATTEMPTS_MAX {
   1037         return Ok(None);
   1038     }
   1039     let failure = match outcome {
   1040         DeliveryAttemptOutcome::Receipt(_) => {
   1041             let retry_at = attempted_at_unix_ms
   1042                 .checked_add(WORK_RETRY_DELAY_MS)
   1043                 .ok_or(Error::DeadlineOverflow)?;
   1044             WorkFailure::new(
   1045                 "delivery_pending",
   1046                 WorkPhase::Delivery,
   1047                 FailureClass::Retryable,
   1048                 Some(retry_at),
   1049                 None,
   1050             )
   1051             .map_err(map_storage_error)?
   1052         }
   1053         DeliveryAttemptOutcome::SinkFailure(failure)
   1054             if failure.retryability() == Retryability::Retryable =>
   1055         {
   1056             WorkFailure::new(
   1057                 failure.code(),
   1058                 WorkPhase::Delivery,
   1059                 FailureClass::Retryable,
   1060                 failure.retry_after_unix_ms(),
   1061                 failure.message().map(str::to_owned),
   1062             )
   1063             .map_err(map_storage_error)?
   1064         }
   1065         DeliveryAttemptOutcome::SinkFailure(_) => return Ok(None),
   1066     };
   1067     let retry_at = failure.retry_after_unix_ms().unwrap_or(
   1068         attempted_at_unix_ms
   1069             .checked_add(WORK_RETRY_DELAY_MS)
   1070             .ok_or(Error::DeadlineOverflow)?,
   1071     );
   1072     RetrySchedule::new(attempt, retry_at, failure)
   1073         .map(Some)
   1074         .map_err(map_storage_error)
   1075 }
   1076 
   1077 fn outbound_admission(
   1078     event: &radroots_event::SignedEvent,
   1079     contract_id: &str,
   1080     observed_at_unix_ms: u64,
   1081 ) -> Result<EventAdmission, Error> {
   1082     let verified = verify::signature(
   1083         verify::id(RawEvent::new(event.envelope().clone()))
   1084             .map_err(|_| Error::InvalidSignerOutput)?,
   1085         &Nip01SignatureVerifier,
   1086     )
   1087     .map_err(|_| Error::InvalidSignerOutput)?;
   1088     let validated = verified
   1089         .validate_contract_for_admission(contract_id)
   1090         .map_err(|_| Error::InvalidSignerOutput)?;
   1091     let policy = RegistryPolicy::visible();
   1092     if policy.decide(&validated) != AdmissionDecision::Visible {
   1093         return Err(Error::InvalidSignerOutput);
   1094     }
   1095     struct Evidence;
   1096     impl radroots_event::admission::AdmissionPolicy for Evidence {
   1097         type Error = core::convert::Infallible;
   1098         fn policy_id(&self) -> &'static str {
   1099             "radroots.registry_v7"
   1100         }
   1101         fn admit(
   1102             &self,
   1103             _: &radroots_event::admission::ContractValidatedEvent,
   1104         ) -> Result<(), Self::Error> {
   1105             Ok(())
   1106         }
   1107     }
   1108     impl radroots_event::admission::VisibilityPolicy for Evidence {
   1109         type Error = core::convert::Infallible;
   1110         fn policy_id(&self) -> &'static str {
   1111             "radroots.registry_v7"
   1112         }
   1113         fn make_visible(
   1114             &self,
   1115             _: &radroots_event::admission::AdmittedEvent,
   1116         ) -> Result<(), Self::Error> {
   1117             Ok(())
   1118         }
   1119     }
   1120     let visible = validated
   1121         .admit_with(&Evidence)
   1122         .and_then(|event| event.make_visible_with(&Evidence))
   1123         .map_err(|never| match never {})?;
   1124     let local = Target::new(TransportId::LOCAL, "local:sync-authored")
   1125         .map_err(|_| Error::InvalidPushRequest)?;
   1126     let provenance = EventProvenance::new(
   1127         TransportId::LOCAL,
   1128         local.fingerprint().clone(),
   1129         observed_at_unix_ms,
   1130     )
   1131     .map_err(|_| Error::InvalidPushRequest)?;
   1132     EventAdmission::visible(ObservedEvent::new(event.clone(), provenance), visible)
   1133         .map_err(map_storage_error)
   1134 }
   1135 
   1136 fn authored_push_input_digest(request: &PushRequest) -> Result<AtomicCommitDigest, Error> {
   1137     let mut hasher = Sha256::new();
   1138     hash_field(&mut hasher, b"radroots.sync.authored-push.v2");
   1139     hash_field(&mut hasher, request.operation_id.as_bytes());
   1140     hash_field(&mut hasher, request.idempotency_key.as_str().as_bytes());
   1141     hash_field(&mut hasher, request.plan.digest().as_bytes());
   1142     hash_field(&mut hasher, request.actor.public_key().as_bytes());
   1143     let source = match request.actor.source() {
   1144         radroots_signing::actor::ActorSource::LocalAccount(_) => 0,
   1145         radroots_signing::actor::ActorSource::ExplicitPublicKey => 1,
   1146         radroots_signing::actor::ActorSource::RemoteSigner(_) => 2,
   1147         radroots_signing::actor::ActorSource::Service(_) => 3,
   1148         _ => return Err(Error::InvalidPushRequest),
   1149     };
   1150     hasher.update([source]);
   1151     if let Some(account_id) = request.actor.account_id() {
   1152         hash_field(&mut hasher, account_id.as_bytes());
   1153     }
   1154     for role in request.actor.roles() {
   1155         hash_field(&mut hasher, role.as_str().as_bytes());
   1156     }
   1157     for target in request.targets.targets() {
   1158         hash_field(&mut hasher, target.fingerprint().as_str().as_bytes());
   1159     }
   1160     hash_satisfaction(&mut hasher, &request.satisfaction);
   1161     hasher.update(request.delivery_deadline_unix_ms.to_be_bytes());
   1162     let cancellation = match request.cancellation {
   1163         CancellationPolicy::PreservePublishedRequest => 0,
   1164         CancellationPolicy::LocalCooperative => 1,
   1165         _ => return Err(Error::InvalidPushRequest),
   1166     };
   1167     hasher.update([cancellation]);
   1168     Ok(AtomicCommitDigest::new(hasher.finalize().into()))
   1169 }
   1170 
   1171 fn hash_satisfaction(hasher: &mut Sha256, policy: &SatisfactionPolicy) {
   1172     hasher.update([match policy.class() {
   1173         radroots_transport::policy::SatisfactionClass::Accepted => 0,
   1174         radroots_transport::policy::SatisfactionClass::Delivered => 1,
   1175     }]);
   1176     let targets = policy.targets();
   1177     if targets.is_any() {
   1178         hasher.update([0]);
   1179     } else if targets.is_all() {
   1180         hasher.update([1]);
   1181     } else if let Some(threshold) = targets.quorum_threshold() {
   1182         hasher.update([2]);
   1183         hasher.update(threshold.to_be_bytes());
   1184     } else if let Some(required) = targets.required_targets() {
   1185         hasher.update([3]);
   1186         for target in required {
   1187             hash_field(hasher, target.as_str().as_bytes());
   1188         }
   1189     }
   1190 }
   1191 
   1192 fn authored_ids(
   1193     operation_id: SyncId,
   1194 ) -> Result<
   1195     (
   1196         OperationInstanceId,
   1197         AuthoredArtifactId,
   1198         AuthoredDeliveryPlanId,
   1199     ),
   1200     Error,
   1201 > {
   1202     let operation =
   1203         OperationInstanceId::new(*operation_id.as_bytes()).map_err(map_storage_error)?;
   1204     let artifact = AuthoredArtifactId::new(derive_child_id(
   1205         b"radroots.sync.authored-artifact.v2",
   1206         operation_id.as_bytes(),
   1207     ))
   1208     .map_err(map_storage_error)?;
   1209     let delivery = AuthoredDeliveryPlanId::new(derive_child_id(
   1210         b"radroots.sync.authored-delivery.v2",
   1211         artifact.as_bytes(),
   1212     ))
   1213     .map_err(map_storage_error)?;
   1214     Ok((operation, artifact, delivery))
   1215 }
   1216 
   1217 fn derive_child_id(domain: &[u8], parent: &[u8; 16]) -> [u8; 16] {
   1218     let mut hasher = Sha256::new();
   1219     hash_field(&mut hasher, domain);
   1220     hash_field(&mut hasher, parent);
   1221     let digest: [u8; 32] = hasher.finalize().into();
   1222     let mut child = [0_u8; 16];
   1223     child.copy_from_slice(&digest[..16]);
   1224     child
   1225 }
   1226 
   1227 fn hash_field(hasher: &mut Sha256, value: &[u8]) {
   1228     hasher.update(u64::try_from(value.len()).unwrap_or(u64::MAX).to_be_bytes());
   1229     hasher.update(value);
   1230 }
   1231 
   1232 fn delivery_request_id(id: SyncId) -> String {
   1233     const HEX: &[u8; 16] = b"0123456789abcdef";
   1234     let mut value = String::from("push-");
   1235     for byte in id.as_bytes() {
   1236         value.push(HEX[(byte >> 4) as usize] as char);
   1237         value.push(HEX[(byte & 0x0f) as usize] as char);
   1238     }
   1239     value
   1240 }
   1241 
   1242 fn map_storage_error(error: radroots_storage::Error) -> Error {
   1243     match error {
   1244         radroots_storage::Error::SpaceInsufficient => Error::StorageSpaceInsufficient,
   1245         radroots_storage::Error::IdempotencyConflict
   1246         | radroots_storage::Error::OperationIdentityMismatch
   1247         | radroots_storage::Error::JournalRevisionConflict
   1248         | radroots_storage::Error::OutboxPlanConflict
   1249         | radroots_storage::Error::AtomicCommitConflict => Error::StorageConflict,
   1250         _ => Error::StorageFailed,
   1251     }
   1252 }
   1253 
   1254 fn map_claim_error(error: radroots_storage::Error) -> Error {
   1255     match error {
   1256         radroots_storage::Error::AtomicCommitConflict
   1257         | radroots_storage::Error::InvalidAuthoredTransition
   1258         | radroots_storage::Error::InvalidWorkClaim
   1259         | radroots_storage::Error::DeliveryPlanClaimConflict => Error::WorkClaimConflict,
   1260         other => map_storage_error(other),
   1261     }
   1262 }
   1263 
   1264 #[cfg(test)]
   1265 #[cfg_attr(coverage_nightly, coverage(off))]
   1266 mod tests {
   1267     use super::*;
   1268     use radroots_event::{GenericEventDraft, contract::AuthorRole};
   1269     use radroots_signing::{
   1270         SignerStatus,
   1271         actor::ActorSource,
   1272         capability::{CancellationSupport, SignerCapability, SignerKind},
   1273         status::SignerAvailability,
   1274     };
   1275     use radroots_transport::{
   1276         DeliveryReceipt,
   1277         outcome::DeliveryOutcome,
   1278         policy::{SatisfactionClass, TargetPolicy},
   1279         sink::{DeliveryPayload, DeliveryTargetReceipt},
   1280     };
   1281 
   1282     const AUTHOR: &str = "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df";
   1283     const OTHER_AUTHOR: &str = "e0266e3cfb0d2886f91c73f5f868f3b98273713e5fcd97c081663f5518a4b3af";
   1284 
   1285     fn request_with(
   1286         target_policy: TargetPolicy,
   1287         deadline: u64,
   1288         cancellation: CancellationPolicy,
   1289     ) -> PushRequest {
   1290         let plan = AuthoredEventPlan::from_generic(
   1291             GenericEventDraft::new(
   1292                 "radroots.social.geochat.v1",
   1293                 20_000,
   1294                 1_800_000_100,
   1295                 Vec::new(),
   1296                 "push helper coverage",
   1297                 AUTHOR,
   1298             )
   1299             .unwrap(),
   1300         )
   1301         .unwrap();
   1302         let actor =
   1303             Actor::from_public_key_hex(AUTHOR, ActorSource::ExplicitPublicKey, [AuthorRole::Any])
   1304                 .unwrap();
   1305         PushRequest::new(
   1306             SyncId::new([31; 16]).unwrap(),
   1307             IdempotencyKey::parse("push-helper-coverage").unwrap(),
   1308             actor,
   1309             plan,
   1310             TargetSet::new(vec![Target::nostr_relay("wss://helper.example").unwrap()]).unwrap(),
   1311             SatisfactionPolicy::new(SatisfactionClass::Accepted, target_policy),
   1312             deadline,
   1313             cancellation,
   1314         )
   1315         .unwrap()
   1316     }
   1317 
   1318     fn delivery_plan(request: &PushRequest) -> AuthoredDeliveryPlan {
   1319         let (_, artifact_id, plan_id) = authored_ids(request.operation_id()).unwrap();
   1320         let delivery_request = DeliveryRequest::new(
   1321             delivery_request_id(request.operation_id()),
   1322             DeliveryPayload::new(
   1323                 radroots_event_codec::Codec::decode_signed_event(
   1324                     r#"{"id":"762bee187e9e645b81ec26ade05a69b5e8398caf527be8de0d9a45311ed0c7a0","pubkey":"585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df","created_at":1800000100,"kind":0,"tags":[],"content":"{\"display_name\":\"Moss Street Farm\",\"bot\":false,\"website\":\"https://mossstreet.example\",\"picture\":42}","sig":"4290da0bb6422986647bc8cd5f63bd52d49f41e7b665d3b47105b8109183e8d596f322c531d4061df53e1d2b70fda12d5d1c14f3720d7a56d9d0a03746af5109"}"#,
   1325                 )
   1326                 .unwrap(),
   1327             ),
   1328             request.targets().clone(),
   1329             request.satisfaction().clone(),
   1330             request.delivery_deadline_unix_ms(),
   1331         )
   1332         .unwrap();
   1333         AuthoredDeliveryPlan::new_bound(plan_id, artifact_id, delivery_request, 10).unwrap()
   1334     }
   1335 
   1336     fn capability(replay: ReplayCapability) -> SignerCapability {
   1337         SignerCapability::new(
   1338             SignerKind::Local,
   1339             replay,
   1340             CancellationSupport::BeforePublication,
   1341             false,
   1342             false,
   1343         )
   1344     }
   1345 
   1346     #[test]
   1347     fn push_request_validates_every_binding_and_exposes_redacted_accessors() {
   1348         let request = request_with(
   1349             TargetPolicy::any(),
   1350             1_800_000_300_000,
   1351             CancellationPolicy::LocalCooperative,
   1352         );
   1353         assert_eq!(request.idempotency_key().as_str(), "push-helper-coverage");
   1354         assert_eq!(request.actor().public_key(), *request.plan().author());
   1355         assert_eq!(request.targets().len(), 1);
   1356         assert!(request.satisfaction().targets().is_any());
   1357         assert_eq!(request.delivery_deadline_unix_ms(), 1_800_000_300_000);
   1358         assert_eq!(request.cancellation(), CancellationPolicy::LocalCooperative);
   1359         let debug = format!("{request:?}");
   1360         assert!(debug.contains("[redacted exact authored plan]"));
   1361         assert!(!debug.contains("push helper coverage"));
   1362 
   1363         let invalid_policy = TargetPolicy::required(vec![
   1364             Target::nostr_relay("wss://absent.example")
   1365                 .unwrap()
   1366                 .fingerprint()
   1367                 .clone(),
   1368         ])
   1369         .unwrap();
   1370         let mut invalid = request.clone();
   1371         invalid.satisfaction = SatisfactionPolicy::new(SatisfactionClass::Accepted, invalid_policy);
   1372         assert!(matches!(
   1373             PushRequest::new(
   1374                 invalid.operation_id,
   1375                 invalid.idempotency_key,
   1376                 invalid.actor,
   1377                 invalid.plan,
   1378                 invalid.targets,
   1379                 invalid.satisfaction,
   1380                 invalid.delivery_deadline_unix_ms,
   1381                 invalid.cancellation,
   1382             ),
   1383             Err(Error::InvalidPushRequest)
   1384         ));
   1385 
   1386         let valid = request_with(
   1387             TargetPolicy::any(),
   1388             1_800_000_300_000,
   1389             CancellationPolicy::PreservePublishedRequest,
   1390         );
   1391         assert!(matches!(
   1392             PushRequest::new(
   1393                 valid.operation_id,
   1394                 valid.idempotency_key.clone(),
   1395                 valid.actor.clone(),
   1396                 valid.plan.clone(),
   1397                 valid.targets.clone(),
   1398                 valid.satisfaction.clone(),
   1399                 0,
   1400                 valid.cancellation,
   1401             ),
   1402             Err(Error::InvalidPushRequest)
   1403         ));
   1404         let wrong_actor = Actor::from_public_key_hex(
   1405             OTHER_AUTHOR,
   1406             ActorSource::ExplicitPublicKey,
   1407             [AuthorRole::Any],
   1408         )
   1409         .unwrap();
   1410         assert!(matches!(
   1411             PushRequest::new(
   1412                 valid.operation_id,
   1413                 valid.idempotency_key,
   1414                 wrong_actor,
   1415                 valid.plan,
   1416                 valid.targets,
   1417                 valid.satisfaction,
   1418                 valid.delivery_deadline_unix_ms,
   1419                 valid.cancellation,
   1420             ),
   1421             Err(Error::InvalidPushRequest)
   1422         ));
   1423     }
   1424 
   1425     #[test]
   1426     fn signer_and_retry_helpers_cover_every_stable_policy_variant() {
   1427         let empty = SignerStatus::new(SignerAvailability::Ready, Vec::new(), None);
   1428         assert_eq!(
   1429             signer_replay_capability(&empty),
   1430             Err(Error::SignerCapabilityUnavailable)
   1431         );
   1432         let exact = SignerStatus::new(
   1433             SignerAvailability::Ready,
   1434             vec![capability(ReplayCapability::ExactReplayByRequestId)],
   1435             None,
   1436         );
   1437         assert_eq!(
   1438             signer_replay_capability(&exact).unwrap(),
   1439             ReplayCapability::ExactReplayByRequestId
   1440         );
   1441         let mixed = SignerStatus::new(
   1442             SignerAvailability::Ready,
   1443             vec![
   1444                 capability(ReplayCapability::ExactReplayByRequestId),
   1445                 capability(ReplayCapability::LocalReplaySafe),
   1446             ],
   1447             None,
   1448         );
   1449         assert_eq!(
   1450             signer_replay_capability(&mixed).unwrap(),
   1451             ReplayCapability::NonReplayable
   1452         );
   1453         assert_eq!(
   1454             signing_claim_owner(ReplayCapability::ExactReplayByRequestId),
   1455             SIGNING_CLAIM_OWNER_EXACT
   1456         );
   1457         assert_eq!(
   1458             signing_claim_owner(ReplayCapability::LocalReplaySafe),
   1459             SIGNING_CLAIM_OWNER_LOCAL
   1460         );
   1461         assert_eq!(
   1462             signing_claim_owner(ReplayCapability::NonReplayable),
   1463             SIGNING_CLAIM_OWNER_NON_REPLAYABLE
   1464         );
   1465 
   1466         let failure = WorkFailure::new(
   1467             "retry",
   1468             WorkPhase::Signing,
   1469             FailureClass::Retryable,
   1470             Some(20),
   1471             None,
   1472         )
   1473         .unwrap();
   1474         assert!(retry_schedule(None, &failure, None).unwrap().is_none());
   1475         let first = retry_schedule(None, &failure, Some(20)).unwrap().unwrap();
   1476         assert_eq!(first.attempt(), NonZeroU32::MIN);
   1477         let next_failure = WorkFailure::new(
   1478             "retry",
   1479             WorkPhase::Signing,
   1480             FailureClass::Retryable,
   1481             Some(21),
   1482             None,
   1483         )
   1484         .unwrap();
   1485         assert_eq!(
   1486             retry_schedule(Some(&first), &next_failure, Some(21))
   1487                 .unwrap()
   1488                 .unwrap()
   1489                 .attempt()
   1490                 .get(),
   1491             2
   1492         );
   1493     }
   1494 
   1495     #[test]
   1496     fn delivery_failure_and_retry_helpers_normalize_all_outcome_classes() {
   1497         let request = request_with(
   1498             TargetPolicy::any(),
   1499             1_800_000_300_000,
   1500             CancellationPolicy::PreservePublishedRequest,
   1501         );
   1502         let plan = delivery_plan(&request);
   1503         let delivery_request = plan.request().unwrap();
   1504         let terminal = SinkFailure::for_request(
   1505             delivery_request,
   1506             "terminal",
   1507             Retryability::Terminal,
   1508             None,
   1509             None,
   1510             Vec::new(),
   1511         )
   1512         .unwrap();
   1513         assert_eq!(
   1514             normalize_sink_failure(delivery_request, terminal.clone(), 100).unwrap(),
   1515             terminal
   1516         );
   1517         let retryable = SinkFailure::for_request(
   1518             delivery_request,
   1519             "retryable",
   1520             Retryability::Retryable,
   1521             None,
   1522             Some("retry".to_owned()),
   1523             Vec::new(),
   1524         )
   1525         .unwrap();
   1526         let normalized = normalize_sink_failure(delivery_request, retryable, 100).unwrap();
   1527         assert_eq!(normalized.retry_after_unix_ms(), Some(1_100));
   1528         assert_eq!(
   1529             normalize_sink_failure(delivery_request, normalized.clone(), 100).unwrap(),
   1530             normalized
   1531         );
   1532         let past = SinkFailure::for_request(
   1533             delivery_request,
   1534             "past",
   1535             Retryability::Retryable,
   1536             Some(99),
   1537             None,
   1538             Vec::new(),
   1539         )
   1540         .unwrap();
   1541         assert_eq!(
   1542             normalize_sink_failure(delivery_request, past, 100)
   1543                 .unwrap()
   1544                 .retry_after_unix_ms(),
   1545             Some(101)
   1546         );
   1547 
   1548         let receipt = DeliveryReceipt::for_request(
   1549             delivery_request,
   1550             delivery_request
   1551                 .target_set()
   1552                 .targets()
   1553                 .iter()
   1554                 .cloned()
   1555                 .map(|target| {
   1556                     DeliveryTargetReceipt::attempted(target, DeliveryOutcome::unavailable())
   1557                 })
   1558                 .collect(),
   1559         )
   1560         .unwrap();
   1561         let receipt_outcome = DeliveryAttemptOutcome::Receipt(receipt);
   1562         assert!(
   1563             delivery_retry_schedule(
   1564                 NonZeroU32::new(plan.attempt_count() + 1).unwrap(),
   1565                 &receipt_outcome,
   1566                 SatisfactionState::Satisfied,
   1567                 100
   1568             )
   1569             .unwrap()
   1570             .is_none()
   1571         );
   1572         let schedule = delivery_retry_schedule(
   1573             NonZeroU32::new(plan.attempt_count() + 1).unwrap(),
   1574             &receipt_outcome,
   1575             SatisfactionState::Pending,
   1576             100,
   1577         )
   1578         .unwrap()
   1579         .unwrap();
   1580         assert_eq!(schedule.not_before_unix_ms(), 1_100);
   1581         let retry_outcome = DeliveryAttemptOutcome::SinkFailure(normalized);
   1582         assert!(
   1583             delivery_retry_schedule(
   1584                 NonZeroU32::new(plan.attempt_count() + 1).unwrap(),
   1585                 &retry_outcome,
   1586                 SatisfactionState::Pending,
   1587                 100
   1588             )
   1589             .unwrap()
   1590             .is_some()
   1591         );
   1592         let terminal_outcome = DeliveryAttemptOutcome::SinkFailure(terminal);
   1593         assert!(
   1594             delivery_retry_schedule(
   1595                 NonZeroU32::new(plan.attempt_count() + 1).unwrap(),
   1596                 &terminal_outcome,
   1597                 SatisfactionState::Pending,
   1598                 100
   1599             )
   1600             .unwrap()
   1601             .is_none()
   1602         );
   1603 
   1604         for policy in [
   1605             TargetPolicy::any(),
   1606             TargetPolicy::all(),
   1607             TargetPolicy::quorum(1).unwrap(),
   1608             TargetPolicy::required(vec![request.targets().targets()[0].fingerprint().clone()])
   1609                 .unwrap(),
   1610         ] {
   1611             let request = request_with(
   1612                 policy,
   1613                 1_800_000_300_000,
   1614                 CancellationPolicy::LocalCooperative,
   1615             );
   1616             assert_ne!(
   1617                 authored_push_input_digest(&request).unwrap().as_bytes(),
   1618                 &[0; 32]
   1619             );
   1620         }
   1621     }
   1622 
   1623     #[test]
   1624     fn storage_error_maps_are_explicit_and_fail_closed() {
   1625         assert_eq!(
   1626             map_storage_error(radroots_storage::Error::SpaceInsufficient),
   1627             Error::StorageSpaceInsufficient
   1628         );
   1629         assert_eq!(
   1630             map_claim_error(radroots_storage::Error::SpaceInsufficient),
   1631             Error::StorageSpaceInsufficient
   1632         );
   1633         assert_eq!(
   1634             map_storage_error(radroots_storage::Error::AtomicCommitConflict),
   1635             Error::StorageConflict
   1636         );
   1637         assert_eq!(
   1638             map_storage_error(radroots_storage::Error::BackendUnavailable),
   1639             Error::StorageFailed
   1640         );
   1641         assert_eq!(
   1642             map_claim_error(radroots_storage::Error::DeliveryPlanClaimConflict),
   1643             Error::WorkClaimConflict
   1644         );
   1645         assert_eq!(
   1646             map_claim_error(radroots_storage::Error::BackendUnavailable),
   1647             Error::StorageFailed
   1648         );
   1649         assert_eq!(
   1650             delivery_request_id(SyncId::new([0xab; 16]).unwrap()).len(),
   1651             37
   1652         );
   1653     }
   1654 }