authored_delivery.rs (29831B)
1 use core::num::{NonZeroU32, NonZeroU64}; 2 use radroots_event::{SignedEvent, wire::v1::Nip01EventWire}; 3 use radroots_storage::{ 4 Error, 5 authored::{ 6 AuthoredArtifactId, FailureClass, RetrySchedule, WorkClaim, WorkFailure, WorkPhase, 7 }, 8 authored_delivery::{ 9 AuthoredDeliveryAttempt, AuthoredDeliveryIntent, AuthoredDeliveryPlan, 10 AuthoredDeliveryPlanId, AuthoredDeliveryState, DELIVERY_PLAN_ATTEMPTS_MAX, 11 DeliveryAttemptOutcome, 12 }, 13 }; 14 use radroots_transport::{ 15 DeliveryReceipt, DeliveryRequest, SinkFailure, Target, TargetSet, 16 outcome::{DeliveryOutcome, Retryability}, 17 policy::{SatisfactionClass, SatisfactionPolicy, SatisfactionState, TargetPolicy}, 18 sink::{DeliveryPayload, DeliveryTargetReceipt}, 19 }; 20 21 fn signed_event() -> SignedEvent { 22 let mut wire = Nip01EventWire { 23 id: "0".repeat(64), 24 pubkey: "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df".to_owned(), 25 created_at: 1_800_000_100, 26 kind: 20_000, 27 tags: Vec::new(), 28 content: "delivery plan".to_owned(), 29 sig: "dd".repeat(64), 30 extra: Default::default(), 31 }; 32 wire.id = wire.computed_event_id().expect("event ID").to_hex(); 33 let raw = serde_json::to_string(&wire).expect("raw event"); 34 SignedEvent::from_wire_verified_id(wire, raw).expect("signed event") 35 } 36 37 fn target_set() -> TargetSet { 38 TargetSet::new(vec![ 39 Target::nostr_relay("wss://one.example").expect("one"), 40 Target::nostr_relay("wss://two.example").expect("two"), 41 ]) 42 .expect("targets") 43 } 44 45 fn request(policy: TargetPolicy) -> DeliveryRequest { 46 DeliveryRequest::new( 47 "authored-delivery", 48 DeliveryPayload::new(signed_event()), 49 target_set(), 50 SatisfactionPolicy::new(SatisfactionClass::Accepted, policy), 51 100, 52 ) 53 .expect("request") 54 } 55 56 fn plan(value: u8, policy: TargetPolicy) -> AuthoredDeliveryPlan { 57 AuthoredDeliveryPlan::new_bound( 58 AuthoredDeliveryPlanId::new([value; 16]).expect("plan ID"), 59 AuthoredArtifactId::new([9; 16]).expect("artifact ID"), 60 request(policy), 61 10, 62 ) 63 .expect("plan") 64 } 65 66 fn bound_request(plan: &AuthoredDeliveryPlan) -> &DeliveryRequest { 67 plan.request().expect("bound request") 68 } 69 70 fn claim(plan: &AuthoredDeliveryPlan, token: u8, acquired: u64) -> WorkClaim { 71 WorkClaim::new( 72 [token; 16], 73 format!("worker-{token}"), 74 NonZeroU64::new(u64::from(token)).expect("generation"), 75 acquired, 76 acquired + 20, 77 plan.revision(), 78 ) 79 .expect("claim") 80 } 81 82 fn receipt(request: &DeliveryRequest, outcomes: Vec<DeliveryOutcome>) -> DeliveryReceipt { 83 DeliveryReceipt::for_request( 84 request, 85 request 86 .target_set() 87 .targets() 88 .iter() 89 .cloned() 90 .zip(outcomes) 91 .map(|(target, outcome)| DeliveryTargetReceipt::attempted(target, outcome)) 92 .collect(), 93 ) 94 .expect("receipt") 95 } 96 97 fn retry(code: &str, attempt: u32, not_before: u64) -> RetrySchedule { 98 let failure = WorkFailure::new( 99 code, 100 WorkPhase::Delivery, 101 FailureClass::Retryable, 102 Some(not_before), 103 None, 104 ) 105 .expect("failure"); 106 RetrySchedule::new( 107 NonZeroU32::new(attempt).expect("attempt"), 108 not_before, 109 failure, 110 ) 111 .expect("retry") 112 } 113 114 #[test] 115 fn independent_plans_claim_and_progress_without_cross_blocking() { 116 let mut first = plan(1, TargetPolicy::any()); 117 let mut second = plan(2, TargetPolicy::all()); 118 let first_claim = claim(&first, 1, 11); 119 let second_claim = claim(&second, 2, 11); 120 first.claim(first_claim.clone(), 11).expect("first claim"); 121 second 122 .claim(second_claim.clone(), 11) 123 .expect("second claim"); 124 125 let first_receipt = receipt( 126 bound_request(&first), 127 vec![DeliveryOutcome::accepted(), DeliveryOutcome::unavailable()], 128 ); 129 first 130 .apply_receipt( 131 first_claim.token(), 132 first_claim.generation(), 133 first_claim.row_revision(), 134 first_receipt, 135 None, 136 12, 137 ) 138 .expect("first satisfied"); 139 assert_eq!(first.state(), AuthoredDeliveryState::Satisfied); 140 assert!(second.claim_evidence().is_some()); 141 142 let second_receipt = receipt( 143 bound_request(&second), 144 vec![DeliveryOutcome::accepted(), DeliveryOutcome::unavailable()], 145 ); 146 second 147 .apply_receipt( 148 second_claim.token(), 149 second_claim.generation(), 150 second_claim.row_revision(), 151 second_receipt, 152 Some(retry("delivery_pending", 1, 20)), 153 12, 154 ) 155 .expect("second retry"); 156 assert_eq!(second.state(), AuthoredDeliveryState::Retryable); 157 assert_eq!(second.retry().expect("retry").not_before_unix_ms(), 20); 158 let before_early_claim = second.clone(); 159 let early = claim(&second, 7, 19); 160 assert_eq!( 161 second.claim(early, 19), 162 Err(Error::DeliveryPlanClaimConflict) 163 ); 164 assert_eq!(second, before_early_claim); 165 } 166 167 #[test] 168 fn retry_schedule_and_partial_evidence_survive_reconstruction() { 169 let mut delivery = plan(3, TargetPolicy::all()); 170 let active = claim(&delivery, 3, 11); 171 delivery.claim(active.clone(), 11).expect("claim"); 172 let partial = DeliveryTargetReceipt::attempted( 173 bound_request(&delivery).target_set().targets()[0].clone(), 174 DeliveryOutcome::accepted(), 175 ); 176 let sink_failure = SinkFailure::for_request( 177 bound_request(&delivery), 178 "relay_batch_unavailable", 179 Retryability::Retryable, 180 Some(20), 181 None, 182 vec![partial], 183 ) 184 .expect("sink failure"); 185 delivery 186 .apply_sink_failure( 187 active.token(), 188 active.generation(), 189 active.row_revision(), 190 sink_failure, 191 Some(retry("relay_batch_unavailable", 1, 20)), 192 12, 193 ) 194 .expect("retryable failure"); 195 196 let json = serde_json::to_string(&delivery).expect("plan json"); 197 let reopened: AuthoredDeliveryPlan = serde_json::from_str(&json).expect("reopen plan"); 198 assert_eq!(reopened, delivery); 199 assert_eq!(reopened.attempt_count(), 1); 200 assert_eq!(reopened.state(), AuthoredDeliveryState::Retryable); 201 assert_eq!( 202 reopened.attempts()[0].satisfaction(), 203 SatisfactionState::Pending 204 ); 205 assert_eq!( 206 reopened.last_failure().expect("failure").code(), 207 "relay_batch_unavailable" 208 ); 209 } 210 211 #[test] 212 fn terminal_partial_success_stale_claim_and_invalid_receipt_fail_closed() { 213 let mut any = plan(4, TargetPolicy::any()); 214 let active = claim(&any, 4, 11); 215 any.claim(active.clone(), 11).expect("claim"); 216 let partial = DeliveryTargetReceipt::attempted( 217 bound_request(&any).target_set().targets()[0].clone(), 218 DeliveryOutcome::accepted(), 219 ); 220 let failure = SinkFailure::for_request( 221 bound_request(&any), 222 "terminal_batch_failure", 223 Retryability::Terminal, 224 None, 225 None, 226 vec![partial], 227 ) 228 .expect("terminal failure"); 229 let other_request = DeliveryRequest::new( 230 "other-request", 231 bound_request(&any).payload().clone(), 232 bound_request(&any).target_set().clone(), 233 bound_request(&any).satisfaction().clone(), 234 bound_request(&any).deadline_unix_ms(), 235 ) 236 .expect("other request"); 237 let invalid_receipt = receipt( 238 &other_request, 239 vec![DeliveryOutcome::accepted(), DeliveryOutcome::accepted()], 240 ); 241 assert_eq!( 242 any.apply_receipt( 243 active.token(), 244 active.generation(), 245 active.row_revision(), 246 invalid_receipt, 247 None, 248 12, 249 ), 250 Err(Error::InvalidAuthoredDeliveryPlan) 251 ); 252 assert_eq!(any.attempt_count(), 0); 253 assert!(any.claim_evidence().is_some()); 254 assert_eq!( 255 any.apply_sink_failure( 256 &[8; 16], 257 active.generation(), 258 active.row_revision(), 259 failure.clone(), 260 None, 261 12, 262 ), 263 Err(Error::DeliveryPlanClaimConflict) 264 ); 265 assert_eq!(any.attempt_count(), 0); 266 any.apply_sink_failure( 267 active.token(), 268 active.generation(), 269 active.row_revision(), 270 failure, 271 None, 272 12, 273 ) 274 .expect("partial satisfaction"); 275 assert_eq!(any.state(), AuthoredDeliveryState::Satisfied); 276 277 let mut all = plan(5, TargetPolicy::all()); 278 let all_claim = claim(&all, 5, 11); 279 all.claim(all_claim.clone(), 11).expect("claim"); 280 let terminal = SinkFailure::for_request( 281 bound_request(&all), 282 "terminal_batch_failure", 283 Retryability::Terminal, 284 None, 285 None, 286 Vec::new(), 287 ) 288 .expect("terminal failure"); 289 all.apply_sink_failure( 290 all_claim.token(), 291 all_claim.generation(), 292 all_claim.row_revision(), 293 terminal, 294 None, 295 12, 296 ) 297 .expect("terminal apply"); 298 assert_eq!(all.state(), AuthoredDeliveryState::FailedTerminal); 299 } 300 301 #[test] 302 fn attempt_limit_is_checked_without_mutating_claimed_state() { 303 let base = plan(6, TargetPolicy::all()); 304 let pending_receipt = receipt( 305 bound_request(&base), 306 vec![ 307 DeliveryOutcome::unavailable(), 308 DeliveryOutcome::unavailable(), 309 ], 310 ); 311 let attempts: Vec<_> = (1..=DELIVERY_PLAN_ATTEMPTS_MAX) 312 .map(|attempt| { 313 AuthoredDeliveryAttempt::reconstruct( 314 NonZeroU32::new(attempt).expect("attempt"), 315 12, 316 DeliveryAttemptOutcome::Receipt(pending_receipt.clone()), 317 SatisfactionState::Pending, 318 ) 319 .expect("attempt record") 320 }) 321 .collect(); 322 let mut value = serde_json::to_value(&base).expect("plan value"); 323 value["state"] = serde_json::json!("retryable"); 324 value["attempts"] = serde_json::to_value(attempts).expect("attempts value"); 325 value["attempt_count"] = serde_json::json!(DELIVERY_PLAN_ATTEMPTS_MAX); 326 let schedule = retry("delivery_pending", DELIVERY_PLAN_ATTEMPTS_MAX, 20); 327 value["retry"] = serde_json::to_value(&schedule).expect("retry value"); 328 value["last_failure"] = serde_json::to_value(schedule.failure()).expect("failure value"); 329 value["updated_at_unix_ms"] = serde_json::json!(12); 330 let mut saturated: AuthoredDeliveryPlan = 331 serde_json::from_value(value).expect("saturated plan"); 332 let active = claim(&saturated, 6, 20); 333 saturated.claim(active.clone(), 20).expect("claim"); 334 let before = saturated.clone(); 335 assert_eq!( 336 saturated.apply_receipt( 337 active.token(), 338 active.generation(), 339 active.row_revision(), 340 pending_receipt, 341 Some(retry("delivery_pending", DELIVERY_PLAN_ATTEMPTS_MAX, 30,)), 342 21, 343 ), 344 Err(Error::DeliveryAttemptOverflow) 345 ); 346 assert_eq!(saturated, before); 347 } 348 349 #[test] 350 fn delivery_identity_intent_state_and_attempt_contracts_are_total() { 351 assert_eq!( 352 AuthoredDeliveryPlanId::new([0; 16]), 353 Err(Error::InvalidAuthoredDeliveryPlan) 354 ); 355 let id = AuthoredDeliveryPlanId::new([7; 16]).unwrap(); 356 assert_eq!(id.as_bytes(), &[7; 16]); 357 assert_eq!(AuthoredDeliveryPlanId::try_from([7; 16]).unwrap(), id); 358 assert_eq!(<[u8; 16]>::from(id), [7; 16]); 359 360 let targets = target_set(); 361 let policy = SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::all()); 362 assert_eq!( 363 AuthoredDeliveryIntent::new("", targets.clone(), policy.clone(), 100), 364 Err(Error::InvalidAuthoredDeliveryPlan) 365 ); 366 assert_eq!( 367 AuthoredDeliveryIntent::new("intent", targets.clone(), policy.clone(), 0), 368 Err(Error::InvalidAuthoredDeliveryPlan) 369 ); 370 let invalid_policy = SatisfactionPolicy::new( 371 SatisfactionClass::Accepted, 372 TargetPolicy::required(vec![ 373 Target::nostr_relay("wss://foreign.example") 374 .unwrap() 375 .fingerprint() 376 .clone(), 377 ]) 378 .unwrap(), 379 ); 380 assert_eq!( 381 AuthoredDeliveryIntent::new("intent", targets.clone(), invalid_policy, 100), 382 Err(Error::InvalidAuthoredDeliveryPlan) 383 ); 384 385 for policy in [ 386 TargetPolicy::any(), 387 TargetPolicy::all(), 388 TargetPolicy::quorum(2).unwrap(), 389 TargetPolicy::required(vec![targets.targets()[0].fingerprint().clone()]).unwrap(), 390 ] { 391 for class in [SatisfactionClass::Accepted, SatisfactionClass::Delivered] { 392 let request = DeliveryRequest::new( 393 "intent", 394 DeliveryPayload::new(signed_event()), 395 targets.clone(), 396 SatisfactionPolicy::new(class, policy.clone()), 397 100, 398 ) 399 .unwrap(); 400 let intent = AuthoredDeliveryIntent::from_request(&request); 401 assert_eq!(intent.request_id(), request.request_id()); 402 assert_eq!(intent.target_set(), request.target_set()); 403 assert_eq!(intent.satisfaction(), request.satisfaction()); 404 assert_eq!(intent.deadline_unix_ms(), 100); 405 assert_eq!( 406 intent.materialize(request.payload().clone()).unwrap(), 407 request 408 ); 409 AuthoredDeliveryPlan::new(id, artifact_id(9), intent, 10).unwrap(); 410 } 411 } 412 413 for (state, terminal) in [ 414 (AuthoredDeliveryState::Pending, false), 415 (AuthoredDeliveryState::Retryable, false), 416 (AuthoredDeliveryState::Satisfied, true), 417 (AuthoredDeliveryState::Exhausted, true), 418 (AuthoredDeliveryState::FailedTerminal, true), 419 (AuthoredDeliveryState::Cancelled, true), 420 ] { 421 assert_eq!(state.is_terminal(), terminal); 422 } 423 424 let outcome = DeliveryAttemptOutcome::Receipt(receipt( 425 &request(TargetPolicy::all()), 426 vec![DeliveryOutcome::accepted(), DeliveryOutcome::unavailable()], 427 )); 428 assert_eq!( 429 AuthoredDeliveryAttempt::reconstruct( 430 NonZeroU32::MIN, 431 0, 432 outcome.clone(), 433 SatisfactionState::Pending, 434 ), 435 Err(Error::InvalidAuthoredDeliveryPlan) 436 ); 437 let attempt = AuthoredDeliveryAttempt::reconstruct( 438 NonZeroU32::MIN, 439 12, 440 outcome.clone(), 441 SatisfactionState::Pending, 442 ) 443 .unwrap(); 444 assert_eq!(attempt.attempt(), NonZeroU32::MIN); 445 assert_eq!(attempt.recorded_at_unix_ms(), 12); 446 assert_eq!(attempt.outcome(), &outcome); 447 assert_eq!(attempt.satisfaction(), SatisfactionState::Pending); 448 } 449 450 fn artifact_id(value: u8) -> AuthoredArtifactId { 451 AuthoredArtifactId::new([value; 16]).unwrap() 452 } 453 454 fn assert_plan_json_rejected( 455 mut value: serde_json::Value, 456 key: &str, 457 replacement: serde_json::Value, 458 ) { 459 value[key] = replacement; 460 assert!( 461 serde_json::from_value::<AuthoredDeliveryPlan>(value).is_err(), 462 "forged field {key} was accepted" 463 ); 464 } 465 466 #[test] 467 fn reconstructed_delivery_plans_fail_closed_for_each_durable_invariant() { 468 let base = plan(8, TargetPolicy::all()); 469 let value = serde_json::to_value(&base).unwrap(); 470 assert_plan_json_rejected(value.clone(), "created_at_unix_ms", serde_json::json!(0)); 471 assert_plan_json_rejected(value.clone(), "updated_at_unix_ms", serde_json::json!(9)); 472 assert_plan_json_rejected( 473 value.clone(), 474 "request_digest", 475 serde_json::to_value([0_u8; 32]).unwrap(), 476 ); 477 assert_plan_json_rejected( 478 value.clone(), 479 "attempt_count", 480 serde_json::json!(DELIVERY_PLAN_ATTEMPTS_MAX + 1), 481 ); 482 assert_plan_json_rejected(value.clone(), "attempt_count", serde_json::json!(1)); 483 assert_plan_json_rejected(value.clone(), "state", serde_json::json!("retryable")); 484 assert_plan_json_rejected(value.clone(), "state", serde_json::json!("failed_terminal")); 485 486 let mut wrong_exhausted_phase = value.clone(); 487 wrong_exhausted_phase["state"] = serde_json::json!("exhausted"); 488 wrong_exhausted_phase["last_failure"] = 489 serde_json::to_value(retry_failure_for_phase(WorkPhase::Signing, 20)).unwrap(); 490 assert!(serde_json::from_value::<AuthoredDeliveryPlan>(wrong_exhausted_phase).is_err()); 491 let mut wrong_exhausted_class = value.clone(); 492 wrong_exhausted_class["state"] = serde_json::json!("exhausted"); 493 wrong_exhausted_class["last_failure"] = 494 serde_json::to_value(retry_failure_for_phase(WorkPhase::Delivery, 20)).unwrap(); 495 assert!(serde_json::from_value::<AuthoredDeliveryPlan>(wrong_exhausted_class).is_err()); 496 497 let mut request_mismatch = value.clone(); 498 request_mismatch["request"]["request_id"] = serde_json::json!("different-request"); 499 assert!(serde_json::from_value::<AuthoredDeliveryPlan>(request_mismatch).is_err()); 500 501 let pending_receipt = receipt( 502 bound_request(&base), 503 vec![ 504 DeliveryOutcome::unavailable(), 505 DeliveryOutcome::unavailable(), 506 ], 507 ); 508 let attempt = AuthoredDeliveryAttempt::reconstruct( 509 NonZeroU32::MIN, 510 12, 511 DeliveryAttemptOutcome::Receipt(pending_receipt), 512 SatisfactionState::Pending, 513 ) 514 .unwrap(); 515 let mut attempt_without_request = value.clone(); 516 attempt_without_request["request"] = serde_json::Value::Null; 517 attempt_without_request["attempts"] = serde_json::json!([attempt]); 518 attempt_without_request["attempt_count"] = serde_json::json!(1); 519 assert!(serde_json::from_value::<AuthoredDeliveryPlan>(attempt_without_request).is_err()); 520 521 let active = claim(&base, 8, 11); 522 let mut claimed = base.clone(); 523 claimed.claim(active, 11).unwrap(); 524 let claimed_value = serde_json::to_value(&claimed).unwrap(); 525 assert_plan_json_rejected( 526 claimed_value.clone(), 527 "state", 528 serde_json::json!("satisfied"), 529 ); 530 for (field, replacement) in [ 531 ("token", serde_json::to_value([0_u8; 16]).unwrap()), 532 ("acquired_at_unix_ms", serde_json::json!(10)), 533 ("row_revision", serde_json::json!(9)), 534 ] { 535 let mut forged = claimed_value.clone(); 536 forged["claim"][field] = replacement; 537 assert!(serde_json::from_value::<AuthoredDeliveryPlan>(forged).is_err()); 538 } 539 540 let mut retryable = plan(9, TargetPolicy::all()); 541 let active = claim(&retryable, 9, 11); 542 retryable.claim(active.clone(), 11).unwrap(); 543 retryable 544 .apply_receipt( 545 active.token(), 546 active.generation(), 547 active.row_revision(), 548 receipt( 549 bound_request(&retryable), 550 vec![DeliveryOutcome::accepted(), DeliveryOutcome::unavailable()], 551 ), 552 Some(retry("delivery_pending", 1, 20)), 553 12, 554 ) 555 .unwrap(); 556 let retryable_value = serde_json::to_value(&retryable).unwrap(); 557 for (key, replacement) in [ 558 ("attempt_count", serde_json::json!(0)), 559 ("retry", serde_json::Value::Null), 560 ("last_failure", serde_json::Value::Null), 561 ] { 562 assert_plan_json_rejected(retryable_value.clone(), key, replacement); 563 } 564 for (field, replacement) in [ 565 ("attempt", serde_json::json!(2)), 566 ("recorded_at_unix_ms", serde_json::json!(9)), 567 ("satisfaction", serde_json::json!("satisfied")), 568 ] { 569 let mut forged = retryable_value.clone(); 570 forged["attempts"][0][field] = replacement; 571 assert!(serde_json::from_value::<AuthoredDeliveryPlan>(forged).is_err()); 572 } 573 for (field, replacement) in [ 574 ("attempt", serde_json::json!(2)), 575 ("not_before_unix_ms", serde_json::json!(12)), 576 ] { 577 let mut forged = retryable_value.clone(); 578 forged["retry"][field] = replacement; 579 assert!(serde_json::from_value::<AuthoredDeliveryPlan>(forged).is_err()); 580 } 581 } 582 583 #[test] 584 fn delivery_claim_binding_cancellation_and_retry_validation_are_transactional() { 585 let mut unbound = AuthoredDeliveryPlan::new( 586 AuthoredDeliveryPlanId::new([10; 16]).unwrap(), 587 artifact_id(10), 588 AuthoredDeliveryIntent::from_request(&request(TargetPolicy::all())), 589 10, 590 ) 591 .unwrap(); 592 assert!(unbound.request().is_none()); 593 assert_eq!(unbound.attempts(), &[]); 594 assert_eq!(unbound.attempt_count(), 0); 595 assert_eq!(unbound.retry(), None); 596 assert_eq!(unbound.claim_evidence(), None); 597 assert_eq!(unbound.last_failure(), None); 598 assert_eq!(unbound.created_at_unix_ms(), 10); 599 assert_eq!(unbound.updated_at_unix_ms(), 10); 600 assert_eq!(unbound.revision(), NonZeroU64::MIN); 601 let impossible_claim = WorkClaim::new( 602 [1; 16], 603 "worker", 604 NonZeroU64::MIN, 605 11, 606 20, 607 unbound.revision(), 608 ) 609 .unwrap(); 610 assert_eq!( 611 unbound.claim(impossible_claim, 11), 612 Err(Error::DeliveryPlanClaimConflict) 613 ); 614 assert_eq!( 615 unbound.evaluate_next_attempt(&DeliveryAttemptOutcome::Receipt(receipt( 616 &request(TargetPolicy::all()), 617 vec![DeliveryOutcome::accepted(), DeliveryOutcome::accepted()], 618 ))), 619 Err(Error::InvalidAuthoredDeliveryPlan) 620 ); 621 assert_eq!( 622 unbound.bind_signed_event(signed_event(), 9), 623 Err(Error::InvalidAuthoredDeliveryPlan) 624 ); 625 unbound.bind_signed_event(signed_event(), 11).unwrap(); 626 assert!(unbound.request().is_some()); 627 assert_eq!( 628 unbound.bind_signed_event(signed_event(), 12), 629 Err(Error::InvalidAuthoredDeliveryPlan) 630 ); 631 632 let active = claim(&unbound, 10, 12); 633 let stale_revision = WorkClaim::new( 634 [11; 16], 635 "worker", 636 NonZeroU64::new(11).unwrap(), 637 12, 638 20, 639 NonZeroU64::MIN, 640 ) 641 .unwrap(); 642 assert_eq!( 643 unbound.claim(stale_revision, 12), 644 Err(Error::DeliveryPlanClaimConflict) 645 ); 646 unbound.claim(active.clone(), 12).unwrap(); 647 assert_eq!( 648 unbound.claim(claim(&unbound, 11, 13), 13), 649 Err(Error::DeliveryPlanClaimConflict) 650 ); 651 let before = unbound.clone(); 652 assert_eq!( 653 unbound.apply_receipt( 654 active.token(), 655 active.generation(), 656 active.row_revision(), 657 receipt( 658 bound_request(&unbound), 659 vec![ 660 DeliveryOutcome::unavailable(), 661 DeliveryOutcome::unavailable() 662 ], 663 ), 664 Some( 665 RetrySchedule::new( 666 NonZeroU32::MIN, 667 20, 668 retry_failure_for_phase(WorkPhase::Signing, 20), 669 ) 670 .unwrap() 671 ), 672 13, 673 ), 674 Err(Error::InvalidRetrySchedule) 675 ); 676 assert_eq!(unbound, before); 677 678 let mut cancelled = plan(11, TargetPolicy::all()); 679 cancelled.cancel(11).unwrap(); 680 assert_eq!(cancelled.state(), AuthoredDeliveryState::Cancelled); 681 assert_eq!( 682 cancelled.cancel(12), 683 Err(Error::InvalidAuthoredDeliveryPlan) 684 ); 685 assert_eq!( 686 cancelled.claim(claim(&cancelled, 12, 12), 12), 687 Err(Error::DeliveryPlanClaimConflict) 688 ); 689 let mut backwards = plan(12, TargetPolicy::all()); 690 let before = backwards.clone(); 691 assert_eq!(backwards.cancel(9), Err(Error::InvalidAuthoredDeliveryPlan)); 692 assert_eq!(backwards, before); 693 } 694 695 fn retry_failure_for_phase(phase: WorkPhase, at: u64) -> WorkFailure { 696 WorkFailure::new( 697 "delivery_pending", 698 phase, 699 FailureClass::Retryable, 700 Some(at), 701 None, 702 ) 703 .unwrap() 704 } 705 706 fn claimed_plan(value: u8) -> (AuthoredDeliveryPlan, WorkClaim) { 707 let mut plan = plan(value, TargetPolicy::all()); 708 let active = claim(&plan, value, 11); 709 plan.claim(active.clone(), 11).unwrap(); 710 (plan, active) 711 } 712 713 #[test] 714 fn delivery_attempt_retry_and_terminal_error_matrix_is_fail_closed() { 715 let (mut satisfied, active) = claimed_plan(20); 716 let accepted = receipt( 717 bound_request(&satisfied), 718 vec![DeliveryOutcome::accepted(), DeliveryOutcome::accepted()], 719 ); 720 assert_eq!( 721 satisfied.apply_receipt( 722 active.token(), 723 active.generation(), 724 active.row_revision(), 725 accepted, 726 Some(retry("unexpected_retry", 1, 20)), 727 12, 728 ), 729 Err(Error::InvalidRetrySchedule) 730 ); 731 732 let (mut exhausted, active) = claimed_plan(21); 733 let rejected = receipt( 734 bound_request(&exhausted), 735 vec![DeliveryOutcome::rejected(), DeliveryOutcome::rejected()], 736 ); 737 assert_eq!( 738 exhausted.apply_receipt( 739 active.token(), 740 active.generation(), 741 active.row_revision(), 742 rejected, 743 Some(retry("unexpected_retry", 1, 20)), 744 12, 745 ), 746 Err(Error::InvalidRetrySchedule) 747 ); 748 749 let (mut validation_rollback, active) = claimed_plan(22); 750 let pending = receipt( 751 bound_request(&validation_rollback), 752 vec![DeliveryOutcome::accepted(), DeliveryOutcome::unavailable()], 753 ); 754 let before = validation_rollback.clone(); 755 assert_eq!( 756 validation_rollback.apply_receipt( 757 active.token(), 758 active.generation(), 759 active.row_revision(), 760 pending, 761 Some(retry("delivery_pending", 2, 12)), 762 12, 763 ), 764 Err(Error::InvalidAuthoredDeliveryPlan) 765 ); 766 assert_eq!(validation_rollback, before); 767 768 for (value, evidence, retryability, retry_schedule, expected) in [ 769 ( 770 23, 771 vec![DeliveryOutcome::accepted(), DeliveryOutcome::accepted()], 772 Retryability::Terminal, 773 Some(retry("sink_failure", 1, 20)), 774 Err(Error::InvalidRetrySchedule), 775 ), 776 ( 777 24, 778 vec![DeliveryOutcome::rejected(), DeliveryOutcome::rejected()], 779 Retryability::Terminal, 780 Some(retry("sink_failure", 1, 20)), 781 Err(Error::InvalidRetrySchedule), 782 ), 783 ( 784 25, 785 Vec::new(), 786 Retryability::Retryable, 787 Some(retry("different_failure", 1, 20)), 788 Err(Error::InvalidRetrySchedule), 789 ), 790 ( 791 26, 792 Vec::new(), 793 Retryability::Terminal, 794 Some(retry("sink_failure", 1, 20)), 795 Err(Error::InvalidRetrySchedule), 796 ), 797 ] { 798 let (mut plan, active) = claimed_plan(value); 799 let request = bound_request(&plan); 800 let partial = request 801 .target_set() 802 .targets() 803 .iter() 804 .cloned() 805 .zip(evidence) 806 .map(|(target, outcome)| DeliveryTargetReceipt::attempted(target, outcome)) 807 .collect(); 808 let failure = SinkFailure::for_request( 809 request, 810 "sink_failure", 811 retryability, 812 (retryability == Retryability::Retryable).then_some(20), 813 None, 814 partial, 815 ) 816 .unwrap(); 817 let before = plan.clone(); 818 assert_eq!( 819 plan.apply_sink_failure( 820 active.token(), 821 active.generation(), 822 active.row_revision(), 823 failure, 824 retry_schedule, 825 12, 826 ), 827 expected 828 ); 829 assert_eq!(plan, before); 830 } 831 832 let (mut exhausted, active) = claimed_plan(27); 833 let request = bound_request(&exhausted); 834 let partial = request 835 .target_set() 836 .targets() 837 .iter() 838 .cloned() 839 .map(|target| DeliveryTargetReceipt::attempted(target, DeliveryOutcome::rejected())) 840 .collect(); 841 let failure = SinkFailure::for_request( 842 request, 843 "sink_failure", 844 Retryability::Terminal, 845 None, 846 None, 847 partial, 848 ) 849 .unwrap(); 850 exhausted 851 .apply_sink_failure( 852 active.token(), 853 active.generation(), 854 active.row_revision(), 855 failure, 856 None, 857 12, 858 ) 859 .unwrap(); 860 assert_eq!(exhausted.state(), AuthoredDeliveryState::Exhausted); 861 } 862 863 #[test] 864 fn delivery_claim_guards_distinguish_generation_time_revision_and_clock() { 865 let mut delivery = plan(28, TargetPolicy::all()); 866 let first = WorkClaim::new( 867 [1; 16], 868 "worker", 869 NonZeroU64::new(2).unwrap(), 870 11, 871 20, 872 delivery.revision(), 873 ) 874 .unwrap(); 875 delivery.claim(first, 11).unwrap(); 876 let lower_generation = WorkClaim::new( 877 [2; 16], 878 "worker", 879 NonZeroU64::MIN, 880 20, 881 30, 882 delivery.revision(), 883 ) 884 .unwrap(); 885 assert_eq!( 886 delivery.claim(lower_generation, 20), 887 Err(Error::DeliveryPlanClaimConflict) 888 ); 889 890 let mut wrong_time = plan(29, TargetPolicy::all()); 891 let acquired_later = WorkClaim::new( 892 [3; 16], 893 "worker", 894 NonZeroU64::MIN, 895 12, 896 20, 897 wrong_time.revision(), 898 ) 899 .unwrap(); 900 assert_eq!( 901 wrong_time.claim(acquired_later, 11), 902 Err(Error::DeliveryPlanClaimConflict) 903 ); 904 905 let mut backwards = plan(30, TargetPolicy::all()); 906 let claim = WorkClaim::new( 907 [4; 16], 908 "worker", 909 NonZeroU64::MIN, 910 9, 911 20, 912 backwards.revision(), 913 ) 914 .unwrap(); 915 let before = backwards.clone(); 916 assert_eq!( 917 backwards.claim(claim, 9), 918 Err(Error::InvalidAuthoredDeliveryPlan) 919 ); 920 assert_eq!(backwards, before); 921 }