lib

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

sync.rs (24283B)


      1 //! Client-scoped access to the canonical synchronization engine.
      2 
      3 #[cfg(feature = "sync")]
      4 use std::sync::Arc;
      5 
      6 #[cfg(feature = "sync")]
      7 use radroots_storage::{authored_delivery::AuthoredDeliveryPlan, projection::ProjectionId};
      8 #[cfg(feature = "sync")]
      9 use radroots_sync::{
     10     Engine, PullReceipt, PullRequest, PushRequest, SyncStatus,
     11     ingest::{AdmissionPolicy, IngestBatchReceipt, IngestReceipt},
     12     policy::Error,
     13     projection::{Reducer, RefreshReceipt, RefreshRequest},
     14     push::{
     15         AdmissionRunReceipt, DeliveryExecutionReceipt, PushCancellationReceipt, PushPreparation,
     16         PushStatus, SigningRunReceipt,
     17     },
     18 };
     19 #[cfg(feature = "sync")]
     20 use radroots_transport::source::ObservedEvent;
     21 
     22 /// Explicit host policy for SDK-composed synchronization.
     23 ///
     24 /// Selecting this policy opts into the system clock and operating-system
     25 /// randomness for operation IDs. It creates no executor, timer, or worker.
     26 #[cfg(feature = "sync")]
     27 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
     28 pub struct HostPolicy {
     29     deadlines: radroots_sync::policy::DeadlinePolicy,
     30 }
     31 
     32 #[cfg(feature = "sync")]
     33 impl HostPolicy {
     34     /// Creates a bounded policy for pull, signing, and delivery calls.
     35     pub fn new(
     36         pull_timeout_ms: u64,
     37         sign_timeout_ms: u64,
     38         delivery_timeout_ms: u64,
     39     ) -> Result<Self, radroots_sync::policy::Error> {
     40         Ok(Self {
     41             deadlines: radroots_sync::policy::DeadlinePolicy::new(
     42                 pull_timeout_ms,
     43                 sign_timeout_ms,
     44                 delivery_timeout_ms,
     45             )?,
     46         })
     47     }
     48 
     49     /// Returns the ordinary bounded native-host policy.
     50     #[must_use]
     51     pub fn standard() -> Self {
     52         Self::new(30_000, 30_000, 30_000).expect("static host deadlines are valid")
     53     }
     54 
     55     pub(crate) fn composition(
     56         self,
     57     ) -> (
     58         Arc<dyn radroots_sync::policy::Clock>,
     59         Arc<dyn radroots_sync::policy::IdSource>,
     60         radroots_sync::policy::DeadlinePolicy,
     61     ) {
     62         (Arc::new(SystemClock), Arc::new(RandomIds), self.deadlines)
     63     }
     64 }
     65 
     66 #[cfg(feature = "sync")]
     67 impl Default for HostPolicy {
     68     fn default() -> Self {
     69         Self::standard()
     70     }
     71 }
     72 
     73 #[cfg(feature = "sync")]
     74 struct SystemClock;
     75 
     76 #[cfg(feature = "sync")]
     77 impl radroots_sync::policy::Clock for SystemClock {
     78     fn now_unix_ms(&self) -> Result<u64, radroots_sync::policy::Error> {
     79         std::time::SystemTime::now()
     80             .duration_since(std::time::UNIX_EPOCH)
     81             .ok()
     82             .and_then(|duration| u64::try_from(duration.as_millis()).ok())
     83             .filter(|value| *value != 0)
     84             .ok_or(radroots_sync::policy::Error::ClockUnavailable)
     85     }
     86 }
     87 
     88 #[cfg(feature = "sync")]
     89 struct RandomIds;
     90 
     91 #[cfg(feature = "sync")]
     92 impl radroots_sync::policy::IdSource for RandomIds {
     93     fn next_id(
     94         &self,
     95         _operation: radroots_sync::policy::OperationKind,
     96     ) -> Result<radroots_sync::policy::SyncId, radroots_sync::policy::Error> {
     97         radroots_sync::policy::SyncId::new(*uuid::Uuid::new_v4().as_bytes())
     98     }
     99 }
    100 
    101 /// Borrowed client operations over one explicitly composed sync engine.
    102 ///
    103 /// This type owns no scheduling, retries, status strings, outbox state, or
    104 /// projection state. Every method delegates once to the canonical engine and
    105 /// returns its native receipt or error.
    106 #[cfg(feature = "sync")]
    107 #[derive(Clone, Copy)]
    108 pub struct Operations<'a> {
    109     engine: &'a Engine,
    110 }
    111 
    112 #[cfg(feature = "sync")]
    113 impl<'a> Operations<'a> {
    114     pub(crate) const fn new(engine: &'a Engine) -> Self {
    115         Self { engine }
    116     }
    117 
    118     /// Runs one caller-bounded pull and canonical ingest sequence.
    119     pub async fn pull(
    120         &self,
    121         request: PullRequest,
    122         admission: &dyn AdmissionPolicy,
    123     ) -> Result<PullReceipt, Error> {
    124         self.engine.pull(request, admission).await
    125     }
    126 
    127     /// Verifies and atomically ingests one observed event.
    128     pub async fn ingest(
    129         &self,
    130         observed: ObservedEvent,
    131         admission: &dyn AdmissionPolicy,
    132     ) -> Result<IngestReceipt, Error> {
    133         self.engine.ingest(observed, admission).await
    134     }
    135 
    136     /// Ingests a bounded caller-owned batch while preserving partial outcomes.
    137     pub async fn ingest_batch(
    138         &self,
    139         observed: Vec<ObservedEvent>,
    140         admission: &dyn AdmissionPolicy,
    141     ) -> IngestBatchReceipt {
    142         self.engine.ingest_batch(observed, admission).await
    143     }
    144 
    145     /// Runs one bounded projection refresh through its owning reducer.
    146     pub async fn refresh_projection(
    147         &self,
    148         request: RefreshRequest,
    149         reducer: &dyn Reducer,
    150     ) -> Result<RefreshReceipt, Error> {
    151         self.engine.refresh_projection(request, reducer).await
    152     }
    153 
    154     /// Atomically persists one complete authored operation before any side effect.
    155     pub async fn prepare_push(&self, request: PushRequest) -> Result<PushPreparation, Error> {
    156         self.engine.prepare_push(request).await
    157     }
    158 
    159     /// Returns the complete durable state of one prepared push operation.
    160     pub async fn push_status(
    161         &self,
    162         operation_id: radroots_sync::policy::SyncId,
    163     ) -> Result<Option<PushStatus>, Error> {
    164         self.engine.push_status(operation_id).await
    165     }
    166 
    167     /// Cancels every remaining local phase while preserving durable evidence.
    168     pub async fn cancel_push(
    169         &self,
    170         operation_id: radroots_sync::policy::SyncId,
    171     ) -> Result<PushCancellationReceipt, Error> {
    172         self.engine.cancel_push(operation_id).await
    173     }
    174 
    175     /// Runs one bounded signing phase for an exactly prepared operation.
    176     pub async fn sign_prepared(&self, request: PushRequest) -> Result<SigningRunReceipt, Error> {
    177         self.engine.sign_prepared(request).await
    178     }
    179 
    180     /// Runs one bounded local-admission phase for a durably signed artifact.
    181     pub async fn admit_signed(
    182         &self,
    183         operation_id: radroots_sync::policy::SyncId,
    184     ) -> Result<AdmissionRunReceipt, Error> {
    185         self.engine.admit_signed(operation_id).await
    186     }
    187 
    188     /// Prepares, signs, and locally admits one authored operation.
    189     ///
    190     /// Preparation commits the complete recoverable intent before the signer
    191     /// can be invoked. Delivery remains an explicit caller-driven phase.
    192     pub async fn submit_push(&self, request: PushRequest) -> Result<PushStatus, Error> {
    193         let operation_id = request.operation_id();
    194         self.sign_prepared(request).await?;
    195         self.admit_signed(operation_id).await?;
    196         self.push_status(operation_id)
    197             .await?
    198             .ok_or(Error::StorageFailed)
    199     }
    200 
    201     /// Runs one bounded delivery attempt for one durable authored plan.
    202     pub async fn deliver_push(
    203         &self,
    204         operation_id: radroots_sync::policy::SyncId,
    205     ) -> Result<DeliveryExecutionReceipt, Error> {
    206         self.engine.deliver_push(operation_id).await
    207     }
    208 
    209     /// Attempts an exact subset while retaining the full durable request binding.
    210     pub async fn deliver_push_selected(
    211         &self,
    212         operation_id: radroots_sync::policy::SyncId,
    213         selected: radroots_transport::TargetSet,
    214     ) -> Result<DeliveryExecutionReceipt, Error> {
    215         self.engine
    216             .deliver_push_selected(operation_id, selected)
    217             .await
    218     }
    219 
    220     /// Returns the native passive sync status without starting recovery work.
    221     pub async fn status(&self, projections: &[ProjectionId]) -> Result<SyncStatus, Error> {
    222         self.engine.status(projections).await
    223     }
    224 
    225     /// Returns the native host scheduling decision for one durable plan.
    226     pub fn retry_decision(
    227         &self,
    228         plan: &AuthoredDeliveryPlan,
    229         now_unix_ms: u64,
    230     ) -> Result<radroots_protocol::runtime::v1::SyncRetryDecision, Error> {
    231         self.engine.retry_decision(plan, now_unix_ms)
    232     }
    233 }
    234 
    235 #[cfg(feature = "sync")]
    236 impl std::fmt::Debug for Operations<'_> {
    237     fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
    238         formatter
    239             .debug_struct("Operations")
    240             .field("engine", &"<borrowed canonical engine>")
    241             .finish()
    242     }
    243 }
    244 
    245 #[cfg(all(test, feature = "sync", feature = "memory"))]
    246 mod tests {
    247     #[cfg(feature = "local-signing")]
    248     mod selected;
    249     use std::sync::{
    250         Arc,
    251         atomic::{AtomicU8, Ordering},
    252     };
    253 
    254     use radroots_event::{
    255         GenericEventDraft, SignedEvent, contract::AuthorRole, wire::Nip01EventWire,
    256     };
    257     use radroots_event_codec::authoring::AuthoredEventPlan;
    258     use radroots_identity::PublicKey;
    259     use radroots_protocol::runtime::v1::SyncCapabilityState;
    260     use radroots_signing::{Actor, actor::ActorSource, request::CancellationPolicy};
    261     use radroots_storage::{
    262         event::{SourceGeneration, StoredVisibleEvent},
    263         journal::IdempotencyKey,
    264         memory::MemoryStorage,
    265         projection::{
    266             ProjectionGeneration, ProjectionId, RawSourceDigest, RebuildFailure, RebuildTicketId,
    267         },
    268     };
    269     use radroots_sync::{
    270         Engine, PullRequest, PushRequest,
    271         ingest::RegistryPolicy,
    272         policy::{Clock, DeadlinePolicy, Error, IdSource, OperationKind, SyncId, SyncStorage},
    273         projection::{Reducer, ReducerError, RefreshRequest, RefreshState},
    274         pull::PullTermination,
    275     };
    276     use radroots_transport::{
    277         Error as TransportError, EventSource, FetchPage, FetchRequest, SourceStatus, Target,
    278         TargetSet, TransportId,
    279         capability::{Availability, Maturity, SourceCapabilities},
    280         policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy},
    281         source::{EventProvenance, NextPage, ObservedEvent},
    282     };
    283 
    284     use crate::{ClientBuilder, error::ErrorKind, sync::HostPolicy};
    285 
    286     const PUBLIC_KEY: &str = "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df";
    287 
    288     struct FixedClock;
    289     struct SequenceIds(AtomicU8);
    290     struct CancelledSource;
    291 
    292     impl Clock for FixedClock {
    293         fn now_unix_ms(&self) -> Result<u64, Error> {
    294             Ok(1_700_000_000_000)
    295         }
    296     }
    297 
    298     impl IdSource for SequenceIds {
    299         fn next_id(&self, _operation: OperationKind) -> Result<SyncId, Error> {
    300             let next = self.0.fetch_add(1, Ordering::Relaxed);
    301             SyncId::new([next; 16])
    302         }
    303     }
    304 
    305     impl EventSource for CancelledSource {
    306         fn status(
    307             &self,
    308         ) -> radroots_transport::BoxFuture<'_, Result<SourceStatus, TransportError>> {
    309             Box::pin(async {
    310                 Ok(SourceStatus::new(
    311                     TransportId::NOSTR,
    312                     true,
    313                     Maturity::Stable,
    314                     Availability::Available,
    315                     SourceCapabilities::FETCH,
    316                     "ready",
    317                 ))
    318             })
    319         }
    320 
    321         fn fetch(
    322             &self,
    323             request: FetchRequest,
    324         ) -> radroots_transport::BoxFuture<'_, Result<FetchPage, TransportError>> {
    325             Box::pin(async move {
    326                 FetchPage::for_request(
    327                     &request,
    328                     Vec::new(),
    329                     Vec::new(),
    330                     NextPage::Cancelled { resume_from: None },
    331                 )
    332             })
    333         }
    334     }
    335 
    336     struct PartialThenCompleteSource(AtomicU8);
    337 
    338     impl EventSource for PartialThenCompleteSource {
    339         fn status(
    340             &self,
    341         ) -> radroots_transport::BoxFuture<'_, Result<SourceStatus, TransportError>> {
    342             Box::pin(async { unreachable!("explicit pull only") })
    343         }
    344 
    345         fn fetch(
    346             &self,
    347             request: FetchRequest,
    348         ) -> radroots_transport::BoxFuture<'_, Result<FetchPage, TransportError>> {
    349             Box::pin(async move {
    350                 use radroots_transport::{
    351                     outcome::{FetchTargetOutcome, FetchTargetState},
    352                     source::FetchCursor,
    353                 };
    354                 let page = self.0.fetch_add(1, Ordering::Relaxed);
    355                 assert!(page < 2);
    356                 let (state, next) = if page == 0 {
    357                     assert!(request.cursor().is_none());
    358                     (
    359                         FetchTargetState::Partial,
    360                         NextPage::Cursor(FetchCursor::parse("second").expect("cursor")),
    361                     )
    362                 } else {
    363                     assert_eq!(request.cursor().map(FetchCursor::as_str), Some("second"));
    364                     (FetchTargetState::Complete, NextPage::Complete)
    365                 };
    366                 let outcome = FetchTargetOutcome::new(
    367                     request.target_set().targets()[0].fingerprint().clone(),
    368                     state,
    369                 );
    370                 FetchPage::for_request(&request, vec![], vec![outcome], next)
    371             })
    372         }
    373     }
    374 
    375     #[tokio::test]
    376     async fn operations_preserve_cumulative_pull_evidence_after_a_later_complete_page() {
    377         let storage = Arc::new(MemoryStorage::new(
    378             SourceGeneration::new([3; 32]).expect("generation"),
    379         ));
    380         let source = Arc::new(PartialThenCompleteSource(AtomicU8::new(0)));
    381         let engine = Engine::builder(
    382             storage.clone(),
    383             Arc::new(FixedClock),
    384             Arc::new(SequenceIds(AtomicU8::new(1))),
    385             DeadlinePolicy::new(1_000, 1_000, 1_000).expect("deadlines"),
    386         )
    387         .source(source.clone())
    388         .build()
    389         .expect("engine");
    390         let client = ClientBuilder::new()
    391             .storage(storage)
    392             .sync_engine(engine)
    393             .build()
    394             .expect("client");
    395         let receipt = client
    396             .sync()
    397             .expect("open client")
    398             .expect("sync")
    399             .pull(
    400                 PullRequest::new(TargetSet::new(vec![target()]).expect("targets"), 1, 2)
    401                     .expect("request"),
    402                 &RegistryPolicy::visible(),
    403             )
    404             .await
    405             .expect("receipt");
    406         assert_eq!(source.0.load(Ordering::Relaxed), 2);
    407         assert_eq!(receipt.termination(), PullTermination::Complete);
    408         assert_eq!(
    409             receipt.target_outcomes()[0].state(),
    410             radroots_transport::outcome::FetchTargetState::Complete
    411         );
    412         let summary = &receipt.target_summaries().expect("measured")[0];
    413         assert_eq!(summary.pages_observed(), 2);
    414         assert_eq!(summary.incomplete_pages(), 1);
    415         assert_eq!(
    416             summary.last_incomplete(),
    417             Some(radroots_transport::outcome::FetchTargetState::Partial)
    418         );
    419         assert!(!summary.all_pages_complete());
    420     }
    421 
    422     struct EmptyReducer {
    423         id: ProjectionId,
    424         generation: ProjectionGeneration,
    425     }
    426 
    427     impl Reducer for EmptyReducer {
    428         fn projection_id(&self) -> &ProjectionId {
    429             &self.id
    430         }
    431 
    432         fn generation(&self) -> ProjectionGeneration {
    433             self.generation
    434         }
    435 
    436         fn begin_rebuild(
    437             &self,
    438             _ticket_id: RebuildTicketId,
    439             _source_generation: SourceGeneration,
    440             _source_digest: RawSourceDigest,
    441         ) -> Result<(), ReducerError> {
    442             Ok(())
    443         }
    444 
    445         fn reduce(
    446             &self,
    447             events: &[StoredVisibleEvent],
    448             prior_projected_rows: u64,
    449             _rebuild_ticket: Option<RebuildTicketId>,
    450         ) -> Result<u64, ReducerError> {
    451             assert!(events.is_empty());
    452             Ok(prior_projected_rows)
    453         }
    454 
    455         fn abort_rebuild(
    456             &self,
    457             _ticket_id: RebuildTicketId,
    458             _failure: RebuildFailure,
    459         ) -> Result<(), ReducerError> {
    460             Ok(())
    461         }
    462     }
    463 
    464     fn target() -> Target {
    465         Target::nostr_relay("wss://sync.example").expect("target")
    466     }
    467 
    468     fn engine(storage: Arc<MemoryStorage>) -> Engine {
    469         let capability: Arc<dyn SyncStorage> = storage;
    470         Engine::builder(
    471             capability,
    472             Arc::new(FixedClock),
    473             Arc::new(SequenceIds(AtomicU8::new(1))),
    474             DeadlinePolicy::new(1_000, 1_000, 1_000).expect("deadlines"),
    475         )
    476         .source(Arc::new(CancelledSource))
    477         .build()
    478         .expect("engine")
    479     }
    480 
    481     fn invalid_observation() -> ObservedEvent {
    482         let mut wire = Nip01EventWire {
    483             id: "0".repeat(64),
    484             pubkey: PUBLIC_KEY.to_owned(),
    485             created_at: 1_800_000_100,
    486             kind: 0,
    487             tags: vec![],
    488             content: "invalid signature fixture".to_owned(),
    489             sig: "42".repeat(64),
    490             extra: Default::default(),
    491         };
    492         wire.id = wire.computed_event_id().expect("event id").to_hex();
    493         let raw = format!(
    494             "{{\"id\":\"{}\",\"pubkey\":\"{}\",\"created_at\":{},\"kind\":0,\"tags\":[],\"content\":\"invalid signature fixture\",\"sig\":\"{}\"}}",
    495             wire.id, wire.pubkey, wire.created_at, wire.sig
    496         );
    497         let event = SignedEvent::from_wire_verified_id(wire, raw).expect("signed event");
    498         let target = target();
    499         let provenance = EventProvenance::new(
    500             TransportId::NOSTR,
    501             target.fingerprint().clone(),
    502             1_700_000_000_000,
    503         )
    504         .expect("provenance");
    505         ObservedEvent::new(event, provenance)
    506     }
    507 
    508     fn push_request() -> PushRequest {
    509         push_request_for_id(8)
    510     }
    511 
    512     fn push_request_for_id(id: u8) -> PushRequest {
    513         let actor = Actor::new(
    514             PublicKey::from_hex(PUBLIC_KEY).expect("public key"),
    515             ActorSource::ExplicitPublicKey,
    516             [AuthorRole::Any],
    517         )
    518         .expect("actor");
    519         let plan = AuthoredEventPlan::from_generic(
    520             GenericEventDraft::new(
    521                 "radroots.social.geochat.v1",
    522                 20_000,
    523                 1_700_000_000,
    524                 Vec::new(),
    525                 "content",
    526                 PUBLIC_KEY,
    527             )
    528             .expect("draft"),
    529         )
    530         .expect("authored plan");
    531         PushRequest::new(
    532             SyncId::new([id; 16]).expect("operation id"),
    533             IdempotencyKey::parse(format!("sdk-sync-wrapper-{id}")).expect("idempotency key"),
    534             actor,
    535             plan,
    536             TargetSet::new(vec![target()]).expect("targets"),
    537             SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::all()),
    538             1_700_000_001_000,
    539             CancellationPolicy::PreservePublishedRequest,
    540         )
    541         .expect("push request")
    542     }
    543 
    544     #[tokio::test]
    545     async fn operations_delegate_native_pull_ingest_projection_push_delivery_and_status() {
    546         let storage = Arc::new(MemoryStorage::new(
    547             SourceGeneration::new([3; 32]).expect("generation"),
    548         ));
    549         let client = ClientBuilder::new()
    550             .storage(storage.clone())
    551             .sync_engine(engine(storage))
    552             .build()
    553             .expect("client");
    554         let operations = client.sync().expect("open client").expect("sync");
    555         let policy = RegistryPolicy::verified();
    556 
    557         let pull = operations
    558             .pull(
    559                 PullRequest::new(TargetSet::new(vec![target()]).expect("targets"), 10, 1)
    560                     .expect("pull request"),
    561                 &policy,
    562             )
    563             .await
    564             .expect("pull");
    565         assert_eq!(pull.termination(), PullTermination::Cancelled);
    566         let summaries = pull.target_summaries().expect("measured target evidence");
    567         assert_eq!(summaries.len(), 1);
    568         assert_eq!(summaries[0].target(), target().fingerprint());
    569         assert_eq!(summaries[0].pages_observed(), 1);
    570         assert_eq!(summaries[0].missing_outcome_pages(), 1);
    571         assert!(!summaries[0].all_pages_complete());
    572 
    573         let ingest = operations
    574             .ingest_batch(vec![invalid_observation(), invalid_observation()], &policy)
    575             .await;
    576         assert_eq!(ingest.accepted(), 0);
    577         assert_eq!(ingest.rejected(), 2);
    578         assert!(
    579             ingest
    580                 .outcomes()
    581                 .iter()
    582                 .all(|outcome| matches!(outcome, Err(Error::VerificationFailed)))
    583         );
    584 
    585         let projection_id = ProjectionId::parse("sdk.sync.test").expect("projection id");
    586         let generation = ProjectionGeneration::new([5; 32]).expect("generation");
    587         let projection = operations
    588             .refresh_projection(
    589                 RefreshRequest::new(projection_id.clone(), generation, 10, 1)
    590                     .expect("refresh request"),
    591                 &EmptyReducer {
    592                     id: projection_id.clone(),
    593                     generation,
    594                 },
    595             )
    596             .await
    597             .expect("projection");
    598         assert_eq!(projection.state(), RefreshState::Complete);
    599 
    600         let push = push_request();
    601         let operation_id = push.operation_id();
    602         let preparation = operations
    603             .prepare_push(push.clone())
    604             .await
    605             .expect("durable preparation");
    606         assert!(!preparation.is_replay());
    607         assert_eq!(
    608             operations.sign_prepared(push).await,
    609             Err(Error::MissingSigner)
    610         );
    611         assert!(
    612             operations
    613                 .push_status(operation_id)
    614                 .await
    615                 .expect("push status")
    616                 .is_some()
    617         );
    618         assert_eq!(
    619             operations.deliver_push(operation_id).await,
    620             Err(Error::InvalidSignerOutput)
    621         );
    622 
    623         let status = operations
    624             .status(std::slice::from_ref(&projection_id))
    625             .await
    626             .expect("status");
    627         assert_eq!(status.source().state(), SyncCapabilityState::Available);
    628         assert_eq!(status.sink().state(), SyncCapabilityState::Unsupported);
    629         assert_eq!(status.projections().len(), 1);
    630 
    631         client.close().await.expect("close");
    632         assert_eq!(
    633             client.sync().expect_err("closed").kind(),
    634             ErrorKind::ClientClosed
    635         );
    636     }
    637 
    638     #[tokio::test]
    639     async fn host_policy_and_remaining_operation_wrappers_preserve_native_results() {
    640         assert!(HostPolicy::new(0, 1, 1).is_err());
    641         let policy = HostPolicy::default();
    642         assert_eq!(policy, HostPolicy::standard());
    643         let (clock, ids, deadlines) = policy.composition();
    644         assert!(clock.now_unix_ms().expect("system clock") > 0);
    645         let first = ids.next_id(OperationKind::Pull).expect("random id");
    646         let second = ids.next_id(OperationKind::Pull).expect("random id");
    647         assert_ne!(first, second);
    648 
    649         let storage = Arc::new(MemoryStorage::new(
    650             SourceGeneration::new([9; 32]).expect("generation"),
    651         ));
    652         let engine = Engine::builder(
    653             storage.clone(),
    654             Arc::new(FixedClock),
    655             Arc::new(SequenceIds(AtomicU8::new(20))),
    656             deadlines,
    657         )
    658         .source(Arc::new(CancelledSource))
    659         .build()
    660         .expect("engine");
    661         let client = ClientBuilder::new()
    662             .storage(storage)
    663             .sync_engine(engine)
    664             .build()
    665             .expect("client");
    666         let operations = client.sync().expect("open client").expect("sync");
    667         let admission = RegistryPolicy::verified();
    668 
    669         assert_eq!(
    670             operations.ingest(invalid_observation(), &admission).await,
    671             Err(Error::VerificationFailed)
    672         );
    673 
    674         let request = push_request();
    675         let operation_id = request.operation_id();
    676         operations
    677             .prepare_push(request.clone())
    678             .await
    679             .expect("durable preparation");
    680         assert_eq!(
    681             operations.admit_signed(operation_id).await,
    682             Err(Error::InvalidSignerOutput)
    683         );
    684         let cancellation = operations
    685             .cancel_push(operation_id)
    686             .await
    687             .expect("cancel prepared push");
    688         assert_eq!(
    689             cancellation.status().operation().operation_id().as_bytes(),
    690             operation_id.as_bytes()
    691         );
    692         assert!(cancellation.changed());
    693 
    694         let submitted = push_request();
    695         assert_eq!(
    696             operations.submit_push(submitted).await,
    697             Err(Error::SigningCancelled)
    698         );
    699         assert_eq!(
    700             operations.submit_push(push_request_for_id(9)).await,
    701             Err(Error::MissingSigner)
    702         );
    703     }
    704 }