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