delivery_evidence.rs (48672B)
1 use super::*; 2 use futures::channel::oneshot; 3 use radroots_storage::{ 4 authored_atomic::{AuthoredAtomicCommand, RecordDeliveryFact}, 5 authored_delivery::DeliveryAttemptOutcome, 6 }; 7 use radroots_transport::policy::SatisfactionState; 8 9 type Pending = ( 10 DeliveryRequest, 11 oneshot::Sender<Result<DeliveryReceipt, SinkFailure>>, 12 ); 13 14 #[derive(Default)] 15 struct HeldSink { 16 pending: Mutex<VecDeque<Pending>>, 17 calls: AtomicUsize, 18 } 19 20 impl HeldSink { 21 fn take(&self) -> Pending { 22 self.pending.lock().unwrap().pop_front().unwrap() 23 } 24 } 25 impl EventSink for HeldSink { 26 fn status(&self) -> radroots_transport::BoxFuture<'_, Result<SinkStatus, TransportError>> { 27 Box::pin(async { panic!("delivery must not probe transport status") }) 28 } 29 fn deliver( 30 &self, 31 request: DeliveryRequest, 32 ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> { 33 Box::pin(async move { 34 self.calls.fetch_add(1, Ordering::Relaxed); 35 let (sender, receiver) = oneshot::channel(); 36 self.pending.lock().unwrap().push_back((request, sender)); 37 receiver.await.expect("host retains admitted future") 38 }) 39 } 40 } 41 fn setup( 42 byte: u8, 43 ) -> ( 44 Engine, 45 Arc<MemoryStorage>, 46 Arc<TestClock>, 47 Arc<HeldSink>, 48 PushRequest, 49 ) { 50 let sink = Arc::new(HeldSink::default()); 51 let ((engine, storage), clock) = setup_engine_with_sink( 52 Arc::new(MockSigner::new(SignBehavior::Success { 53 completed_at_unix_ms: 1_800_000_200_500, 54 })), 55 sink.clone(), 56 ); 57 let push = request(byte, "wss://relay.example"); 58 execute_to_admitted(&engine, &push); 59 (engine, storage, clock, sink, push) 60 } 61 pub(super) fn source_only(storage: Arc<dyn SyncStorage>, clock: Arc<dyn Clock>) -> Engine { 62 Engine::builder( 63 storage, 64 clock, 65 Arc::new(TestIds(AtomicU64::new(230))), 66 DeadlinePolicy::new(10_000, 10_000, 10_000).unwrap(), 67 ) 68 .source(Arc::new(MockSource)) 69 .build() 70 .unwrap() 71 } 72 fn accept((request, sender): Pending) -> DeliveryReceipt { 73 let result = receipt( 74 &request, 75 vec![DeliveryOutcome::accepted(); request.target_set().len()], 76 ) 77 .unwrap(); 78 sender.send(Ok(result.clone())).unwrap(); 79 result 80 } 81 82 #[test] 83 fn capacity_delivery_receipt_failure_retains_claim_and_known_or_unknown_effects() { 84 for after in [false, true] { 85 let storage = Arc::new(FaultStorage::new(185)); 86 let sink = Arc::new(HeldSink::default()); 87 let signer = Arc::new(MockSigner::new(SignBehavior::Success { 88 completed_at_unix_ms: 1_800_000_200_500, 89 })); 90 let engine = fault_engine(storage.clone(), signer.clone(), sink.clone()); 91 let push = request(185, "wss://capacity.example"); 92 execute_to_admitted(&engine, &push); 93 let before = block_on(engine.push_status(push.operation_id())) 94 .unwrap() 95 .unwrap(); 96 let mut future = Box::pin(engine.deliver_push(push.operation_id())); 97 assert!( 98 future 99 .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref())) 100 .is_pending() 101 ); 102 let pending = sink.take(); 103 assert_eq!(&pending.0, before.delivery_plan().request().unwrap()); 104 *storage.capacity.lock().unwrap() = Some(capacity::Fault { 105 phase: capacity::Phase::DeliveryFact, 106 after, 107 }); 108 let expected = accept(pending); 109 assert_eq!(block_on(future), Err(Error::StorageSpaceInsufficient)); 110 let failed = block_on(engine.push_status(push.operation_id())) 111 .unwrap() 112 .unwrap(); 113 assert_eq!(failed.artifact().signed(), before.artifact().signed()); 114 assert_eq!( 115 failed.delivery_plan().request(), 116 before.delivery_plan().request() 117 ); 118 assert!(!failed.delivery_history().proves_no_issued_attempt()); 119 assert_eq!(failed.delivery_history().has_unresolved_claims(), !after); 120 assert_eq!( 121 failed.delivery_plan().delivery_facts().len(), 122 usize::from(after) 123 ); 124 if after { 125 assert_eq!( 126 failed.delivery_plan().delivery_facts()[0].outcome(), 127 &DeliveryAttemptOutcome::Receipt(expected) 128 ); 129 assert_eq!( 130 block_on(engine.deliver_push(push.operation_id())), 131 Err(Error::WorkClaimConflict) 132 ); 133 let expires_at = failed 134 .delivery_plan() 135 .claim_evidence() 136 .unwrap() 137 .expires_at_unix_ms(); 138 let recovery = source_only( 139 storage.clone(), 140 Arc::new(TestClock(AtomicU64::new(expires_at))), 141 ); 142 let reconciled = block_on(recovery.deliver_push(push.operation_id())).unwrap(); 143 assert_eq!(reconciled.plan().state(), AuthoredDeliveryState::Satisfied); 144 } 145 assert_eq!(sink.calls.load(Ordering::Relaxed), 1); 146 assert_eq!(signer.calls.load(Ordering::Relaxed), 1); 147 } 148 } 149 150 #[test] 151 fn late_accepted_result_survives_stop_and_terminal_replay_without_sink() { 152 let (engine, storage, clock, sink, push) = setup(121); 153 let status = block_on(engine.push_status(push.operation_id())) 154 .unwrap() 155 .unwrap(); 156 assert!(status.delivery_history().proves_no_issued_attempt()); 157 let mut future = Box::pin(engine.deliver_push(push.operation_id())); 158 assert!( 159 future 160 .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref())) 161 .is_pending() 162 ); 163 let pending = sink.take(); 164 let stopped = block_on(engine.cancel_push(push.operation_id())).unwrap(); 165 assert!(stopped.changed()); 166 assert!(stopped.status().delivery_history().has_unresolved_claims()); 167 assert!( 168 !stopped 169 .status() 170 .delivery_history() 171 .proves_no_issued_attempt() 172 ); 173 let stop_at = stopped.status().delivery_plan().stop_requested_at_unix_ms(); 174 let expected = accept(pending); 175 let delivered = block_on(future).unwrap(); 176 assert_eq!(delivered.plan().state(), AuthoredDeliveryState::Cancelled); 177 assert_eq!(delivered.plan().attempt_count(), 0); 178 assert_eq!( 179 delivered.plan().delivery_satisfaction().unwrap(), 180 SatisfactionState::Satisfied 181 ); 182 assert_eq!( 183 delivered.plan().delivery_facts()[0].outcome(), 184 &DeliveryAttemptOutcome::Receipt(expected) 185 ); 186 let status = block_on(engine.push_status(push.operation_id())) 187 .unwrap() 188 .unwrap(); 189 assert!(!status.delivery_history().has_unresolved_claims()); 190 assert_eq!(status.delivery_plan().stop_requested_at_unix_ms(), stop_at); 191 assert!( 192 !block_on(engine.cancel_push(push.operation_id())) 193 .unwrap() 194 .changed() 195 ); 196 let replay = block_on(source_only(storage, clock).deliver_push(push.operation_id())).unwrap(); 197 assert!(replay.is_replay()); 198 assert_eq!(replay.plan(), status.delivery_plan()); 199 assert_eq!(sink.calls.load(Ordering::Relaxed), 1); 200 } 201 202 #[test] 203 fn expired_callback_retains_fact_but_only_fresh_invocation_reconciles() { 204 let (engine, storage, clock, sink, push) = setup(122); 205 let mut future = Box::pin(engine.deliver_push(push.operation_id())); 206 assert!( 207 future 208 .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref())) 209 .is_pending() 210 ); 211 let before = block_on(engine.push_status(push.operation_id())) 212 .unwrap() 213 .unwrap(); 214 clock.0.store( 215 before 216 .delivery_plan() 217 .claim_evidence() 218 .unwrap() 219 .expires_at_unix_ms(), 220 Ordering::Relaxed, 221 ); 222 accept(sink.take()); 223 let late = block_on(future).unwrap(); 224 assert_eq!(late.plan().revision(), before.delivery_plan().revision()); 225 assert_eq!(late.plan().attempt_count(), 0); 226 let recovery = source_only(storage, clock); 227 let recovered = block_on(recovery.deliver_push(push.operation_id())).unwrap(); 228 assert!(recovered.is_replay()); 229 assert_eq!(recovered.plan().state(), AuthoredDeliveryState::Satisfied); 230 assert_eq!(recovered.plan().attempt_count(), 1); 231 assert_eq!( 232 recovered.plan().attempts()[0].claim_evidence(), 233 before.delivery_plan().claim_evidence() 234 ); 235 assert_eq!(sink.calls.load(Ordering::Relaxed), 1); 236 } 237 238 #[test] 239 fn superseded_callback_cannot_clear_newer_claim_and_acceptance_never_regresses() { 240 let (engine, _, clock, sink, push) = setup(123); 241 let mut first = Box::pin(engine.deliver_push(push.operation_id())); 242 assert!( 243 first 244 .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref())) 245 .is_pending() 246 ); 247 let first_pending = sink.take(); 248 let before = block_on(engine.push_status(push.operation_id())) 249 .unwrap() 250 .unwrap(); 251 clock.0.store( 252 before 253 .delivery_plan() 254 .claim_evidence() 255 .unwrap() 256 .expires_at_unix_ms(), 257 Ordering::Relaxed, 258 ); 259 let mut second = Box::pin(engine.deliver_push(push.operation_id())); 260 assert!( 261 second 262 .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref())) 263 .is_pending() 264 ); 265 let (second_request, second_sender) = sink.take(); 266 assert_eq!(first_pending.0, second_request); 267 let newer = block_on(engine.push_status(push.operation_id())) 268 .unwrap() 269 .unwrap(); 270 accept(first_pending); 271 let late = block_on(first).unwrap(); 272 assert_eq!( 273 late.plan().claim_evidence(), 274 newer.delivery_plan().claim_evidence() 275 ); 276 assert_eq!(late.plan().revision(), newer.delivery_plan().revision()); 277 assert_eq!( 278 block_on(engine.deliver_push(push.operation_id())), 279 Err(Error::WorkClaimConflict) 280 ); 281 second_sender 282 .send(Err(SinkFailure::for_request( 283 &second_request, 284 "provider_failed", 285 Retryability::Terminal, 286 None, 287 None, 288 Vec::new(), 289 ) 290 .unwrap())) 291 .unwrap(); 292 let complete = block_on(second).unwrap(); 293 assert_eq!(complete.plan().state(), AuthoredDeliveryState::Satisfied); 294 assert_eq!(complete.plan().attempt_count(), 2); 295 assert_eq!(complete.plan().delivery_facts().len(), 2); 296 assert_eq!(sink.calls.load(Ordering::Relaxed), 2); 297 } 298 299 #[test] 300 fn clock_loss_retains_raw_result_before_reporting_unavailable() { 301 for accepted in [true, false] { 302 let (engine, storage, clock, sink, push) = setup(124 + u8::from(accepted)); 303 let mut future = Box::pin(engine.deliver_push(push.operation_id())); 304 assert!( 305 future 306 .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref())) 307 .is_pending() 308 ); 309 let (request, sender) = sink.take(); 310 let before = block_on(engine.push_status(push.operation_id())) 311 .unwrap() 312 .unwrap(); 313 let outcome = if accepted { 314 DeliveryAttemptOutcome::Receipt( 315 receipt(&request, vec![DeliveryOutcome::accepted()]).unwrap(), 316 ) 317 } else { 318 DeliveryAttemptOutcome::SinkFailure( 319 SinkFailure::for_request( 320 &request, 321 "raw_failure", 322 Retryability::Retryable, 323 None, 324 Some("raw diagnostic".into()), 325 Vec::new(), 326 ) 327 .unwrap(), 328 ) 329 }; 330 sender 331 .send(match &outcome { 332 DeliveryAttemptOutcome::Receipt(value) => Ok(value.clone()), 333 DeliveryAttemptOutcome::SinkFailure(value) => Err(value.clone()), 334 }) 335 .unwrap(); 336 clock.0.store(0, Ordering::Relaxed); 337 assert_eq!(block_on(future), Err(Error::ClockUnavailable)); 338 let status = block_on(engine.push_status(push.operation_id())) 339 .unwrap() 340 .unwrap(); 341 assert_eq!( 342 status.delivery_plan().delivery_facts()[0].outcome(), 343 &outcome 344 ); 345 assert_eq!( 346 status.delivery_plan().revision(), 347 before.delivery_plan().revision() 348 ); 349 assert_eq!(status.delivery_plan().attempt_count(), 0); 350 let expires = status 351 .delivery_plan() 352 .claim_evidence() 353 .unwrap() 354 .expires_at_unix_ms(); 355 clock.0.store(expires, Ordering::Relaxed); 356 let recovered = 357 block_on(source_only(storage, clock).deliver_push(push.operation_id())).unwrap(); 358 assert_eq!(recovered.plan().attempt_count(), 1); 359 assert_eq!(recovered.plan().attempts()[0].outcome(), &outcome); 360 if !accepted { 361 assert!(recovered.plan().retry().unwrap().not_before_unix_ms() >= expires + 1_000); 362 } 363 assert_eq!(sink.calls.load(Ordering::Relaxed), 1); 364 } 365 } 366 367 #[test] 368 fn stop_before_delivery_proves_no_issued_attempt_and_stop_after_acceptance_keeps_success() { 369 let (engine, _, _, sink, push) = setup(126); 370 let stopped = block_on(engine.cancel_push(push.operation_id())).unwrap(); 371 assert!( 372 stopped 373 .status() 374 .delivery_history() 375 .proves_no_issued_attempt() 376 ); 377 assert!( 378 block_on(engine.deliver_push(push.operation_id())) 379 .unwrap() 380 .is_replay() 381 ); 382 assert_eq!(sink.calls.load(Ordering::Relaxed), 0); 383 let (engine, _, _, sink, push) = setup(127); 384 let mut future = Box::pin(engine.deliver_push(push.operation_id())); 385 assert!( 386 future 387 .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref())) 388 .is_pending() 389 ); 390 accept(sink.take()); 391 assert_eq!( 392 block_on(future).unwrap().plan().state(), 393 AuthoredDeliveryState::Satisfied 394 ); 395 let stopped = block_on(engine.cancel_push(push.operation_id())).unwrap(); 396 assert!(stopped.changed()); 397 assert!( 398 stopped 399 .status() 400 .delivery_plan() 401 .stop_requested_at_unix_ms() 402 .is_some() 403 ); 404 assert_eq!( 405 stopped.status().delivery_plan().state(), 406 AuthoredDeliveryState::Satisfied 407 ); 408 assert_eq!( 409 stopped 410 .status() 411 .delivery_plan() 412 .delivery_satisfaction() 413 .unwrap(), 414 SatisfactionState::Satisfied 415 ); 416 assert!( 417 !block_on(engine.cancel_push(push.operation_id())) 418 .unwrap() 419 .changed() 420 ); 421 } 422 423 #[test] 424 fn equal_fact_replay_preserves_first_observation_and_conflict_cannot_erase_acceptance() { 425 let (engine, storage, _, sink, push) = setup(128); 426 let mut future = Box::pin(engine.deliver_push(push.operation_id())); 427 assert!( 428 future 429 .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref())) 430 .is_pending() 431 ); 432 accept(sink.take()); 433 let delivered = block_on(future).unwrap(); 434 let plan = delivered.plan(); 435 let fact = &plan.delivery_facts()[0]; 436 let repeated = AuthoredAtomicCommand::RecordDelivery( 437 RecordDeliveryFact::new( 438 plan.plan_id(), 439 plan.artifact_id(), 440 fact.claim().clone(), 441 fact.outcome().clone(), 442 fact.observed_at_unix_ms() + 10_000, 443 ) 444 .unwrap(), 445 ); 446 block_on(storage.execute_authored(repeated)).unwrap(); 447 let changed = DeliveryAttemptOutcome::Receipt( 448 receipt( 449 plan.request().unwrap(), 450 vec![DeliveryOutcome::unavailable()], 451 ) 452 .unwrap(), 453 ); 454 let conflicting = AuthoredAtomicCommand::RecordDelivery( 455 RecordDeliveryFact::new( 456 plan.plan_id(), 457 plan.artifact_id(), 458 fact.claim().clone(), 459 changed, 460 fact.observed_at_unix_ms() + 10_000, 461 ) 462 .unwrap(), 463 ); 464 assert!(block_on(storage.execute_authored(conflicting)).is_err()); 465 let status = block_on(engine.push_status(push.operation_id())) 466 .unwrap() 467 .unwrap(); 468 assert_eq!(status.delivery_plan(), plan); 469 } 470 #[test] 471 fn every_delivery_storage_receipt_is_checked_and_durable_results_remain_recoverable() { 472 for nth in 1..=3 { 473 let storage = Arc::new(FaultStorage::new(150 + nth as u8)); 474 let sink = Arc::new(ScriptedSink::new([DeliveryBehavior::Outcomes(vec![ 475 DeliveryOutcome::accepted(), 476 ])])); 477 let engine = fault_engine( 478 storage.clone(), 479 Arc::new(MockSigner::new(SignBehavior::Success { 480 completed_at_unix_ms: 1_800_000_200_500, 481 })), 482 sink.clone(), 483 ); 484 let push = request(150 + nth as u8, "wss://relay.example"); 485 execute_to_admitted(&engine, &push); 486 storage.fault_nth_plan(nth); 487 assert_eq!( 488 block_on(engine.deliver_push(push.operation_id())), 489 Err(Error::StorageFailed) 490 ); 491 let status = block_on(engine.push_status(push.operation_id())) 492 .unwrap() 493 .unwrap(); 494 assert_eq!(sink.requests.lock().unwrap().len(), usize::from(nth > 1)); 495 assert_eq!( 496 status.delivery_plan().delivery_facts().len(), 497 usize::from(nth > 1) 498 ); 499 if nth > 1 { 500 let clock = Arc::new(TestClock(AtomicU64::new(1_800_000_250_000))); 501 let recovered = 502 block_on(source_only(storage, clock).deliver_push(push.operation_id())).unwrap(); 503 assert_eq!(recovered.plan().state(), AuthoredDeliveryState::Satisfied); 504 assert_eq!(recovered.plan().attempt_count(), 1); 505 } 506 } 507 } 508 509 #[test] 510 fn legacy_applied_result_gets_exact_marker_without_counting_a_second_attempt() { 511 legacy_reconciliation(false); 512 } 513 514 #[test] 515 fn conflicting_legacy_result_retains_both_observations_without_reconciliation() { 516 legacy_reconciliation(true); 517 } 518 519 fn legacy_reconciliation(conflict: bool) { 520 use core::num::NonZeroU32; 521 use radroots_storage::{ 522 authored::{FailureClass, RetrySchedule, WorkFailure, WorkPhase}, 523 authored_atomic::{ApplyDeliveryAttempt, WorkFence}, 524 }; 525 let (engine, storage, clock, sink, push) = setup(154); 526 let mut future = Box::pin(engine.deliver_push(push.operation_id())); 527 assert!( 528 future 529 .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref())) 530 .is_pending() 531 ); 532 let (request, _sender) = sink.take(); 533 drop(future); 534 let status = block_on(engine.push_status(push.operation_id())) 535 .unwrap() 536 .unwrap(); 537 let plan = status.delivery_plan(); 538 let claim = plan.claim_evidence().unwrap(); 539 let at = claim.acquired_at_unix_ms() + 100; 540 let outcome = DeliveryAttemptOutcome::Receipt( 541 receipt(&request, vec![DeliveryOutcome::unavailable()]).unwrap(), 542 ); 543 let retry = RetrySchedule::new( 544 NonZeroU32::MIN, 545 at + 1000, 546 WorkFailure::new( 547 "delivery_pending", 548 WorkPhase::Delivery, 549 FailureClass::Retryable, 550 Some(at + 1000), 551 None, 552 ) 553 .unwrap(), 554 ) 555 .unwrap(); 556 block_on( 557 storage.execute_authored(AuthoredAtomicCommand::ApplyDelivery( 558 ApplyDeliveryAttempt::new( 559 plan.plan_id(), 560 WorkFence::new(*claim.token(), claim.generation(), claim.row_revision()).unwrap(), 561 outcome.clone(), 562 Some(retry), 563 at, 564 ) 565 .unwrap(), 566 )), 567 ) 568 .unwrap(); 569 block_on( 570 storage.execute_authored(AuthoredAtomicCommand::RecordDelivery( 571 RecordDeliveryFact::new( 572 plan.plan_id(), 573 plan.artifact_id(), 574 claim.clone(), 575 if conflict { 576 DeliveryAttemptOutcome::Receipt( 577 receipt(&request, vec![DeliveryOutcome::accepted()]).unwrap(), 578 ) 579 } else { 580 outcome.clone() 581 }, 582 at + 1, 583 ) 584 .unwrap(), 585 )), 586 ) 587 .unwrap(); 588 clock.0.store(at + 2, Ordering::Relaxed); 589 let result = block_on(source_only(storage, clock).deliver_push(push.operation_id())); 590 if conflict { 591 assert_eq!(result, Err(Error::StorageConflict)); 592 let current = block_on(engine.push_status(push.operation_id())) 593 .unwrap() 594 .unwrap(); 595 assert_eq!(current.delivery_plan().attempt_count(), 1); 596 assert_eq!(current.delivery_plan().attempts()[0].outcome(), &outcome); 597 assert_eq!(current.delivery_plan().delivery_facts().len(), 1); 598 assert_eq!( 599 current.delivery_plan().delivery_satisfaction().unwrap(), 600 SatisfactionState::Satisfied 601 ); 602 assert_eq!(sink.calls.load(Ordering::Relaxed), 1); 603 return; 604 } 605 let reconciled = result.unwrap(); 606 assert_eq!(reconciled.plan().attempt_count(), 1); 607 assert_eq!(reconciled.plan().attempts()[0].recorded_at_unix_ms(), at); 608 assert_eq!(reconciled.plan().attempts()[0].outcome(), &outcome); 609 assert_eq!( 610 reconciled.plan().attempts()[0].claim_evidence(), 611 Some(claim) 612 ); 613 assert_eq!(sink.calls.load(Ordering::Relaxed), 1); 614 } 615 616 #[tokio::test] 617 async fn sqlite_late_facts_reopen_after_stop_or_expiry_and_reconcile_once() { 618 use radroots_storage::authored_atomic::{CancelAuthoredTarget, CancelAuthoredWork}; 619 struct BoundarySink { 620 storage: Arc<SqliteStorage>, 621 clock: Arc<TestClock>, 622 id: radroots_storage::authored_delivery::AuthoredDeliveryPlanId, 623 stop: bool, 624 } 625 impl EventSink for BoundarySink { 626 fn status(&self) -> radroots_transport::BoxFuture<'_, Result<SinkStatus, TransportError>> { 627 Box::pin(async { panic!("no probe") }) 628 } 629 fn deliver( 630 &self, 631 request: DeliveryRequest, 632 ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> { 633 Box::pin(async move { 634 let history = self 635 .storage 636 .authored_delivery_history(self.id) 637 .await 638 .unwrap() 639 .unwrap(); 640 let plan = history.plan(); 641 let at = plan.claim_evidence().unwrap().expires_at_unix_ms(); 642 self.clock.0.store(at, Ordering::Relaxed); 643 if self.stop { 644 self.storage 645 .execute_authored(AuthoredAtomicCommand::Cancel( 646 CancelAuthoredWork::new( 647 CancelAuthoredTarget::DeliveryPlan(self.id), 648 plan.revision(), 649 at, 650 ) 651 .unwrap(), 652 )) 653 .await 654 .unwrap(); 655 } 656 Ok(receipt(&request, vec![DeliveryOutcome::accepted()]).unwrap()) 657 }) 658 } 659 } 660 for stop in [true, false] { 661 let directory = tempfile::tempdir().unwrap(); 662 let paths = Paths::from_directory(directory.path()).unwrap(); 663 let storage = Arc::new( 664 SqliteStorage::open( 665 OpenOptions::new(paths.clone(), OpenMode::Create) 666 .with_source_generation(SourceGeneration::new([155; 32]).unwrap(), 1) 667 .unwrap(), 668 ) 669 .await 670 .unwrap(), 671 ); 672 let clock = Arc::new(TestClock(AtomicU64::new(1_800_000_200_000))); 673 let push = request(155, "wss://relay.example"); 674 let id = push 675 .authored_preparation(1_800_000_200_000) 676 .unwrap() 677 .delivery_plans()[0] 678 .plan_id(); 679 let engine = Engine::builder( 680 storage.clone(), 681 clock.clone(), 682 Arc::new(TestIds(AtomicU64::new(10))), 683 DeadlinePolicy::new(10_000, 10_000, 10_000).unwrap(), 684 ) 685 .sink(Arc::new(BoundarySink { 686 storage: storage.clone(), 687 clock: clock.clone(), 688 id, 689 stop, 690 })) 691 .signer(Arc::new(MockSigner::new(SignBehavior::Success { 692 completed_at_unix_ms: 1_800_000_200_500, 693 }))) 694 .build() 695 .unwrap(); 696 engine.sign_prepared(push.clone()).await.unwrap(); 697 engine.admit_signed(push.operation_id()).await.unwrap(); 698 let late = engine.deliver_push(push.operation_id()).await.unwrap(); 699 assert_eq!(late.plan().attempt_count(), 0); 700 assert_eq!( 701 late.plan().delivery_satisfaction().unwrap(), 702 SatisfactionState::Satisfied 703 ); 704 drop(engine); 705 storage.close().await.unwrap(); 706 let storage = Arc::new( 707 SqliteStorage::open(OpenOptions::new(paths.clone(), OpenMode::ReadWriteExisting)) 708 .await 709 .unwrap(), 710 ); 711 let recovery = source_only(storage.clone(), clock.clone()); 712 let recovered = recovery.deliver_push(push.operation_id()).await.unwrap(); 713 assert_eq!(recovered.plan().attempt_count(), u32::from(!stop)); 714 assert_eq!( 715 recovered.plan().delivery_facts(), 716 late.plan().delivery_facts() 717 ); 718 let status = recovery 719 .push_status(push.operation_id()) 720 .await 721 .unwrap() 722 .unwrap(); 723 assert!(!status.delivery_history().has_unresolved_claims()); 724 drop(recovery); 725 storage.close().await.unwrap(); 726 let storage = Arc::new( 727 SqliteStorage::open(OpenOptions::new(paths, OpenMode::ReadOnly)) 728 .await 729 .unwrap(), 730 ); 731 let readonly = source_only(storage.clone(), clock); 732 let replay = readonly.deliver_push(push.operation_id()).await.unwrap(); 733 assert_eq!(replay.plan(), recovered.plan()); 734 drop(readonly); 735 storage.close().await.unwrap(); 736 } 737 } 738 pub(super) enum Injection { 739 FactBeforeClaim(Box<AuthoredAtomicCommand>), 740 StopAfterClaim, 741 ReplaceAfterClaim, 742 StopAfterFact, 743 StopAfterReconcile, 744 } 745 746 pub(super) async fn inject(storage: &FaultStorage, command: &AuthoredAtomicCommand, before: bool) { 747 use radroots_storage::{ 748 authored::WorkClaim, 749 authored_atomic::{ 750 CancelAuthoredTarget, CancelAuthoredWork, ClaimAuthoredTarget, ClaimAuthoredWork, 751 }, 752 }; 753 let claim = matches!(command, AuthoredAtomicCommand::Claim(value) 754 if matches!(value.target(), ClaimAuthoredTarget::DeliveryPlan(_))); 755 let injection = { 756 let mut armed = storage.delivery_injection.lock().unwrap(); 757 let ready = match armed.as_ref() { 758 Some(Injection::FactBeforeClaim(_)) => before && claim, 759 Some(Injection::StopAfterClaim | Injection::ReplaceAfterClaim) => !before && claim, 760 Some(Injection::StopAfterFact) => { 761 !before && matches!(command, AuthoredAtomicCommand::RecordDelivery(_)) 762 } 763 Some(Injection::StopAfterReconcile) => { 764 !before && matches!(command, AuthoredAtomicCommand::ReconcileDelivery(_)) 765 } 766 None => false, 767 }; 768 if ready { armed.take() } else { None } 769 }; 770 let Some(injection) = injection else { 771 return; 772 }; 773 if let Injection::FactBeforeClaim(fact) = injection { 774 storage.inner.execute_authored(*fact).await.unwrap(); 775 return; 776 } 777 let id = match command { 778 AuthoredAtomicCommand::Claim(value) => match value.target() { 779 ClaimAuthoredTarget::DeliveryPlan(id) => *id, 780 _ => panic!("delivery-only injection"), 781 }, 782 AuthoredAtomicCommand::RecordDelivery(value) => value.plan_id(), 783 AuthoredAtomicCommand::ReconcileDelivery(value) => value.plan_id(), 784 _ => panic!("delivery-only injection"), 785 }; 786 let history = storage 787 .inner 788 .authored_delivery_history(id) 789 .await 790 .unwrap() 791 .unwrap(); 792 let plan = history.plan(); 793 let mutation = if matches!(injection, Injection::ReplaceAfterClaim) { 794 let original = plan.claim_evidence().unwrap(); 795 let at = original.expires_at_unix_ms(); 796 AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( 797 ClaimAuthoredTarget::DeliveryPlan(id), 798 WorkClaim::new( 799 [222; 16], 800 "competing-worker", 801 core::num::NonZeroU64::new(2).unwrap(), 802 at, 803 at + 10_000, 804 plan.revision(), 805 ) 806 .unwrap(), 807 )) 808 } else { 809 AuthoredAtomicCommand::Cancel( 810 CancelAuthoredWork::new( 811 CancelAuthoredTarget::DeliveryPlan(id), 812 plan.revision(), 813 plan.updated_at_unix_ms(), 814 ) 815 .unwrap(), 816 ) 817 }; 818 storage.inner.execute_authored(mutation).await.unwrap(); 819 } 820 821 #[test] 822 fn observed_stop_or_replacement_after_claim_prevents_sink_admission() { 823 for stop in [true, false] { 824 let storage = Arc::new(FaultStorage::new(160)); 825 let sink = Arc::new(HeldSink::default()); 826 let engine = fault_engine( 827 storage.clone(), 828 Arc::new(MockSigner::new(SignBehavior::Success { 829 completed_at_unix_ms: 1_800_000_200_500, 830 })), 831 sink.clone(), 832 ); 833 let push = request(160, "wss://relay.example"); 834 execute_to_admitted(&engine, &push); 835 *storage.delivery_injection.lock().unwrap() = Some(if stop { 836 Injection::StopAfterClaim 837 } else { 838 Injection::ReplaceAfterClaim 839 }); 840 let result = block_on(engine.deliver_push(push.operation_id())); 841 if stop { 842 assert_eq!( 843 result.unwrap().plan().state(), 844 AuthoredDeliveryState::Cancelled 845 ); 846 } else { 847 assert_eq!(result, Err(Error::WorkClaimConflict)); 848 } 849 assert_eq!(sink.calls.load(Ordering::Relaxed), 0); 850 } 851 } 852 853 #[test] 854 fn fresh_claim_reconciles_fact_arriving_after_initial_status_without_another_effect() { 855 let storage = Arc::new(FaultStorage::new(161)); 856 let sink = Arc::new(HeldSink::default()); 857 let clock = Arc::new(TestClock(AtomicU64::new(1_800_000_200_000))); 858 let engine = Engine::builder( 859 storage.clone(), 860 clock.clone(), 861 Arc::new(TestIds(AtomicU64::new(10))), 862 DeadlinePolicy::new(10_000, 10_000, 10_000).unwrap(), 863 ) 864 .sink(sink.clone()) 865 .signer(Arc::new(MockSigner::new(SignBehavior::Success { 866 completed_at_unix_ms: 1_800_000_200_500, 867 }))) 868 .build() 869 .unwrap(); 870 let push = request(161, "wss://relay.example"); 871 execute_to_admitted(&engine, &push); 872 let mut first = Box::pin(engine.deliver_push(push.operation_id())); 873 assert!( 874 first 875 .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref())) 876 .is_pending() 877 ); 878 let (request, _sender) = sink.take(); 879 drop(first); 880 let status = block_on(engine.push_status(push.operation_id())) 881 .unwrap() 882 .unwrap(); 883 let plan = status.delivery_plan(); 884 let claim = plan.claim_evidence().unwrap(); 885 let at = claim.expires_at_unix_ms(); 886 let fact = AuthoredAtomicCommand::RecordDelivery( 887 RecordDeliveryFact::new( 888 plan.plan_id(), 889 plan.artifact_id(), 890 claim.clone(), 891 DeliveryAttemptOutcome::Receipt( 892 receipt(&request, vec![DeliveryOutcome::accepted()]).unwrap(), 893 ), 894 at, 895 ) 896 .unwrap(), 897 ); 898 *storage.delivery_injection.lock().unwrap() = Some(Injection::FactBeforeClaim(Box::new(fact))); 899 clock.0.store(at, Ordering::Relaxed); 900 let reconciled = block_on(engine.deliver_push(push.operation_id())).unwrap(); 901 assert!(reconciled.is_replay()); 902 assert_eq!(reconciled.plan().state(), AuthoredDeliveryState::Satisfied); 903 assert_eq!(reconciled.plan().attempt_count(), 1); 904 assert_eq!(sink.calls.load(Ordering::Relaxed), 1); 905 } 906 907 #[test] 908 fn receipt_return_reloads_stop_after_fact_or_reconciliation_commit() { 909 for before_reconcile in [true, false] { 910 let storage = Arc::new(FaultStorage::new(162)); 911 let sink = Arc::new(ScriptedSink::new([DeliveryBehavior::Outcomes(vec![ 912 DeliveryOutcome::accepted(), 913 ])])); 914 let engine = fault_engine( 915 storage.clone(), 916 Arc::new(MockSigner::new(SignBehavior::Success { 917 completed_at_unix_ms: 1_800_000_200_500, 918 })), 919 sink.clone(), 920 ); 921 let push = request(162, "wss://relay.example"); 922 execute_to_admitted(&engine, &push); 923 *storage.delivery_injection.lock().unwrap() = Some(if before_reconcile { 924 Injection::StopAfterFact 925 } else { 926 Injection::StopAfterReconcile 927 }); 928 let delivered = block_on(engine.deliver_push(push.operation_id())).unwrap(); 929 assert!(delivered.plan().stop_requested_at_unix_ms().is_some()); 930 assert_eq!( 931 delivered.plan().attempt_count(), 932 u32::from(!before_reconcile) 933 ); 934 assert_eq!( 935 delivered.plan().delivery_satisfaction().unwrap(), 936 SatisfactionState::Satisfied 937 ); 938 assert_eq!( 939 engine 940 .retry_decision(delivered.plan(), 1_800_000_250_000) 941 .unwrap(), 942 SyncRetryDecision::Satisfied 943 ); 944 assert_eq!(sink.requests.lock().unwrap().len(), 1); 945 } 946 } 947 948 #[test] 949 fn forged_stop_receipt_is_not_reported_as_success() { 950 let storage = Arc::new(FaultStorage::new(163)); 951 let engine = fault_engine( 952 storage.clone(), 953 Arc::new(MockSigner::new(SignBehavior::Success { 954 completed_at_unix_ms: 1_800_000_200_500, 955 })), 956 Arc::new(MockSink), 957 ); 958 let push = request(163, "wss://relay.example"); 959 execute_to_admitted(&engine, &push); 960 storage.fault_nth_plan(1); 961 assert_eq!( 962 block_on(engine.cancel_push(push.operation_id())), 963 Err(Error::StorageFailed) 964 ); 965 let status = block_on(engine.push_status(push.operation_id())) 966 .unwrap() 967 .unwrap(); 968 assert!(status.delivery_plan().stop_requested_at_unix_ms().is_some()); 969 assert!(status.delivery_history().proves_no_issued_attempt()); 970 } 971 use radroots_storage::authored_atomic::{AuthoredAtomicOutcome, AuthoredAtomicReceipt}; 972 use radroots_storage::authored_delivery::{AuthoredDeliveryHistory, AuthoredDeliveryPlan}; 973 974 #[derive(Clone, Copy, Debug)] 975 pub(super) enum ReceiptMutation { 976 Identity, 977 Plan, 978 Artifact, 979 Request, 980 Created, 981 Claim, 982 CommitTime, 983 Attempts, 984 State, 985 Retry, 986 StopTime, 987 Revision, 988 } 989 pub(super) fn mutate_receipt( 990 receipt: AuthoredAtomicReceipt, 991 mutation: ReceiptMutation, 992 ) -> AuthoredAtomicReceipt { 993 let AuthoredAtomicOutcome::DeliveryPlan(plan) = receipt.outcome() else { 994 panic!("delivery receipt"); 995 }; 996 let mut wire = serde_json::to_value(plan).unwrap(); 997 match mutation { 998 ReceiptMutation::Identity | ReceiptMutation::CommitTime => {} 999 ReceiptMutation::Plan => wire["plan_id"] = serde_json::json!(vec![241; 16]), 1000 ReceiptMutation::Artifact => wire["artifact_id"] = serde_json::json!(vec![242; 16]), 1001 ReceiptMutation::Request => { 1002 let original = plan.request().unwrap(); 1003 let changed = DeliveryRequest::new( 1004 "different-request", 1005 original.payload().clone(), 1006 original.target_set().clone(), 1007 original.satisfaction().clone(), 1008 original.deadline_unix_ms(), 1009 ) 1010 .unwrap(); 1011 let other = AuthoredDeliveryPlan::new_bound( 1012 plan.plan_id(), 1013 plan.artifact_id(), 1014 changed, 1015 plan.created_at_unix_ms(), 1016 ) 1017 .unwrap(); 1018 let other = serde_json::to_value(other).unwrap(); 1019 for key in ["request", "intent", "request_digest"] { 1020 wire[key] = other[key].clone(); 1021 } 1022 } 1023 ReceiptMutation::Created => { 1024 wire["created_at_unix_ms"] = serde_json::json!(plan.created_at_unix_ms() - 1) 1025 } 1026 ReceiptMutation::Claim => wire["claim"]["owner"] = serde_json::json!("another-worker"), 1027 ReceiptMutation::Attempts => { 1028 let outcome = DeliveryAttemptOutcome::Receipt(receipt_for_pending(plan)); 1029 wire["attempts"] = serde_json::json!([{"attempt":1,"recorded_at_unix_ms":plan.created_at_unix_ms(),"outcome":outcome,"satisfaction":"pending"}]); 1030 wire["attempt_count"] = serde_json::json!(1); 1031 } 1032 ReceiptMutation::State => { 1033 wire["state"] = serde_json::json!("pending"); 1034 wire["retry"] = serde_json::Value::Null; 1035 wire["last_failure"] = serde_json::Value::Null; 1036 } 1037 ReceiptMutation::Retry => { 1038 let at = plan.retry().unwrap().not_before_unix_ms() + 1000; 1039 wire["retry"]["not_before_unix_ms"] = serde_json::json!(at); 1040 wire["retry"]["failure"]["retry_after_unix_ms"] = serde_json::json!(at); 1041 wire["last_failure"] = wire["retry"]["failure"].clone(); 1042 } 1043 ReceiptMutation::StopTime => { 1044 wire["stop_requested_at_unix_ms"] = 1045 serde_json::json!(plan.stop_requested_at_unix_ms().unwrap() - 1) 1046 } 1047 ReceiptMutation::Revision => { 1048 wire["revision"] = serde_json::json!(plan.revision().get() + 1) 1049 } 1050 } 1051 let changed: AuthoredDeliveryPlan = 1052 serde_json::from_value(wire).expect("valid but incorrectly bound projection"); 1053 AuthoredAtomicReceipt::from_durable_parts( 1054 if matches!(mutation, ReceiptMutation::Identity) { 1055 radroots_storage::atomic::AtomicCommitId::new([243; 16]).unwrap() 1056 } else { 1057 receipt.commit_id() 1058 }, 1059 receipt.digest(), 1060 receipt.disposition(), 1061 receipt.committed_at_unix_ms() + u64::from(matches!(mutation, ReceiptMutation::CommitTime)), 1062 AuthoredAtomicOutcome::DeliveryPlan(changed), 1063 ) 1064 .unwrap() 1065 } 1066 fn receipt_for_pending(plan: &AuthoredDeliveryPlan) -> DeliveryReceipt { 1067 receipt( 1068 plan.request().unwrap(), 1069 vec![DeliveryOutcome::unavailable()], 1070 ) 1071 .unwrap() 1072 } 1073 1074 #[test] 1075 fn claim_projection_must_match_every_original_binding_before_transport() { 1076 for mutation in [ 1077 ReceiptMutation::Identity, 1078 ReceiptMutation::Plan, 1079 ReceiptMutation::Artifact, 1080 ReceiptMutation::Request, 1081 ReceiptMutation::Created, 1082 ReceiptMutation::Claim, 1083 ReceiptMutation::CommitTime, 1084 ReceiptMutation::Attempts, 1085 ] { 1086 let storage = Arc::new(FaultStorage::new(171)); 1087 let sink = Arc::new(HeldSink::default()); 1088 let engine = fault_engine( 1089 storage.clone(), 1090 Arc::new(MockSigner::new(SignBehavior::Success { 1091 completed_at_unix_ms: 1_800_000_200_500, 1092 })), 1093 sink.clone(), 1094 ); 1095 let push = request(171, "wss://relay.example"); 1096 execute_to_admitted(&engine, &push); 1097 *storage.receipt_mutation.lock().unwrap() = Some(mutation); 1098 assert_eq!( 1099 block_on(engine.deliver_push(push.operation_id())), 1100 Err(Error::StorageFailed), 1101 "{mutation:?}" 1102 ); 1103 assert_eq!(sink.calls.load(Ordering::Relaxed), 0); 1104 } 1105 } 1106 1107 #[test] 1108 fn stop_projection_must_match_identity_request_time_and_revision() { 1109 for mutation in [ 1110 ReceiptMutation::Identity, 1111 ReceiptMutation::Plan, 1112 ReceiptMutation::Artifact, 1113 ReceiptMutation::Request, 1114 ReceiptMutation::StopTime, 1115 ReceiptMutation::Revision, 1116 ] { 1117 let storage = Arc::new(FaultStorage::new(172)); 1118 let engine = fault_engine( 1119 storage.clone(), 1120 Arc::new(MockSigner::new(SignBehavior::Success { 1121 completed_at_unix_ms: 1_800_000_200_500, 1122 })), 1123 Arc::new(MockSink), 1124 ); 1125 let push = request(172, "wss://relay.example"); 1126 execute_to_admitted(&engine, &push); 1127 *storage.receipt_mutation.lock().unwrap() = Some(mutation); 1128 assert_eq!( 1129 block_on(engine.cancel_push(push.operation_id())), 1130 Err(Error::StorageFailed), 1131 "{mutation:?}" 1132 ); 1133 let current = block_on(engine.push_status(push.operation_id())) 1134 .unwrap() 1135 .unwrap(); 1136 assert!(current.delivery_history().proves_no_issued_attempt()); 1137 assert!( 1138 current 1139 .delivery_plan() 1140 .stop_requested_at_unix_ms() 1141 .is_some() 1142 ); 1143 } 1144 } 1145 1146 #[test] 1147 fn claim_projection_cannot_replace_retained_retry_state_or_backoff() { 1148 for mutation in [ReceiptMutation::State, ReceiptMutation::Retry] { 1149 let storage = Arc::new(FaultStorage::new(173)); 1150 let clock = Arc::new(TestClock(AtomicU64::new(1_800_000_200_000))); 1151 let sink = Arc::new(ScriptedSink::new([DeliveryBehavior::Failure( 1152 Retryability::Retryable, 1153 )])); 1154 let engine = Engine::builder( 1155 storage.clone(), 1156 clock.clone(), 1157 Arc::new(TestIds(AtomicU64::new(10))), 1158 DeadlinePolicy::new(10_000, 10_000, 10_000).unwrap(), 1159 ) 1160 .sink(sink.clone()) 1161 .signer(Arc::new(MockSigner::new(SignBehavior::Success { 1162 completed_at_unix_ms: 1_800_000_200_500, 1163 }))) 1164 .build() 1165 .unwrap(); 1166 let push = request(173, "wss://relay.example"); 1167 execute_to_admitted(&engine, &push); 1168 let first = block_on(engine.deliver_push(push.operation_id())).unwrap(); 1169 clock.0.store( 1170 first.plan().retry().unwrap().not_before_unix_ms(), 1171 Ordering::Relaxed, 1172 ); 1173 *storage.receipt_mutation.lock().unwrap() = Some(mutation); 1174 assert_eq!( 1175 block_on(engine.deliver_push(push.operation_id())), 1176 Err(Error::StorageFailed), 1177 "{mutation:?}" 1178 ); 1179 assert_eq!(sink.requests.lock().unwrap().len(), 1); 1180 } 1181 } 1182 1183 #[derive(Clone, Copy)] 1184 pub(super) enum HistoryMutation { 1185 Plan, 1186 Artifact, 1187 MissingFact, 1188 MissingStop, 1189 } 1190 pub(super) fn mutate_history( 1191 history: AuthoredDeliveryHistory, 1192 mutation: HistoryMutation, 1193 ) -> AuthoredDeliveryHistory { 1194 let mut wire = serde_json::to_value(history.plan()).unwrap(); 1195 match mutation { 1196 HistoryMutation::Plan => wire["plan_id"] = serde_json::json!(vec![244; 16]), 1197 HistoryMutation::Artifact => wire["artifact_id"] = serde_json::json!(vec![245; 16]), 1198 HistoryMutation::MissingFact => wire["delivery_facts"] = serde_json::json!([]), 1199 HistoryMutation::MissingStop => { 1200 wire["stop_requested_at_unix_ms"] = serde_json::Value::Null; 1201 wire["state"] = serde_json::json!("pending"); 1202 } 1203 } 1204 AuthoredDeliveryHistory::new(serde_json::from_value(wire).unwrap(), None).unwrap() 1205 } 1206 1207 #[test] 1208 fn wrong_history_and_lost_committed_fact_or_stop_never_become_success() { 1209 for (nth, mutation, stop) in [ 1210 (1, HistoryMutation::Plan, false), 1211 (1, HistoryMutation::Artifact, false), 1212 (2, HistoryMutation::Plan, false), 1213 (3, HistoryMutation::MissingFact, false), 1214 (2, HistoryMutation::MissingStop, true), 1215 ] { 1216 let storage = Arc::new(FaultStorage::new(174)); 1217 let sink = Arc::new(ScriptedSink::new([DeliveryBehavior::Outcomes(vec![ 1218 DeliveryOutcome::accepted(), 1219 ])])); 1220 let engine = fault_engine( 1221 storage.clone(), 1222 Arc::new(MockSigner::new(SignBehavior::Success { 1223 completed_at_unix_ms: 1_800_000_200_500, 1224 })), 1225 sink.clone(), 1226 ); 1227 let push = request(174, "wss://relay.example"); 1228 execute_to_admitted(&engine, &push); 1229 *storage.history_mutation.lock().unwrap() = Some((nth, mutation)); 1230 if stop { 1231 assert_eq!( 1232 block_on(engine.cancel_push(push.operation_id())), 1233 Err(Error::StorageFailed) 1234 ); 1235 } else { 1236 assert_eq!( 1237 block_on(engine.deliver_push(push.operation_id())), 1238 Err(Error::StorageFailed) 1239 ); 1240 } 1241 let actual = block_on(engine.push_status(push.operation_id())) 1242 .unwrap() 1243 .unwrap(); 1244 assert_eq!( 1245 actual.delivery_plan().delivery_facts().len(), 1246 usize::from(nth == 3) 1247 ); 1248 } 1249 } 1250 1251 #[test] 1252 fn clock_expiry_before_sink_and_future_fact_before_reconciliation_fail_closed() { 1253 struct ExpiringClock(AtomicUsize); 1254 impl Clock for ExpiringClock { 1255 fn now_unix_ms(&self) -> Result<u64, Error> { 1256 Ok(1_800_000_220_000 + 10_000 * self.0.fetch_add(1, Ordering::Relaxed) as u64) 1257 } 1258 } 1259 let (engine, storage, _, sink, push) = setup(175); 1260 let boundary = Engine::builder( 1261 storage, 1262 Arc::new(ExpiringClock(AtomicUsize::new(0))), 1263 Arc::new(TestIds(AtomicU64::new(230))), 1264 DeadlinePolicy::new(10_000, 10_000, 10_000).unwrap(), 1265 ) 1266 .sink(sink.clone()) 1267 .build() 1268 .unwrap(); 1269 assert_eq!( 1270 block_on(boundary.deliver_push(push.operation_id())), 1271 Err(Error::WorkClaimConflict) 1272 ); 1273 assert_eq!(sink.calls.load(Ordering::Relaxed), 0); 1274 assert_eq!( 1275 block_on(engine.deliver_push(SyncId::new([231; 16]).unwrap())), 1276 Err(Error::StorageFailed) 1277 ); 1278 let (engine, storage, clock, sink, push) = setup(176); 1279 let mut future = Box::pin(engine.deliver_push(push.operation_id())); 1280 assert!( 1281 future 1282 .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref())) 1283 .is_pending() 1284 ); 1285 let (request, _sender) = sink.take(); 1286 drop(future); 1287 let status = block_on(engine.push_status(push.operation_id())) 1288 .unwrap() 1289 .unwrap(); 1290 let plan = status.delivery_plan(); 1291 let claim = plan.claim_evidence().unwrap(); 1292 let at = claim.expires_at_unix_ms(); 1293 block_on( 1294 storage.execute_authored(AuthoredAtomicCommand::RecordDelivery( 1295 RecordDeliveryFact::new( 1296 plan.plan_id(), 1297 plan.artifact_id(), 1298 claim.clone(), 1299 DeliveryAttemptOutcome::Receipt( 1300 receipt(&request, vec![DeliveryOutcome::accepted()]).unwrap(), 1301 ), 1302 at + 1000, 1303 ) 1304 .unwrap(), 1305 )), 1306 ) 1307 .unwrap(); 1308 clock.0.store(at, Ordering::Relaxed); 1309 assert_eq!( 1310 block_on(engine.deliver_push(push.operation_id())), 1311 Err(Error::ClockUnavailable) 1312 ); 1313 let current = block_on(engine.push_status(push.operation_id())) 1314 .unwrap() 1315 .unwrap(); 1316 assert_eq!(current.delivery_plan().attempt_count(), 0); 1317 assert_eq!(current.delivery_plan().delivery_facts().len(), 1); 1318 } 1319 1320 #[test] 1321 fn unbound_preparation_and_stop_keep_honest_retry_and_no_issued_proof() { 1322 let signer = Arc::new(MockSigner::new(SignBehavior::Success { 1323 completed_at_unix_ms: 1_800_000_200_500, 1324 })); 1325 let (engine, _) = setup_engine(signer); 1326 let push = request(177, "wss://relay.example"); 1327 block_on(engine.prepare_push(push.clone())).unwrap(); 1328 let status = block_on(engine.push_status(push.operation_id())) 1329 .unwrap() 1330 .unwrap(); 1331 assert!(status.delivery_plan().request().is_none()); 1332 assert!(status.delivery_history().proves_no_issued_attempt()); 1333 assert_eq!( 1334 engine 1335 .retry_decision(status.delivery_plan(), 1_800_000_200_001) 1336 .unwrap(), 1337 SyncRetryDecision::Ready 1338 ); 1339 let stopped = block_on(engine.cancel_push(push.operation_id())).unwrap(); 1340 assert!( 1341 stopped 1342 .status() 1343 .delivery_history() 1344 .proves_no_issued_attempt() 1345 ); 1346 assert_eq!( 1347 engine 1348 .retry_decision(stopped.status().delivery_plan(), 1_800_000_200_010) 1349 .unwrap(), 1350 SyncRetryDecision::Exhausted 1351 ); 1352 }