lib

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

push_enqueue.rs (84323B)


      1 use std::{
      2     collections::VecDeque,
      3     sync::{
      4         Arc, Mutex,
      5         atomic::{AtomicU64, AtomicUsize, Ordering},
      6     },
      7 };
      8 
      9 use futures::{FutureExt, task::noop_waker_ref};
     10 use futures_executor::block_on;
     11 use radroots_event::{
     12     GenericEventDraft, SignedEvent,
     13     contract::AuthorRole,
     14     draft::SignedEventParts,
     15     food::availability::{
     16         FoodAvailabilityDetails, FoodAvailabilityDetailsParts, FoodAvailabilityStatus, FoodContent,
     17         FoodCurrency, FoodIdentifier, FoodPrice, FoodPublishedAt, FoodText, FoodUnit,
     18     },
     19 };
     20 use radroots_event_codec::authoring::AuthoredEventPlan;
     21 use radroots_protocol::runtime::v1::SyncRetryDecision;
     22 use radroots_signing::{
     23     Actor, Error as SigningError, SignReceipt, SignRequest, Signer, SignerStatus,
     24     actor::ActorSource,
     25     capability::{CancellationSupport, SignerCapability, SignerKind},
     26     error::Kind as SigningErrorKind,
     27     recovery::ReplayCapability,
     28     request::CancellationPolicy,
     29     status::SignerAvailability,
     30 };
     31 use radroots_storage::{
     32     EventStore, Journal, Outbox, ProjectionStore,
     33     atomic::AtomicStorage,
     34     authored_atomic::AuthoredAtomicStorage,
     35     authored_delivery::AuthoredDeliveryState,
     36     event::{EventQuery, EventQueryBounds, SourceGeneration},
     37     journal::{IdempotencyKey, OperationInstanceId},
     38     memory::MemoryStorage,
     39     status::StorageStatusProvider,
     40 };
     41 use radroots_storage_sqlite::{OpenMode, OpenOptions, Paths, SqliteStorage};
     42 use radroots_sync::{
     43     Engine, PushRequest,
     44     policy::{Clock, DeadlinePolicy, Error, IdSource, OperationKind, SyncId, SyncStorage},
     45 };
     46 use radroots_transport::{
     47     DeliveryReceipt, DeliveryRequest, Error as TransportError, EventSink, EventSource, FetchPage,
     48     FetchRequest, SinkFailure, SinkStatus, SourceStatus, Target, TargetSet, TransportId,
     49     outcome::{DeliveryOutcome, Retryability},
     50     policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy},
     51     sink::DeliveryTargetReceipt,
     52 };
     53 use secp256k1::{Keypair, Message, Secp256k1, SecretKey};
     54 
     55 const CONTENT: &str = "frozen-content";
     56 
     57 #[path = "push_enqueue/signing_evidence.rs"]
     58 mod signing_evidence;
     59 
     60 #[path = "push_enqueue/delivery_evidence.rs"]
     61 mod delivery_evidence;
     62 
     63 #[path = "push_enqueue/delivery_selection.rs"]
     64 mod delivery_selection;
     65 
     66 #[path = "push_enqueue/capacity.rs"]
     67 mod capacity;
     68 
     69 struct MockSink;
     70 
     71 struct FaultStorage {
     72     capacity: Mutex<Option<capacity::Fault>>,
     73     inner: Arc<MemoryStorage>,
     74     delivery_injection: Mutex<Option<delivery_evidence::Injection>>,
     75     receipt_mutation: Mutex<Option<delivery_evidence::ReceiptMutation>>,
     76     history_mutation: Mutex<Option<(usize, delivery_evidence::HistoryMutation)>>,
     77     fault_kind: AtomicUsize,
     78     remaining: AtomicUsize,
     79     admit_error: AtomicUsize,
     80     prepared: Mutex<Option<radroots_storage::authored_atomic::AuthoredAtomicOutcome>>,
     81 }
     82 
     83 impl FaultStorage {
     84     fn new(generation: u8) -> Self {
     85         Self {
     86             capacity: Mutex::new(None),
     87             delivery_injection: Mutex::new(None),
     88             receipt_mutation: Mutex::new(None),
     89             history_mutation: Mutex::new(None),
     90             inner: Arc::new(MemoryStorage::new(
     91                 SourceGeneration::new([generation; 32]).expect("generation"),
     92             )),
     93             fault_kind: AtomicUsize::new(0),
     94             remaining: AtomicUsize::new(0),
     95             admit_error: AtomicUsize::new(0),
     96             prepared: Mutex::new(None),
     97         }
     98     }
     99 
    100     fn fault_next_prepared(&self) {
    101         self.fault_kind.store(1, Ordering::Relaxed);
    102         self.remaining.store(1, Ordering::Relaxed);
    103     }
    104 
    105     fn fault_nth_artifact(&self, nth: usize) {
    106         self.fault_kind.store(2, Ordering::Relaxed);
    107         self.remaining.store(nth, Ordering::Relaxed);
    108     }
    109 
    110     fn fault_nth_plan(&self, nth: usize) {
    111         self.fault_kind.store(3, Ordering::Relaxed);
    112         self.remaining.store(nth, Ordering::Relaxed);
    113     }
    114 
    115     fn fault_prepared_identity(&self) {
    116         self.fault_kind.store(4, Ordering::Relaxed);
    117         self.remaining.store(1, Ordering::Relaxed);
    118     }
    119 
    120     fn fail_admission_with(&self, value: usize) {
    121         self.admit_error.store(value, Ordering::Relaxed);
    122     }
    123 }
    124 
    125 impl EventStore for FaultStorage {
    126     fn status(
    127         &self,
    128     ) -> radroots_transport::BoxFuture<
    129         '_,
    130         Result<radroots_storage::status::EventStoreStatus, radroots_storage::Error>,
    131     > {
    132         EventStore::status(self.inner.as_ref())
    133     }
    134     fn admit(
    135         &self,
    136         value: radroots_storage::event::EventAdmission,
    137     ) -> radroots_transport::BoxFuture<
    138         '_,
    139         Result<radroots_storage::event::AdmissionReceipt, radroots_storage::Error>,
    140     > {
    141         match self.admit_error.load(Ordering::Relaxed) {
    142             1 => Box::pin(async { Err(radroots_storage::Error::EventConflict) }),
    143             2 => Box::pin(async { Err(radroots_storage::Error::BackendUnavailable) }),
    144             3 => Box::pin(async { Err(radroots_storage::Error::SpaceInsufficient) }),
    145             _ => EventStore::admit(self.inner.as_ref(), value),
    146         }
    147     }
    148     fn query_raw(
    149         &self,
    150         value: radroots_storage::event::EventQuery,
    151     ) -> radroots_transport::BoxFuture<
    152         '_,
    153         Result<
    154             radroots_storage::event::EventPage<radroots_storage::event::StoredRawEvent>,
    155             radroots_storage::Error,
    156         >,
    157     > {
    158         EventStore::query_raw(self.inner.as_ref(), value)
    159     }
    160     fn query_verified(
    161         &self,
    162         value: radroots_storage::event::EventQuery,
    163     ) -> radroots_transport::BoxFuture<
    164         '_,
    165         Result<
    166             radroots_storage::event::EventPage<radroots_storage::event::StoredVerifiedEvent>,
    167             radroots_storage::Error,
    168         >,
    169     > {
    170         EventStore::query_verified(self.inner.as_ref(), value)
    171     }
    172     fn query_visible(
    173         &self,
    174         value: radroots_storage::event::EventQuery,
    175     ) -> radroots_transport::BoxFuture<
    176         '_,
    177         Result<
    178             radroots_storage::event::EventPage<radroots_storage::event::StoredVisibleEvent>,
    179             radroots_storage::Error,
    180         >,
    181     > {
    182         EventStore::query_visible(self.inner.as_ref(), value)
    183     }
    184     fn rebuild_visibility(
    185         &self,
    186     ) -> radroots_transport::BoxFuture<
    187         '_,
    188         Result<radroots_storage::event::VisibilitySnapshot, radroots_storage::Error>,
    189     > {
    190         EventStore::rebuild_visibility(self.inner.as_ref())
    191     }
    192     fn query_provenance(
    193         &self,
    194         id: radroots_storage::event::EventId,
    195         bounds: radroots_storage::event::EventQueryBounds,
    196     ) -> radroots_transport::BoxFuture<
    197         '_,
    198         Result<
    199             radroots_storage::event::EventPage<radroots_storage::event::StoredEventProvenance>,
    200             radroots_storage::Error,
    201         >,
    202     > {
    203         EventStore::query_provenance(self.inner.as_ref(), id, bounds)
    204     }
    205 }
    206 
    207 impl Journal for FaultStorage {
    208     fn prepare(
    209         &self,
    210         value: radroots_storage::journal::PrepareOperation,
    211     ) -> radroots_transport::BoxFuture<
    212         '_,
    213         Result<radroots_storage::journal::PrepareReceipt, radroots_storage::Error>,
    214     > {
    215         Journal::prepare(self.inner.as_ref(), value)
    216     }
    217     fn operation(
    218         &self,
    219         id: OperationInstanceId,
    220     ) -> radroots_transport::BoxFuture<
    221         '_,
    222         Result<Option<radroots_storage::journal::OperationRecord>, radroots_storage::Error>,
    223     > {
    224         Journal::operation(self.inner.as_ref(), id)
    225     }
    226     fn by_idempotency_key(
    227         &self,
    228         id: radroots_storage::journal::OperationId,
    229         key: IdempotencyKey,
    230     ) -> radroots_transport::BoxFuture<
    231         '_,
    232         Result<Option<radroots_storage::journal::OperationRecord>, radroots_storage::Error>,
    233     > {
    234         Journal::by_idempotency_key(self.inner.as_ref(), id, key)
    235     }
    236     fn transition(
    237         &self,
    238         value: radroots_storage::journal::JournalTransition,
    239     ) -> radroots_transport::BoxFuture<
    240         '_,
    241         Result<radroots_storage::journal::OperationRecord, radroots_storage::Error>,
    242     > {
    243         Journal::transition(self.inner.as_ref(), value)
    244     }
    245     fn recoverable(
    246         &self,
    247         limit: u16,
    248     ) -> radroots_transport::BoxFuture<
    249         '_,
    250         Result<Vec<radroots_storage::journal::OperationRecord>, radroots_storage::Error>,
    251     > {
    252         Journal::recoverable(self.inner.as_ref(), limit)
    253     }
    254 }
    255 
    256 impl Outbox for FaultStorage {
    257     fn enqueue(
    258         &self,
    259         value: radroots_storage::outbox::EnqueueOutboxItem,
    260     ) -> radroots_transport::BoxFuture<
    261         '_,
    262         Result<radroots_storage::outbox::EnqueueReceipt, radroots_storage::Error>,
    263     > {
    264         Outbox::enqueue(self.inner.as_ref(), value)
    265     }
    266     fn item(
    267         &self,
    268         id: radroots_storage::outbox::OutboxItemId,
    269     ) -> radroots_transport::BoxFuture<
    270         '_,
    271         Result<Option<radroots_storage::outbox::OutboxRecord>, radroots_storage::Error>,
    272     > {
    273         Outbox::item(self.inner.as_ref(), id)
    274     }
    275     fn claim(
    276         &self,
    277         value: radroots_storage::outbox::ClaimOutboxItems,
    278     ) -> radroots_transport::BoxFuture<
    279         '_,
    280         Result<Vec<radroots_storage::outbox::ClaimedOutboxItem>, radroots_storage::Error>,
    281     > {
    282         Outbox::claim(self.inner.as_ref(), value)
    283     }
    284     fn record_attempt(
    285         &self,
    286         value: radroots_storage::outbox::DeliveryAttemptEvidence,
    287     ) -> radroots_transport::BoxFuture<
    288         '_,
    289         Result<radroots_storage::outbox::OutboxRecord, radroots_storage::Error>,
    290     > {
    291         Outbox::record_attempt(self.inner.as_ref(), value)
    292     }
    293     fn release(
    294         &self,
    295         id: radroots_storage::outbox::OutboxItemId,
    296         lease: radroots_storage::outbox::LeaseId,
    297         revision: radroots_storage::outbox::OutboxRevision,
    298         at: u64,
    299         retry: Option<u64>,
    300     ) -> radroots_transport::BoxFuture<
    301         '_,
    302         Result<radroots_storage::outbox::OutboxRecord, radroots_storage::Error>,
    303     > {
    304         Outbox::release(self.inner.as_ref(), id, lease, revision, at, retry)
    305     }
    306     fn status(
    307         &self,
    308     ) -> radroots_transport::BoxFuture<
    309         '_,
    310         Result<radroots_storage::outbox::OutboxStatus, radroots_storage::Error>,
    311     > {
    312         Outbox::status(self.inner.as_ref())
    313     }
    314 }
    315 
    316 impl ProjectionStore for FaultStorage {
    317     fn status(
    318         &self,
    319         id: radroots_storage::projection::ProjectionId,
    320     ) -> radroots_transport::BoxFuture<
    321         '_,
    322         Result<Option<radroots_storage::projection::ProjectionStatus>, radroots_storage::Error>,
    323     > {
    324         ProjectionStore::status(self.inner.as_ref(), id)
    325     }
    326     fn checkpoint(
    327         &self,
    328         value: radroots_storage::projection::ProjectionCheckpoint,
    329     ) -> radroots_transport::BoxFuture<
    330         '_,
    331         Result<radroots_storage::projection::ProjectionStatus, radroots_storage::Error>,
    332     > {
    333         ProjectionStore::checkpoint(self.inner.as_ref(), value)
    334     }
    335     fn invalidate(
    336         &self,
    337         value: radroots_storage::projection::ProjectionInvalidation,
    338     ) -> radroots_transport::BoxFuture<
    339         '_,
    340         Result<radroots_storage::projection::ProjectionStatus, radroots_storage::Error>,
    341     > {
    342         ProjectionStore::invalidate(self.inner.as_ref(), value)
    343     }
    344     fn invalidation(
    345         &self,
    346         id: radroots_storage::projection::ProjectionId,
    347         generation: radroots_storage::projection::ProjectionGeneration,
    348     ) -> radroots_transport::BoxFuture<
    349         '_,
    350         Result<
    351             Option<radroots_storage::projection::ProjectionInvalidation>,
    352             radroots_storage::Error,
    353         >,
    354     > {
    355         ProjectionStore::invalidation(self.inner.as_ref(), id, generation)
    356     }
    357     fn request_rebuild(
    358         &self,
    359         value: radroots_storage::projection::RebuildTicket,
    360     ) -> radroots_transport::BoxFuture<
    361         '_,
    362         Result<radroots_storage::projection::RebuildTicket, radroots_storage::Error>,
    363     > {
    364         ProjectionStore::request_rebuild(self.inner.as_ref(), value)
    365     }
    366     fn rebuild(
    367         &self,
    368         id: radroots_storage::projection::RebuildTicketId,
    369     ) -> radroots_transport::BoxFuture<
    370         '_,
    371         Result<Option<radroots_storage::projection::RebuildTicket>, radroots_storage::Error>,
    372     > {
    373         ProjectionStore::rebuild(self.inner.as_ref(), id)
    374     }
    375     fn transition_rebuild(
    376         &self,
    377         value: radroots_storage::projection::RebuildTransition,
    378     ) -> radroots_transport::BoxFuture<
    379         '_,
    380         Result<radroots_storage::projection::RebuildTicket, radroots_storage::Error>,
    381     > {
    382         ProjectionStore::transition_rebuild(self.inner.as_ref(), value)
    383     }
    384     fn event_index_manifest(
    385         &self,
    386         generation: radroots_storage::projection::ProjectionGeneration,
    387     ) -> radroots_transport::BoxFuture<
    388         '_,
    389         Result<Option<radroots_storage::projection::EventIndexManifest>, radroots_storage::Error>,
    390     > {
    391         ProjectionStore::event_index_manifest(self.inner.as_ref(), generation)
    392     }
    393     fn put_event_index_manifest(
    394         &self,
    395         value: radroots_storage::projection::EventIndexManifest,
    396     ) -> radroots_transport::BoxFuture<'_, Result<(), radroots_storage::Error>> {
    397         ProjectionStore::put_event_index_manifest(self.inner.as_ref(), value)
    398     }
    399     fn event_index_checkpoint(
    400         &self,
    401         generation: radroots_storage::projection::ProjectionGeneration,
    402     ) -> radroots_transport::BoxFuture<
    403         '_,
    404         Result<Option<radroots_storage::projection::EventIndexCheckpoint>, radroots_storage::Error>,
    405     > {
    406         ProjectionStore::event_index_checkpoint(self.inner.as_ref(), generation)
    407     }
    408     fn put_event_index_checkpoint(
    409         &self,
    410         value: radroots_storage::projection::EventIndexCheckpoint,
    411     ) -> radroots_transport::BoxFuture<'_, Result<(), radroots_storage::Error>> {
    412         ProjectionStore::put_event_index_checkpoint(self.inner.as_ref(), value)
    413     }
    414     fn put_projection_document(
    415         &self,
    416         projection_id: radroots_storage::projection::ProjectionId,
    417         generation: radroots_storage::projection::ProjectionGeneration,
    418         document: radroots_storage::projection::ProjectionDocument,
    419     ) -> radroots_transport::BoxFuture<'_, Result<(), radroots_storage::Error>> {
    420         ProjectionStore::put_projection_document(
    421             self.inner.as_ref(),
    422             projection_id,
    423             generation,
    424             document,
    425         )
    426     }
    427     fn projection_document(
    428         &self,
    429         projection_id: radroots_storage::projection::ProjectionId,
    430         generation: radroots_storage::projection::ProjectionGeneration,
    431         key: String,
    432     ) -> radroots_transport::BoxFuture<
    433         '_,
    434         Result<Option<radroots_storage::projection::ProjectionDocument>, radroots_storage::Error>,
    435     > {
    436         ProjectionStore::projection_document(self.inner.as_ref(), projection_id, generation, key)
    437     }
    438     fn put_projection_snapshot(
    439         &self,
    440         snapshot: radroots_storage::projection::ProjectionSnapshot,
    441     ) -> radroots_transport::BoxFuture<'_, Result<(), radroots_storage::Error>> {
    442         ProjectionStore::put_projection_snapshot(self.inner.as_ref(), snapshot)
    443     }
    444     fn projection_snapshot(
    445         &self,
    446         projection_id: radroots_storage::projection::ProjectionId,
    447         snapshot_id: [u8; 32],
    448     ) -> radroots_transport::BoxFuture<
    449         '_,
    450         Result<Option<radroots_storage::projection::ProjectionSnapshot>, radroots_storage::Error>,
    451     > {
    452         ProjectionStore::projection_snapshot(self.inner.as_ref(), projection_id, snapshot_id)
    453     }
    454 }
    455 
    456 impl AtomicStorage for FaultStorage {
    457     fn commit(
    458         &self,
    459         value: radroots_storage::atomic::AtomicCommit,
    460     ) -> radroots_transport::BoxFuture<
    461         '_,
    462         Result<radroots_storage::atomic::AtomicCommitReceipt, radroots_storage::Error>,
    463     > {
    464         AtomicStorage::commit(self.inner.as_ref(), value)
    465     }
    466     fn receipt(
    467         &self,
    468         id: radroots_storage::atomic::AtomicCommitId,
    469     ) -> radroots_transport::BoxFuture<
    470         '_,
    471         Result<Option<radroots_storage::atomic::AtomicCommitReceipt>, radroots_storage::Error>,
    472     > {
    473         AtomicStorage::receipt(self.inner.as_ref(), id)
    474     }
    475 }
    476 
    477 impl StorageStatusProvider for FaultStorage {
    478     fn storage_status(
    479         &self,
    480     ) -> radroots_transport::BoxFuture<
    481         '_,
    482         Result<radroots_storage::status::StorageStatus, radroots_storage::Error>,
    483     > {
    484         StorageStatusProvider::storage_status(self.inner.as_ref())
    485     }
    486 }
    487 
    488 impl AuthoredAtomicStorage for FaultStorage {
    489     fn execute_authored(
    490         &self,
    491         command: radroots_storage::authored_atomic::AuthoredAtomicCommand,
    492     ) -> radroots_transport::BoxFuture<
    493         '_,
    494         Result<radroots_storage::authored_atomic::AuthoredAtomicReceipt, radroots_storage::Error>,
    495     > {
    496         Box::pin(async move {
    497             delivery_evidence::inject(self, &command, true).await;
    498             capacity::inject(self, &command, false)?;
    499             let receipt =
    500                 AuthoredAtomicStorage::execute_authored(self.inner.as_ref(), command.clone())
    501                     .await?;
    502             capacity::inject(self, &command, true)?;
    503             delivery_evidence::inject(self, &command, false).await;
    504             if matches!(
    505                 receipt.outcome(),
    506                 radroots_storage::authored_atomic::AuthoredAtomicOutcome::DeliveryPlan(_)
    507             ) && let Some(mutation) = self.receipt_mutation.lock().unwrap().take()
    508             {
    509                 return Ok(delivery_evidence::mutate_receipt(receipt, mutation));
    510             }
    511             if let radroots_storage::authored_atomic::AuthoredAtomicOutcome::Prepared { .. } =
    512                 receipt.outcome()
    513             {
    514                 let mut cache = self.prepared.lock().expect("prepared cache");
    515                 if cache.is_none() {
    516                     *cache = Some(receipt.outcome().clone());
    517                 }
    518             }
    519             let kind = match receipt.outcome() {
    520                 radroots_storage::authored_atomic::AuthoredAtomicOutcome::Prepared { .. } => 1,
    521                 radroots_storage::authored_atomic::AuthoredAtomicOutcome::Artifact(_) => 2,
    522                 radroots_storage::authored_atomic::AuthoredAtomicOutcome::DeliveryPlan(_) => 3,
    523                 radroots_storage::authored_atomic::AuthoredAtomicOutcome::Submitted(_) => 5,
    524             };
    525             let fault = self.fault_kind.load(Ordering::Relaxed);
    526             if (fault != kind && !(fault == 4 && kind == 1))
    527                 || self.remaining.fetch_sub(1, Ordering::Relaxed) != 1
    528             {
    529                 return Ok(receipt);
    530             }
    531             let cached = self
    532                 .prepared
    533                 .lock()
    534                 .expect("prepared cache")
    535                 .clone()
    536                 .expect("prepared outcome");
    537             let forged = match (fault, cached) {
    538                 (4, cached) => cached,
    539                 (
    540                     1,
    541                     radroots_storage::authored_atomic::AuthoredAtomicOutcome::Prepared {
    542                         artifacts,
    543                         ..
    544                     },
    545                 ) => radroots_storage::authored_atomic::AuthoredAtomicOutcome::Artifact(
    546                     artifacts[0].clone(),
    547                 ),
    548                 (
    549                     2,
    550                     radroots_storage::authored_atomic::AuthoredAtomicOutcome::Prepared {
    551                         delivery_plans,
    552                         ..
    553                     },
    554                 ) => radroots_storage::authored_atomic::AuthoredAtomicOutcome::DeliveryPlan(
    555                     delivery_plans[0].clone(),
    556                 ),
    557                 (
    558                     3,
    559                     radroots_storage::authored_atomic::AuthoredAtomicOutcome::Prepared {
    560                         artifacts,
    561                         ..
    562                     },
    563                 ) => radroots_storage::authored_atomic::AuthoredAtomicOutcome::Artifact(
    564                     artifacts[0].clone(),
    565                 ),
    566                 _ => unreachable!(),
    567             };
    568             radroots_storage::authored_atomic::AuthoredAtomicReceipt::from_durable_parts(
    569                 receipt.commit_id(),
    570                 receipt.digest(),
    571                 receipt.disposition(),
    572                 receipt.committed_at_unix_ms(),
    573                 forged,
    574             )
    575         })
    576     }
    577     fn authored_receipt(
    578         &self,
    579         id: radroots_storage::atomic::AtomicCommitId,
    580     ) -> radroots_transport::BoxFuture<
    581         '_,
    582         Result<
    583             Option<radroots_storage::authored_atomic::AuthoredAtomicReceipt>,
    584             radroots_storage::Error,
    585         >,
    586     > {
    587         AuthoredAtomicStorage::authored_receipt(self.inner.as_ref(), id)
    588     }
    589     fn authored_operation(
    590         &self,
    591         id: OperationInstanceId,
    592     ) -> radroots_transport::BoxFuture<
    593         '_,
    594         Result<Option<radroots_storage::authored::AuthoredOperation>, radroots_storage::Error>,
    595     > {
    596         AuthoredAtomicStorage::authored_operation(self.inner.as_ref(), id)
    597     }
    598     fn authored_artifact(
    599         &self,
    600         id: radroots_storage::authored::AuthoredArtifactId,
    601     ) -> radroots_transport::BoxFuture<
    602         '_,
    603         Result<Option<radroots_storage::authored::AuthoredArtifact>, radroots_storage::Error>,
    604     > {
    605         AuthoredAtomicStorage::authored_artifact(self.inner.as_ref(), id)
    606     }
    607     fn authored_delivery_plan(
    608         &self,
    609         id: radroots_storage::authored_delivery::AuthoredDeliveryPlanId,
    610     ) -> radroots_transport::BoxFuture<
    611         '_,
    612         Result<
    613             Option<radroots_storage::authored_delivery::AuthoredDeliveryPlan>,
    614             radroots_storage::Error,
    615         >,
    616     > {
    617         AuthoredAtomicStorage::authored_delivery_plan(self.inner.as_ref(), id)
    618     }
    619     fn authored_delivery_history(
    620         &self,
    621         id: radroots_storage::authored_delivery::AuthoredDeliveryPlanId,
    622     ) -> radroots_transport::BoxFuture<
    623         '_,
    624         Result<
    625             Option<radroots_storage::authored_delivery::AuthoredDeliveryHistory>,
    626             radroots_storage::Error,
    627         >,
    628     > {
    629         Box::pin(async move {
    630             let history = self.inner.authored_delivery_history(id).await?;
    631             let mutation = {
    632                 let mut armed = self.history_mutation.lock().unwrap();
    633                 if let Some((remaining, _)) = armed.as_mut() {
    634                     *remaining -= 1;
    635                 }
    636                 if armed.as_ref().is_some_and(|(remaining, _)| *remaining == 0) {
    637                     armed.take().map(|(_, mutation)| mutation)
    638                 } else {
    639                     None
    640                 }
    641             };
    642             Ok(match (history, mutation) {
    643                 (Some(history), Some(mutation)) => {
    644                     Some(delivery_evidence::mutate_history(history, mutation))
    645                 }
    646                 (history, _) => history,
    647             })
    648         })
    649     }
    650 }
    651 
    652 struct MockSource;
    653 
    654 impl EventSource for MockSource {
    655     fn status(&self) -> radroots_transport::BoxFuture<'_, Result<SourceStatus, TransportError>> {
    656         Box::pin(async { unreachable!("missing-sink check does not inspect source") })
    657     }
    658 
    659     fn fetch(
    660         &self,
    661         _request: FetchRequest,
    662     ) -> radroots_transport::BoxFuture<'_, Result<FetchPage, TransportError>> {
    663         Box::pin(async { unreachable!("missing-sink check does not fetch") })
    664     }
    665 }
    666 
    667 impl EventSink for MockSink {
    668     fn status(&self) -> radroots_transport::BoxFuture<'_, Result<SinkStatus, TransportError>> {
    669         Box::pin(async { unreachable!("enqueue does not inspect sink") })
    670     }
    671     fn deliver(
    672         &self,
    673         _request: DeliveryRequest,
    674     ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> {
    675         Box::pin(async { unreachable!("enqueue does not deliver") })
    676     }
    677 }
    678 
    679 enum DeliveryBehavior {
    680     Outcomes(Vec<DeliveryOutcome>),
    681     MismatchedRequest,
    682     Failure(Retryability),
    683     MismatchedFailure,
    684     Pending,
    685 }
    686 
    687 struct ScriptedSink {
    688     behaviors: Mutex<VecDeque<DeliveryBehavior>>,
    689     requests: Mutex<Vec<DeliveryRequest>>,
    690 }
    691 
    692 impl ScriptedSink {
    693     fn new(behaviors: impl IntoIterator<Item = DeliveryBehavior>) -> Self {
    694         Self {
    695             behaviors: Mutex::new(behaviors.into_iter().collect()),
    696             requests: Mutex::new(Vec::new()),
    697         }
    698     }
    699 }
    700 
    701 impl EventSink for ScriptedSink {
    702     fn status(&self) -> radroots_transport::BoxFuture<'_, Result<SinkStatus, TransportError>> {
    703         Box::pin(async { unreachable!("delivery does not inspect sink status") })
    704     }
    705 
    706     fn deliver(
    707         &self,
    708         request: DeliveryRequest,
    709     ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> {
    710         self.requests
    711             .lock()
    712             .expect("scripted request lock")
    713             .push(request.clone());
    714         let behavior = self
    715             .behaviors
    716             .lock()
    717             .expect("scripted behavior lock")
    718             .pop_front()
    719             .expect("scripted delivery behavior");
    720         Box::pin(async move {
    721             match behavior {
    722                 DeliveryBehavior::Outcomes(outcomes) => {
    723                     Ok(receipt(&request, outcomes).expect("valid scripted receipt"))
    724                 }
    725                 DeliveryBehavior::MismatchedRequest => {
    726                     let mismatched = DeliveryRequest::new(
    727                         "mismatched-request",
    728                         request.payload().clone(),
    729                         request.target_set().clone(),
    730                         request.satisfaction().clone(),
    731                         request.deadline_unix_ms(),
    732                     )
    733                     .expect("mismatched request");
    734                     Ok(receipt(
    735                         &mismatched,
    736                         vec![DeliveryOutcome::accepted(); mismatched.target_set().len()],
    737                     )
    738                     .expect("mismatched receipt"))
    739                 }
    740                 DeliveryBehavior::Failure(retryability) => Err(SinkFailure::for_request(
    741                     &request,
    742                     "scripted_failure",
    743                     retryability,
    744                     None,
    745                     Some("scripted failure".to_owned()),
    746                     Vec::new(),
    747                 )
    748                 .expect("valid scripted failure")),
    749                 DeliveryBehavior::MismatchedFailure => {
    750                     let mismatched = DeliveryRequest::new(
    751                         "mismatched-failure",
    752                         request.payload().clone(),
    753                         request.target_set().clone(),
    754                         request.satisfaction().clone(),
    755                         request.deadline_unix_ms(),
    756                     )
    757                     .expect("mismatched request");
    758                     Err(SinkFailure::for_request(
    759                         &mismatched,
    760                         "mismatched_failure",
    761                         Retryability::Terminal,
    762                         None,
    763                         None,
    764                         Vec::new(),
    765                     )
    766                     .expect("mismatched failure"))
    767                 }
    768                 DeliveryBehavior::Pending => std::future::pending().await,
    769             }
    770         })
    771     }
    772 }
    773 
    774 struct TestClock(AtomicU64);
    775 
    776 impl Clock for TestClock {
    777     fn now_unix_ms(&self) -> Result<u64, Error> {
    778         Ok(self.0.fetch_add(1, Ordering::Relaxed))
    779     }
    780 }
    781 
    782 struct TestIds(AtomicU64);
    783 
    784 impl IdSource for TestIds {
    785     fn next_id(&self, _operation: OperationKind) -> Result<SyncId, Error> {
    786         let value = self.0.fetch_add(1, Ordering::Relaxed);
    787         let byte = u8::try_from(value).map_err(|_| Error::InvalidSyncId)?;
    788         SyncId::new([byte; 16])
    789     }
    790 }
    791 
    792 #[derive(Clone, Copy)]
    793 enum SignBehavior {
    794     Success { completed_at_unix_ms: u64 },
    795     Error(SigningErrorKind),
    796     Uncertain(SigningErrorKind),
    797     Pending,
    798 }
    799 
    800 struct MockSigner {
    801     behavior: SignBehavior,
    802     replay: ReplayCapability,
    803     calls: AtomicUsize,
    804 }
    805 
    806 impl MockSigner {
    807     fn new(behavior: SignBehavior) -> Self {
    808         Self {
    809             behavior,
    810             replay: ReplayCapability::LocalReplaySafe,
    811             calls: AtomicUsize::new(0),
    812         }
    813     }
    814 
    815     fn with_replay(behavior: SignBehavior, replay: ReplayCapability) -> Self {
    816         Self {
    817             behavior,
    818             replay,
    819             calls: AtomicUsize::new(0),
    820         }
    821     }
    822 }
    823 
    824 impl Signer for MockSigner {
    825     fn status(
    826         &self,
    827     ) -> radroots_signing::signer::BoxFuture<'_, Result<SignerStatus, SigningError>> {
    828         let replay = self.replay;
    829         Box::pin(async move {
    830             Ok(SignerStatus::new(
    831                 SignerAvailability::Ready,
    832                 vec![SignerCapability::new(
    833                     if replay == ReplayCapability::LocalReplaySafe {
    834                         SignerKind::Local
    835                     } else {
    836                         SignerKind::Remote
    837                     },
    838                     replay,
    839                     CancellationSupport::BeforePublication,
    840                     false,
    841                     false,
    842                 )],
    843                 None,
    844             ))
    845         })
    846     }
    847 
    848     fn sign(
    849         &self,
    850         request: SignRequest,
    851     ) -> radroots_signing::signer::BoxFuture<'_, Result<SignReceipt, SigningError>> {
    852         self.calls.fetch_add(1, Ordering::Relaxed);
    853         Box::pin(async move {
    854             match self.behavior {
    855                 SignBehavior::Success {
    856                     completed_at_unix_ms,
    857                 } => SignReceipt::from_signed_event(
    858                     &request,
    859                     signed_event(&request),
    860                     completed_at_unix_ms,
    861                 ),
    862                 SignBehavior::Error(kind) => Err(SigningError::new(kind)),
    863                 SignBehavior::Uncertain(kind) => {
    864                     Err(SigningError::new(kind).with_possible_remote_effect())
    865                 }
    866                 SignBehavior::Pending => std::future::pending().await,
    867             }
    868         })
    869     }
    870 }
    871 
    872 #[derive(Clone, Copy)]
    873 enum BoundaryViolation {
    874     CompletesAfterDeadline,
    875     CancelsBeforeReturn,
    876 }
    877 
    878 struct BoundaryViolatingSigner {
    879     violation: BoundaryViolation,
    880     clock: Arc<TestClock>,
    881 }
    882 
    883 impl Signer for BoundaryViolatingSigner {
    884     fn status(
    885         &self,
    886     ) -> radroots_signing::signer::BoxFuture<'_, Result<SignerStatus, SigningError>> {
    887         Box::pin(async {
    888             Ok(SignerStatus::new(
    889                 SignerAvailability::Ready,
    890                 vec![SignerCapability::new(
    891                     SignerKind::Remote,
    892                     ReplayCapability::ExactReplayByRequestId,
    893                     CancellationSupport::BeforeAndAfterPublication,
    894                     false,
    895                     false,
    896                 )],
    897                 None,
    898             ))
    899         })
    900     }
    901 
    902     fn sign(
    903         &self,
    904         request: SignRequest,
    905     ) -> radroots_signing::signer::BoxFuture<'_, Result<SignReceipt, SigningError>> {
    906         Box::pin(async move {
    907             let deadline = request.policy().deadline_unix_ms();
    908             let receipt =
    909                 SignReceipt::from_signed_event(&request, signed_event(&request), deadline - 1)?;
    910             match self.violation {
    911                 BoundaryViolation::CompletesAfterDeadline => {
    912                     self.clock.0.store(deadline, Ordering::Release);
    913                 }
    914                 BoundaryViolation::CancelsBeforeReturn => request.cancellation_signal().cancel(),
    915             }
    916             Ok(receipt)
    917         })
    918     }
    919 }
    920 
    921 fn signing_keypair() -> Keypair {
    922     let secret = SecretKey::from_slice(&[1; 32]).expect("secret key");
    923     Keypair::from_secret_key(&Secp256k1::new(), &secret)
    924 }
    925 
    926 fn public_key_hex() -> String {
    927     signing_keypair().x_only_public_key().0.to_string()
    928 }
    929 
    930 fn signed_event(request: &SignRequest) -> SignedEvent {
    931     let id = request.expected_event_id().to_hex();
    932     let pubkey = public_key_hex();
    933     let signature = Secp256k1::new()
    934         .sign_schnorr_no_aux_rand(
    935             &Message::from_digest(*request.expected_event_id().as_bytes()),
    936             &signing_keypair(),
    937         )
    938         .to_string();
    939     let raw_json = format!(
    940         "{{\"id\":\"{id}\",\"pubkey\":\"{pubkey}\",\"created_at\":{},\"kind\":{},\"tags\":{:?},\"content\":{content:?},\"sig\":\"{signature}\"}}",
    941         request.created_at(),
    942         request.kind(),
    943         request.tags(),
    944         content = request.content(),
    945     );
    946     SignedEvent::new(SignedEventParts {
    947         id,
    948         pubkey,
    949         created_at: request.created_at(),
    950         kind: request.kind(),
    951         tags: request.tags().to_vec(),
    952         content: request.content().to_owned(),
    953         sig: signature,
    954         raw_json,
    955     })
    956     .expect("signed event")
    957 }
    958 
    959 fn request(operation_byte: u8, relay: &str) -> PushRequest {
    960     request_with_policy(
    961         operation_byte,
    962         &[relay],
    963         SatisfactionClass::Accepted,
    964         TargetPolicy::any(),
    965     )
    966 }
    967 
    968 fn food_request(operation_byte: u8, relay: &str) -> PushRequest {
    969     let created_at = 1_800_000_100;
    970     let details = FoodAvailabilityDetails::new(FoodAvailabilityDetailsParts {
    971         content: FoodContent::new("Carrots available this week.").expect("food content"),
    972         identifier: FoodIdentifier::parse("nantes-carrots").expect("food identifier"),
    973         title: FoodText::new("Nantes Carrots").expect("food title"),
    974         summary: FoodText::new("Fresh bunches").expect("food summary"),
    975         published_at: FoodPublishedAt::new(created_at).expect("published at"),
    976         location: FoodText::new("Central Saanich, BC").expect("food location"),
    977         price: FoodPrice::new(
    978             "3",
    979             FoodCurrency::parse("CAD").expect("currency"),
    980             FoodUnit::Pound,
    981         )
    982         .expect("food price"),
    983         quantity: None,
    984         status: FoodAvailabilityStatus::Active,
    985         images: Vec::new(),
    986     })
    987     .expect("food availability");
    988     let plan = AuthoredEventPlan::from_food_availability(&details, created_at, public_key_hex())
    989         .expect("typed food plan");
    990     let actor = Actor::new(
    991         *plan.author(),
    992         ActorSource::ExplicitPublicKey,
    993         [AuthorRole::Seller],
    994     )
    995     .expect("food actor");
    996     PushRequest::new(
    997         SyncId::new([operation_byte; 16]).expect("operation id"),
    998         IdempotencyKey::parse(format!("food-push-{operation_byte}")).expect("idempotency key"),
    999         actor,
   1000         plan,
   1001         TargetSet::new(vec![
   1002             Target::new(TransportId::NOSTR, relay).expect("target"),
   1003         ])
   1004         .expect("targets"),
   1005         SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::any()),
   1006         1_800_000_300_000,
   1007         CancellationPolicy::PreservePublishedRequest,
   1008     )
   1009     .expect("food push request")
   1010 }
   1011 
   1012 fn request_with_policy(
   1013     operation_byte: u8,
   1014     relays: &[&str],
   1015     class: SatisfactionClass,
   1016     target_policy: TargetPolicy,
   1017 ) -> PushRequest {
   1018     let pubkey = public_key_hex();
   1019     let plan = AuthoredEventPlan::from_generic(
   1020         GenericEventDraft::new(
   1021             "radroots.social.geochat.v1",
   1022             20_000,
   1023             1_800_000_100,
   1024             vec![],
   1025             CONTENT,
   1026             pubkey,
   1027         )
   1028         .expect("draft"),
   1029     )
   1030     .expect("authored plan");
   1031     let actor = Actor::new(
   1032         *plan.author(),
   1033         ActorSource::ExplicitPublicKey,
   1034         [AuthorRole::Any],
   1035     )
   1036     .expect("actor");
   1037     PushRequest::new(
   1038         SyncId::new([operation_byte; 16]).expect("operation id"),
   1039         IdempotencyKey::parse(format!("push-{operation_byte}")).expect("idempotency key"),
   1040         actor,
   1041         plan,
   1042         TargetSet::new(
   1043             relays
   1044                 .iter()
   1045                 .map(|relay| Target::new(TransportId::NOSTR, *relay).expect("target"))
   1046                 .collect(),
   1047         )
   1048         .expect("targets"),
   1049         SatisfactionPolicy::new(class, target_policy),
   1050         1_800_000_300_000,
   1051         CancellationPolicy::PreservePublishedRequest,
   1052     )
   1053     .expect("push request")
   1054 }
   1055 
   1056 fn receipt(
   1057     request: &DeliveryRequest,
   1058     outcomes: Vec<DeliveryOutcome>,
   1059 ) -> Result<DeliveryReceipt, TransportError> {
   1060     let targets = request
   1061         .target_set()
   1062         .targets()
   1063         .iter()
   1064         .cloned()
   1065         .zip(outcomes)
   1066         .map(|(target, outcome)| DeliveryTargetReceipt::attempted(target, outcome))
   1067         .collect();
   1068     DeliveryReceipt::for_request(request, targets)
   1069 }
   1070 
   1071 fn setup_engine(signer: Arc<MockSigner>) -> (Engine, Arc<MemoryStorage>) {
   1072     setup_engine_with_sink(signer, Arc::new(MockSink)).0
   1073 }
   1074 
   1075 fn setup_engine_with_sink(
   1076     signer: Arc<MockSigner>,
   1077     sink: Arc<dyn EventSink>,
   1078 ) -> ((Engine, Arc<MemoryStorage>), Arc<TestClock>) {
   1079     let storage = Arc::new(MemoryStorage::new(
   1080         SourceGeneration::new([6; 32]).expect("generation"),
   1081     ));
   1082     let capability: Arc<dyn SyncStorage> = storage.clone();
   1083     let clock = Arc::new(TestClock(AtomicU64::new(1_800_000_200_000)));
   1084     let engine = Engine::builder(
   1085         capability,
   1086         clock.clone(),
   1087         Arc::new(TestIds(AtomicU64::new(10))),
   1088         DeadlinePolicy::new(10_000, 10_000, 10_000).expect("deadlines"),
   1089     )
   1090     .sink(sink)
   1091     .signer(signer)
   1092     .build()
   1093     .expect("engine");
   1094     ((engine, storage), clock)
   1095 }
   1096 
   1097 fn fault_engine(
   1098     storage: Arc<FaultStorage>,
   1099     signer: Arc<MockSigner>,
   1100     sink: Arc<dyn EventSink>,
   1101 ) -> Engine {
   1102     let capability: Arc<dyn SyncStorage> = storage;
   1103     Engine::builder(
   1104         capability,
   1105         Arc::new(TestClock(AtomicU64::new(1_800_000_200_000))),
   1106         Arc::new(TestIds(AtomicU64::new(10))),
   1107         DeadlinePolicy::new(10_000, 10_000, 10_000).unwrap(),
   1108     )
   1109     .sink(sink)
   1110     .signer(signer)
   1111     .build()
   1112     .unwrap()
   1113 }
   1114 
   1115 fn execute_to_admitted(engine: &Engine, push: &PushRequest) {
   1116     block_on(engine.sign_prepared(push.clone())).expect("sign prepared push");
   1117     block_on(engine.admit_signed(push.operation_id())).expect("admit signed push");
   1118 }
   1119 
   1120 #[test]
   1121 fn typed_only_authored_contract_identity_survives_signing_and_local_admission() {
   1122     let signer = Arc::new(MockSigner::new(SignBehavior::Success {
   1123         completed_at_unix_ms: 1_800_000_200_500,
   1124     }));
   1125     let (engine, storage) = setup_engine(signer);
   1126     let push = food_request(9, "wss://relay.example");
   1127 
   1128     block_on(engine.sign_prepared(push.clone())).expect("sign typed food plan");
   1129     let admitted = block_on(engine.admit_signed(push.operation_id()))
   1130         .expect("admit with the frozen typed contract identity");
   1131     assert!(admitted.artifact().admission_state().is_admitted());
   1132     let events = block_on(storage.query_visible(EventQuery::all(
   1133         EventQueryBounds::first(10).expect("bounds"),
   1134     )))
   1135     .expect("visible authored events");
   1136     assert_eq!(events.items().len(), 1);
   1137     assert_eq!(events.items()[0].event().kind(), 30_402);
   1138 }
   1139 
   1140 #[test]
   1141 fn caller_retains_late_and_cancelled_signer_facts_before_reporting_wait_outcome() {
   1142     for (byte, violation, expected) in [
   1143         (
   1144             55,
   1145             BoundaryViolation::CompletesAfterDeadline,
   1146             Error::SignerDeadlineExceeded,
   1147         ),
   1148         (
   1149             56,
   1150             BoundaryViolation::CancelsBeforeReturn,
   1151             Error::SigningCancelled,
   1152         ),
   1153     ] {
   1154         let storage = Arc::new(MemoryStorage::new(
   1155             SourceGeneration::new([byte; 32]).expect("generation"),
   1156         ));
   1157         let capability: Arc<dyn SyncStorage> = storage;
   1158         let clock = Arc::new(TestClock(AtomicU64::new(1_800_000_200_000)));
   1159         let signer: Arc<dyn Signer> = Arc::new(BoundaryViolatingSigner {
   1160             violation,
   1161             clock: Arc::clone(&clock),
   1162         });
   1163         let engine = Engine::builder(
   1164             capability,
   1165             clock,
   1166             Arc::new(TestIds(AtomicU64::new(10))),
   1167             DeadlinePolicy::new(10_000, 10_000, 10_000).expect("deadlines"),
   1168         )
   1169         .sink(Arc::new(MockSink))
   1170         .signer(signer)
   1171         .build()
   1172         .expect("engine");
   1173         let push = request(byte, "wss://relay.example");
   1174 
   1175         assert_eq!(block_on(engine.sign_prepared(push.clone())), Err(expected));
   1176         let status = block_on(engine.push_status(push.operation_id()))
   1177             .expect("status")
   1178             .expect("prepared operation");
   1179         let signed = status
   1180             .artifact()
   1181             .signed()
   1182             .expect("retain verified signer facts");
   1183         assert_eq!(signed.event().id(), push.plan().expected_event_id());
   1184         assert_eq!(signed.event().created_at(), push.plan().created_at());
   1185         assert_eq!(
   1186             status.artifact().admission_state(),
   1187             if matches!(violation, BoundaryViolation::CancelsBeforeReturn) {
   1188                 radroots_storage::authored::AdmissionState::Cancelled
   1189             } else {
   1190                 radroots_storage::authored::AdmissionState::Pending
   1191             }
   1192         );
   1193         assert!(status.delivery_plan().attempts().is_empty());
   1194     }
   1195 }
   1196 
   1197 #[test]
   1198 fn preparation_is_atomic_status_visible_and_replays_without_external_effects() {
   1199     let signer = Arc::new(MockSigner::new(SignBehavior::Pending));
   1200     let (engine, storage) = setup_engine(signer.clone());
   1201     let push = request(6, "wss://relay.example");
   1202 
   1203     let never_polled = engine.prepare_push(push.clone());
   1204     drop(never_polled);
   1205     assert!(
   1206         block_on(engine.push_status(push.operation_id()))
   1207             .expect("status before prepare")
   1208             .is_none()
   1209     );
   1210 
   1211     let prepared = block_on(engine.prepare_push(push.clone())).expect("prepare");
   1212     assert!(!prepared.is_replay());
   1213     assert_eq!(
   1214         prepared.operation().artifact_ids(),
   1215         &[prepared.artifact().artifact_id()]
   1216     );
   1217     assert_eq!(
   1218         prepared.delivery_plan().artifact_id(),
   1219         prepared.artifact().artifact_id()
   1220     );
   1221     assert!(prepared.delivery_plan().request().is_none());
   1222     assert!(
   1223         block_on(Journal::operation(
   1224             &*storage,
   1225             OperationInstanceId::new(*push.operation_id().as_bytes()).expect("operation"),
   1226         ))
   1227         .expect("legacy journal lookup")
   1228         .is_none()
   1229     );
   1230 
   1231     let status = block_on(engine.push_status(push.operation_id()))
   1232         .expect("status after prepare")
   1233         .expect("prepared status");
   1234     assert_eq!(status.operation(), prepared.operation());
   1235     assert_eq!(status.artifact(), prepared.artifact());
   1236     assert_eq!(status.delivery_plan(), prepared.delivery_plan());
   1237 
   1238     let replay = block_on(engine.prepare_push(push)).expect("exact replay");
   1239     assert!(replay.is_replay());
   1240     assert_eq!(replay.operation(), prepared.operation());
   1241     assert_eq!(replay.artifact(), prepared.artifact());
   1242     assert_eq!(replay.delivery_plan(), prepared.delivery_plan());
   1243     assert_eq!(signer.calls.load(Ordering::Relaxed), 0);
   1244 
   1245     assert_eq!(
   1246         block_on(engine.prepare_push(request(6, "wss://other.example"))),
   1247         Err(Error::StorageConflict)
   1248     );
   1249     assert_eq!(signer.calls.load(Ordering::Relaxed), 0);
   1250 }
   1251 
   1252 #[test]
   1253 fn phase_aware_cancellation_preserves_signed_artifacts_and_is_reentrant() {
   1254     let unsigned_signer = Arc::new(MockSigner::new(SignBehavior::Pending));
   1255     let (unsigned_engine, _) = setup_engine(unsigned_signer.clone());
   1256     let unsigned = request(61, "wss://relay.example");
   1257     block_on(unsigned_engine.prepare_push(unsigned.clone())).expect("prepare unsigned");
   1258     let cancelled = block_on(unsigned_engine.cancel_push(unsigned.operation_id()))
   1259         .expect("cancel unsigned operation");
   1260     assert!(cancelled.changed());
   1261     assert_eq!(
   1262         cancelled.status().artifact().signing_state(),
   1263         radroots_storage::authored::SigningState::Cancelled
   1264     );
   1265     assert_eq!(
   1266         cancelled.status().delivery_plan().state(),
   1267         AuthoredDeliveryState::Cancelled
   1268     );
   1269     assert!(cancelled.status().settlement().is_settled());
   1270     let replay = block_on(unsigned_engine.cancel_push(unsigned.operation_id()))
   1271         .expect("replay cancellation");
   1272     assert!(!replay.changed());
   1273     assert_eq!(replay.status(), cancelled.status());
   1274     assert_eq!(unsigned_signer.calls.load(Ordering::Relaxed), 0);
   1275 
   1276     let signed_signer = Arc::new(MockSigner::new(SignBehavior::Success {
   1277         completed_at_unix_ms: 1_800_000_200_010,
   1278     }));
   1279     let (signed_engine, _) = setup_engine(signed_signer);
   1280     let signed = request(62, "wss://relay.example");
   1281     block_on(signed_engine.sign_prepared(signed.clone())).expect("sign operation");
   1282     let before = block_on(signed_engine.push_status(signed.operation_id()))
   1283         .unwrap()
   1284         .unwrap();
   1285     let raw = before
   1286         .artifact()
   1287         .signed()
   1288         .expect("exact signed artifact")
   1289         .event()
   1290         .raw_json()
   1291         .to_owned();
   1292     let cancelled = block_on(signed_engine.cancel_push(signed.operation_id()))
   1293         .expect("cancel signed operation");
   1294     assert_eq!(
   1295         cancelled.status().artifact().admission_state(),
   1296         radroots_storage::authored::AdmissionState::Cancelled
   1297     );
   1298     assert_eq!(
   1299         cancelled
   1300             .status()
   1301             .artifact()
   1302             .signed()
   1303             .expect("signed bytes retained")
   1304             .event()
   1305             .raw_json(),
   1306         raw
   1307     );
   1308     assert_eq!(
   1309         cancelled.status().delivery_plan().state(),
   1310         AuthoredDeliveryState::Cancelled
   1311     );
   1312 
   1313     let admitted_signer = Arc::new(MockSigner::new(SignBehavior::Success {
   1314         completed_at_unix_ms: 1_800_000_200_010,
   1315     }));
   1316     let (admitted_engine, _) = setup_engine(admitted_signer);
   1317     let admitted = request(63, "wss://relay.example");
   1318     execute_to_admitted(&admitted_engine, &admitted);
   1319     let cancelled = block_on(admitted_engine.cancel_push(admitted.operation_id()))
   1320         .expect("cancel admitted delivery");
   1321     assert!(
   1322         cancelled
   1323             .status()
   1324             .artifact()
   1325             .admission_state()
   1326             .is_admitted()
   1327     );
   1328     assert_eq!(
   1329         cancelled.status().delivery_plan().state(),
   1330         AuthoredDeliveryState::Cancelled
   1331     );
   1332 }
   1333 
   1334 #[test]
   1335 fn authored_storage_outcome_contracts_fail_closed_at_each_orchestration_phase() {
   1336     let storage = Arc::new(FaultStorage::new(84));
   1337     storage.fault_next_prepared();
   1338     let engine = fault_engine(
   1339         storage,
   1340         Arc::new(MockSigner::new(SignBehavior::Pending)),
   1341         Arc::new(MockSink),
   1342     );
   1343     assert_eq!(
   1344         block_on(engine.prepare_push(request(84, "wss://fault.example"))),
   1345         Err(Error::StorageFailed)
   1346     );
   1347 
   1348     let storage = Arc::new(FaultStorage::new(85));
   1349     let engine = fault_engine(
   1350         storage.clone(),
   1351         Arc::new(MockSigner::new(SignBehavior::Pending)),
   1352         Arc::new(MockSink),
   1353     );
   1354     block_on(engine.prepare_push(request(85, "wss://fault.example"))).unwrap();
   1355     storage.fault_prepared_identity();
   1356     assert_eq!(
   1357         block_on(engine.prepare_push(request(86, "wss://fault.example"))),
   1358         Err(Error::StorageFailed)
   1359     );
   1360 
   1361     let storage = Arc::new(FaultStorage::new(87));
   1362     let engine = fault_engine(
   1363         storage.clone(),
   1364         Arc::new(MockSigner::new(SignBehavior::Success {
   1365             completed_at_unix_ms: 1_800_000_200_500,
   1366         })),
   1367         Arc::new(MockSink),
   1368     );
   1369     let push = request(87, "wss://fault.example");
   1370     block_on(engine.prepare_push(push.clone())).unwrap();
   1371     storage.fault_nth_artifact(1);
   1372     assert_eq!(
   1373         block_on(engine.sign_prepared(push)),
   1374         Err(Error::StorageFailed)
   1375     );
   1376 
   1377     let storage = Arc::new(FaultStorage::new(88));
   1378     let engine = fault_engine(
   1379         storage.clone(),
   1380         Arc::new(MockSigner::new(SignBehavior::Success {
   1381             completed_at_unix_ms: 1_800_000_200_500,
   1382         })),
   1383         Arc::new(MockSink),
   1384     );
   1385     let push = request(88, "wss://fault.example");
   1386     block_on(engine.prepare_push(push.clone())).unwrap();
   1387     storage.fault_nth_artifact(2);
   1388     assert_eq!(
   1389         block_on(engine.sign_prepared(push)),
   1390         Err(Error::StorageFailed)
   1391     );
   1392 
   1393     let storage = Arc::new(FaultStorage::new(89));
   1394     let engine = fault_engine(
   1395         storage.clone(),
   1396         Arc::new(MockSigner::new(SignBehavior::Success {
   1397             completed_at_unix_ms: 1_800_000_200_500,
   1398         })),
   1399         Arc::new(MockSink),
   1400     );
   1401     let push = request(89, "wss://fault.example");
   1402     block_on(engine.sign_prepared(push.clone())).unwrap();
   1403     storage.fault_nth_artifact(2);
   1404     assert_eq!(
   1405         block_on(engine.admit_signed(push.operation_id())),
   1406         Err(Error::StorageFailed)
   1407     );
   1408 
   1409     let storage = Arc::new(FaultStorage::new(90));
   1410     let engine = fault_engine(
   1411         storage.clone(),
   1412         Arc::new(MockSigner::new(SignBehavior::Error(
   1413             SigningErrorKind::SignerRejected,
   1414         ))),
   1415         Arc::new(MockSink),
   1416     );
   1417     let push = request(90, "wss://fault.example");
   1418     block_on(engine.prepare_push(push.clone())).unwrap();
   1419     storage.fault_nth_artifact(2);
   1420     assert_eq!(
   1421         block_on(engine.sign_prepared(push)),
   1422         Err(Error::StorageFailed)
   1423     );
   1424 
   1425     for (byte, nth) in [(91, 1), (92, 2)] {
   1426         let storage = Arc::new(FaultStorage::new(byte));
   1427         let sink = Arc::new(ScriptedSink::new([DeliveryBehavior::Outcomes(vec![
   1428             DeliveryOutcome::accepted(),
   1429         ])]));
   1430         let engine = fault_engine(
   1431             storage.clone(),
   1432             Arc::new(MockSigner::new(SignBehavior::Success {
   1433                 completed_at_unix_ms: 1_800_000_200_500,
   1434             })),
   1435             sink,
   1436         );
   1437         let push = request(byte, "wss://fault.example");
   1438         execute_to_admitted(&engine, &push);
   1439         storage.fault_nth_plan(nth);
   1440         assert_eq!(
   1441             block_on(engine.deliver_push(push.operation_id())),
   1442             Err(Error::StorageFailed)
   1443         );
   1444     }
   1445 }
   1446 
   1447 #[test]
   1448 fn signing_claim_recovery_respects_exact_and_non_replayable_capabilities() {
   1449     let exact_storage = Arc::new(MemoryStorage::new(
   1450         SourceGeneration::new([8; 32]).expect("generation"),
   1451     ));
   1452     let exact_pending = Arc::new(MockSigner::with_replay(
   1453         SignBehavior::Pending,
   1454         ReplayCapability::ExactReplayByRequestId,
   1455     ));
   1456     let exact_engine = Engine::builder(
   1457         exact_storage.clone(),
   1458         Arc::new(TestClock(AtomicU64::new(1_800_000_200_000))),
   1459         Arc::new(TestIds(AtomicU64::new(100))),
   1460         DeadlinePolicy::new(10_000, 10_000, 10_000).expect("deadlines"),
   1461     )
   1462     .sink(Arc::new(MockSink))
   1463     .signer(exact_pending.clone())
   1464     .build()
   1465     .expect("exact engine");
   1466     let exact_push = request(8, "wss://relay.example");
   1467     let mut exact_future = Box::pin(exact_engine.sign_prepared(exact_push.clone())).fuse();
   1468     let mut context = std::task::Context::from_waker(noop_waker_ref());
   1469     assert!(exact_future.poll_unpin(&mut context).is_pending());
   1470     drop(exact_future);
   1471     let claimed = block_on(exact_engine.push_status(exact_push.operation_id()))
   1472         .expect("status")
   1473         .expect("claimed status");
   1474     assert!(claimed.artifact().signing_claim().is_some());
   1475 
   1476     let exact_success = Arc::new(MockSigner::with_replay(
   1477         SignBehavior::Success {
   1478             completed_at_unix_ms: 1_800_000_220_000,
   1479         },
   1480         ReplayCapability::ExactReplayByRequestId,
   1481     ));
   1482     let exact_recovery = Engine::builder(
   1483         exact_storage,
   1484         Arc::new(TestClock(AtomicU64::new(1_800_000_211_000))),
   1485         Arc::new(TestIds(AtomicU64::new(110))),
   1486         DeadlinePolicy::new(10_000, 10_000, 10_000).expect("deadlines"),
   1487     )
   1488     .sink(Arc::new(MockSink))
   1489     .signer(exact_success.clone())
   1490     .build()
   1491     .expect("recovery engine");
   1492     let recovered =
   1493         block_on(exact_recovery.sign_prepared(exact_push.clone())).expect("exact replay recovery");
   1494     assert_eq!(
   1495         recovered.artifact().signing_state(),
   1496         radroots_storage::authored::SigningState::Signed
   1497     );
   1498     assert_eq!(exact_success.calls.load(Ordering::Relaxed), 1);
   1499     let replay = block_on(exact_recovery.sign_prepared(exact_push)).expect("signed replay");
   1500     assert!(replay.is_replay());
   1501     assert_eq!(exact_success.calls.load(Ordering::Relaxed), 1);
   1502 
   1503     let unsafe_storage = Arc::new(MemoryStorage::new(
   1504         SourceGeneration::new([9; 32]).expect("generation"),
   1505     ));
   1506     let unsafe_pending = Arc::new(MockSigner::with_replay(
   1507         SignBehavior::Pending,
   1508         ReplayCapability::NonReplayable,
   1509     ));
   1510     let unsafe_engine = Engine::builder(
   1511         unsafe_storage.clone(),
   1512         Arc::new(TestClock(AtomicU64::new(1_800_000_200_000))),
   1513         Arc::new(TestIds(AtomicU64::new(120))),
   1514         DeadlinePolicy::new(10_000, 10_000, 10_000).expect("deadlines"),
   1515     )
   1516     .sink(Arc::new(MockSink))
   1517     .signer(unsafe_pending)
   1518     .build()
   1519     .expect("unsafe engine");
   1520     let unsafe_push = request(9, "wss://relay.example");
   1521     let mut unsafe_future = Box::pin(unsafe_engine.sign_prepared(unsafe_push.clone())).fuse();
   1522     let mut context = std::task::Context::from_waker(noop_waker_ref());
   1523     assert!(unsafe_future.poll_unpin(&mut context).is_pending());
   1524     drop(unsafe_future);
   1525 
   1526     let unsafe_success = Arc::new(MockSigner::with_replay(
   1527         SignBehavior::Success {
   1528             completed_at_unix_ms: 1_800_000_220_000,
   1529         },
   1530         ReplayCapability::NonReplayable,
   1531     ));
   1532     let unsafe_recovery = Engine::builder(
   1533         unsafe_storage,
   1534         Arc::new(TestClock(AtomicU64::new(1_800_000_211_000))),
   1535         Arc::new(TestIds(AtomicU64::new(130))),
   1536         DeadlinePolicy::new(10_000, 10_000, 10_000).expect("deadlines"),
   1537     )
   1538     .sink(Arc::new(MockSink))
   1539     .signer(unsafe_success.clone())
   1540     .build()
   1541     .expect("unsafe recovery engine");
   1542     assert_eq!(
   1543         block_on(unsafe_recovery.sign_prepared(unsafe_push.clone())),
   1544         Err(Error::SigningIndeterminate)
   1545     );
   1546     assert_eq!(unsafe_success.calls.load(Ordering::Relaxed), 0);
   1547     let unsafe_status = block_on(unsafe_recovery.push_status(unsafe_push.operation_id()))
   1548         .expect("unsafe status")
   1549         .expect("unsafe operation");
   1550     assert_eq!(
   1551         unsafe_status.artifact().signing_state(),
   1552         radroots_storage::authored::SigningState::Indeterminate
   1553     );
   1554 }
   1555 
   1556 #[test]
   1557 fn signing_failures_persist_retry_indeterminate_terminal_and_cancelled_states() {
   1558     let exact = Arc::new(MockSigner::with_replay(
   1559         SignBehavior::Uncertain(SigningErrorKind::SignerTimeout),
   1560         ReplayCapability::ExactReplayByRequestId,
   1561     ));
   1562     let (engine, storage) = setup_engine(exact);
   1563     let retryable = request(10, "wss://relay.example");
   1564     assert_eq!(
   1565         block_on(engine.sign_prepared(retryable.clone())),
   1566         Err(Error::SignerFailed)
   1567     );
   1568     let retryable_status = block_on(engine.push_status(retryable.operation_id()))
   1569         .expect("retryable status")
   1570         .expect("retryable operation");
   1571     assert_eq!(
   1572         retryable_status.artifact().signing_state(),
   1573         radroots_storage::authored::SigningState::Retryable
   1574     );
   1575     let retry_at = retryable_status
   1576         .artifact()
   1577         .signing_retry()
   1578         .expect("retry schedule")
   1579         .not_before_unix_ms();
   1580     let success = Arc::new(MockSigner::with_replay(
   1581         SignBehavior::Success {
   1582             completed_at_unix_ms: retry_at + 2,
   1583         },
   1584         ReplayCapability::ExactReplayByRequestId,
   1585     ));
   1586     let recovery = Engine::builder(
   1587         storage,
   1588         Arc::new(TestClock(AtomicU64::new(retry_at + 1))),
   1589         Arc::new(TestIds(AtomicU64::new(140))),
   1590         DeadlinePolicy::new(10_000, 10_000, 10_000).expect("deadlines"),
   1591     )
   1592     .sink(Arc::new(MockSink))
   1593     .signer(success.clone())
   1594     .build()
   1595     .expect("recovery engine");
   1596     assert_eq!(
   1597         block_on(recovery.sign_prepared(retryable))
   1598             .expect("retry exact request")
   1599             .artifact()
   1600             .signing_state(),
   1601         radroots_storage::authored::SigningState::Signed
   1602     );
   1603     assert_eq!(success.calls.load(Ordering::Relaxed), 1);
   1604 
   1605     let non_replayable = Arc::new(MockSigner::with_replay(
   1606         SignBehavior::Uncertain(SigningErrorKind::SignerTimeout),
   1607         ReplayCapability::NonReplayable,
   1608     ));
   1609     let (engine, _) = setup_engine(non_replayable);
   1610     let uncertain = request(11, "wss://relay.example");
   1611     assert_eq!(
   1612         block_on(engine.sign_prepared(uncertain.clone())),
   1613         Err(Error::SigningIndeterminate)
   1614     );
   1615     assert_eq!(
   1616         block_on(engine.push_status(uncertain.operation_id()))
   1617             .expect("indeterminate status")
   1618             .expect("indeterminate operation")
   1619             .artifact()
   1620             .signing_state(),
   1621         radroots_storage::authored::SigningState::Indeterminate
   1622     );
   1623     assert_eq!(
   1624         block_on(engine.sign_prepared(uncertain)),
   1625         Err(Error::SigningIndeterminate)
   1626     );
   1627 
   1628     let cancelled = Arc::new(MockSigner::new(SignBehavior::Error(
   1629         SigningErrorKind::SignerCancelled,
   1630     )));
   1631     let (engine, _) = setup_engine(cancelled);
   1632     let cancelled_push = request(12, "wss://relay.example");
   1633     assert_eq!(
   1634         block_on(engine.sign_prepared(cancelled_push.clone())),
   1635         Err(Error::SigningCancelled)
   1636     );
   1637     assert_eq!(
   1638         block_on(engine.push_status(cancelled_push.operation_id()))
   1639             .expect("cancelled status")
   1640             .expect("cancelled operation")
   1641             .artifact()
   1642             .signing_state(),
   1643         radroots_storage::authored::SigningState::Cancelled
   1644     );
   1645     assert_eq!(
   1646         block_on(engine.sign_prepared(cancelled_push)),
   1647         Err(Error::SignerFailed)
   1648     );
   1649 
   1650     let terminal = Arc::new(MockSigner::new(SignBehavior::Error(
   1651         SigningErrorKind::SignerRejected,
   1652     )));
   1653     let (engine, _) = setup_engine(terminal);
   1654     let terminal_push = request(13, "wss://relay.example");
   1655     assert_eq!(
   1656         block_on(engine.sign_prepared(terminal_push.clone())),
   1657         Err(Error::SignerFailed)
   1658     );
   1659     assert_eq!(
   1660         block_on(engine.push_status(terminal_push.operation_id()))
   1661             .expect("terminal status")
   1662             .expect("terminal operation")
   1663             .artifact()
   1664             .signing_state(),
   1665         radroots_storage::authored::SigningState::FailedTerminal
   1666     );
   1667     assert_eq!(
   1668         block_on(engine.sign_prepared(terminal_push)),
   1669         Err(Error::SignerFailed)
   1670     );
   1671 
   1672     let deadline = Arc::new(MockSigner::new(SignBehavior::Error(
   1673         SigningErrorKind::DeadlineExceeded,
   1674     )));
   1675     let (engine, _) = setup_engine(deadline);
   1676     let deadline_push = request(14, "wss://relay.example");
   1677     assert_eq!(
   1678         block_on(engine.sign_prepared(deadline_push)),
   1679         Err(Error::SignerDeadlineExceeded)
   1680     );
   1681 }
   1682 
   1683 #[tokio::test]
   1684 async fn sqlite_signing_and_admission_recover_across_every_reopen_boundary() {
   1685     let directory = tempfile::tempdir().expect("database directory");
   1686     let paths = Paths::from_directory(directory.path()).expect("paths");
   1687     let push = request(7, "wss://relay.example");
   1688     let original = {
   1689         let store = Arc::new(
   1690             SqliteStorage::open(
   1691                 OpenOptions::new(paths.clone(), OpenMode::Create)
   1692                     .with_source_generation(SourceGeneration::new([7; 32]).expect("generation"), 1)
   1693                     .expect("source generation"),
   1694             )
   1695             .await
   1696             .expect("open SQLite"),
   1697         );
   1698         let capability: Arc<dyn SyncStorage> = store;
   1699         let signer = Arc::new(MockSigner::new(SignBehavior::Pending));
   1700         let engine = Engine::builder(
   1701             capability,
   1702             Arc::new(TestClock(AtomicU64::new(1_800_000_200_000))),
   1703             Arc::new(TestIds(AtomicU64::new(80))),
   1704             DeadlinePolicy::new(10_000, 10_000, 10_000).expect("deadlines"),
   1705         )
   1706         .sink(Arc::new(MockSink))
   1707         .signer(signer.clone())
   1708         .build()
   1709         .expect("engine");
   1710         let prepared = engine.prepare_push(push.clone()).await.expect("prepare");
   1711         assert!(!prepared.is_replay());
   1712         assert_eq!(signer.calls.load(Ordering::Relaxed), 0);
   1713         prepared
   1714     };
   1715 
   1716     {
   1717         let store = Arc::new(
   1718             SqliteStorage::open(OpenOptions::new(paths.clone(), OpenMode::ReadWriteExisting))
   1719                 .await
   1720                 .expect("reopen for signing"),
   1721         );
   1722         let capability: Arc<dyn SyncStorage> = store;
   1723         let signer = Arc::new(MockSigner::new(SignBehavior::Success {
   1724             completed_at_unix_ms: 1_800_000_210_500,
   1725         }));
   1726         let engine = Engine::builder(
   1727             capability,
   1728             Arc::new(TestClock(AtomicU64::new(1_800_000_210_000))),
   1729             Arc::new(TestIds(AtomicU64::new(90))),
   1730             DeadlinePolicy::new(10_000, 10_000, 10_000).expect("deadlines"),
   1731         )
   1732         .sink(Arc::new(MockSink))
   1733         .signer(signer.clone())
   1734         .build()
   1735         .expect("engine");
   1736         let status = engine
   1737             .push_status(push.operation_id())
   1738             .await
   1739             .expect("status")
   1740             .expect("durable preparation");
   1741         assert_eq!(status.operation(), original.operation());
   1742         assert_eq!(status.artifact(), original.artifact());
   1743         assert_eq!(status.delivery_plan(), original.delivery_plan());
   1744         let replay = engine
   1745             .prepare_push(push.clone())
   1746             .await
   1747             .expect("prepare replay");
   1748         assert!(replay.is_replay());
   1749         assert_eq!(signer.calls.load(Ordering::Relaxed), 0);
   1750         let signed = engine
   1751             .sign_prepared(push.clone())
   1752             .await
   1753             .expect("sign after reopen");
   1754         assert_eq!(
   1755             signed.artifact().signing_state(),
   1756             radroots_storage::authored::SigningState::Signed
   1757         );
   1758         assert_eq!(
   1759             signed.artifact().admission_state(),
   1760             radroots_storage::authored::AdmissionState::Pending
   1761         );
   1762         assert!(signed.artifact().signed().is_some());
   1763         assert_eq!(signer.calls.load(Ordering::Relaxed), 1);
   1764         assert_eq!(
   1765             engine.deliver_push(push.operation_id()).await,
   1766             Err(Error::AdmissionFailed)
   1767         );
   1768     }
   1769 
   1770     {
   1771         let store = Arc::new(
   1772             SqliteStorage::open(OpenOptions::new(paths.clone(), OpenMode::ReadWriteExisting))
   1773                 .await
   1774                 .expect("reopen for admission"),
   1775         );
   1776         let capability: Arc<dyn SyncStorage> = store.clone();
   1777         let signer = Arc::new(MockSigner::new(SignBehavior::Pending));
   1778         let engine = Engine::builder(
   1779             capability,
   1780             Arc::new(TestClock(AtomicU64::new(1_800_000_211_000))),
   1781             Arc::new(TestIds(AtomicU64::new(100))),
   1782             DeadlinePolicy::new(10_000, 10_000, 10_000).expect("deadlines"),
   1783         )
   1784         .sink(Arc::new(MockSink))
   1785         .signer(signer.clone())
   1786         .build()
   1787         .expect("engine");
   1788         let admitted = engine
   1789             .admit_signed(push.operation_id())
   1790             .await
   1791             .expect("admit after reopen");
   1792         assert!(!admitted.is_replay());
   1793         assert!(admitted.artifact().admission_state().is_admitted());
   1794         let replay = engine
   1795             .admit_signed(push.operation_id())
   1796             .await
   1797             .expect("admission replay");
   1798         assert!(replay.is_replay());
   1799         assert_eq!(signer.calls.load(Ordering::Relaxed), 0);
   1800         let admitted_events = store
   1801             .query_raw(EventQuery::all(
   1802                 EventQueryBounds::first(10).expect("bounds"),
   1803             ))
   1804             .await
   1805             .expect("admitted events");
   1806         assert_eq!(admitted_events.items().len(), 1);
   1807         assert_eq!(
   1808             admitted_events.items()[0].stage(),
   1809             radroots_storage::event::AdmissionStage::Visible
   1810         );
   1811     }
   1812 
   1813     let store = Arc::new(
   1814         SqliteStorage::open(OpenOptions::new(paths, OpenMode::ReadOnly))
   1815             .await
   1816             .expect("final read-only reopen"),
   1817     );
   1818     let engine = Engine::builder(
   1819         store,
   1820         Arc::new(TestClock(AtomicU64::new(1_800_000_212_000))),
   1821         Arc::new(TestIds(AtomicU64::new(110))),
   1822         DeadlinePolicy::new(10_000, 10_000, 10_000).expect("deadlines"),
   1823     )
   1824     .sink(Arc::new(MockSink))
   1825     .build()
   1826     .expect("read-only engine");
   1827     let final_status = engine
   1828         .push_status(push.operation_id())
   1829         .await
   1830         .expect("final status")
   1831         .expect("durable operation");
   1832     assert_eq!(
   1833         final_status.artifact().signing_state(),
   1834         radroots_storage::authored::SigningState::Signed
   1835     );
   1836     assert!(final_status.artifact().admission_state().is_admitted());
   1837     assert!(final_status.delivery_plan().request().is_some());
   1838 }
   1839 
   1840 #[test]
   1841 fn authored_delivery_persists_retry_evidence_and_complete_settlement() {
   1842     let sink = Arc::new(ScriptedSink::new([
   1843         DeliveryBehavior::Outcomes(vec![
   1844             DeliveryOutcome::accepted(),
   1845             DeliveryOutcome::unavailable(),
   1846         ]),
   1847         DeliveryBehavior::Outcomes(vec![
   1848             DeliveryOutcome::accepted(),
   1849             DeliveryOutcome::accepted(),
   1850         ]),
   1851     ]));
   1852     let signer = Arc::new(MockSigner::new(SignBehavior::Success {
   1853         completed_at_unix_ms: 1_800_000_200_500,
   1854     }));
   1855     let ((engine, _), clock) = setup_engine_with_sink(signer, sink.clone());
   1856     let push = request_with_policy(
   1857         73,
   1858         &["wss://one.example", "wss://two.example"],
   1859         SatisfactionClass::Accepted,
   1860         TargetPolicy::all(),
   1861     );
   1862     execute_to_admitted(&engine, &push);
   1863     let admitted = block_on(engine.push_status(push.operation_id()))
   1864         .expect("admitted status")
   1865         .expect("admitted operation");
   1866     assert_eq!(
   1867         engine.retry_decision(admitted.delivery_plan(), 0),
   1868         Err(Error::ClockUnavailable)
   1869     );
   1870     assert_eq!(
   1871         engine.retry_decision(admitted.delivery_plan(), 1_800_000_200_100),
   1872         Ok(SyncRetryDecision::Ready)
   1873     );
   1874     assert_eq!(
   1875         engine.retry_decision(admitted.delivery_plan(), 1_800_000_300_000),
   1876         Ok(SyncRetryDecision::Expired)
   1877     );
   1878 
   1879     let first = block_on(engine.deliver_push(push.operation_id())).expect("first attempt");
   1880     assert!(!first.is_replay());
   1881     assert_eq!(first.plan().state(), AuthoredDeliveryState::Retryable);
   1882     assert_eq!(first.plan().attempt_count(), 1);
   1883     let retry_at = first
   1884         .plan()
   1885         .retry()
   1886         .expect("durable retry")
   1887         .not_before_unix_ms();
   1888     assert_eq!(
   1889         engine.retry_decision(first.plan(), retry_at - 1),
   1890         Ok(SyncRetryDecision::DeferredUntil { unix_ms: retry_at })
   1891     );
   1892     assert_eq!(
   1893         engine.retry_decision(first.plan(), retry_at),
   1894         Ok(SyncRetryDecision::Ready)
   1895     );
   1896     assert_eq!(
   1897         block_on(engine.deliver_push(push.operation_id())),
   1898         Err(Error::DeliveryDeferred)
   1899     );
   1900     let pending = block_on(engine.push_status(push.operation_id()))
   1901         .expect("pending status")
   1902         .expect("pending operation");
   1903     assert_eq!(pending.settlement().artifacts(), 1);
   1904     assert_eq!(pending.settlement().signed(), 1);
   1905     assert_eq!(pending.settlement().admitted(), 1);
   1906     assert_eq!(pending.settlement().delivery_plans(), 1);
   1907     assert_eq!(pending.settlement().delivery_retryable(), 1);
   1908     assert!(!pending.settlement().is_settled());
   1909 
   1910     clock.0.store(retry_at, Ordering::Relaxed);
   1911     let second = block_on(engine.deliver_push(push.operation_id())).expect("retry attempt");
   1912     assert_eq!(second.plan().state(), AuthoredDeliveryState::Satisfied);
   1913     assert_eq!(
   1914         engine.retry_decision(second.plan(), retry_at),
   1915         Ok(SyncRetryDecision::Satisfied)
   1916     );
   1917     assert_eq!(second.plan().attempt_count(), 2);
   1918     assert_eq!(second.plan().attempts().len(), 2);
   1919     let replay = block_on(engine.deliver_push(push.operation_id())).expect("terminal replay");
   1920     assert!(replay.is_replay());
   1921     assert_eq!(replay.plan(), second.plan());
   1922     let settled = block_on(engine.push_status(push.operation_id()))
   1923         .expect("settled status")
   1924         .expect("settled operation");
   1925     assert_eq!(settled.settlement().delivery_satisfied(), 1);
   1926     assert!(settled.settlement().is_settled());
   1927     assert!(!settled.settlement().has_failures());
   1928     assert!(settled.settlement().is_successful());
   1929     let requests = sink.requests.lock().expect("request log");
   1930     assert_eq!(requests.len(), 2);
   1931     assert_eq!(requests[0], requests[1]);
   1932 }
   1933 
   1934 #[test]
   1935 fn invalid_authored_delivery_receipt_terminalizes_without_hot_loop() {
   1936     let sink = Arc::new(ScriptedSink::new([DeliveryBehavior::MismatchedRequest]));
   1937     let signer = Arc::new(MockSigner::new(SignBehavior::Success {
   1938         completed_at_unix_ms: 1_800_000_200_500,
   1939     }));
   1940     let ((engine, _), _) = setup_engine_with_sink(signer, sink.clone());
   1941     let push = request(74, "wss://one.example");
   1942     execute_to_admitted(&engine, &push);
   1943 
   1944     let failed = block_on(engine.deliver_push(push.operation_id())).expect("durable failure");
   1945     assert_eq!(failed.plan().state(), AuthoredDeliveryState::FailedTerminal);
   1946     assert_eq!(
   1947         engine.retry_decision(failed.plan(), 1_800_000_200_100),
   1948         Ok(SyncRetryDecision::Exhausted)
   1949     );
   1950     assert_eq!(failed.plan().attempt_count(), 1);
   1951     assert_eq!(
   1952         failed.plan().last_failure().expect("failure").code(),
   1953         "invalid_transport_contract"
   1954     );
   1955     let replay = block_on(engine.deliver_push(push.operation_id())).expect("failure replay");
   1956     assert!(replay.is_replay());
   1957     assert_eq!(replay.plan(), failed.plan());
   1958     assert_eq!(sink.requests.lock().expect("request log").len(), 1);
   1959     let status = block_on(engine.push_status(push.operation_id()))
   1960         .expect("failed status")
   1961         .expect("failed operation");
   1962     assert_eq!(status.settlement().delivery_failed_terminal(), 1);
   1963     assert!(status.settlement().is_settled());
   1964     assert!(status.settlement().has_failures());
   1965     assert!(!status.settlement().is_successful());
   1966 }
   1967 
   1968 #[test]
   1969 fn sink_failures_deadlines_and_missing_capabilities_are_durable_and_bounded() {
   1970     for (byte, behavior, expected) in [
   1971         (
   1972             77,
   1973             DeliveryBehavior::Failure(Retryability::Retryable),
   1974             AuthoredDeliveryState::Retryable,
   1975         ),
   1976         (
   1977             78,
   1978             DeliveryBehavior::Failure(Retryability::Terminal),
   1979             AuthoredDeliveryState::FailedTerminal,
   1980         ),
   1981         (
   1982             79,
   1983             DeliveryBehavior::MismatchedFailure,
   1984             AuthoredDeliveryState::FailedTerminal,
   1985         ),
   1986     ] {
   1987         let sink = Arc::new(ScriptedSink::new([behavior]));
   1988         let signer = Arc::new(MockSigner::new(SignBehavior::Success {
   1989             completed_at_unix_ms: 1_800_000_200_500,
   1990         }));
   1991         let ((engine, _), _) = setup_engine_with_sink(signer, sink);
   1992         let push = request(byte, "wss://failure.example");
   1993         execute_to_admitted(&engine, &push);
   1994         let delivered =
   1995             block_on(engine.deliver_push(push.operation_id())).expect("durable failure");
   1996         assert_eq!(delivered.plan().state(), expected);
   1997     }
   1998 
   1999     let signer = Arc::new(MockSigner::new(SignBehavior::Success {
   2000         completed_at_unix_ms: 1_800_000_200_500,
   2001     }));
   2002     let ((engine, _), clock) =
   2003         setup_engine_with_sink(signer, Arc::new(ScriptedSink::new(std::iter::empty())));
   2004     let push = request(80, "wss://deadline.example");
   2005     execute_to_admitted(&engine, &push);
   2006     clock
   2007         .0
   2008         .store(push.delivery_deadline_unix_ms(), Ordering::Relaxed);
   2009     let expired = block_on(engine.deliver_push(push.operation_id())).expect("deadline evidence");
   2010     assert_eq!(
   2011         expired.plan().state(),
   2012         AuthoredDeliveryState::FailedTerminal
   2013     );
   2014     assert_eq!(
   2015         expired.plan().last_failure().expect("failure").code(),
   2016         "delivery_deadline_exceeded"
   2017     );
   2018 
   2019     let storage = Arc::new(MemoryStorage::new(
   2020         SourceGeneration::new([81; 32]).expect("generation"),
   2021     ));
   2022     let capability: Arc<dyn SyncStorage> = storage;
   2023     let no_signer = Engine::builder(
   2024         capability,
   2025         Arc::new(TestClock(AtomicU64::new(1_800_000_200_000))),
   2026         Arc::new(TestIds(AtomicU64::new(10))),
   2027         DeadlinePolicy::new(10_000, 10_000, 10_000).unwrap(),
   2028     )
   2029     .sink(Arc::new(MockSink))
   2030     .build()
   2031     .unwrap();
   2032     assert_eq!(
   2033         block_on(no_signer.sign_prepared(request(81, "wss://missing.example"))),
   2034         Err(Error::MissingSigner)
   2035     );
   2036 
   2037     let storage = Arc::new(MemoryStorage::new(
   2038         SourceGeneration::new([82; 32]).expect("generation"),
   2039     ));
   2040     let capability: Arc<dyn SyncStorage> = storage;
   2041     let no_sink = Engine::builder(
   2042         capability,
   2043         Arc::new(TestClock(AtomicU64::new(1_800_000_200_000))),
   2044         Arc::new(TestIds(AtomicU64::new(10))),
   2045         DeadlinePolicy::new(10_000, 10_000, 10_000).unwrap(),
   2046     )
   2047     .source(Arc::new(MockSource))
   2048     .build()
   2049     .unwrap();
   2050     assert_eq!(
   2051         block_on(no_sink.deliver_push(SyncId::new([82; 16]).unwrap())),
   2052         Err(Error::MissingSink)
   2053     );
   2054 
   2055     let signer = Arc::new(MockSigner::new(SignBehavior::Success {
   2056         completed_at_unix_ms: 1_800_000_200_500,
   2057     }));
   2058     let (engine, _) = setup_engine(signer);
   2059     let unsigned = request(83, "wss://unsigned.example");
   2060     block_on(engine.prepare_push(unsigned.clone())).unwrap();
   2061     assert_eq!(
   2062         block_on(engine.admit_signed(unsigned.operation_id())),
   2063         Err(Error::InvalidSignerOutput)
   2064     );
   2065     assert_eq!(
   2066         block_on(engine.deliver_push(unsigned.operation_id())),
   2067         Err(Error::InvalidSignerOutput)
   2068     );
   2069 }
   2070 
   2071 #[test]
   2072 fn active_authored_claims_and_admission_storage_failures_are_fenced() {
   2073     let storage = Arc::new(FaultStorage::new(93));
   2074     let engine = fault_engine(
   2075         storage.clone(),
   2076         Arc::new(MockSigner::new(SignBehavior::Success {
   2077             completed_at_unix_ms: 1_800_000_200_500,
   2078         })),
   2079         Arc::new(MockSink),
   2080     );
   2081     let push = request(93, "wss://claim.example");
   2082     block_on(engine.prepare_push(push.clone())).unwrap();
   2083     let artifact = block_on(engine.push_status(push.operation_id()))
   2084         .unwrap()
   2085         .unwrap()
   2086         .artifact()
   2087         .clone();
   2088     let claim = radroots_storage::authored::WorkClaim::new(
   2089         [93; 16],
   2090         "external-signer",
   2091         std::num::NonZeroU64::MIN,
   2092         1_800_000_200_000,
   2093         1_800_000_210_000,
   2094         artifact.revision(),
   2095     )
   2096     .unwrap();
   2097     block_on(AuthoredAtomicStorage::execute_authored(
   2098         storage.as_ref(),
   2099         radroots_storage::authored_atomic::AuthoredAtomicCommand::Claim(
   2100             radroots_storage::authored_atomic::ClaimAuthoredWork::new(
   2101                 radroots_storage::authored_atomic::ClaimAuthoredTarget::ArtifactSigning(
   2102                     artifact.artifact_id(),
   2103                 ),
   2104                 claim,
   2105             ),
   2106         ),
   2107     ))
   2108     .unwrap();
   2109     assert_eq!(
   2110         block_on(engine.sign_prepared(push)),
   2111         Err(Error::WorkClaimConflict)
   2112     );
   2113 
   2114     let storage = Arc::new(FaultStorage::new(94));
   2115     let engine = fault_engine(
   2116         storage.clone(),
   2117         Arc::new(MockSigner::new(SignBehavior::Success {
   2118             completed_at_unix_ms: 1_800_000_200_500,
   2119         })),
   2120         Arc::new(MockSink),
   2121     );
   2122     let push = request(94, "wss://claim.example");
   2123     block_on(engine.sign_prepared(push.clone())).unwrap();
   2124     let artifact = block_on(engine.push_status(push.operation_id()))
   2125         .unwrap()
   2126         .unwrap()
   2127         .artifact()
   2128         .clone();
   2129     let claim = radroots_storage::authored::WorkClaim::new(
   2130         [94; 16],
   2131         "external-admission",
   2132         std::num::NonZeroU64::MIN,
   2133         1_800_000_200_500,
   2134         1_800_000_210_500,
   2135         artifact.revision(),
   2136     )
   2137     .unwrap();
   2138     block_on(AuthoredAtomicStorage::execute_authored(
   2139         storage.as_ref(),
   2140         radroots_storage::authored_atomic::AuthoredAtomicCommand::Claim(
   2141             radroots_storage::authored_atomic::ClaimAuthoredWork::new(
   2142                 radroots_storage::authored_atomic::ClaimAuthoredTarget::ArtifactAdmission(
   2143                     artifact.artifact_id(),
   2144                 ),
   2145                 claim,
   2146             ),
   2147         ),
   2148     ))
   2149     .unwrap();
   2150     assert_eq!(
   2151         block_on(engine.admit_signed(push.operation_id())),
   2152         Err(Error::WorkClaimConflict)
   2153     );
   2154 
   2155     for (byte, fault, expected) in [
   2156         (95, 1, radroots_storage::authored::AdmissionState::Rejected),
   2157         (96, 2, radroots_storage::authored::AdmissionState::Retryable),
   2158     ] {
   2159         let storage = Arc::new(FaultStorage::new(byte));
   2160         let engine = fault_engine(
   2161             storage.clone(),
   2162             Arc::new(MockSigner::new(SignBehavior::Success {
   2163                 completed_at_unix_ms: 1_800_000_200_500,
   2164             })),
   2165             Arc::new(MockSink),
   2166         );
   2167         let push = request(byte, "wss://admission-error.example");
   2168         block_on(engine.sign_prepared(push.clone())).unwrap();
   2169         storage.fail_admission_with(fault);
   2170         assert_eq!(
   2171             block_on(engine.admit_signed(push.operation_id())),
   2172             Err(Error::AdmissionFailed),
   2173             "admission fault {fault}"
   2174         );
   2175         assert_eq!(
   2176             block_on(engine.push_status(push.operation_id()))
   2177                 .unwrap()
   2178                 .unwrap()
   2179                 .artifact()
   2180                 .admission_state(),
   2181             expected
   2182         );
   2183     }
   2184 }
   2185 
   2186 #[test]
   2187 fn authored_delivery_claims_fence_concurrent_and_stale_workers() {
   2188     let pending_sink = Arc::new(ScriptedSink::new([DeliveryBehavior::Pending]));
   2189     let signer = Arc::new(MockSigner::new(SignBehavior::Success {
   2190         completed_at_unix_ms: 1_800_000_200_500,
   2191     }));
   2192     let ((engine, storage), clock) = setup_engine_with_sink(signer, pending_sink.clone());
   2193     let push = request(75, "wss://one.example");
   2194     execute_to_admitted(&engine, &push);
   2195 
   2196     let mut pending = Box::pin(engine.deliver_push(push.operation_id())).fuse();
   2197     let mut context = std::task::Context::from_waker(noop_waker_ref());
   2198     assert!(pending.poll_unpin(&mut context).is_pending());
   2199     let claimed = block_on(engine.push_status(push.operation_id()))
   2200         .expect("claimed status")
   2201         .expect("claimed operation");
   2202     let expires_at = claimed
   2203         .delivery_plan()
   2204         .claim_evidence()
   2205         .expect("delivery claim")
   2206         .expires_at_unix_ms();
   2207     assert_eq!(
   2208         engine.retry_decision(claimed.delivery_plan(), expires_at - 1),
   2209         Ok(SyncRetryDecision::InFlightUntil {
   2210             unix_ms: expires_at,
   2211         })
   2212     );
   2213     assert_eq!(
   2214         engine.retry_decision(claimed.delivery_plan(), expires_at),
   2215         Ok(SyncRetryDecision::Ready)
   2216     );
   2217     assert_eq!(
   2218         block_on(engine.deliver_push(push.operation_id())),
   2219         Err(Error::WorkClaimConflict)
   2220     );
   2221     drop(pending);
   2222 
   2223     clock.0.store(expires_at + 1, Ordering::Relaxed);
   2224     let recovery_sink = Arc::new(ScriptedSink::new([DeliveryBehavior::Outcomes(vec![
   2225         DeliveryOutcome::accepted(),
   2226     ])]));
   2227     let capability: Arc<dyn SyncStorage> = storage;
   2228     let recovery = Engine::builder(
   2229         capability,
   2230         clock,
   2231         Arc::new(TestIds(AtomicU64::new(220))),
   2232         DeadlinePolicy::new(10_000, 10_000, 10_000).expect("deadlines"),
   2233     )
   2234     .sink(recovery_sink.clone())
   2235     .build()
   2236     .expect("recovery engine");
   2237     let recovered = block_on(recovery.deliver_push(push.operation_id())).expect("stale recovery");
   2238     assert_eq!(recovered.plan().state(), AuthoredDeliveryState::Satisfied);
   2239     assert_eq!(
   2240         recovered.plan().claim_evidence(),
   2241         None,
   2242         "settlement clears the claim"
   2243     );
   2244     assert_eq!(recovery_sink.requests.lock().expect("request log").len(), 1);
   2245 }
   2246 
   2247 #[tokio::test]
   2248 async fn sqlite_authored_delivery_retry_survives_reopen() {
   2249     let directory = tempfile::tempdir().expect("database directory");
   2250     let paths = Paths::from_directory(directory.path()).expect("paths");
   2251     let push = request_with_policy(
   2252         76,
   2253         &["wss://one.example", "wss://two.example"],
   2254         SatisfactionClass::Accepted,
   2255         TargetPolicy::all(),
   2256     );
   2257     let retry_at = {
   2258         let store = Arc::new(
   2259             SqliteStorage::open(
   2260                 OpenOptions::new(paths.clone(), OpenMode::Create)
   2261                     .with_source_generation(SourceGeneration::new([76; 32]).expect("generation"), 1)
   2262                     .expect("source generation"),
   2263             )
   2264             .await
   2265             .expect("open SQLite"),
   2266         );
   2267         let sink = Arc::new(ScriptedSink::new([DeliveryBehavior::Outcomes(vec![
   2268             DeliveryOutcome::accepted(),
   2269             DeliveryOutcome::unavailable(),
   2270         ])]));
   2271         let engine = Engine::builder(
   2272             store,
   2273             Arc::new(TestClock(AtomicU64::new(1_800_000_200_000))),
   2274             Arc::new(TestIds(AtomicU64::new(230))),
   2275             DeadlinePolicy::new(10_000, 10_000, 10_000).expect("deadlines"),
   2276         )
   2277         .sink(sink)
   2278         .signer(Arc::new(MockSigner::new(SignBehavior::Success {
   2279             completed_at_unix_ms: 1_800_000_200_500,
   2280         })))
   2281         .build()
   2282         .expect("engine");
   2283         engine
   2284             .sign_prepared(push.clone())
   2285             .await
   2286             .expect("sign prepared push");
   2287         engine
   2288             .admit_signed(push.operation_id())
   2289             .await
   2290             .expect("admit signed push");
   2291         let first = engine
   2292             .deliver_push(push.operation_id())
   2293             .await
   2294             .expect("first attempt");
   2295         assert_eq!(first.plan().attempt_count(), 1);
   2296         first
   2297             .plan()
   2298             .retry()
   2299             .expect("retry schedule")
   2300             .not_before_unix_ms()
   2301     };
   2302 
   2303     let store = Arc::new(
   2304         SqliteStorage::open(OpenOptions::new(paths, OpenMode::ReadWriteExisting))
   2305             .await
   2306             .expect("reopen SQLite"),
   2307     );
   2308     let sink = Arc::new(ScriptedSink::new([DeliveryBehavior::Outcomes(vec![
   2309         DeliveryOutcome::accepted(),
   2310         DeliveryOutcome::accepted(),
   2311     ])]));
   2312     let clock = Arc::new(TestClock(AtomicU64::new(retry_at - 1)));
   2313     let engine = Engine::builder(
   2314         store,
   2315         clock.clone(),
   2316         Arc::new(TestIds(AtomicU64::new(240))),
   2317         DeadlinePolicy::new(10_000, 10_000, 10_000).expect("deadlines"),
   2318     )
   2319     .sink(sink.clone())
   2320     .build()
   2321     .expect("reopen engine");
   2322     let reopened = engine
   2323         .push_status(push.operation_id())
   2324         .await
   2325         .expect("reopened status")
   2326         .expect("reopened operation");
   2327     assert_eq!(reopened.delivery_plan().attempt_count(), 1);
   2328     assert_eq!(
   2329         reopened.delivery_plan().state(),
   2330         AuthoredDeliveryState::Retryable
   2331     );
   2332     assert_eq!(
   2333         engine.deliver_push(push.operation_id()).await,
   2334         Err(Error::DeliveryDeferred)
   2335     );
   2336     clock.0.store(retry_at, Ordering::Relaxed);
   2337     let completed = engine
   2338         .deliver_push(push.operation_id())
   2339         .await
   2340         .expect("retry after reopen");
   2341     assert_eq!(completed.plan().attempt_count(), 2);
   2342     assert_eq!(completed.plan().state(), AuthoredDeliveryState::Satisfied);
   2343     assert_eq!(sink.requests.lock().expect("request log").len(), 1);
   2344 }
   2345 
   2346 #[test]
   2347 fn captured_preparation_matches_engine_without_changing_ordinary_replay_identity() {
   2348     use radroots_storage::authored_atomic::AuthoredAtomicCommand;
   2349     let signer = Arc::new(MockSigner::new(SignBehavior::Pending));
   2350     let (engine, _) = setup_engine(signer.clone());
   2351     let push = request(91, "wss://captured.example");
   2352     let captured = push.authored_preparation(1_800_000_200_000).unwrap();
   2353     assert!(push.authored_preparation(0).is_err());
   2354     assert!(
   2355         block_on(engine.push_status(push.operation_id()))
   2356             .unwrap()
   2357             .is_none()
   2358     );
   2359     assert_eq!(signer.calls.load(Ordering::Relaxed), 0);
   2360     let later = push.authored_preparation(1_800_000_200_001).unwrap();
   2361     assert_ne!(
   2362         captured, later,
   2363         "the composite command compares captured time exactly"
   2364     );
   2365     let original_command = AuthoredAtomicCommand::Prepare(captured.clone());
   2366     let later_command = AuthoredAtomicCommand::Prepare(later);
   2367     assert_eq!(original_command.commit_id(), later_command.commit_id());
   2368     assert_eq!(original_command.digest(), later_command.digest());
   2369     let prepared = block_on(engine.prepare_push(push)).unwrap();
   2370     assert_eq!(prepared.operation(), captured.operation());
   2371     assert_eq!(
   2372         [prepared.artifact()],
   2373         captured.artifacts().iter().collect::<Vec<_>>().as_slice()
   2374     );
   2375     assert_eq!(
   2376         [prepared.delivery_plan()],
   2377         captured
   2378             .delivery_plans()
   2379             .iter()
   2380             .collect::<Vec<_>>()
   2381             .as_slice()
   2382     );
   2383     assert_eq!(signer.calls.load(Ordering::Relaxed), 0);
   2384 }