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 }