lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

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 }