reconciliation_tests.rs (37438B)
1 use super::*; 2 use radroots_storage::{ 3 authored::{FailureClass, RetrySchedule, WorkFailure, WorkPhase}, 4 authored_atomic::ReconcileDeliveryFacts, 5 authored_delivery::AuthoredDeliveryHistory, 6 }; 7 use std::num::NonZeroU32; 8 9 fn history(storage: &MemoryStorage) -> AuthoredDeliveryHistory { 10 block_on(storage.authored_delivery_history(ids().2)) 11 .unwrap() 12 .unwrap() 13 } 14 15 fn fence(claim: &WorkClaim) -> WorkFence { 16 WorkFence::new(*claim.token(), claim.generation(), claim.row_revision()).unwrap() 17 } 18 19 fn reconcile( 20 plan: &AuthoredDeliveryPlan, 21 authority: Option<WorkFence>, 22 retry: Option<RetrySchedule>, 23 at: u64, 24 ) -> AuthoredAtomicCommand { 25 AuthoredAtomicCommand::ReconcileDelivery( 26 ReconcileDeliveryFacts::new(plan, authority, retry, at).unwrap(), 27 ) 28 } 29 30 fn retry(attempt: u32, at: u64) -> RetrySchedule { 31 RetrySchedule::new( 32 NonZeroU32::new(attempt).unwrap(), 33 at, 34 WorkFailure::new( 35 "delivery_pending", 36 WorkPhase::Delivery, 37 FailureClass::Retryable, 38 Some(at), 39 None, 40 ) 41 .unwrap(), 42 ) 43 .unwrap() 44 } 45 46 #[test] 47 fn history_distinguishes_no_issued_work_from_stopped_unresolved_work() { 48 let untouched = MemoryStorage::new(SourceGeneration::new([71; 32]).unwrap()); 49 assert!( 50 block_on(untouched.authored_delivery_history(ids().2)) 51 .unwrap() 52 .is_none() 53 ); 54 block_on(untouched.execute_authored(prepare().0)).unwrap(); 55 let before = history(&untouched); 56 assert!(before.is_complete()); 57 assert!(before.proves_no_issued_attempt()); 58 assert!(!before.has_unresolved_claims()); 59 stop(&untouched, 20); 60 assert!(history(&untouched).proves_no_issued_attempt()); 61 62 let (storage, active, _) = prepared(); 63 let issued = history(&storage); 64 assert_eq!(issued.claims().len(), 1); 65 assert_eq!(issued.claims()[0].claim(), &active); 66 assert_eq!(issued.claims()[0].prior_attempt_count(), 0); 67 assert!(issued.has_unresolved_claims()); 68 assert!(!issued.proves_no_issued_attempt()); 69 stop(&storage, 20); 70 assert!(history(&storage).has_unresolved_claims()); 71 block_on(storage.execute_authored(fact(&plan(&storage), active, true, 50))).unwrap(); 72 let observed = history(&storage); 73 assert!(!observed.has_unresolved_claims()); 74 assert!(!observed.proves_no_issued_attempt()); 75 assert_eq!(observed.plan().state(), AuthoredDeliveryState::Cancelled); 76 assert_eq!( 77 observed.plan().delivery_satisfaction().unwrap(), 78 SatisfactionState::Satisfied 79 ); 80 } 81 82 #[test] 83 fn exact_current_fence_reconciles_once_and_preserves_raw_fact_provenance() { 84 let (storage, active, original) = prepared(); 85 block_on(storage.execute_authored(fact(&plan(&storage), active.clone(), true, 14))).unwrap(); 86 let before = plan(&storage); 87 let command = reconcile(&before, Some(fence(&active)), None, 15); 88 let receipt = block_on(storage.execute_authored(command.clone())).unwrap(); 89 let after = plan(&storage); 90 assert!(receipt.matches_command(&command)); 91 assert_eq!(after.state(), AuthoredDeliveryState::Satisfied); 92 assert_eq!(after.attempt_count(), 1); 93 assert_eq!(after.attempts()[0].claim_evidence(), Some(&active)); 94 assert_eq!( 95 after.attempts()[0].outcome(), 96 before.delivery_facts()[0].outcome() 97 ); 98 assert_eq!(after.attempts()[0].recorded_at_unix_ms(), 15); 99 assert_eq!(after.delivery_facts(), before.delivery_facts()); 100 assert_eq!(after.pending_delivery_facts().count(), 0); 101 assert!(!history(&storage).has_unresolved_claims()); 102 let replay = block_on(storage.execute_authored(command)).unwrap(); 103 assert_eq!(replay.disposition(), AtomicCommitDisposition::Replay); 104 assert_eq!(plan(&storage), after); 105 assert_eq!( 106 block_on(storage.authored_receipt(original.commit_id())) 107 .unwrap() 108 .unwrap(), 109 original 110 ); 111 assert_eq!( 112 ReconcileDeliveryFacts::new(&after, None, None, 16), 113 Err(Error::AtomicWorkflowMismatch) 114 ); 115 } 116 117 #[test] 118 fn distinct_current_reconciliation_survives_legacy_generation_identity_reuse() { 119 let (storage, first, _) = prepared(); 120 let first_plan = plan(&storage); 121 block_on(storage.execute_authored(fact(&first_plan, first.clone(), false, 14))).unwrap(); 122 let first_command = reconcile(&plan(&storage), Some(fence(&first)), Some(retry(1, 18)), 15); 123 block_on(storage.execute_authored(first_command.clone())).unwrap(); 124 let second = claim(plan(&storage).revision(), 2, 20); 125 block_on( 126 storage.execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( 127 ClaimAuthoredTarget::DeliveryPlan(ids().2), 128 second.clone(), 129 ))), 130 ) 131 .unwrap(); 132 block_on(storage.execute_authored(fact(&plan(&storage), second.clone(), false, 21))).unwrap(); 133 let second_command = reconcile( 134 &plan(&storage), 135 Some(fence(&second)), 136 Some(retry(2, 25)), 137 22, 138 ); 139 assert_ne!(first_command.commit_id(), second_command.commit_id()); 140 let legacy = |claim: &WorkClaim, at| { 141 AuthoredAtomicCommand::ApplyDelivery( 142 ApplyDeliveryAttempt::new(ids().2, fence(claim), outcome(&first_plan, false), None, at) 143 .unwrap(), 144 ) 145 }; 146 assert_eq!( 147 legacy(&first, 15).commit_id(), 148 legacy(&second, 22).commit_id() 149 ); 150 block_on(storage.execute_authored(second_command)).unwrap(); 151 let after = plan(&storage); 152 assert_eq!(after.state(), AuthoredDeliveryState::Retryable); 153 assert_eq!(after.attempt_count(), 2); 154 assert_eq!(after.retry().unwrap().not_before_unix_ms(), 25); 155 assert_eq!(after.attempts()[0].claim_evidence(), Some(&first)); 156 assert_eq!(after.attempts()[1].claim_evidence(), Some(&second)); 157 assert!(!history(&storage).has_unresolved_claims()); 158 } 159 160 #[test] 161 fn late_fact_cannot_take_a_newer_lease_and_its_marker_cannot_resolve_that_lease() { 162 let (storage, old, _) = prepared(); 163 let newer = claim(plan(&storage).revision(), 3, 40); 164 block_on( 165 storage.execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( 166 ClaimAuthoredTarget::DeliveryPlan(ids().2), 167 newer.clone(), 168 ))), 169 ) 170 .unwrap(); 171 block_on(storage.execute_authored(fact(&plan(&storage), old.clone(), true, 50))).unwrap(); 172 let before = plan(&storage); 173 for authority in [None, Some(fence(&old))] { 174 let command = reconcile(&before, authority, None, 51); 175 assert_eq!( 176 block_on(storage.execute_authored(command)), 177 Err(Error::DeliveryPlanClaimConflict) 178 ); 179 assert_eq!(plan(&storage), before); 180 } 181 block_on(storage.execute_authored(reconcile(&before, Some(fence(&newer)), None, 51))).unwrap(); 182 let after = history(&storage); 183 assert_eq!(after.plan().attempts()[0].claim_evidence(), Some(&old)); 184 assert!( 185 after.has_unresolved_claims(), 186 "the new worker has no final result yet" 187 ); 188 assert_eq!( 189 after.plan().delivery_satisfaction().unwrap(), 190 SatisfactionState::Satisfied 191 ); 192 } 193 194 #[test] 195 fn fresh_observer_can_reconcile_after_expiry_but_expired_worker_cannot() { 196 let (storage, old, _) = prepared(); 197 block_on(storage.execute_authored(fact(&plan(&storage), old.clone(), true, 50))).unwrap(); 198 let before = plan(&storage); 199 assert_eq!( 200 block_on(storage.execute_authored(reconcile(&before, Some(fence(&old)), None, 51))), 201 Err(Error::DeliveryPlanClaimConflict) 202 ); 203 block_on(storage.execute_authored(reconcile(&before, None, None, 51))).unwrap(); 204 assert_eq!(plan(&storage).state(), AuthoredDeliveryState::Satisfied); 205 assert_eq!(plan(&storage).attempts()[0].claim_evidence(), Some(&old)); 206 } 207 208 #[test] 209 fn exact_fact_set_and_revision_races_fail_without_partial_scheduling() { 210 let (storage, old, _) = prepared(); 211 let newer = claim(plan(&storage).revision(), 3, 40); 212 block_on( 213 storage.execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( 214 ClaimAuthoredTarget::DeliveryPlan(ids().2), 215 newer.clone(), 216 ))), 217 ) 218 .unwrap(); 219 block_on(storage.execute_authored(fact(&plan(&storage), old.clone(), true, 50))).unwrap(); 220 let before = plan(&storage); 221 let stale = reconcile(&before, Some(fence(&newer)), None, 55); 222 block_on(storage.execute_authored(fact(&before, newer.clone(), false, 51))).unwrap(); 223 let current = plan(&storage); 224 assert_eq!(current.revision(), before.revision()); 225 assert_eq!( 226 block_on(storage.execute_authored(stale)), 227 Err(Error::DeliveryPlanClaimConflict) 228 ); 229 assert_eq!(plan(&storage), current); 230 block_on(storage.execute_authored(reconcile(¤t, Some(fence(&newer)), None, 55))).unwrap(); 231 let settled = plan(&storage); 232 assert_eq!(settled.attempt_count(), 2); 233 assert_eq!(settled.state(), AuthoredDeliveryState::Satisfied); 234 assert_eq!(settled.attempts()[0].claim_evidence(), Some(&old)); 235 assert_eq!(settled.attempts()[1].claim_evidence(), Some(&newer)); 236 assert!( 237 settled 238 .attempts() 239 .iter() 240 .all(|attempt| attempt.recorded_at_unix_ms() == 55) 241 ); 242 243 let (stopped, active, _) = prepared(); 244 block_on(stopped.execute_authored(fact(&plan(&stopped), active, true, 50))).unwrap(); 245 let stale = reconcile(&plan(&stopped), None, None, 51); 246 stop(&stopped, 52); 247 let before = plan(&stopped); 248 assert_eq!( 249 block_on(stopped.execute_authored(stale)), 250 Err(Error::DeliveryPlanClaimConflict) 251 ); 252 assert_eq!( 253 block_on(stopped.execute_authored(reconcile(&before, None, None, 53))), 254 Err(Error::InvalidAuthoredTransition) 255 ); 256 assert_eq!(plan(&stopped), before); 257 } 258 259 #[test] 260 fn historical_application_is_resolved_without_rewriting_its_snapshot_shape() { 261 let (storage, active, _) = prepared(); 262 let command = AuthoredAtomicCommand::ApplyDelivery( 263 ApplyDeliveryAttempt::new( 264 ids().2, 265 fence(&active), 266 outcome(&plan(&storage), true), 267 None, 268 14, 269 ) 270 .unwrap(), 271 ); 272 block_on(storage.execute_authored(command)).unwrap(); 273 let legacy = history(&storage); 274 assert!(!legacy.has_unresolved_claims()); 275 let attempt = &legacy.plan().attempts()[0]; 276 assert!(attempt.claim_evidence().is_none()); 277 assert!( 278 !serde_json::to_value(attempt) 279 .unwrap() 280 .as_object() 281 .unwrap() 282 .contains_key("claim") 283 ); 284 let restored: AuthoredDeliveryPlan = 285 serde_json::from_slice(&serde_json::to_vec(legacy.plan()).unwrap()).unwrap(); 286 assert_eq!(restored, *legacy.plan()); 287 } 288 289 #[test] 290 fn missing_truncated_or_malformed_history_never_proves_absence() { 291 let storage = MemoryStorage::new(SourceGeneration::new([72; 32]).unwrap()); 292 let original = block_on(storage.execute_authored(prepare().0)).unwrap(); 293 let initial = plan(&storage); 294 let unknown = AuthoredDeliveryHistory::new(initial.clone(), None).unwrap(); 295 assert!(!unknown.is_complete()); 296 assert!(!unknown.proves_no_issued_attempt()); 297 assert!(unknown.has_unresolved_claims()); 298 let mut bounded = AuthoredDeliveryHistory::new(initial, Some(&original)).unwrap(); 299 bounded.mark_truncated(); 300 assert!(bounded.is_truncated()); 301 assert!(!bounded.proves_no_issued_attempt()); 302 assert!(bounded.has_unresolved_claims()); 303 assert_eq!( 304 bounded.require_pending_fact_provenance(), 305 Err(Error::DeliveryAttemptOverflow) 306 ); 307 assert_eq!( 308 bounded.push_claim(&original), 309 Err(Error::AtomicWorkflowMismatch) 310 ); 311 312 let (storage, _, issued) = prepared(); 313 let mut duplicate = history(&storage); 314 assert_eq!( 315 duplicate.push_claim(&issued), 316 Err(Error::AtomicWorkflowMismatch) 317 ); 318 assert!(AuthoredDeliveryHistory::new(plan(&storage), Some(&issued)).is_err()); 319 let incomplete = AuthoredDeliveryHistory::new(plan(&storage), Some(&original)).unwrap(); 320 assert_eq!(incomplete.validate(), Err(Error::AtomicWorkflowMismatch)); 321 } 322 323 #[test] 324 fn forged_or_duplicate_reconciliation_markers_fail_snapshot_validation() { 325 let (storage, active, _) = prepared(); 326 block_on(storage.execute_authored(fact(&plan(&storage), active.clone(), true, 14))).unwrap(); 327 block_on(storage.execute_authored(reconcile(&plan(&storage), Some(fence(&active)), None, 15))) 328 .unwrap(); 329 let valid = serde_json::to_value(plan(&storage)).unwrap(); 330 let mut missing_fact = valid.clone(); 331 missing_fact["delivery_facts"] = serde_json::json!([]); 332 assert!(serde_json::from_value::<AuthoredDeliveryPlan>(missing_fact).is_err()); 333 let mut forged = valid.clone(); 334 forged["attempts"][0]["claim"] = serde_json::to_value(claim(NonZeroU64::MIN, 9, 13)).unwrap(); 335 assert!(serde_json::from_value::<AuthoredDeliveryPlan>(forged).is_err()); 336 let mut duplicate = valid; 337 let mut second = duplicate["attempts"][0].clone(); 338 second["attempt"] = serde_json::json!(2); 339 duplicate["attempts"].as_array_mut().unwrap().push(second); 340 duplicate["attempt_count"] = serde_json::json!(2); 341 assert!(serde_json::from_value::<AuthoredDeliveryPlan>(duplicate).is_err()); 342 } 343 344 #[test] 345 fn issued_claim_limit_rejects_new_work_atomically_and_keeps_original_replay() { 346 let (storage, _, first) = prepared(); 347 for index in 1..DELIVERY_PLAN_ATTEMPTS_MAX { 348 let at = 40 + u64::from(index) * 21; 349 let active = WorkClaim::new( 350 [7; 16], 351 "bounded-worker", 352 NonZeroU64::new(u64::from(index) + 2).unwrap(), 353 at, 354 at + 20, 355 plan(&storage).revision(), 356 ) 357 .unwrap(); 358 block_on( 359 storage.execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( 360 ClaimAuthoredTarget::DeliveryPlan(ids().2), 361 active, 362 ))), 363 ) 364 .unwrap(); 365 } 366 let before = plan(&storage); 367 let at = 40 + u64::from(DELIVERY_PLAN_ATTEMPTS_MAX) * 21; 368 let active = WorkClaim::new( 369 [8; 16], 370 "overflow-worker", 371 NonZeroU64::new(2048).unwrap(), 372 at, 373 at + 20, 374 before.revision(), 375 ) 376 .unwrap(); 377 let command = AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( 378 ClaimAuthoredTarget::DeliveryPlan(ids().2), 379 active, 380 )); 381 assert_eq!( 382 block_on(storage.execute_authored(command.clone())), 383 Err(Error::DeliveryAttemptOverflow) 384 ); 385 assert_eq!(plan(&storage), before); 386 assert!( 387 block_on(storage.authored_receipt(command.commit_id())) 388 .unwrap() 389 .is_none() 390 ); 391 let history = history(&storage); 392 assert_eq!(history.claims().len(), DELIVERY_PLAN_ATTEMPTS_MAX as usize); 393 assert!(history.is_complete()); 394 assert!(history.has_unresolved_claims()); 395 assert_eq!( 396 history.clone().push_claim(&first), 397 Err(Error::DeliveryAttemptOverflow) 398 ); 399 let AuthoredAtomicOutcome::DeliveryPlan(original) = first.outcome() else { 400 unreachable!() 401 }; 402 let original = AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( 403 ClaimAuthoredTarget::DeliveryPlan(ids().2), 404 original.claim_evidence().unwrap().clone(), 405 )); 406 assert_eq!( 407 block_on(storage.execute_authored(original)) 408 .unwrap() 409 .disposition(), 410 AtomicCommitDisposition::Replay 411 ); 412 assert_eq!(plan(&storage), before); 413 } 414 415 #[test] 416 fn reconciliation_rejects_future_facts_zero_time_and_invalid_retry_without_mutation() { 417 for accepted in [false, true] { 418 let (storage, active, _) = prepared(); 419 block_on(storage.execute_authored(fact(&plan(&storage), active.clone(), accepted, 30))) 420 .unwrap(); 421 let before = plan(&storage); 422 assert_eq!( 423 ReconcileDeliveryFacts::new(&before, None, None, 0), 424 Err(Error::AtomicWorkflowMismatch) 425 ); 426 assert_eq!( 427 block_on(storage.execute_authored(reconcile(&before, Some(fence(&active)), None, 29))), 428 Err(Error::AtomicWorkflowMismatch) 429 ); 430 let invalid = if accepted { Some(retry(1, 60)) } else { None }; 431 assert_eq!( 432 block_on(storage.execute_authored(reconcile(&before, None, invalid, 50))), 433 Err(Error::InvalidRetrySchedule) 434 ); 435 if !accepted { 436 assert_eq!( 437 block_on(storage.execute_authored(reconcile( 438 &before, 439 None, 440 Some(retry(2, 60)), 441 50 442 ))), 443 Err(Error::InvalidRetrySchedule) 444 ); 445 } 446 assert_eq!(plan(&storage), before); 447 } 448 } 449 450 #[test] 451 fn raw_sink_failure_controls_retry_diagnostic_and_provider_backoff() { 452 let (storage, active, _) = prepared(); 453 let failure = SinkFailure::for_request( 454 plan(&storage).request().unwrap(), 455 "sink_lost", 456 Retryability::Retryable, 457 Some(65), 458 Some("connection lost".into()), 459 vec![], 460 ) 461 .unwrap(); 462 block_on( 463 storage.execute_authored(AuthoredAtomicCommand::RecordDelivery( 464 RecordDeliveryFact::new( 465 ids().2, 466 ids().1, 467 active, 468 DeliveryAttemptOutcome::SinkFailure(failure.clone()), 469 50, 470 ) 471 .unwrap(), 472 )), 473 ) 474 .unwrap(); 475 let before = plan(&storage); 476 for (code, diagnostic, at) in [ 477 ("other", Some("connection lost"), 65), 478 ("sink_lost", None, 65), 479 ("sink_lost", Some("connection lost"), 64), 480 ] { 481 let retry = RetrySchedule::new( 482 NonZeroU32::MIN, 483 at, 484 WorkFailure::new( 485 code, 486 WorkPhase::Delivery, 487 FailureClass::Retryable, 488 None, 489 diagnostic.map(str::to_owned), 490 ) 491 .unwrap(), 492 ) 493 .unwrap(); 494 assert_eq!( 495 block_on(storage.execute_authored(reconcile(&before, None, Some(retry), 51))), 496 Err(Error::InvalidRetrySchedule) 497 ); 498 assert_eq!(plan(&storage), before); 499 } 500 let retry = RetrySchedule::new( 501 NonZeroU32::MIN, 502 65, 503 WorkFailure::new( 504 "sink_lost", 505 WorkPhase::Delivery, 506 FailureClass::Retryable, 507 Some(65), 508 Some("connection lost".into()), 509 ) 510 .unwrap(), 511 ) 512 .unwrap(); 513 block_on(storage.execute_authored(reconcile(&before, None, Some(retry.clone()), 51))).unwrap(); 514 let after = plan(&storage); 515 assert_eq!(after.state(), AuthoredDeliveryState::Retryable); 516 assert_eq!(after.retry(), Some(&retry)); 517 assert_eq!( 518 after.attempts()[0].outcome(), 519 &DeliveryAttemptOutcome::SinkFailure(failure) 520 ); 521 } 522 523 #[test] 524 fn matching_late_fact_marks_legacy_attempt_without_inventing_another_attempt() { 525 for accepted in [false, true] { 526 let (storage, active, _) = prepared(); 527 let command = AuthoredAtomicCommand::ApplyDelivery( 528 ApplyDeliveryAttempt::new( 529 ids().2, 530 fence(&active), 531 outcome(&plan(&storage), false), 532 Some(retry(1, 18)), 533 14, 534 ) 535 .unwrap(), 536 ); 537 let original = block_on(storage.execute_authored(command)).unwrap(); 538 block_on(storage.execute_authored(fact(&plan(&storage), active.clone(), accepted, 50))) 539 .unwrap(); 540 let before = plan(&storage); 541 let command = reconcile(&before, None, Some(retry(1, 60)), 51); 542 if accepted { 543 assert_eq!( 544 block_on(storage.execute_authored(command)), 545 Err(Error::AtomicWorkflowMismatch) 546 ); 547 assert_eq!(plan(&storage), before); 548 // Conflicting raw facts remain factual evidence, not a replacement 549 // for the immutable earlier attempt or authority for another effect. 550 assert_eq!( 551 before.delivery_satisfaction().unwrap(), 552 SatisfactionState::Satisfied 553 ); 554 } else { 555 block_on(storage.execute_authored(command.clone())).unwrap(); 556 let after = plan(&storage); 557 assert_eq!(after.attempt_count(), 1); 558 assert_eq!(after.attempts()[0].recorded_at_unix_ms(), 14); 559 assert_eq!(after.attempts()[0].claim_evidence(), Some(&active)); 560 assert_eq!(after.pending_delivery_facts().count(), 0); 561 assert_eq!( 562 block_on(storage.execute_authored(command)) 563 .unwrap() 564 .disposition(), 565 AtomicCommitDisposition::Replay 566 ); 567 } 568 assert_eq!( 569 block_on(storage.authored_receipt(original.commit_id())) 570 .unwrap() 571 .unwrap(), 572 original 573 ); 574 } 575 } 576 577 #[test] 578 fn partial_acceptance_and_terminal_evidence_settle_without_retry_authority() { 579 for (retryability, partial, expected) in [ 580 ( 581 Retryability::Terminal, 582 None, 583 AuthoredDeliveryState::FailedTerminal, 584 ), 585 ( 586 Retryability::Retryable, 587 Some(true), 588 AuthoredDeliveryState::Satisfied, 589 ), 590 ( 591 Retryability::Retryable, 592 Some(false), 593 AuthoredDeliveryState::Exhausted, 594 ), 595 ( 596 Retryability::Terminal, 597 Some(false), 598 AuthoredDeliveryState::Exhausted, 599 ), 600 ] { 601 let (storage, active, _) = prepared(); 602 let before = plan(&storage); 603 let evidence = partial 604 .map(|accepted| { 605 DeliveryTargetReceipt::attempted( 606 before.request().unwrap().target_set().targets()[0].clone(), 607 if accepted { 608 DeliveryOutcome::accepted() 609 } else { 610 DeliveryOutcome::rejected() 611 }, 612 ) 613 }) 614 .into_iter() 615 .collect(); 616 let failure = SinkFailure::for_request( 617 before.request().unwrap(), 618 "sink_lost", 619 retryability, 620 None, 621 None, 622 evidence, 623 ) 624 .unwrap(); 625 block_on( 626 storage.execute_authored(AuthoredAtomicCommand::RecordDelivery( 627 RecordDeliveryFact::new( 628 ids().2, 629 ids().1, 630 active, 631 DeliveryAttemptOutcome::SinkFailure(failure), 632 50, 633 ) 634 .unwrap(), 635 )), 636 ) 637 .unwrap(); 638 let before = plan(&storage); 639 assert_eq!( 640 block_on(storage.execute_authored(reconcile(&before, None, Some(retry(1, 60)), 51))), 641 Err(Error::InvalidRetrySchedule) 642 ); 643 assert_eq!(plan(&storage), before); 644 block_on(storage.execute_authored(reconcile(&before, None, None, 51))).unwrap(); 645 assert_eq!(plan(&storage).state(), expected); 646 assert!(plan(&storage).retry().is_none()); 647 } 648 } 649 650 #[test] 651 fn earlier_unresolved_claim_cannot_adopt_a_new_workers_legacy_attempt() { 652 let (storage, first, _) = prepared(); 653 let second = claim(plan(&storage).revision(), 3, 40); 654 block_on( 655 storage.execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( 656 ClaimAuthoredTarget::DeliveryPlan(ids().2), 657 second.clone(), 658 ))), 659 ) 660 .unwrap(); 661 block_on(storage.execute_authored(fact(&plan(&storage), first, false, 41))).unwrap(); 662 block_on(storage.execute_authored(reconcile( 663 &plan(&storage), 664 Some(fence(&second)), 665 Some(retry(1, 43)), 666 42, 667 ))) 668 .unwrap(); 669 let third = claim(plan(&storage).revision(), 4, 44); 670 block_on( 671 storage.execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( 672 ClaimAuthoredTarget::DeliveryPlan(ids().2), 673 third.clone(), 674 ))), 675 ) 676 .unwrap(); 677 let legacy = AuthoredAtomicCommand::ApplyDelivery( 678 ApplyDeliveryAttempt::new( 679 ids().2, 680 fence(&third), 681 outcome(&plan(&storage), false), 682 Some(retry(2, 47)), 683 45, 684 ) 685 .unwrap(), 686 ); 687 block_on(storage.execute_authored(legacy)).unwrap(); 688 assert!(history(&storage).has_unresolved_claims()); 689 block_on(storage.execute_authored(fact(&plan(&storage), second.clone(), false, 50))).unwrap(); 690 block_on(storage.execute_authored(reconcile(&plan(&storage), None, Some(retry(3, 65)), 61))) 691 .unwrap(); 692 let after = plan(&storage); 693 assert_eq!(after.attempt_count(), 3); 694 assert!(after.attempts()[1].claim_evidence().is_none()); 695 assert_eq!(after.attempts()[2].claim_evidence(), Some(&second)); 696 assert!(!history(&storage).has_unresolved_claims()); 697 } 698 699 fn restored_receipt( 700 original: &AuthoredAtomicReceipt, 701 plan: AuthoredDeliveryPlan, 702 at: u64, 703 ) -> AuthoredAtomicReceipt { 704 AuthoredAtomicReceipt::from_durable_parts( 705 original.commit_id(), 706 original.digest(), 707 AtomicCommitDisposition::Committed, 708 at, 709 AuthoredAtomicOutcome::DeliveryPlan(plan), 710 ) 711 .unwrap() 712 } 713 714 #[test] 715 fn original_claim_history_rejects_each_mismatched_binding() { 716 let (storage, _, original) = prepared(); 717 let current = plan(&storage); 718 for field in ["plan_id", "artifact_id", "request", "created_at_unix_ms"] { 719 let mut wire = serde_json::to_value(¤t).unwrap(); 720 match field { 721 "plan_id" => { 722 wire[field] = serde_json::to_value( 723 radroots_storage::authored_delivery::AuthoredDeliveryPlanId::new([99; 16]) 724 .unwrap(), 725 ) 726 .unwrap() 727 } 728 "artifact_id" => { 729 wire[field] = serde_json::to_value( 730 radroots_storage::authored::AuthoredArtifactId::new([99; 16]).unwrap(), 731 ) 732 .unwrap() 733 } 734 "request" => wire[field] = serde_json::Value::Null, 735 _ => wire[field] = serde_json::json!(9), 736 } 737 let altered = serde_json::from_value(wire).unwrap(); 738 let forged = restored_receipt(&original, altered, 13); 739 let mut history = AuthoredDeliveryHistory::new(current.clone(), None).unwrap(); 740 assert_eq!( 741 history.push_claim(&forged), 742 Err(Error::AtomicWorkflowMismatch), 743 "{field}" 744 ); 745 assert!(history.claims().is_empty()); 746 } 747 let mut rebound = serde_json::to_value(¤t).unwrap(); 748 rebound["request"] = serde_json::to_value( 749 current 750 .intent() 751 .materialize(radroots_transport::sink::DeliveryPayload::new(event( 752 OTHER_RAW, 753 ))) 754 .unwrap(), 755 ) 756 .unwrap(); 757 let mut history = 758 AuthoredDeliveryHistory::new(serde_json::from_value(rebound).unwrap(), None).unwrap(); 759 assert_eq!( 760 history.push_claim(&original), 761 Err(Error::AtomicWorkflowMismatch) 762 ); 763 let mut history = AuthoredDeliveryHistory::new(current.clone(), None).unwrap(); 764 assert_eq!( 765 history.push_claim(&restored_receipt(&original, current.clone(), 14)), 766 Err(Error::AtomicWorkflowMismatch) 767 ); 768 let corrupt_digest = AuthoredAtomicReceipt::from_durable_parts( 769 original.commit_id(), 770 AtomicCommitDigest::new([99; 32]), 771 AtomicCommitDisposition::Committed, 772 13, 773 original.outcome().clone(), 774 ) 775 .unwrap(); 776 assert_eq!( 777 history.push_claim(&corrupt_digest), 778 Err(Error::AtomicWorkflowMismatch) 779 ); 780 // A structurally valid future original cannot explain the current row. 781 for (revision, at) in [ 782 (current.revision().get() + 1, 13), 783 (current.revision().get(), 14), 784 ] { 785 let mut wire = serde_json::to_value(¤t).unwrap(); 786 let active = claim(NonZeroU64::new(revision - 1).unwrap(), 8, at); 787 wire["revision"] = serde_json::json!(revision); 788 wire["updated_at_unix_ms"] = serde_json::json!(at); 789 wire["claim"] = serde_json::to_value(&active).unwrap(); 790 let future = serde_json::from_value(wire).unwrap(); 791 let command = AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( 792 ClaimAuthoredTarget::DeliveryPlan(ids().2), 793 active, 794 )); 795 let receipt = AuthoredAtomicReceipt::new( 796 &command, 797 AtomicCommitDisposition::Committed, 798 at, 799 AuthoredAtomicOutcome::DeliveryPlan(future), 800 ) 801 .unwrap(); 802 assert_eq!( 803 history.push_claim(&receipt), 804 Err(Error::AtomicWorkflowMismatch) 805 ); 806 } 807 } 808 809 #[test] 810 fn preparation_history_requires_exact_initial_identity_and_monotonic_row() { 811 let storage = MemoryStorage::new(SourceGeneration::new([73; 32]).unwrap()); 812 let original = block_on(storage.execute_authored(prepare().0)).unwrap(); 813 let initial = plan(&storage); 814 for field in ["plan_id", "artifact_id", "created_at_unix_ms"] { 815 let mut wire = serde_json::to_value(&initial).unwrap(); 816 match field { 817 "plan_id" => { 818 wire[field] = serde_json::to_value( 819 radroots_storage::authored_delivery::AuthoredDeliveryPlanId::new([99; 16]) 820 .unwrap(), 821 ) 822 .unwrap() 823 } 824 "artifact_id" => { 825 wire[field] = serde_json::to_value( 826 radroots_storage::authored::AuthoredArtifactId::new([99; 16]).unwrap(), 827 ) 828 .unwrap() 829 } 830 _ => wire[field] = serde_json::json!(9), 831 } 832 assert!( 833 AuthoredDeliveryHistory::new(serde_json::from_value(wire).unwrap(), Some(&original)) 834 .is_err(), 835 "{field}" 836 ); 837 } 838 let intent = radroots_storage::authored_delivery::AuthoredDeliveryIntent::new( 839 "different-intent", 840 initial.intent().target_set().clone(), 841 initial.intent().satisfaction().clone(), 842 100, 843 ) 844 .unwrap(); 845 let changed = AuthoredDeliveryPlan::new(ids().2, ids().1, intent, 10).unwrap(); 846 assert!(AuthoredDeliveryHistory::new(changed, Some(&original)).is_err()); 847 for field in ["revision", "updated_at_unix_ms"] { 848 let mut wire = serde_json::to_value(&initial).unwrap(); 849 wire[field] = serde_json::json!(12); 850 let altered = serde_json::from_value(wire).unwrap(); 851 let AuthoredAtomicOutcome::Prepared { 852 operation, 853 artifacts, 854 .. 855 } = original.outcome() 856 else { 857 unreachable!() 858 }; 859 let receipt = AuthoredAtomicReceipt::from_durable_parts( 860 original.commit_id(), 861 original.digest(), 862 AtomicCommitDisposition::Committed, 863 12, 864 AuthoredAtomicOutcome::Prepared { 865 operation: operation.clone(), 866 artifacts: artifacts.clone(), 867 delivery_plans: vec![altered], 868 }, 869 ) 870 .unwrap(); 871 assert!( 872 AuthoredDeliveryHistory::new(initial.clone(), Some(&receipt)).is_err(), 873 "{field}" 874 ); 875 } 876 } 877 878 #[test] 879 fn reconciliation_receipts_cannot_substitute_a_different_plan_or_result() { 880 let (storage, active, _) = prepared(); 881 block_on(storage.execute_authored(fact(&plan(&storage), active.clone(), false, 14))).unwrap(); 882 let command = reconcile( 883 &plan(&storage), 884 Some(fence(&active)), 885 Some(retry(1, 20)), 886 15, 887 ); 888 let receipt = block_on(storage.execute_authored(command.clone())).unwrap(); 889 let after = plan(&storage); 890 for field in [ 891 "plan_id", 892 "revision", 893 "updated_at_unix_ms", 894 "claim", 895 "stop", 896 "pending", 897 "facts", 898 "retry", 899 ] { 900 let mut wire = serde_json::to_value(&after).unwrap(); 901 match field { 902 "plan_id" => { 903 wire[field] = serde_json::to_value( 904 radroots_storage::authored_delivery::AuthoredDeliveryPlanId::new([99; 16]) 905 .unwrap(), 906 ) 907 .unwrap() 908 } 909 "revision" => wire[field] = serde_json::json!(after.revision().get() + 1), 910 "updated_at_unix_ms" => wire[field] = serde_json::json!(16), 911 "claim" => { 912 wire[field] = serde_json::to_value(claim( 913 NonZeroU64::new(after.revision().get() - 1).unwrap(), 914 9, 915 15, 916 )) 917 .unwrap(); 918 } 919 "stop" => { 920 wire["stop_requested_at_unix_ms"] = serde_json::json!(15); 921 wire["state"] = serde_json::json!("cancelled"); 922 wire["retry"] = serde_json::Value::Null; 923 wire["last_failure"] = serde_json::Value::Null; 924 } 925 "pending" => wire["attempts"][0]["claim"] = serde_json::Value::Null, 926 "facts" => wire["delivery_facts"][0]["observed_at_unix_ms"] = serde_json::json!(15), 927 _ => { 928 let other = retry(1, 21); 929 wire["retry"] = serde_json::to_value(&other).unwrap(); 930 wire["last_failure"] = serde_json::to_value(other.failure()).unwrap(); 931 } 932 } 933 let altered: AuthoredDeliveryPlan = serde_json::from_value(wire).unwrap(); 934 let wrong = restored_receipt(&receipt, altered.clone(), 15); 935 assert!(!wrong.matches_command(&command), "{field}"); 936 assert_eq!( 937 AuthoredAtomicReceipt::new( 938 &command, 939 AtomicCommitDisposition::Committed, 940 15, 941 AuthoredAtomicOutcome::DeliveryPlan(altered) 942 ), 943 Err(Error::AtomicWorkflowMismatch), 944 "{field}" 945 ); 946 } 947 let artifact = block_on(storage.authored_artifact(ids().1)) 948 .unwrap() 949 .unwrap(); 950 let wrong = AuthoredAtomicReceipt::from_durable_parts( 951 receipt.commit_id(), 952 receipt.digest(), 953 AtomicCommitDisposition::Committed, 954 15, 955 AuthoredAtomicOutcome::Artifact(artifact.clone()), 956 ) 957 .unwrap(); 958 assert!(!wrong.matches_command(&command)); 959 assert!( 960 AuthoredAtomicReceipt::new( 961 &command, 962 AtomicCommitDisposition::Committed, 963 15, 964 AuthoredAtomicOutcome::Artifact(artifact) 965 ) 966 .is_err() 967 ); 968 } 969 970 #[test] 971 fn missing_fact_claims_and_prepopulated_preparation_remain_uncertain() { 972 let (storage, active, _) = prepared(); 973 let original = block_on(storage.authored_receipt(prepare().0.commit_id())) 974 .unwrap() 975 .unwrap(); 976 let AuthoredAtomicOutcome::Prepared { 977 operation, 978 artifacts, 979 .. 980 } = original.outcome() 981 else { 982 unreachable!() 983 }; 984 let prepopulated = AuthoredAtomicReceipt::from_durable_parts( 985 original.commit_id(), 986 original.digest(), 987 AtomicCommitDisposition::Committed, 988 13, 989 AuthoredAtomicOutcome::Prepared { 990 operation: operation.clone(), 991 artifacts: artifacts.clone(), 992 delivery_plans: vec![plan(&storage)], 993 }, 994 ) 995 .unwrap(); 996 let legacy = AuthoredDeliveryHistory::new(plan(&storage), Some(&prepopulated)).unwrap(); 997 assert!(!legacy.is_complete()); 998 assert!(!legacy.proves_no_issued_attempt()); 999 assert!(legacy.has_unresolved_claims()); 1000 block_on(storage.execute_authored(fact(&plan(&storage), active, true, 50))).unwrap(); 1001 stop(&storage, 51); 1002 let unknown = AuthoredDeliveryHistory::new(plan(&storage), None).unwrap(); 1003 assert_eq!( 1004 unknown.require_pending_fact_provenance(), 1005 Err(Error::AtomicWorkflowMismatch) 1006 ); 1007 let incomplete = AuthoredDeliveryHistory::new(plan(&storage), Some(&original)).unwrap(); 1008 assert_eq!(incomplete.validate(), Err(Error::AtomicWorkflowMismatch)); 1009 } 1010 1011 #[test] 1012 fn later_legacy_attempt_outside_old_lease_cannot_resolve_old_claim() { 1013 let (storage, _, _) = prepared(); 1014 let newer = claim(plan(&storage).revision(), 3, 40); 1015 block_on( 1016 storage.execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( 1017 ClaimAuthoredTarget::DeliveryPlan(ids().2), 1018 newer.clone(), 1019 ))), 1020 ) 1021 .unwrap(); 1022 block_on( 1023 storage.execute_authored(AuthoredAtomicCommand::ApplyDelivery( 1024 ApplyDeliveryAttempt::new( 1025 ids().2, 1026 fence(&newer), 1027 outcome(&plan(&storage), true), 1028 None, 1029 41, 1030 ) 1031 .unwrap(), 1032 )), 1033 ) 1034 .unwrap(); 1035 let history = history(&storage); 1036 assert_eq!(history.plan().attempt_count(), 1); 1037 assert_eq!(history.plan().state(), AuthoredDeliveryState::Satisfied); 1038 assert!(history.has_unresolved_claims()); 1039 }