authored_delivery_facts.rs (27097B)
1 use futures_executor::block_on; 2 use radroots_storage::{ 3 Error, 4 atomic::{AtomicCommitDigest, AtomicCommitDisposition}, 5 authored::WorkClaim, 6 authored_atomic::{ 7 ApplyDeliveryAttempt, AuthoredAtomicCommand, AuthoredAtomicOutcome, AuthoredAtomicReceipt, 8 AuthoredAtomicStorage, CancelAuthoredTarget, CancelAuthoredWork, ClaimAuthoredTarget, 9 ClaimAuthoredWork, RecordDeliveryFact, WorkFence, 10 }, 11 authored_delivery::{ 12 AuthoredDeliveryPlan, AuthoredDeliveryState, DELIVERY_PLAN_ATTEMPTS_MAX, 13 DeliveryAttemptOutcome, 14 }, 15 event::SourceGeneration, 16 memory::MemoryStorage, 17 }; 18 use radroots_transport::{ 19 DeliveryReceipt, SinkFailure, 20 outcome::{DeliveryOutcome, Retryability}, 21 policy::SatisfactionState, 22 sink::DeliveryTargetReceipt, 23 }; 24 use std::num::NonZeroU64; 25 26 #[path = "authored_signing/fixture.rs"] 27 mod fixture; 28 use fixture::*; 29 30 #[path = "authored_delivery/reconciliation_tests.rs"] 31 mod reconciliation; 32 33 fn plan(storage: &MemoryStorage) -> AuthoredDeliveryPlan { 34 block_on(storage.authored_delivery_plan(ids().2)) 35 .unwrap() 36 .unwrap() 37 } 38 39 fn prepared() -> (MemoryStorage, WorkClaim, AuthoredAtomicReceipt) { 40 let storage = MemoryStorage::new(SourceGeneration::new([1; 32]).unwrap()); 41 let (command, event) = prepare(); 42 block_on(storage.execute_authored(command)).unwrap(); 43 let signing = claim(NonZeroU64::MIN, 1, 11); 44 block_on( 45 storage.execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( 46 ClaimAuthoredTarget::ArtifactSigning(ids().1), 47 signing.clone(), 48 ))), 49 ) 50 .unwrap(); 51 block_on(storage.execute_authored(record(event, signing, 12))).unwrap(); 52 let delivery = claim(plan(&storage).revision(), 2, 13); 53 let receipt = block_on(storage.execute_authored(AuthoredAtomicCommand::Claim( 54 ClaimAuthoredWork::new(ClaimAuthoredTarget::DeliveryPlan(ids().2), delivery.clone()), 55 ))) 56 .unwrap(); 57 (storage, delivery, receipt) 58 } 59 60 fn outcome(plan: &AuthoredDeliveryPlan, accepted: bool) -> DeliveryAttemptOutcome { 61 let request = plan.request().unwrap(); 62 DeliveryAttemptOutcome::Receipt( 63 DeliveryReceipt::for_request( 64 request, 65 request 66 .target_set() 67 .targets() 68 .iter() 69 .cloned() 70 .map(|target| { 71 DeliveryTargetReceipt::attempted( 72 target, 73 if accepted { 74 DeliveryOutcome::accepted() 75 } else { 76 DeliveryOutcome::unavailable() 77 }, 78 ) 79 }) 80 .collect(), 81 ) 82 .unwrap(), 83 ) 84 } 85 86 fn fact( 87 plan: &AuthoredDeliveryPlan, 88 claim: WorkClaim, 89 accepted: bool, 90 at: u64, 91 ) -> AuthoredAtomicCommand { 92 AuthoredAtomicCommand::RecordDelivery( 93 RecordDeliveryFact::new( 94 plan.plan_id(), 95 plan.artifact_id(), 96 claim, 97 outcome(plan, accepted), 98 at, 99 ) 100 .unwrap(), 101 ) 102 } 103 104 fn stop(storage: &MemoryStorage, at: u64) { 105 block_on( 106 storage.execute_authored(AuthoredAtomicCommand::Cancel( 107 CancelAuthoredWork::new( 108 CancelAuthoredTarget::DeliveryPlan(ids().2), 109 plan(storage).revision(), 110 at, 111 ) 112 .unwrap(), 113 )), 114 ) 115 .unwrap(); 116 } 117 118 #[test] 119 fn expired_cancelled_delivery_retains_exact_evidence_and_idempotent_first_time() { 120 let (storage, active, original) = prepared(); 121 stop(&storage, 20); 122 let before = plan(&storage); 123 let command = fact(&before, active.clone(), true, 50); 124 let receipt = block_on(storage.execute_authored(command.clone())).unwrap(); 125 assert!(receipt.matches_command(&command)); 126 let after = plan(&storage); 127 assert_eq!(after.state(), AuthoredDeliveryState::Cancelled); 128 assert_eq!(after.stop_requested_at_unix_ms(), Some(20)); 129 assert_eq!(after.request().unwrap().payload().event().raw_json(), RAW); 130 assert_eq!( 131 after.delivery_satisfaction().unwrap(), 132 SatisfactionState::Satisfied 133 ); 134 assert_eq!(after.revision(), before.revision()); 135 assert_eq!(after.updated_at_unix_ms(), before.updated_at_unix_ms()); 136 assert_eq!(after.attempt_count(), 0); 137 assert_eq!(after.delivery_facts()[0].claim(), &active); 138 assert_eq!(after.delivery_facts()[0].observed_at_unix_ms(), 50); 139 let replay_command = fact(&after, active.clone(), true, 90); 140 assert_eq!(replay_command.commit_id(), command.commit_id()); 141 let replay = block_on(storage.execute_authored(replay_command)).unwrap(); 142 assert_eq!(replay.disposition(), AtomicCommitDisposition::Replay); 143 assert_eq!(replay.committed_at_unix_ms(), 50); 144 assert_eq!(plan(&storage), after); 145 assert_eq!( 146 block_on(storage.authored_receipt(original.commit_id())) 147 .unwrap() 148 .unwrap(), 149 original 150 ); 151 stop(&storage, 99); 152 assert_eq!(plan(&storage), after); 153 assert!(block_on(storage.execute_authored(fact(&after, active.clone(), false, 95))).is_err()); 154 let fenced = AuthoredAtomicCommand::ApplyDelivery( 155 ApplyDeliveryAttempt::new( 156 ids().2, 157 WorkFence::new(*active.token(), active.generation(), active.row_revision()).unwrap(), 158 outcome(&after, true), 159 None, 160 50, 161 ) 162 .unwrap(), 163 ); 164 assert_eq!( 165 block_on(storage.execute_authored(fenced)), 166 Err(Error::DeliveryPlanClaimConflict) 167 ); 168 assert!( 169 block_on( 170 storage.execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( 171 ClaimAuthoredTarget::DeliveryPlan(ids().2), 172 claim(after.revision(), 9, 100) 173 ))) 174 ) 175 .is_err() 176 ); 177 } 178 179 #[test] 180 fn superseded_result_does_not_steal_new_claim_or_regress_acceptance() { 181 let (storage, old, _) = prepared(); 182 let newer = claim(plan(&storage).revision(), 3, 40); 183 block_on( 184 storage.execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( 185 ClaimAuthoredTarget::DeliveryPlan(ids().2), 186 newer.clone(), 187 ))), 188 ) 189 .unwrap(); 190 let before = plan(&storage); 191 block_on(storage.execute_authored(fact(&before, old, true, 30))).unwrap(); 192 let after = plan(&storage); 193 assert_eq!(after.claim_evidence(), Some(&newer)); 194 assert_eq!(after.revision(), before.revision()); 195 assert_eq!(after.updated_at_unix_ms(), 40); 196 let partial = SinkFailure::for_request( 197 after.request().unwrap(), 198 "sink_lost", 199 Retryability::Retryable, 200 Some(65), 201 Some("connection lost".into()), 202 vec![], 203 ) 204 .unwrap(); 205 let command = AuthoredAtomicCommand::RecordDelivery( 206 RecordDeliveryFact::new( 207 ids().2, 208 ids().1, 209 newer.clone(), 210 DeliveryAttemptOutcome::SinkFailure(partial.clone()), 211 45, 212 ) 213 .unwrap(), 214 ); 215 block_on(storage.execute_authored(command)).unwrap(); 216 assert_eq!( 217 plan(&storage).delivery_satisfaction().unwrap(), 218 SatisfactionState::Satisfied 219 ); 220 assert_eq!(plan(&storage).claim_evidence(), Some(&newer)); 221 let retry = radroots_storage::authored::RetrySchedule::new( 222 std::num::NonZeroU32::MIN, 223 65, 224 radroots_storage::authored::WorkFailure::new( 225 "sink_lost", 226 radroots_storage::authored::WorkPhase::Delivery, 227 radroots_storage::authored::FailureClass::Retryable, 228 Some(65), 229 Some("connection lost".into()), 230 ) 231 .unwrap(), 232 ) 233 .unwrap(); 234 let fenced = AuthoredAtomicCommand::ApplyDelivery( 235 ApplyDeliveryAttempt::new( 236 ids().2, 237 WorkFence::new(*newer.token(), newer.generation(), newer.row_revision()).unwrap(), 238 DeliveryAttemptOutcome::SinkFailure(partial), 239 Some(retry.clone()), 240 46, 241 ) 242 .unwrap(), 243 ); 244 block_on(storage.execute_authored(fenced)).unwrap(); 245 let transitioned = plan(&storage); 246 assert_eq!(transitioned.state(), AuthoredDeliveryState::Retryable); 247 assert_eq!(transitioned.retry(), Some(&retry)); 248 assert_eq!(transitioned.attempt_count(), 1); 249 assert_eq!( 250 transitioned.delivery_satisfaction().unwrap(), 251 SatisfactionState::Satisfied 252 ); 253 stop(&storage, 47); 254 assert_eq!( 255 plan(&storage).delivery_satisfaction().unwrap(), 256 SatisfactionState::Satisfied 257 ); 258 } 259 260 #[test] 261 fn backend_provenance_rejects_every_forged_complete_claim_and_wrong_artifact() { 262 let (storage, original, _) = prepared(); 263 let before = plan(&storage); 264 for changed in 0..6 { 265 let forged = WorkClaim::new( 266 if changed == 0 { 267 [9; 16] 268 } else { 269 *original.token() 270 }, 271 if changed == 1 { 272 "forged" 273 } else { 274 original.owner() 275 }, 276 if changed == 2 { 277 NonZeroU64::new(99).unwrap() 278 } else { 279 original.generation() 280 }, 281 if changed == 3 { 282 14 283 } else { 284 original.acquired_at_unix_ms() 285 }, 286 if changed == 4 { 287 90 288 } else { 289 original.expires_at_unix_ms() 290 }, 291 if changed == 5 { 292 NonZeroU64::new(99).unwrap() 293 } else { 294 original.row_revision() 295 }, 296 ) 297 .unwrap(); 298 let command = fact(&before, forged, true, 100); 299 assert_eq!( 300 block_on(storage.execute_authored(command.clone())), 301 Err(Error::AtomicWorkflowMismatch) 302 ); 303 assert!( 304 block_on(storage.authored_receipt(command.commit_id())) 305 .unwrap() 306 .is_none() 307 ); 308 assert_eq!(plan(&storage), before); 309 } 310 let wrong = AuthoredAtomicCommand::RecordDelivery( 311 RecordDeliveryFact::new( 312 ids().2, 313 radroots_storage::authored::AuthoredArtifactId::new([99; 16]).unwrap(), 314 original.clone(), 315 outcome(&before, true), 316 50, 317 ) 318 .unwrap(), 319 ); 320 assert_eq!( 321 block_on(storage.execute_authored(wrong)), 322 Err(Error::AtomicWorkflowMismatch) 323 ); 324 assert!( 325 RecordDeliveryFact::new(ids().2, ids().1, original, outcome(&before, true), 12).is_err() 326 ); 327 } 328 329 #[test] 330 fn late_fact_rejects_rebound_raw_request_and_caller_receipt_with_wrong_time() { 331 let (storage, active, original) = prepared(); 332 let before = plan(&storage); 333 let value = 334 RecordDeliveryFact::new(ids().2, ids().1, active.clone(), outcome(&before, true), 50) 335 .unwrap(); 336 let request = before 337 .intent() 338 .materialize(radroots_transport::sink::DeliveryPayload::new(event( 339 OTHER_RAW, 340 ))) 341 .unwrap(); 342 let mut wire = serde_json::to_value(&before).unwrap(); 343 wire["request"] = serde_json::to_value(request).unwrap(); 344 let mut rebound: AuthoredDeliveryPlan = serde_json::from_value(wire).unwrap(); 345 let saved = rebound.clone(); 346 assert_eq!( 347 value.apply_to(&mut rebound, &original), 348 Err(Error::AtomicWorkflowMismatch) 349 ); 350 assert_eq!(rebound, saved); 351 let wrong = AuthoredAtomicReceipt::from_durable_parts( 352 original.commit_id(), 353 original.digest(), 354 AtomicCommitDisposition::Committed, 355 14, 356 original.outcome().clone(), 357 ) 358 .unwrap(); 359 let mut unchanged = before.clone(); 360 assert_eq!( 361 value.apply_to(&mut unchanged, &wrong), 362 Err(Error::AtomicWorkflowMismatch) 363 ); 364 assert_eq!(unchanged, before); 365 } 366 367 #[test] 368 fn invalid_result_binding_rolls_back_and_changed_failure_details_have_distinct_ids() { 369 let (storage, active, _) = prepared(); 370 let before = plan(&storage); 371 let request = before.request().unwrap(); 372 let wrong = radroots_transport::DeliveryRequest::new( 373 "wrong-request", 374 request.payload().clone(), 375 request.target_set().clone(), 376 request.satisfaction().clone(), 377 request.deadline_unix_ms(), 378 ) 379 .unwrap(); 380 let wrong = DeliveryReceipt::for_request( 381 &wrong, 382 wrong 383 .target_set() 384 .targets() 385 .iter() 386 .cloned() 387 .map(|target| DeliveryTargetReceipt::attempted(target, DeliveryOutcome::accepted())) 388 .collect(), 389 ) 390 .unwrap(); 391 let command = AuthoredAtomicCommand::RecordDelivery( 392 RecordDeliveryFact::new( 393 ids().2, 394 ids().1, 395 active.clone(), 396 DeliveryAttemptOutcome::Receipt(wrong), 397 50, 398 ) 399 .unwrap(), 400 ); 401 assert_eq!( 402 block_on(storage.execute_authored(command.clone())), 403 Err(Error::InvalidAuthoredDeliveryPlan) 404 ); 405 assert_eq!(plan(&storage), before); 406 assert!( 407 block_on(storage.authored_receipt(command.commit_id())) 408 .unwrap() 409 .is_none() 410 ); 411 let mut identities = std::collections::BTreeSet::new(); 412 for retry in [None, Some(65), Some(66)] { 413 for message in [ 414 None, 415 Some("connection lost".to_owned()), 416 Some("connection closed".to_owned()), 417 ] { 418 let failure = SinkFailure::for_request( 419 request, 420 "sink_lost", 421 Retryability::Retryable, 422 retry, 423 message, 424 vec![], 425 ) 426 .unwrap(); 427 let command = AuthoredAtomicCommand::RecordDelivery( 428 RecordDeliveryFact::new( 429 ids().2, 430 ids().1, 431 active.clone(), 432 DeliveryAttemptOutcome::SinkFailure(failure), 433 50, 434 ) 435 .unwrap(), 436 ); 437 assert!(identities.insert(*command.commit_id().as_bytes())); 438 } 439 } 440 assert_eq!(identities.len(), 9); 441 } 442 443 #[test] 444 fn stop_after_satisfaction_and_legacy_cancelled_snapshot_preserve_intent() { 445 let (storage, active, _) = prepared(); 446 let before = plan(&storage); 447 let command = AuthoredAtomicCommand::ApplyDelivery( 448 ApplyDeliveryAttempt::new( 449 ids().2, 450 WorkFence::new(*active.token(), active.generation(), active.row_revision()).unwrap(), 451 outcome(&before, true), 452 None, 453 14, 454 ) 455 .unwrap(), 456 ); 457 block_on(storage.execute_authored(command)).unwrap(); 458 stop(&storage, 15); 459 let after = plan(&storage); 460 assert_eq!(after.state(), AuthoredDeliveryState::Satisfied); 461 assert_eq!(after.stop_requested_at_unix_ms(), Some(15)); 462 assert_eq!( 463 after.delivery_satisfaction().unwrap(), 464 SatisfactionState::Satisfied 465 ); 466 let (storage, _, _) = prepared(); 467 stop(&storage, 20); 468 let expected = plan(&storage); 469 let mut legacy = serde_json::to_value(&expected).unwrap(); 470 legacy.as_object_mut().unwrap().remove("delivery_facts"); 471 legacy 472 .as_object_mut() 473 .unwrap() 474 .remove("stop_requested_at_unix_ms"); 475 assert_eq!( 476 serde_json::from_value::<AuthoredDeliveryPlan>(legacy).unwrap(), 477 expected 478 ); 479 let mut pending = before.clone(); 480 assert!(pending.request_stop(9).is_err()); 481 assert_eq!(pending, before); 482 assert!(pending.request_stop(12).is_err()); 483 assert_eq!(pending, before); 484 } 485 486 #[test] 487 fn receipt_binding_rejects_wrong_outcome_digest_and_observation_time() { 488 let (storage, active, _) = prepared(); 489 let before = plan(&storage); 490 let command = fact(&before, active, true, 50); 491 assert!( 492 AuthoredAtomicReceipt::new( 493 &command, 494 AtomicCommitDisposition::Committed, 495 50, 496 AuthoredAtomicOutcome::DeliveryPlan(before.clone()) 497 ) 498 .is_err() 499 ); 500 let receipt = block_on(storage.execute_authored(command.clone())).unwrap(); 501 assert!( 502 AuthoredAtomicReceipt::new( 503 &command, 504 AtomicCommitDisposition::Committed, 505 49, 506 receipt.outcome().clone() 507 ) 508 .is_err() 509 ); 510 let wrong = AuthoredAtomicReceipt::from_durable_parts( 511 receipt.commit_id(), 512 AtomicCommitDigest::new([99; 32]), 513 AtomicCommitDisposition::Committed, 514 50, 515 receipt.outcome().clone(), 516 ) 517 .unwrap(); 518 assert!(!wrong.matches_command(&command)); 519 let wrong = AuthoredAtomicReceipt::from_durable_parts( 520 receipt.commit_id(), 521 receipt.digest(), 522 AtomicCommitDisposition::Committed, 523 50, 524 AuthoredAtomicOutcome::Artifact( 525 block_on(storage.authored_artifact(ids().1)) 526 .unwrap() 527 .unwrap(), 528 ), 529 ) 530 .unwrap(); 531 assert!(!wrong.matches_command(&command)); 532 } 533 534 #[test] 535 fn fact_capacity_and_structural_corruption_fail_without_evicting_evidence() { 536 let (storage, active, original) = prepared(); 537 let before = plan(&storage); 538 let command = fact(&before, active.clone(), true, 50); 539 block_on(storage.execute_authored(command)).unwrap(); 540 let after = plan(&storage); 541 let base = serde_json::to_value(&after).unwrap(); 542 for (path, value) in [ 543 ("observed_at_unix_ms", serde_json::json!(1)), 544 ("claim", serde_json::json!(null)), 545 ("outcome", serde_json::json!(null)), 546 ] { 547 let mut corrupt = base.clone(); 548 corrupt["delivery_facts"][0][path] = value; 549 assert!(serde_json::from_value::<AuthoredDeliveryPlan>(corrupt).is_err()); 550 } 551 let mut duplicate = base.clone(); 552 duplicate["delivery_facts"] 553 .as_array_mut() 554 .unwrap() 555 .push(base["delivery_facts"][0].clone()); 556 assert!(serde_json::from_value::<AuthoredDeliveryPlan>(duplicate).is_err()); 557 let mut full = base.clone(); 558 full["revision"] = serde_json::json!(2048); 559 full["claim"] = serde_json::Value::Null; 560 let entries = (1..=DELIVERY_PLAN_ATTEMPTS_MAX) 561 .map(|index| { 562 let mut entry = base["delivery_facts"][0].clone(); 563 let claim = WorkClaim::new( 564 [7; 16], 565 "capacity", 566 NonZeroU64::new(u64::from(index)).unwrap(), 567 13, 568 33, 569 NonZeroU64::new(u64::from(index)).unwrap(), 570 ) 571 .unwrap(); 572 entry["claim"] = serde_json::to_value(claim).unwrap(); 573 entry 574 }) 575 .collect::<Vec<_>>(); 576 full["delivery_facts"] = serde_json::to_value(entries).unwrap(); 577 let mut full_plan: AuthoredDeliveryPlan = serde_json::from_value(full.clone()).unwrap(); 578 let checkpoint = full_plan.clone(); 579 let value = 580 RecordDeliveryFact::new(ids().2, ids().1, active, outcome(&before, true), 50).unwrap(); 581 assert_eq!( 582 value.apply_to(&mut full_plan, &original), 583 Err(Error::DeliveryAttemptOverflow) 584 ); 585 assert_eq!(full_plan, checkpoint); 586 assert!( 587 full_plan 588 .claim(claim(full_plan.revision(), 9, 60), 60) 589 .is_err() 590 ); 591 let mut replay = full.clone(); 592 replay["delivery_facts"][0] = base["delivery_facts"][0].clone(); 593 let mut replay_plan: AuthoredDeliveryPlan = serde_json::from_value(replay).unwrap(); 594 let checkpoint = replay_plan.clone(); 595 value.apply_to(&mut replay_plan, &original).unwrap(); 596 assert_eq!(replay_plan, checkpoint); 597 full["delivery_facts"] 598 .as_array_mut() 599 .unwrap() 600 .push(base["delivery_facts"][0].clone()); 601 assert!(serde_json::from_value::<AuthoredDeliveryPlan>(full).is_err()); 602 } 603 604 #[test] 605 fn corrupt_original_claim_receipts_and_rebound_current_rows_are_rejected() { 606 let (storage, active, original) = prepared(); 607 let before = plan(&storage); 608 let value = 609 RecordDeliveryFact::new(ids().2, ids().1, active, outcome(&before, true), 50).unwrap(); 610 let AuthoredAtomicOutcome::DeliveryPlan(original_plan) = original.outcome() else { 611 unreachable!() 612 }; 613 for updates in [ 614 vec![("plan_id", serde_json::json!(vec![9_u8; 16]))], 615 vec![("artifact_id", serde_json::json!(vec![9_u8; 16]))], 616 vec![("created_at_unix_ms", serde_json::json!(9))], 617 vec![("request", serde_json::Value::Null)], 618 vec![("claim", serde_json::Value::Null)], 619 ] { 620 let mut wire = serde_json::to_value(original_plan).unwrap(); 621 for (key, replacement) in updates { 622 wire[key] = replacement; 623 } 624 let altered: AuthoredDeliveryPlan = serde_json::from_value(wire).unwrap(); 625 let forged = AuthoredAtomicReceipt::from_durable_parts( 626 original.commit_id(), 627 original.digest(), 628 AtomicCommitDisposition::Committed, 629 original.committed_at_unix_ms(), 630 AuthoredAtomicOutcome::DeliveryPlan(altered), 631 ) 632 .unwrap(); 633 let mut current = before.clone(); 634 assert_eq!( 635 value.apply_to(&mut current, &forged), 636 Err(Error::AtomicWorkflowMismatch) 637 ); 638 assert_eq!(current, before); 639 } 640 for updates in [ 641 vec![("plan_id", serde_json::json!(vec![9_u8; 16]))], 642 vec![("artifact_id", serde_json::json!(vec![9_u8; 16]))], 643 vec![("created_at_unix_ms", serde_json::json!(9))], 644 vec![ 645 ("claim", serde_json::Value::Null), 646 ("revision", serde_json::json!(2)), 647 ], 648 vec![ 649 ("claim", serde_json::Value::Null), 650 ("updated_at_unix_ms", serde_json::json!(12)), 651 ], 652 ] { 653 let mut wire = serde_json::to_value(&before).unwrap(); 654 for (key, replacement) in updates { 655 wire[key] = replacement; 656 } 657 let mut current: AuthoredDeliveryPlan = serde_json::from_value(wire).unwrap(); 658 let unchanged = current.clone(); 659 assert_eq!( 660 value.apply_to(&mut current, &original), 661 Err(Error::AtomicWorkflowMismatch) 662 ); 663 assert_eq!(current, unchanged); 664 } 665 let wrong_kind = AuthoredAtomicReceipt::from_durable_parts( 666 original.commit_id(), 667 original.digest(), 668 AtomicCommitDisposition::Committed, 669 original.committed_at_unix_ms(), 670 AuthoredAtomicOutcome::Artifact( 671 block_on(storage.authored_artifact(ids().1)) 672 .unwrap() 673 .unwrap(), 674 ), 675 ) 676 .unwrap(); 677 let mut current = before.clone(); 678 assert_eq!( 679 value.apply_to(&mut current, &wrong_kind), 680 Err(Error::AtomicWorkflowMismatch) 681 ); 682 let wrong_digest = AuthoredAtomicReceipt::from_durable_parts( 683 original.commit_id(), 684 AtomicCommitDigest::new([99; 32]), 685 AtomicCommitDisposition::Committed, 686 original.committed_at_unix_ms(), 687 original.outcome().clone(), 688 ) 689 .unwrap(); 690 assert_eq!( 691 value.apply_to(&mut current, &wrong_digest), 692 Err(Error::AtomicWorkflowMismatch) 693 ); 694 assert_eq!(current, before); 695 } 696 697 #[test] 698 fn forged_fact_snapshots_cannot_bypass_stop_time_claim_or_request_invariants() { 699 let (storage, active, _) = prepared(); 700 let before = plan(&storage); 701 block_on(storage.execute_authored(fact(&before, active, true, 50))).unwrap(); 702 let current = plan(&storage); 703 let base = serde_json::to_value(¤t).unwrap(); 704 for at in [9, 12, 99] { 705 let mut wire = base.clone(); 706 wire["stop_requested_at_unix_ms"] = serde_json::json!(at); 707 assert!(serde_json::from_value::<AuthoredDeliveryPlan>(wire).is_err()); 708 } 709 for (at, revision) in [(9, 2), (15, 2), (13, 3)] { 710 let forged = WorkClaim::new( 711 [8; 16], 712 "invalid-fact", 713 NonZeroU64::MIN, 714 at, 715 at + 20, 716 NonZeroU64::new(revision).unwrap(), 717 ) 718 .unwrap(); 719 let mut wire = base.clone(); 720 wire["delivery_facts"][0]["claim"] = serde_json::to_value(forged).unwrap(); 721 assert!(serde_json::from_value::<AuthoredDeliveryPlan>(wire).is_err()); 722 } 723 let mut wire = base; 724 wire["request"] = serde_json::Value::Null; 725 assert!(serde_json::from_value::<AuthoredDeliveryPlan>(wire).is_err()); 726 } 727 728 #[test] 729 fn receipt_matching_checks_each_plan_fact_and_monotonic_row_time() { 730 let (storage, active, _) = prepared(); 731 let before = plan(&storage); 732 let command = fact(&before, active, true, 50); 733 let receipt = block_on(storage.execute_authored(command.clone())).unwrap(); 734 let current = plan(&storage); 735 let base = serde_json::to_value(¤t).unwrap(); 736 for key in ["plan_id", "artifact_id"] { 737 let mut wire = base.clone(); 738 wire[key] = serde_json::json!(vec![9_u8; 16]); 739 let altered: AuthoredDeliveryPlan = serde_json::from_value(wire).unwrap(); 740 let forged = AuthoredAtomicReceipt::from_durable_parts( 741 receipt.commit_id(), 742 receipt.digest(), 743 AtomicCommitDisposition::Committed, 744 50, 745 AuthoredAtomicOutcome::DeliveryPlan(altered.clone()), 746 ) 747 .unwrap(); 748 assert!(!forged.matches_command(&command)); 749 assert!( 750 AuthoredAtomicReceipt::new( 751 &command, 752 AtomicCommitDisposition::Committed, 753 50, 754 AuthoredAtomicOutcome::DeliveryPlan(altered) 755 ) 756 .is_err() 757 ); 758 } 759 for change_claim in [false, true] { 760 let mut wire = base.clone(); 761 if change_claim { 762 let forged = WorkClaim::new( 763 [8; 16], 764 "different", 765 NonZeroU64::MIN, 766 13, 767 33, 768 NonZeroU64::new(2).unwrap(), 769 ) 770 .unwrap(); 771 wire["delivery_facts"][0]["claim"] = serde_json::to_value(forged).unwrap(); 772 } else { 773 wire["delivery_facts"][0]["outcome"] = 774 serde_json::to_value(outcome(¤t, false)).unwrap(); 775 } 776 let altered: AuthoredDeliveryPlan = serde_json::from_value(wire).unwrap(); 777 let forged = AuthoredAtomicReceipt::from_durable_parts( 778 receipt.commit_id(), 779 receipt.digest(), 780 AtomicCommitDisposition::Committed, 781 50, 782 AuthoredAtomicOutcome::DeliveryPlan(altered), 783 ) 784 .unwrap(); 785 assert!(!forged.matches_command(&command)); 786 } 787 let mut stopped = current; 788 stopped.request_stop(100).unwrap(); 789 assert!( 790 AuthoredAtomicReceipt::new( 791 &command, 792 AtomicCommitDisposition::Committed, 793 60, 794 AuthoredAtomicOutcome::DeliveryPlan(stopped) 795 ) 796 .is_err() 797 ); 798 }