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 }