lib

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

delivery_contract.rs (19841B)


      1 use radroots_event::{SignedEvent, wire::v1::Nip01EventWire};
      2 use radroots_transport::{
      3     DeliveryReceipt, DeliveryRequest, Error, SinkFailure, Target, TargetSet,
      4     outcome::{DeliveryOutcome, DeliveryOutcomeKind, Retryability},
      5     policy::{
      6         SatisfactionClass, SatisfactionPolicy, SatisfactionState, TargetPolicy,
      7         evaluate_satisfaction,
      8     },
      9     sink::{DeliveryPayload, DeliveryTargetReceipt},
     10 };
     11 
     12 fn payload() -> DeliveryPayload {
     13     let raw = r#"{"id":"56bfc78223bb2221bad82b539efdec1ade0f56d0eb0e1f592fd387df4b2ceee0","pubkey":"585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df","created_at":1700000001,"kind":0,"tags":[],"content":"{}","sig":"dddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddd"}"#;
     14     let wire = Nip01EventWire::parse_json(raw).expect("wire event");
     15     DeliveryPayload::new(
     16         SignedEvent::from_wire_verified_id(wire, raw).expect("signed delivery event"),
     17     )
     18 }
     19 
     20 fn targets() -> TargetSet {
     21     TargetSet::new(vec![
     22         Target::nostr_relay("wss://one.example").expect("first"),
     23         Target::nostr_relay("wss://two.example").expect("second"),
     24         Target::nostr_relay("wss://three.example").expect("third"),
     25     ])
     26     .expect("targets")
     27 }
     28 
     29 fn request(policy: SatisfactionPolicy) -> DeliveryRequest {
     30     DeliveryRequest::new(
     31         "delivery-request",
     32         payload(),
     33         targets(),
     34         policy,
     35         1_700_000_100_000,
     36     )
     37     .expect("delivery request")
     38 }
     39 
     40 fn mixed_receipt(request: &DeliveryRequest) -> DeliveryReceipt {
     41     let targets = request.target_set().targets();
     42     DeliveryReceipt::for_request(
     43         request,
     44         vec![
     45             DeliveryTargetReceipt::attempted(
     46                 targets[2].clone(),
     47                 DeliveryOutcome::unavailable()
     48                     .with_detail("offline", "relay unavailable")
     49                     .expect("normalized detail"),
     50             ),
     51             DeliveryTargetReceipt::attempted(targets[0].clone(), DeliveryOutcome::delivered()),
     52             DeliveryTargetReceipt::attempted(targets[1].clone(), DeliveryOutcome::accepted()),
     53         ],
     54     )
     55     .expect("mixed receipt")
     56 }
     57 
     58 #[test]
     59 fn target_selection_is_exact_and_default_adapter_never_widens_it() {
     60     use radroots_transport::{BoxFuture, EventSink, SinkStatus};
     61     use std::sync::atomic::{AtomicUsize, Ordering};
     62     struct Sink(AtomicUsize);
     63     impl EventSink for Sink {
     64         fn status(&self) -> BoxFuture<'_, Result<SinkStatus, Error>> {
     65             Box::pin(async { panic!("no status probe") })
     66         }
     67         fn deliver(
     68             &self,
     69             request: DeliveryRequest,
     70         ) -> BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> {
     71             Box::pin(async move {
     72                 self.0.fetch_add(1, Ordering::SeqCst);
     73                 Ok(mixed_receipt(&request))
     74             })
     75         }
     76     }
     77     let sink = Sink(AtomicUsize::new(0));
     78     let request = request(SatisfactionPolicy::new(
     79         SatisfactionClass::Accepted,
     80         TargetPolicy::all(),
     81     ));
     82     let original = request.clone();
     83     let labelled = Target::new_with_metadata(
     84         radroots_transport::TransportId::NOSTR,
     85         "wss://one.example",
     86         None,
     87         Some(radroots_transport::target::TargetLabel::parse("changed label").unwrap()),
     88     )
     89     .unwrap();
     90     assert_eq!(
     91         labelled.fingerprint(),
     92         request.target_set().targets()[0].fingerprint()
     93     );
     94     assert_eq!(
     95         request.validate_target_selection(&TargetSet::new(vec![labelled]).unwrap()),
     96         Err(Error::InvalidDeliveryTargetSelection)
     97     );
     98     let subset = TargetSet::new(vec![request.target_set().targets()[0].clone()]).unwrap();
     99     request.validate_target_selection(&subset).unwrap();
    100     let unsupported =
    101         futures::executor::block_on(sink.deliver_selected(request.clone(), subset)).unwrap_err();
    102     unsupported.validate_for_request(&request).unwrap();
    103     assert_eq!(unsupported.code(), "target_selection_unsupported");
    104     let foreign =
    105         TargetSet::new(vec![Target::nostr_relay("wss://foreign.example").unwrap()]).unwrap();
    106     assert_eq!(
    107         request.validate_target_selection(&foreign),
    108         Err(Error::InvalidDeliveryTargetSelection)
    109     );
    110     assert_eq!(
    111         Error::InvalidDeliveryTargetSelection.to_string(),
    112         "transport delivery target selection is not an exact subset"
    113     );
    114     let invalid =
    115         futures::executor::block_on(sink.deliver_selected(request.clone(), foreign)).unwrap_err();
    116     invalid.validate_for_request(&request).unwrap();
    117     assert_eq!(invalid.code(), "invalid_transport_contract");
    118     assert_eq!(sink.0.load(Ordering::SeqCst), 0);
    119     let mut reversed = request.target_set().targets().to_vec();
    120     reversed.reverse();
    121     let receipt = futures::executor::block_on(
    122         sink.deliver_selected(request.clone(), TargetSet::new(reversed).unwrap()),
    123     )
    124     .unwrap();
    125     receipt.validate_for_request(&request).unwrap();
    126     assert_eq!(sink.0.load(Ordering::SeqCst), 1);
    127     assert_eq!(request, original);
    128 }
    129 
    130 #[test]
    131 fn any_all_quorum_and_required_targets_are_exact() {
    132     let any = request(SatisfactionPolicy::new(
    133         SatisfactionClass::Accepted,
    134         TargetPolicy::any(),
    135     ));
    136     assert!(mixed_receipt(&any).is_satisfied(&any).expect("any"));
    137 
    138     let all = request(SatisfactionPolicy::new(
    139         SatisfactionClass::Accepted,
    140         TargetPolicy::all(),
    141     ));
    142     assert!(!mixed_receipt(&all).is_satisfied(&all).expect("all"));
    143 
    144     let quorum = request(SatisfactionPolicy::new(
    145         SatisfactionClass::Accepted,
    146         TargetPolicy::quorum(2).expect("quorum"),
    147     ));
    148     assert!(
    149         mixed_receipt(&quorum)
    150             .is_satisfied(&quorum)
    151             .expect("quorum")
    152     );
    153 
    154     let delivered = request(SatisfactionPolicy::new(
    155         SatisfactionClass::Delivered,
    156         TargetPolicy::quorum(2).expect("quorum"),
    157     ));
    158     assert!(
    159         !mixed_receipt(&delivered)
    160             .is_satisfied(&delivered)
    161             .expect("delivered quorum")
    162     );
    163 
    164     let selected = targets();
    165     let required = TargetPolicy::required(vec![
    166         selected.targets()[0].fingerprint().clone(),
    167         selected.targets()[1].fingerprint().clone(),
    168     ])
    169     .expect("required targets");
    170     let required_request = request(SatisfactionPolicy::new(
    171         SatisfactionClass::Accepted,
    172         required,
    173     ));
    174     assert!(
    175         mixed_receipt(&required_request)
    176             .is_satisfied(&required_request)
    177             .expect("required")
    178     );
    179 
    180     let unsatisfied = TargetPolicy::required(vec![
    181         selected.targets()[0].fingerprint().clone(),
    182         selected.targets()[2].fingerprint().clone(),
    183     ])
    184     .expect("required targets");
    185     let unsatisfied_request = request(SatisfactionPolicy::new(
    186         SatisfactionClass::Accepted,
    187         unsatisfied,
    188     ));
    189     assert!(
    190         !mixed_receipt(&unsatisfied_request)
    191             .is_satisfied(&unsatisfied_request)
    192             .expect("required unsatisfied")
    193     );
    194 }
    195 
    196 #[test]
    197 fn policies_and_requests_reject_empty_duplicate_and_impossible_inputs() {
    198     assert_eq!(
    199         TargetPolicy::quorum(0).expect_err("zero quorum"),
    200         Error::InvalidSatisfactionPolicy
    201     );
    202     assert_eq!(
    203         TargetPolicy::required(Vec::new()).expect_err("empty required"),
    204         Error::EmptyRequiredTargetSet
    205     );
    206     let set = targets();
    207     let duplicate = set.targets()[0].fingerprint().clone();
    208     assert_eq!(
    209         TargetPolicy::required(vec![duplicate.clone(), duplicate]).expect_err("duplicate required"),
    210         Error::DuplicateRequiredTargetFingerprint
    211     );
    212     assert_eq!(
    213         DeliveryRequest::new(
    214             "request",
    215             payload(),
    216             set.clone(),
    217             SatisfactionPolicy::new(
    218                 SatisfactionClass::Accepted,
    219                 TargetPolicy::quorum(4).expect("nonzero quorum"),
    220             ),
    221             1,
    222         )
    223         .expect_err("impossible quorum"),
    224         Error::InvalidSatisfactionPolicy
    225     );
    226     assert_eq!(
    227         DeliveryRequest::new(
    228             "",
    229             payload(),
    230             set.clone(),
    231             SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::any()),
    232             1,
    233         )
    234         .expect_err("empty id"),
    235         Error::EmptyDeliveryRequestId
    236     );
    237     assert_eq!(
    238         DeliveryRequest::new(
    239             "request",
    240             payload(),
    241             set,
    242             SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::any()),
    243             0,
    244         )
    245         .expect_err("zero deadline"),
    246         Error::InvalidDeliveryDeadline
    247     );
    248     assert_eq!(
    249         TargetSet::new(Vec::new()).expect_err("empty target set"),
    250         Error::EmptyTargetSet
    251     );
    252 }
    253 
    254 #[test]
    255 fn receipts_reject_duplicate_missing_unexpected_and_false_attempts() {
    256     let request = request(SatisfactionPolicy::new(
    257         SatisfactionClass::Accepted,
    258         TargetPolicy::all(),
    259     ));
    260     let first = request.target_set().targets()[0].clone();
    261     let second = request.target_set().targets()[1].clone();
    262     let third = request.target_set().targets()[2].clone();
    263     let accepted = DeliveryTargetReceipt::attempted(first.clone(), DeliveryOutcome::accepted());
    264     assert_eq!(
    265         DeliveryTargetReceipt::skipped(first.clone(), DeliveryOutcome::accepted())
    266             .expect_err("unattempted success"),
    267         Error::DeliveryTargetReceiptAttemptMismatch
    268     );
    269     assert_eq!(
    270         DeliveryReceipt::for_request(&request, vec![accepted.clone(), accepted])
    271             .expect_err("duplicate receipt"),
    272         Error::DuplicateDeliveryTargetReceipt
    273     );
    274     assert_eq!(
    275         DeliveryReceipt::for_request(
    276             &request,
    277             vec![
    278                 DeliveryTargetReceipt::attempted(first, DeliveryOutcome::accepted()),
    279                 DeliveryTargetReceipt::attempted(second, DeliveryOutcome::accepted()),
    280             ],
    281         )
    282         .expect_err("missing receipt"),
    283         Error::MissingDeliveryTargetReceipt
    284     );
    285     let foreign = Target::nostr_relay("wss://foreign.example").expect("foreign");
    286     assert_eq!(
    287         DeliveryReceipt::for_request(
    288             &request,
    289             vec![
    290                 DeliveryTargetReceipt::attempted(foreign, DeliveryOutcome::accepted()),
    291                 DeliveryTargetReceipt::attempted(third, DeliveryOutcome::accepted()),
    292             ],
    293         )
    294         .expect_err("unexpected receipt"),
    295         Error::UnexpectedDeliveryTargetReceipt
    296     );
    297 }
    298 
    299 #[test]
    300 fn retryability_and_terminality_are_explicit_normalized_data() {
    301     let unavailable = DeliveryOutcome::unavailable();
    302     assert_eq!(unavailable.kind(), DeliveryOutcomeKind::Unavailable);
    303     assert_eq!(unavailable.retryability(), Retryability::Retryable);
    304     assert!(unavailable.is_retryable());
    305     assert!(!unavailable.is_terminal());
    306 
    307     let rejected = DeliveryOutcome::rejected();
    308     assert!(rejected.is_terminal());
    309     assert!(!rejected.is_retryable());
    310     assert_eq!(
    311         DeliveryOutcome::failed(Retryability::NotApplicable).expect_err("unclassified failure"),
    312         Error::InvalidDeliveryOutcome
    313     );
    314     assert!(
    315         DeliveryOutcome::failed(Retryability::Retryable)
    316             .expect("retryable failure")
    317             .is_retryable()
    318     );
    319     assert!(
    320         DeliveryOutcome::failed(Retryability::Terminal)
    321             .expect("terminal failure")
    322             .is_terminal()
    323     );
    324     assert_eq!(
    325         DeliveryOutcome::unavailable()
    326             .with_detail("INVALID", "relay unavailable")
    327             .expect_err("invalid code"),
    328         Error::InvalidDeliveryOutcome
    329     );
    330 }
    331 
    332 #[test]
    333 fn canonical_evaluator_covers_pending_exhausted_and_historical_evidence() {
    334     let set = targets();
    335     let policy = SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::all());
    336     assert_eq!(
    337         evaluate_satisfaction(&policy, &set, core::iter::empty()).expect("empty evidence"),
    338         SatisfactionState::Pending
    339     );
    340 
    341     let first = set.targets()[0].fingerprint();
    342     let second = set.targets()[1].fingerprint();
    343     let third = set.targets()[2].fingerprint();
    344     let accepted = DeliveryOutcome::accepted();
    345     let terminal = DeliveryOutcome::rejected();
    346     let retryable = DeliveryOutcome::unavailable();
    347     assert_eq!(
    348         evaluate_satisfaction(
    349             &policy,
    350             &set,
    351             [(first, &accepted), (second, &terminal), (third, &retryable),],
    352         )
    353         .expect("mixed evidence"),
    354         SatisfactionState::Exhausted
    355     );
    356     assert_eq!(
    357         evaluate_satisfaction(
    358             &policy,
    359             &set,
    360             [
    361                 (first, &accepted),
    362                 (second, &accepted),
    363                 (third, &accepted),
    364                 (first, &terminal),
    365             ],
    366         )
    367         .expect("historical success"),
    368         SatisfactionState::Satisfied
    369     );
    370 
    371     let foreign = Target::nostr_relay("wss://foreign.example").expect("foreign");
    372     assert_eq!(
    373         evaluate_satisfaction(&policy, &set, [(foreign.fingerprint(), &accepted)])
    374             .expect_err("foreign evidence"),
    375         Error::UnexpectedDeliveryTargetReceipt
    376     );
    377 }
    378 
    379 #[test]
    380 fn policy_introspection_validation_and_wire_variants_are_complete() {
    381     let set = targets();
    382     let any = TargetPolicy::any();
    383     assert!(any.is_any());
    384     assert!(!any.is_all());
    385     assert_eq!(any.quorum_threshold(), None);
    386     assert_eq!(any.required_targets(), None);
    387 
    388     let all = TargetPolicy::all();
    389     assert!(!all.is_any());
    390     assert!(all.is_all());
    391     assert_eq!(all.quorum_threshold(), None);
    392     assert_eq!(all.required_targets(), None);
    393 
    394     let quorum = TargetPolicy::quorum(2).expect("quorum");
    395     assert!(!quorum.is_any());
    396     assert!(!quorum.is_all());
    397     assert_eq!(quorum.quorum_threshold(), Some(2));
    398     assert_eq!(quorum.required_targets(), None);
    399 
    400     let fingerprint = set.targets()[0].fingerprint().clone();
    401     let required = TargetPolicy::required(vec![fingerprint.clone()]).expect("required");
    402     assert!(!required.is_any());
    403     assert!(!required.is_all());
    404     assert_eq!(required.quorum_threshold(), None);
    405     assert_eq!(
    406         required.required_targets(),
    407         Some(core::slice::from_ref(&fingerprint))
    408     );
    409 
    410     let policy = SatisfactionPolicy::new(SatisfactionClass::Delivered, required);
    411     assert_eq!(policy.class(), SatisfactionClass::Delivered);
    412     assert_eq!(
    413         policy.targets().required_targets(),
    414         Some(core::slice::from_ref(&fingerprint))
    415     );
    416     policy
    417         .validate_for(&set)
    418         .expect("required target is requested");
    419 
    420     let foreign = Target::nostr_relay("wss://foreign.example").expect("foreign");
    421     let impossible = SatisfactionPolicy::new(
    422         SatisfactionClass::Accepted,
    423         TargetPolicy::required(vec![foreign.fingerprint().clone()]).expect("required"),
    424     );
    425     assert_eq!(
    426         impossible.validate_for(&set).unwrap_err(),
    427         Error::RequiredTargetNotRequested
    428     );
    429 
    430     #[cfg(feature = "serde")]
    431     for (wire, assertion) in [
    432         (serde_json::json!("any"), 0_u8),
    433         (serde_json::json!("all"), 1),
    434         (serde_json::json!({"quorum": 2}), 2),
    435         (serde_json::json!({"required": [fingerprint.as_str()]}), 3),
    436     ] {
    437         let decoded: TargetPolicy = serde_json::from_value(wire).expect("target policy");
    438         match assertion {
    439             0 => assert!(decoded.is_any()),
    440             1 => assert!(decoded.is_all()),
    441             2 => assert_eq!(decoded.quorum_threshold(), Some(2)),
    442             3 => assert_eq!(
    443                 decoded.required_targets(),
    444                 Some(core::slice::from_ref(&fingerprint))
    445             ),
    446             _ => unreachable!(),
    447         }
    448     }
    449 
    450     #[cfg(feature = "serde")]
    451     for invalid in [
    452         serde_json::json!({"quorum": 0}),
    453         serde_json::json!({"required": []}),
    454         serde_json::json!({"required": [fingerprint.as_str(), fingerprint.as_str()]}),
    455     ] {
    456         assert!(serde_json::from_value::<TargetPolicy>(invalid).is_err());
    457     }
    458 }
    459 
    460 #[test]
    461 fn sink_failures_retain_validated_retry_timing_and_partial_evidence() {
    462     let request = request(SatisfactionPolicy::new(
    463         SatisfactionClass::Accepted,
    464         TargetPolicy::all(),
    465     ));
    466     let partial = DeliveryTargetReceipt::attempted(
    467         request.target_set().targets()[0].clone(),
    468         DeliveryOutcome::accepted(),
    469     );
    470     let failure = SinkFailure::for_request(
    471         &request,
    472         "relay_batch_unavailable",
    473         Retryability::Retryable,
    474         Some(1_700_000_200_000),
    475         Some("relay batch unavailable".to_owned()),
    476         vec![partial.clone()],
    477     )
    478     .expect("sink failure");
    479     assert_eq!(failure.code(), "relay_batch_unavailable");
    480     assert_eq!(failure.retryability(), Retryability::Retryable);
    481     assert_eq!(failure.retry_after_unix_ms(), Some(1_700_000_200_000));
    482     assert_eq!(failure.message(), Some("relay batch unavailable"));
    483     assert_eq!(failure.partial_evidence(), core::slice::from_ref(&partial));
    484     assert_eq!(
    485         SinkFailure::for_request(
    486             &request,
    487             "terminal_failure",
    488             Retryability::Terminal,
    489             Some(1),
    490             None,
    491             Vec::new(),
    492         )
    493         .expect_err("terminal retry timing"),
    494         Error::InvalidDeliveryOutcome
    495     );
    496     assert_eq!(
    497         SinkFailure::for_request(
    498             &request,
    499             "invalid_retry_time",
    500             Retryability::Retryable,
    501             Some(0),
    502             None,
    503             Vec::new(),
    504         )
    505         .expect_err("zero retry timing"),
    506         Error::InvalidDeliveryOutcome
    507     );
    508     assert_eq!(
    509         SinkFailure::for_request(
    510             &request,
    511             "duplicate_evidence",
    512             Retryability::Retryable,
    513             None,
    514             None,
    515             vec![partial.clone(), partial],
    516         )
    517         .expect_err("duplicate evidence"),
    518         Error::DuplicateDeliveryTargetReceipt
    519     );
    520 }
    521 
    522 #[test]
    523 #[cfg(feature = "serde")]
    524 fn serde_revalidates_policy_outcome_and_receipt_invariants() {
    525     let request = request(SatisfactionPolicy::new(
    526         SatisfactionClass::Accepted,
    527         TargetPolicy::all(),
    528     ));
    529     let receipt = mixed_receipt(&request);
    530     let encoded_request = serde_json::to_string(&request).expect("request json");
    531     assert_eq!(
    532         serde_json::from_str::<DeliveryRequest>(&encoded_request).expect("request round trip"),
    533         request
    534     );
    535     let encoded_receipt = serde_json::to_string(&receipt).expect("receipt json");
    536     assert_eq!(
    537         serde_json::from_str::<DeliveryReceipt>(&encoded_receipt).expect("receipt round trip"),
    538         receipt
    539     );
    540 
    541     let mut forged_outcome =
    542         serde_json::to_value(DeliveryOutcome::accepted()).expect("outcome value");
    543     forged_outcome["retryability"] = serde_json::json!("retryable");
    544     assert!(serde_json::from_value::<DeliveryOutcome>(forged_outcome).is_err());
    545 
    546     let mut forged_attempt = serde_json::to_value(&receipt).expect("receipt value");
    547     forged_attempt["target_receipts"][0]["attempted"] = false.into();
    548     assert!(serde_json::from_value::<DeliveryReceipt>(forged_attempt).is_err());
    549 
    550     let mut missing = serde_json::to_value(&receipt).expect("receipt value");
    551     missing["target_receipts"]
    552         .as_array_mut()
    553         .expect("receipts array")
    554         .pop();
    555     assert!(serde_json::from_value::<DeliveryReceipt>(missing).is_err());
    556 
    557     let failure = SinkFailure::for_request(
    558         &request,
    559         "relay_unavailable",
    560         Retryability::Retryable,
    561         Some(1_700_000_200_000),
    562         None,
    563         vec![receipt.target_receipts()[0].clone()],
    564     )
    565     .expect("sink failure");
    566     let encoded_failure = serde_json::to_string(&failure).expect("failure json");
    567     assert_eq!(
    568         serde_json::from_str::<SinkFailure>(&encoded_failure).expect("failure round trip"),
    569         failure
    570     );
    571 }