authored_atomic.rs (32584B)
1 use core::num::{NonZeroU32, NonZeroU64}; 2 use futures_executor::block_on; 3 use radroots_event::{GenericEventDraft, SignedEvent, wire::v1::Nip01EventWire}; 4 use radroots_event_codec::authoring::AuthoredEventPlan; 5 use radroots_storage::{ 6 Error, 7 atomic::{AtomicCommitDigest, AtomicCommitDisposition}, 8 authored::{ 9 AdmissionState, AuthoredArtifact, AuthoredArtifactId, AuthoredOperation, FailureClass, 10 RetrySchedule, SigningState, WorkClaim, WorkFailure, WorkPhase, 11 }, 12 authored_atomic::{ 13 ApplyAdmissionResult, ApplyDeliveryAttempt, ApplySignedArtifact, ApplyWorkFailure, 14 AuthoredAtomicCommand, AuthoredAtomicOutcome, AuthoredAtomicReceipt, AuthoredAtomicStorage, 15 AuthoredWorkTarget, CancelAuthoredTarget, CancelAuthoredWork, ClaimAuthoredTarget, 16 ClaimAuthoredWork, PrepareAuthoredOperation, WorkFence, 17 }, 18 authored_delivery::{ 19 AuthoredDeliveryIntent, AuthoredDeliveryPlan, AuthoredDeliveryPlanId, 20 AuthoredDeliveryState, DeliveryAttemptOutcome, 21 }, 22 event::SourceGeneration, 23 journal::OperationInstanceId, 24 memory::MemoryStorage, 25 }; 26 use radroots_transport::{ 27 DeliveryReceipt, SinkFailure, Target, TargetSet, 28 outcome::{DeliveryOutcome, Retryability}, 29 policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy}, 30 sink::DeliveryTargetReceipt, 31 }; 32 33 const AUTHOR: &str = "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df"; 34 35 fn authored_plan() -> AuthoredEventPlan { 36 AuthoredEventPlan::from_generic( 37 GenericEventDraft::new( 38 "radroots.social.geochat.v1", 39 20_000, 40 1_800_000_100, 41 Vec::new(), 42 "atomic authored plan", 43 AUTHOR, 44 ) 45 .expect("draft"), 46 ) 47 .expect("plan") 48 } 49 50 fn signed(plan: &AuthoredEventPlan) -> SignedEvent { 51 let wire = Nip01EventWire { 52 id: plan.expected_event_id().to_hex(), 53 pubkey: plan.author().to_hex(), 54 created_at: plan.created_at(), 55 kind: plan.body().kind(), 56 tags: plan.body().tags().to_vec(), 57 content: plan.body().content().to_owned(), 58 sig: "dd".repeat(64), 59 extra: Default::default(), 60 }; 61 let raw = serde_json::to_string(&wire).expect("raw event"); 62 SignedEvent::from_wire_verified_id(wire, raw).expect("signed event") 63 } 64 65 fn ids() -> ( 66 OperationInstanceId, 67 AuthoredArtifactId, 68 AuthoredDeliveryPlanId, 69 ) { 70 ( 71 OperationInstanceId::new([1; 16]).expect("operation"), 72 AuthoredArtifactId::new([2; 16]).expect("artifact"), 73 AuthoredDeliveryPlanId::new([3; 16]).expect("delivery"), 74 ) 75 } 76 77 fn intent() -> AuthoredDeliveryIntent { 78 AuthoredDeliveryIntent::new( 79 "atomic-delivery", 80 TargetSet::new(vec![ 81 Target::nostr_relay("wss://one.example").expect("one"), 82 Target::nostr_relay("wss://two.example").expect("two"), 83 ]) 84 .expect("targets"), 85 SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::any()), 86 100, 87 ) 88 .expect("intent") 89 } 90 91 fn prepare(input: u8) -> (AuthoredAtomicCommand, AuthoredEventPlan) { 92 let (operation_id, artifact_id, plan_id) = ids(); 93 let authored_plan = authored_plan(); 94 let artifact = AuthoredArtifact::planned(artifact_id, operation_id, 0, &authored_plan, 10) 95 .expect("artifact"); 96 let operation = AuthoredOperation::new(operation_id, vec![artifact_id], 10).expect("operation"); 97 let delivery = 98 AuthoredDeliveryPlan::new(plan_id, artifact_id, intent(), 10).expect("delivery plan"); 99 let prepare = PrepareAuthoredOperation::new( 100 operation, 101 vec![artifact], 102 vec![delivery], 103 AtomicCommitDigest::new([input; 32]), 104 10, 105 ) 106 .expect("prepare"); 107 (AuthoredAtomicCommand::Prepare(prepare), authored_plan) 108 } 109 110 fn claim( 111 target: ClaimAuthoredTarget, 112 revision: NonZeroU64, 113 token: u8, 114 at: u64, 115 ) -> (AuthoredAtomicCommand, WorkClaim) { 116 let claim = WorkClaim::new( 117 [token; 16], 118 format!("worker-{token}"), 119 NonZeroU64::new(u64::from(token)).expect("generation"), 120 at, 121 at + 20, 122 revision, 123 ) 124 .expect("claim"); 125 ( 126 AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new(target, claim.clone())), 127 claim, 128 ) 129 } 130 131 fn fence(claim: &WorkClaim) -> WorkFence { 132 WorkFence::new(*claim.token(), claim.generation(), claim.row_revision()).unwrap() 133 } 134 135 fn prepared_storage() -> (MemoryStorage, AuthoredEventPlan) { 136 let storage = MemoryStorage::new(SourceGeneration::new([1; 32]).unwrap()); 137 let (prepare, plan) = prepare(7); 138 block_on(storage.execute_authored(prepare)).unwrap(); 139 (storage, plan) 140 } 141 142 fn signed_storage() -> MemoryStorage { 143 let (storage, plan) = prepared_storage(); 144 let artifact = block_on(storage.authored_artifact(ids().1)) 145 .unwrap() 146 .unwrap(); 147 let (command, active) = claim( 148 ClaimAuthoredTarget::ArtifactSigning(ids().1), 149 artifact.revision(), 150 4, 151 11, 152 ); 153 block_on(storage.execute_authored(command)).unwrap(); 154 let apply = ApplySignedArtifact::new(ids().1, fence(&active), signed(&plan), 12).unwrap(); 155 block_on(storage.execute_authored(AuthoredAtomicCommand::ApplySigned(apply))).unwrap(); 156 storage 157 } 158 159 fn claim_admission(storage: &MemoryStorage, token: u8, at: u64) -> WorkClaim { 160 let artifact = block_on(storage.authored_artifact(ids().1)) 161 .unwrap() 162 .unwrap(); 163 let (command, active) = claim( 164 ClaimAuthoredTarget::ArtifactAdmission(ids().1), 165 artifact.revision(), 166 token, 167 at, 168 ); 169 block_on(storage.execute_authored(command)).unwrap(); 170 active 171 } 172 173 fn claim_delivery(storage: &MemoryStorage, token: u8, at: u64) -> WorkClaim { 174 let plan = block_on(storage.authored_delivery_plan(ids().2)) 175 .unwrap() 176 .unwrap(); 177 let (command, active) = claim( 178 ClaimAuthoredTarget::DeliveryPlan(ids().2), 179 plan.revision(), 180 token, 181 at, 182 ); 183 block_on(storage.execute_authored(command)).unwrap(); 184 active 185 } 186 187 #[test] 188 fn preparation_is_atomic_deterministic_and_exactly_replayable() { 189 let storage = MemoryStorage::new(SourceGeneration::new([1; 32]).expect("generation")); 190 let (command, _) = prepare(7); 191 assert_eq!(command.commit_id(), command.clone().commit_id()); 192 assert_eq!(command.digest(), command.clone().digest()); 193 194 let committed = block_on(storage.execute_authored(command.clone())).expect("commit"); 195 assert_eq!(committed.disposition(), AtomicCommitDisposition::Committed); 196 let replay = block_on(storage.execute_authored(command.clone())).expect("replay"); 197 assert_eq!(replay.disposition(), AtomicCommitDisposition::Replay); 198 assert_eq!(replay.outcome(), committed.outcome()); 199 assert_eq!( 200 block_on(storage.authored_receipt(command.commit_id())) 201 .expect("receipt") 202 .expect("stored receipt") 203 .outcome(), 204 committed.outcome() 205 ); 206 207 let (conflict, _) = prepare(8); 208 assert_eq!( 209 block_on(storage.execute_authored(conflict)), 210 Err(Error::AtomicCommitConflict) 211 ); 212 let operation = block_on(storage.authored_operation(ids().0)) 213 .expect("operation query") 214 .expect("operation"); 215 assert_eq!(operation.artifact_ids(), &[ids().1]); 216 let delivery = block_on(storage.authored_delivery_plan(ids().2)) 217 .expect("delivery query") 218 .expect("delivery"); 219 assert!(delivery.request().is_none()); 220 } 221 222 #[test] 223 fn signing_atomically_binds_exact_delivery_requests_and_rejects_stale_fences() { 224 let storage = MemoryStorage::new(SourceGeneration::new([1; 32]).expect("generation")); 225 let (prepare, plan) = prepare(7); 226 block_on(storage.execute_authored(prepare)).expect("prepare"); 227 let artifact = block_on(storage.authored_artifact(ids().1)) 228 .expect("artifact query") 229 .expect("artifact"); 230 let (claim_command, active) = claim( 231 ClaimAuthoredTarget::ArtifactSigning(ids().1), 232 artifact.revision(), 233 4, 234 11, 235 ); 236 block_on(storage.execute_authored(claim_command)).expect("claim"); 237 238 let stale = ApplySignedArtifact::new( 239 ids().1, 240 WorkFence::new([8; 16], active.generation(), active.row_revision()).expect("fence"), 241 signed(&plan), 242 12, 243 ) 244 .expect("stale apply"); 245 assert_eq!( 246 block_on(storage.execute_authored(AuthoredAtomicCommand::ApplySigned(stale))), 247 Err(Error::DeliveryPlanClaimConflict) 248 ); 249 assert_eq!( 250 block_on(storage.authored_artifact(ids().1)) 251 .expect("artifact query") 252 .expect("artifact") 253 .signing_state(), 254 SigningState::Planned 255 ); 256 257 let apply = ApplySignedArtifact::new( 258 ids().1, 259 WorkFence::new(*active.token(), active.generation(), active.row_revision()).expect("fence"), 260 signed(&plan), 261 12, 262 ) 263 .expect("apply"); 264 let receipt = block_on(storage.execute_authored(AuthoredAtomicCommand::ApplySigned(apply))) 265 .expect("signed"); 266 assert!(matches!( 267 receipt.outcome(), 268 AuthoredAtomicOutcome::Artifact(_) 269 )); 270 let artifact = block_on(storage.authored_artifact(ids().1)) 271 .expect("artifact query") 272 .expect("artifact"); 273 assert_eq!(artifact.signing_state(), SigningState::Signed); 274 let delivery = block_on(storage.authored_delivery_plan(ids().2)) 275 .expect("delivery query") 276 .expect("delivery"); 277 assert_eq!( 278 delivery 279 .request() 280 .expect("bound delivery") 281 .payload() 282 .event() 283 .raw_json(), 284 artifact 285 .signed() 286 .expect("signed artifact") 287 .event() 288 .raw_json() 289 ); 290 } 291 292 #[test] 293 fn delivery_attempt_and_work_failure_commands_preserve_atomic_state() { 294 let storage = MemoryStorage::new(SourceGeneration::new([1; 32]).expect("generation")); 295 let (prepare, plan) = prepare(7); 296 block_on(storage.execute_authored(prepare)).expect("prepare"); 297 let artifact = block_on(storage.authored_artifact(ids().1)) 298 .expect("artifact") 299 .expect("artifact"); 300 let (sign_claim, active) = claim( 301 ClaimAuthoredTarget::ArtifactSigning(ids().1), 302 artifact.revision(), 303 4, 304 11, 305 ); 306 block_on(storage.execute_authored(sign_claim)).expect("claim signing"); 307 block_on( 308 storage.execute_authored(AuthoredAtomicCommand::ApplySigned( 309 ApplySignedArtifact::new( 310 ids().1, 311 WorkFence::new(*active.token(), active.generation(), active.row_revision()) 312 .expect("fence"), 313 signed(&plan), 314 12, 315 ) 316 .expect("apply"), 317 )), 318 ) 319 .expect("sign"); 320 321 let delivery = block_on(storage.authored_delivery_plan(ids().2)) 322 .expect("delivery") 323 .expect("delivery"); 324 let (delivery_claim, active) = claim( 325 ClaimAuthoredTarget::DeliveryPlan(ids().2), 326 delivery.revision(), 327 5, 328 13, 329 ); 330 block_on(storage.execute_authored(delivery_claim)).expect("claim delivery"); 331 let request = block_on(storage.authored_delivery_plan(ids().2)) 332 .expect("delivery") 333 .expect("delivery") 334 .request() 335 .expect("bound request") 336 .clone(); 337 let receipt = DeliveryReceipt::for_request( 338 &request, 339 request 340 .target_set() 341 .targets() 342 .iter() 343 .cloned() 344 .map(|target| DeliveryTargetReceipt::attempted(target, DeliveryOutcome::accepted())) 345 .collect(), 346 ) 347 .expect("receipt"); 348 let apply = ApplyDeliveryAttempt::new( 349 ids().2, 350 WorkFence::new(*active.token(), active.generation(), active.row_revision()).expect("fence"), 351 DeliveryAttemptOutcome::Receipt(receipt), 352 None, 353 14, 354 ) 355 .expect("apply delivery"); 356 block_on(storage.execute_authored(AuthoredAtomicCommand::ApplyDelivery(apply))) 357 .expect("deliver"); 358 assert_eq!( 359 block_on(storage.authored_delivery_plan(ids().2)) 360 .expect("delivery") 361 .expect("delivery") 362 .state(), 363 AuthoredDeliveryState::Satisfied 364 ); 365 366 let retry_failure = WorkFailure::new( 367 "temporary_signer_failure", 368 WorkPhase::Signing, 369 FailureClass::Retryable, 370 Some(30), 371 None, 372 ) 373 .expect("failure"); 374 let retry = RetrySchedule::new(NonZeroU32::MIN, 30, retry_failure.clone()).expect("retry"); 375 let invalid = radroots_storage::authored_atomic::ApplyWorkFailure::new( 376 AuthoredWorkTarget::Artifact(ids().1), 377 WorkFence::new([9; 16], NonZeroU64::MIN, NonZeroU64::MIN).expect("fence"), 378 retry_failure, 379 Some(retry), 380 15, 381 ) 382 .expect("failure command"); 383 let before = block_on(storage.authored_artifact(ids().1)) 384 .expect("artifact") 385 .expect("artifact"); 386 assert!( 387 block_on(storage.execute_authored(AuthoredAtomicCommand::ApplyFailure(invalid))).is_err() 388 ); 389 assert_eq!( 390 block_on(storage.authored_artifact(ids().1)) 391 .expect("artifact") 392 .expect("artifact"), 393 before 394 ); 395 } 396 397 #[test] 398 fn authored_atomic_command_models_cover_every_phase_identity_and_durable_outcome() { 399 assert_eq!( 400 WorkFence::new([0; 16], NonZeroU64::MIN, NonZeroU64::MIN), 401 Err(Error::InvalidWorkClaim) 402 ); 403 let fence = 404 WorkFence::new([7; 16], NonZeroU64::new(2).unwrap(), NonZeroU64::MIN).expect("fence"); 405 assert_eq!(fence.token(), &[7; 16]); 406 assert_eq!(fence.generation(), NonZeroU64::new(2).unwrap()); 407 assert_eq!(fence.row_revision(), NonZeroU64::MIN); 408 409 let (prepared_command, authored_plan) = prepare(7); 410 let AuthoredAtomicCommand::Prepare(prepared) = prepared_command.clone() else { 411 unreachable!() 412 }; 413 assert_eq!(prepared.operation().operation_id(), ids().0); 414 assert_eq!(prepared.artifacts().len(), 1); 415 assert_eq!(prepared.delivery_plans().len(), 1); 416 assert_eq!(prepared.input_digest(), AtomicCommitDigest::new([7; 32])); 417 assert_eq!(prepared.requested_at_unix_ms(), 10); 418 419 assert_eq!( 420 PrepareAuthoredOperation::new( 421 prepared.operation().clone(), 422 Vec::new(), 423 Vec::new(), 424 AtomicCommitDigest::new([1; 32]), 425 10, 426 ), 427 Err(Error::AtomicWorkflowMismatch) 428 ); 429 assert_eq!( 430 PrepareAuthoredOperation::new( 431 prepared.operation().clone(), 432 prepared.artifacts().to_vec(), 433 prepared.delivery_plans().to_vec(), 434 AtomicCommitDigest::new([1; 32]), 435 0, 436 ), 437 Err(Error::AtomicWorkflowMismatch) 438 ); 439 440 let work_claim = 441 WorkClaim::new([8; 16], "worker", NonZeroU64::MIN, 11, 20, NonZeroU64::MIN).unwrap(); 442 let claim_commands = [ 443 AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( 444 ClaimAuthoredTarget::ArtifactSigning(ids().1), 445 work_claim.clone(), 446 )), 447 AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( 448 ClaimAuthoredTarget::ArtifactAdmission(ids().1), 449 work_claim.clone(), 450 )), 451 AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( 452 ClaimAuthoredTarget::DeliveryPlan(ids().2), 453 work_claim.clone(), 454 )), 455 ]; 456 for command in &claim_commands { 457 let AuthoredAtomicCommand::Claim(value) = command else { 458 unreachable!() 459 }; 460 assert_eq!(value.claim(), &work_claim); 461 let _ = value.target(); 462 } 463 464 let apply_signed = ApplySignedArtifact::new(ids().1, fence.clone(), signed(&authored_plan), 12) 465 .expect("signed command"); 466 assert_eq!(apply_signed.artifact_id(), ids().1); 467 assert_eq!(apply_signed.fence(), &fence); 468 assert_eq!(apply_signed.event().id(), authored_plan.expected_event_id()); 469 assert_eq!(apply_signed.applied_at_unix_ms(), 12); 470 assert_eq!( 471 ApplySignedArtifact::new(ids().1, fence.clone(), signed(&authored_plan), 0), 472 Err(Error::AtomicWorkflowMismatch) 473 ); 474 475 let retry_failure = WorkFailure::new( 476 "temporary_admission", 477 WorkPhase::Admission, 478 FailureClass::Retryable, 479 Some(30), 480 Some("try another worker".to_owned()), 481 ) 482 .unwrap(); 483 let retry = RetrySchedule::new(NonZeroU32::MIN, 30, retry_failure.clone()).unwrap(); 484 let apply_admission = ApplyAdmissionResult::new( 485 ids().1, 486 fence.clone(), 487 AdmissionState::Retryable, 488 Some(retry_failure.clone()), 489 Some(retry.clone()), 490 13, 491 ) 492 .unwrap(); 493 assert_eq!(apply_admission.artifact_id(), ids().1); 494 assert_eq!(apply_admission.fence(), &fence); 495 assert_eq!(apply_admission.state(), AdmissionState::Retryable); 496 assert_eq!(apply_admission.failure(), Some(&retry_failure)); 497 assert_eq!(apply_admission.retry(), Some(&retry)); 498 assert_eq!(apply_admission.applied_at_unix_ms(), 13); 499 assert_eq!( 500 ApplyAdmissionResult::new( 501 ids().1, 502 fence.clone(), 503 AdmissionState::Inserted, 504 None, 505 None, 506 0, 507 ), 508 Err(Error::AtomicWorkflowMismatch) 509 ); 510 511 let request = intent() 512 .materialize(radroots_transport::sink::DeliveryPayload::new(signed( 513 &authored_plan, 514 ))) 515 .unwrap(); 516 let receipt = DeliveryReceipt::for_request( 517 &request, 518 request 519 .target_set() 520 .targets() 521 .iter() 522 .cloned() 523 .map(|target| { 524 DeliveryTargetReceipt::attempted( 525 target, 526 DeliveryOutcome::accepted() 527 .with_detail("accepted", "accepted by relay") 528 .unwrap(), 529 ) 530 }) 531 .collect(), 532 ) 533 .unwrap(); 534 let apply_delivery = ApplyDeliveryAttempt::new( 535 ids().2, 536 fence.clone(), 537 DeliveryAttemptOutcome::Receipt(receipt.clone()), 538 None, 539 14, 540 ) 541 .unwrap(); 542 assert_eq!(apply_delivery.plan_id(), ids().2); 543 assert_eq!(apply_delivery.fence(), &fence); 544 assert!(matches!( 545 apply_delivery.outcome(), 546 DeliveryAttemptOutcome::Receipt(_) 547 )); 548 assert_eq!(apply_delivery.retry(), None); 549 assert_eq!(apply_delivery.applied_at_unix_ms(), 14); 550 assert_eq!( 551 ApplyDeliveryAttempt::new( 552 ids().2, 553 fence.clone(), 554 DeliveryAttemptOutcome::Receipt(receipt), 555 None, 556 0, 557 ), 558 Err(Error::AtomicWorkflowMismatch) 559 ); 560 561 let sink_failure = SinkFailure::for_request( 562 &request, 563 "relay_unavailable", 564 Retryability::Retryable, 565 Some(30), 566 None, 567 Vec::new(), 568 ) 569 .unwrap(); 570 let sink_command = AuthoredAtomicCommand::ApplyDelivery( 571 ApplyDeliveryAttempt::new( 572 ids().2, 573 fence.clone(), 574 DeliveryAttemptOutcome::SinkFailure(sink_failure), 575 Some(retry.clone()), 576 14, 577 ) 578 .unwrap(), 579 ); 580 581 let failure_commands = [ 582 (AuthoredWorkTarget::Artifact(ids().1), WorkPhase::Signing), 583 (AuthoredWorkTarget::Artifact(ids().1), WorkPhase::Admission), 584 ( 585 AuthoredWorkTarget::DeliveryPlan(ids().2), 586 WorkPhase::Delivery, 587 ), 588 ] 589 .map(|(target, phase)| { 590 let failure = WorkFailure::new( 591 "terminal_failure", 592 phase, 593 FailureClass::Terminal, 594 None, 595 Some("terminal".to_owned()), 596 ) 597 .unwrap(); 598 AuthoredAtomicCommand::ApplyFailure( 599 ApplyWorkFailure::new(target, fence.clone(), failure, None, 15).unwrap(), 600 ) 601 }); 602 for command in &failure_commands { 603 let AuthoredAtomicCommand::ApplyFailure(value) = command else { 604 unreachable!() 605 }; 606 let _ = value.target(); 607 assert_eq!(value.fence(), &fence); 608 assert_eq!(value.failure().class(), FailureClass::Terminal); 609 assert_eq!(value.retry(), None); 610 assert_eq!(value.applied_at_unix_ms(), 15); 611 } 612 assert_eq!( 613 ApplyWorkFailure::new( 614 AuthoredWorkTarget::Artifact(ids().1), 615 fence.clone(), 616 retry_failure, 617 Some(retry), 618 0, 619 ), 620 Err(Error::AtomicWorkflowMismatch) 621 ); 622 623 let cancel_commands = [ 624 CancelAuthoredTarget::ArtifactSigning(ids().1), 625 CancelAuthoredTarget::ArtifactAdmission(ids().1), 626 CancelAuthoredTarget::DeliveryPlan(ids().2), 627 ] 628 .map(|target| { 629 AuthoredAtomicCommand::Cancel(CancelAuthoredWork::new(target, NonZeroU64::MIN, 16).unwrap()) 630 }); 631 for command in &cancel_commands { 632 let AuthoredAtomicCommand::Cancel(value) = command else { 633 unreachable!() 634 }; 635 let _ = value.target(); 636 assert_eq!(value.expected_revision(), NonZeroU64::MIN); 637 assert_eq!(value.cancelled_at_unix_ms(), 16); 638 } 639 assert_eq!( 640 CancelAuthoredWork::new( 641 CancelAuthoredTarget::ArtifactSigning(ids().1), 642 NonZeroU64::MIN, 643 0, 644 ), 645 Err(Error::AtomicWorkflowMismatch) 646 ); 647 648 let mut commands = vec![ 649 prepared_command, 650 AuthoredAtomicCommand::ApplySigned(apply_signed), 651 AuthoredAtomicCommand::ApplyAdmission(apply_admission), 652 AuthoredAtomicCommand::ApplyDelivery(apply_delivery), 653 sink_command, 654 ]; 655 commands.extend(claim_commands); 656 commands.extend(failure_commands); 657 commands.extend(cancel_commands); 658 for command in commands { 659 assert_ne!(command.commit_id().as_bytes(), &[0; 16]); 660 assert_ne!(command.digest().as_bytes(), &[0; 32]); 661 assert_ne!(command.requested_at_unix_ms(), 0); 662 } 663 664 let outcome = AuthoredAtomicOutcome::Prepared { 665 operation: prepared.operation().clone(), 666 artifacts: prepared.artifacts().to_vec(), 667 delivery_plans: prepared.delivery_plans().to_vec(), 668 }; 669 let command = AuthoredAtomicCommand::Prepare(prepared.clone()); 670 assert_eq!( 671 AuthoredAtomicReceipt::new( 672 &command, 673 AtomicCommitDisposition::Committed, 674 9, 675 outcome.clone(), 676 ), 677 Err(Error::AtomicWorkflowMismatch) 678 ); 679 let receipt = AuthoredAtomicReceipt::new( 680 &command, 681 AtomicCommitDisposition::Committed, 682 10, 683 outcome.clone(), 684 ) 685 .unwrap(); 686 assert_eq!(receipt.commit_id(), command.commit_id()); 687 assert_eq!(receipt.digest(), command.digest()); 688 assert_eq!(receipt.disposition(), AtomicCommitDisposition::Committed); 689 assert_eq!(receipt.committed_at_unix_ms(), 10); 690 assert_eq!(receipt.outcome(), &outcome); 691 assert_eq!( 692 AuthoredAtomicReceipt::from_durable_parts( 693 receipt.commit_id(), 694 receipt.digest(), 695 AtomicCommitDisposition::Replay, 696 0, 697 outcome.clone(), 698 ), 699 Err(Error::AtomicWorkflowMismatch) 700 ); 701 assert!( 702 AuthoredAtomicReceipt::from_durable_parts( 703 receipt.commit_id(), 704 receipt.digest(), 705 AtomicCommitDisposition::Replay, 706 11, 707 outcome, 708 ) 709 .is_ok() 710 ); 711 let invalid_outcome = AuthoredAtomicOutcome::Prepared { 712 operation: prepared.operation().clone(), 713 artifacts: Vec::new(), 714 delivery_plans: Vec::new(), 715 }; 716 assert_eq!( 717 AuthoredAtomicReceipt::from_durable_parts( 718 receipt.commit_id(), 719 receipt.digest(), 720 AtomicCommitDisposition::Replay, 721 11, 722 invalid_outcome, 723 ), 724 Err(Error::AtomicWorkflowMismatch) 725 ); 726 } 727 728 #[test] 729 fn memory_executes_admission_results_and_every_authored_failure_phase() { 730 let storage = signed_storage(); 731 let active = claim_admission(&storage, 5, 13); 732 let inserted = ApplyAdmissionResult::new( 733 ids().1, 734 fence(&active), 735 AdmissionState::Inserted, 736 None, 737 None, 738 14, 739 ) 740 .unwrap(); 741 block_on(storage.execute_authored(AuthoredAtomicCommand::ApplyAdmission(inserted))).unwrap(); 742 assert_eq!( 743 block_on(storage.authored_artifact(ids().1)) 744 .unwrap() 745 .unwrap() 746 .admission_state(), 747 AdmissionState::Inserted 748 ); 749 750 let (storage, _) = prepared_storage(); 751 let artifact = block_on(storage.authored_artifact(ids().1)) 752 .unwrap() 753 .unwrap(); 754 let (command, active) = claim( 755 ClaimAuthoredTarget::ArtifactSigning(ids().1), 756 artifact.revision(), 757 4, 758 11, 759 ); 760 block_on(storage.execute_authored(command)).unwrap(); 761 let failure = WorkFailure::new( 762 "terminal_signer_failure", 763 WorkPhase::Signing, 764 FailureClass::Terminal, 765 None, 766 None, 767 ) 768 .unwrap(); 769 let command = ApplyWorkFailure::new( 770 AuthoredWorkTarget::Artifact(ids().1), 771 fence(&active), 772 failure, 773 None, 774 12, 775 ) 776 .unwrap(); 777 block_on(storage.execute_authored(AuthoredAtomicCommand::ApplyFailure(command))).unwrap(); 778 assert_eq!( 779 block_on(storage.authored_artifact(ids().1)) 780 .unwrap() 781 .unwrap() 782 .signing_state(), 783 SigningState::FailedTerminal 784 ); 785 786 let storage = signed_storage(); 787 let active = claim_admission(&storage, 5, 13); 788 let failure = WorkFailure::new( 789 "temporary_admission_failure", 790 WorkPhase::Admission, 791 FailureClass::Retryable, 792 Some(30), 793 None, 794 ) 795 .unwrap(); 796 let retry = RetrySchedule::new(NonZeroU32::MIN, 30, failure.clone()).unwrap(); 797 let command = ApplyWorkFailure::new( 798 AuthoredWorkTarget::Artifact(ids().1), 799 fence(&active), 800 failure, 801 Some(retry), 802 14, 803 ) 804 .unwrap(); 805 block_on(storage.execute_authored(AuthoredAtomicCommand::ApplyFailure(command))).unwrap(); 806 assert_eq!( 807 block_on(storage.authored_artifact(ids().1)) 808 .unwrap() 809 .unwrap() 810 .admission_state(), 811 AdmissionState::Retryable 812 ); 813 814 let storage = signed_storage(); 815 let active = claim_admission(&storage, 5, 13); 816 let failure = WorkFailure::new( 817 "terminal_admission_failure", 818 WorkPhase::Admission, 819 FailureClass::Terminal, 820 None, 821 None, 822 ) 823 .unwrap(); 824 let command = ApplyWorkFailure::new( 825 AuthoredWorkTarget::Artifact(ids().1), 826 fence(&active), 827 failure, 828 None, 829 14, 830 ) 831 .unwrap(); 832 block_on(storage.execute_authored(AuthoredAtomicCommand::ApplyFailure(command))).unwrap(); 833 assert_eq!( 834 block_on(storage.authored_artifact(ids().1)) 835 .unwrap() 836 .unwrap() 837 .admission_state(), 838 AdmissionState::Rejected 839 ); 840 } 841 842 #[test] 843 fn memory_executes_delivery_failures_and_rejects_phase_mismatches_atomically() { 844 let storage = signed_storage(); 845 let active = claim_delivery(&storage, 5, 13); 846 let failure = WorkFailure::new( 847 "relay_unavailable", 848 WorkPhase::Delivery, 849 FailureClass::Retryable, 850 Some(30), 851 Some("relay unavailable".to_owned()), 852 ) 853 .unwrap(); 854 let retry = RetrySchedule::new(NonZeroU32::MIN, 30, failure.clone()).unwrap(); 855 let command = ApplyWorkFailure::new( 856 AuthoredWorkTarget::DeliveryPlan(ids().2), 857 fence(&active), 858 failure, 859 Some(retry), 860 14, 861 ) 862 .unwrap(); 863 block_on(storage.execute_authored(AuthoredAtomicCommand::ApplyFailure(command))).unwrap(); 864 assert_eq!( 865 block_on(storage.authored_delivery_plan(ids().2)) 866 .unwrap() 867 .unwrap() 868 .state(), 869 AuthoredDeliveryState::Retryable 870 ); 871 872 let storage = signed_storage(); 873 let active = claim_delivery(&storage, 5, 13); 874 let request = block_on(storage.authored_delivery_plan(ids().2)) 875 .unwrap() 876 .unwrap() 877 .request() 878 .unwrap() 879 .clone(); 880 let sink_failure = SinkFailure::for_request( 881 &request, 882 "relay_terminal", 883 Retryability::Terminal, 884 None, 885 Some("terminal".to_owned()), 886 Vec::new(), 887 ) 888 .unwrap(); 889 let command = ApplyDeliveryAttempt::new( 890 ids().2, 891 fence(&active), 892 DeliveryAttemptOutcome::SinkFailure(sink_failure), 893 None, 894 14, 895 ) 896 .unwrap(); 897 block_on(storage.execute_authored(AuthoredAtomicCommand::ApplyDelivery(command))).unwrap(); 898 assert_eq!( 899 block_on(storage.authored_delivery_plan(ids().2)) 900 .unwrap() 901 .unwrap() 902 .state(), 903 AuthoredDeliveryState::FailedTerminal 904 ); 905 906 for (target, phase, class) in [ 907 ( 908 AuthoredWorkTarget::Artifact(ids().1), 909 WorkPhase::Delivery, 910 FailureClass::Terminal, 911 ), 912 ( 913 AuthoredWorkTarget::DeliveryPlan(ids().2), 914 WorkPhase::Admission, 915 FailureClass::Terminal, 916 ), 917 ( 918 AuthoredWorkTarget::DeliveryPlan(ids().2), 919 WorkPhase::Delivery, 920 FailureClass::Indeterminate, 921 ), 922 ] { 923 let (storage, _) = prepared_storage(); 924 let failure = WorkFailure::new("mismatch", phase, class, None, None).unwrap(); 925 let command = ApplyWorkFailure::new( 926 target, 927 WorkFence::new([9; 16], NonZeroU64::MIN, NonZeroU64::MIN).unwrap(), 928 failure, 929 None, 930 12, 931 ) 932 .unwrap(); 933 assert!( 934 block_on(storage.execute_authored(AuthoredAtomicCommand::ApplyFailure(command))) 935 .is_err() 936 ); 937 } 938 } 939 940 #[test] 941 fn memory_executes_all_cancellation_targets_and_revision_fences() { 942 let (storage, _) = prepared_storage(); 943 let artifact = block_on(storage.authored_artifact(ids().1)) 944 .unwrap() 945 .unwrap(); 946 let stale = CancelAuthoredWork::new( 947 CancelAuthoredTarget::ArtifactSigning(ids().1), 948 NonZeroU64::new(9).unwrap(), 949 11, 950 ) 951 .unwrap(); 952 assert_eq!( 953 block_on(storage.execute_authored(AuthoredAtomicCommand::Cancel(stale))), 954 Err(Error::InvalidAuthoredTransition) 955 ); 956 let cancel = CancelAuthoredWork::new( 957 CancelAuthoredTarget::ArtifactSigning(ids().1), 958 artifact.revision(), 959 11, 960 ) 961 .unwrap(); 962 block_on(storage.execute_authored(AuthoredAtomicCommand::Cancel(cancel))).unwrap(); 963 assert_eq!( 964 block_on(storage.authored_artifact(ids().1)) 965 .unwrap() 966 .unwrap() 967 .signing_state(), 968 SigningState::Cancelled 969 ); 970 971 let storage = signed_storage(); 972 let artifact = block_on(storage.authored_artifact(ids().1)) 973 .unwrap() 974 .unwrap(); 975 let cancel = CancelAuthoredWork::new( 976 CancelAuthoredTarget::ArtifactAdmission(ids().1), 977 artifact.revision(), 978 13, 979 ) 980 .unwrap(); 981 block_on(storage.execute_authored(AuthoredAtomicCommand::Cancel(cancel))).unwrap(); 982 assert_eq!( 983 block_on(storage.authored_artifact(ids().1)) 984 .unwrap() 985 .unwrap() 986 .admission_state(), 987 AdmissionState::Cancelled 988 ); 989 990 let storage = signed_storage(); 991 let plan = block_on(storage.authored_delivery_plan(ids().2)) 992 .unwrap() 993 .unwrap(); 994 let stale = CancelAuthoredWork::new( 995 CancelAuthoredTarget::DeliveryPlan(ids().2), 996 NonZeroU64::MIN, 997 13, 998 ) 999 .unwrap(); 1000 assert_eq!( 1001 block_on(storage.execute_authored(AuthoredAtomicCommand::Cancel(stale))), 1002 Err(Error::InvalidAuthoredDeliveryPlan) 1003 ); 1004 let cancel = CancelAuthoredWork::new( 1005 CancelAuthoredTarget::DeliveryPlan(ids().2), 1006 plan.revision(), 1007 13, 1008 ) 1009 .unwrap(); 1010 block_on(storage.execute_authored(AuthoredAtomicCommand::Cancel(cancel))).unwrap(); 1011 assert_eq!( 1012 block_on(storage.authored_delivery_plan(ids().2)) 1013 .unwrap() 1014 .unwrap() 1015 .state(), 1016 AuthoredDeliveryState::Cancelled 1017 ); 1018 } 1019 1020 #[path = "authored_atomic/draft_submission.rs"] 1021 mod draft_submission;