outbox.rs (57084B)
1 //! Durable outbox and delivery-evidence contracts. 2 //! 3 //! This module stores delivery intent and normalized evidence. It deliberately 4 //! owns no transport adapter and performs no transport I/O. 5 6 use core::fmt; 7 pub use radroots_transport::{ 8 BoxFuture, DeliveryReceipt, DeliveryRequest, TransportId, 9 outcome::{DeliveryOutcome, DeliveryOutcomeKind, Retryability}, 10 policy::{ 11 SatisfactionClass, SatisfactionPolicy, SatisfactionState, TargetPolicy, 12 evaluate_satisfaction as evaluate_transport_satisfaction, 13 }, 14 sink::{DeliveryPayload, DeliveryTargetReceipt}, 15 target::{ 16 TARGET_SET_MAX_ITEMS, Target, TargetFingerprint, TargetLabel, TargetScope, TargetSet, 17 }, 18 }; 19 20 use crate::{Error, journal::OperationInstanceId}; 21 22 /// Maximum records claimed by one bounded outbox query. 23 pub const OUTBOX_CLAIM_LIMIT_MAX: u16 = 256; 24 /// Maximum UTF-8 bytes in a lease owner identity. 25 pub const LEASE_OWNER_MAX_BYTES: usize = 128; 26 27 /// Stable host-generated identity for one durable delivery plan. 28 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 29 #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] 30 pub struct OutboxItemId([u8; 16]); 31 32 impl OutboxItemId { 33 pub const fn new(bytes: [u8; 16]) -> Result<Self, Error> { 34 if bytes_are_zero(&bytes) { 35 return Err(Error::InvalidOutboxItemId); 36 } 37 Ok(Self(bytes)) 38 } 39 40 pub const fn as_bytes(&self) -> &[u8; 16] { 41 &self.0 42 } 43 } 44 45 /// Digest of the canonical delivery request, computed by its owning workflow. 46 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 47 #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] 48 pub struct DeliveryPlanDigest([u8; 32]); 49 50 impl DeliveryPlanDigest { 51 pub const fn new(bytes: [u8; 32]) -> Self { 52 Self(bytes) 53 } 54 55 pub const fn as_bytes(&self) -> &[u8; 32] { 56 &self.0 57 } 58 } 59 60 /// Non-zero optimistic outbox record revision. 61 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 62 #[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)] 63 pub struct OutboxRevision(u64); 64 65 impl OutboxRevision { 66 pub const INITIAL: Self = Self(1); 67 68 pub const fn new(value: u64) -> Result<Self, Error> { 69 if value == 0 { 70 return Err(Error::InvalidOutboxRevision); 71 } 72 Ok(Self(value)) 73 } 74 75 pub const fn get(self) -> u64 { 76 self.0 77 } 78 79 fn next(self) -> Result<Self, Error> { 80 self.0 81 .checked_add(1) 82 .map(Self) 83 .ok_or(Error::CorruptOutboxRecord) 84 } 85 } 86 87 /// Non-zero delivery attempt sequence. 88 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 89 #[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)] 90 pub struct DeliveryAttempt(u32); 91 92 impl DeliveryAttempt { 93 pub const FIRST: Self = Self(1); 94 95 pub const fn new(value: u32) -> Result<Self, Error> { 96 if value == 0 { 97 return Err(Error::InvalidDeliveryAttempt); 98 } 99 Ok(Self(value)) 100 } 101 102 pub const fn get(self) -> u32 { 103 self.0 104 } 105 106 fn next(self) -> Result<Self, Error> { 107 self.0 108 .checked_add(1) 109 .map(Self) 110 .ok_or(Error::CorruptOutboxRecord) 111 } 112 } 113 114 /// Opaque, caller-generated lease identity. 115 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 116 #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] 117 pub struct LeaseId([u8; 16]); 118 119 impl LeaseId { 120 pub const fn new(bytes: [u8; 16]) -> Result<Self, Error> { 121 if bytes_are_zero(&bytes) { 122 return Err(Error::InvalidOutboxLease); 123 } 124 Ok(Self(bytes)) 125 } 126 127 pub const fn as_bytes(&self) -> &[u8; 16] { 128 &self.0 129 } 130 } 131 132 /// Validated worker identity recorded with a lease. 133 #[derive(Clone, Eq, Hash, Ord, PartialEq, PartialOrd)] 134 pub struct LeaseOwner(String); 135 136 impl LeaseOwner { 137 pub fn parse(value: impl Into<String>) -> Result<Self, Error> { 138 let value = value.into(); 139 let invalid = [ 140 value.is_empty(), 141 value.len() > LEASE_OWNER_MAX_BYTES, 142 value != value.trim(), 143 value.chars().any(char::is_control), 144 ]; 145 if invalid.contains(&true) { 146 return Err(Error::InvalidOutboxLeaseOwner); 147 } 148 Ok(Self(value)) 149 } 150 151 pub fn as_str(&self) -> &str { 152 self.0.as_str() 153 } 154 } 155 156 impl fmt::Debug for LeaseOwner { 157 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 158 formatter 159 .debug_struct("LeaseOwner") 160 .field("value", &"[REDACTED]") 161 .field("bytes", &self.0.len()) 162 .finish() 163 } 164 } 165 166 #[cfg(feature = "serde")] 167 impl serde::Serialize for LeaseOwner { 168 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error> 169 where 170 S: serde::Serializer, 171 { 172 serializer.serialize_str(self.as_str()) 173 } 174 } 175 176 #[cfg(feature = "serde")] 177 impl<'de> serde::Deserialize<'de> for LeaseOwner { 178 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error> 179 where 180 D: serde::Deserializer<'de>, 181 { 182 let value = <String as serde::Deserialize>::deserialize(deserializer)?; 183 Self::parse(value).map_err(serde::de::Error::custom) 184 } 185 } 186 187 /// Exclusive, expiring authority to mutate one outbox item. 188 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 189 #[derive(Clone, Debug, Eq, PartialEq)] 190 pub struct OutboxLease { 191 id: LeaseId, 192 owner: LeaseOwner, 193 acquired_at_unix_ms: u64, 194 expires_at_unix_ms: u64, 195 } 196 197 impl OutboxLease { 198 pub fn new( 199 id: LeaseId, 200 owner: LeaseOwner, 201 acquired_at_unix_ms: u64, 202 expires_at_unix_ms: u64, 203 ) -> Result<Self, Error> { 204 if [ 205 acquired_at_unix_ms == 0, 206 expires_at_unix_ms <= acquired_at_unix_ms, 207 ] 208 .contains(&true) 209 { 210 return Err(Error::InvalidOutboxLease); 211 } 212 Ok(Self { 213 id, 214 owner, 215 acquired_at_unix_ms, 216 expires_at_unix_ms, 217 }) 218 } 219 220 pub const fn id(&self) -> LeaseId { 221 self.id 222 } 223 224 pub const fn owner(&self) -> &LeaseOwner { 225 &self.owner 226 } 227 228 pub const fn acquired_at_unix_ms(&self) -> u64 { 229 self.acquired_at_unix_ms 230 } 231 232 pub const fn expires_at_unix_ms(&self) -> u64 { 233 self.expires_at_unix_ms 234 } 235 236 pub const fn is_active_at(&self, unix_ms: u64) -> bool { 237 unix_ms >= self.acquired_at_unix_ms && unix_ms < self.expires_at_unix_ms 238 } 239 } 240 241 /// Durable lifecycle of one delivery plan. 242 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 243 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] 244 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 245 pub enum OutboxStage { 246 Pending, 247 Leased, 248 Retryable, 249 Satisfied, 250 Exhausted, 251 } 252 253 impl OutboxStage { 254 pub const fn is_terminal(self) -> bool { 255 matches!(self, Self::Satisfied | Self::Exhausted) 256 } 257 } 258 259 /// Latest durable evidence for one requested target. 260 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 261 #[derive(Clone, Debug, Eq, PartialEq)] 262 pub struct TargetDeliveryEvidence { 263 target: TargetFingerprint, 264 attempt: DeliveryAttempt, 265 attempted: bool, 266 outcome: DeliveryOutcome, 267 recorded_at_unix_ms: u64, 268 } 269 270 impl TargetDeliveryEvidence { 271 pub fn new( 272 target: TargetFingerprint, 273 attempt: DeliveryAttempt, 274 attempted: bool, 275 outcome: DeliveryOutcome, 276 recorded_at_unix_ms: u64, 277 ) -> Result<Self, Error> { 278 if recorded_at_unix_ms == 0 { 279 return Err(Error::InvalidDeliveryEvidence); 280 } 281 Ok(Self { 282 target, 283 attempt, 284 attempted, 285 outcome, 286 recorded_at_unix_ms, 287 }) 288 } 289 290 pub const fn target(&self) -> &TargetFingerprint { 291 &self.target 292 } 293 294 pub const fn attempt(&self) -> DeliveryAttempt { 295 self.attempt 296 } 297 298 pub const fn was_attempted(&self) -> bool { 299 self.attempted 300 } 301 302 pub const fn outcome(&self) -> &DeliveryOutcome { 303 &self.outcome 304 } 305 306 pub const fn recorded_at_unix_ms(&self) -> u64 { 307 self.recorded_at_unix_ms 308 } 309 } 310 311 /// Result of evaluating the latest complete transport receipt. 312 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 313 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] 314 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 315 pub enum SatisfactionResult { 316 Pending, 317 Satisfied, 318 Exhausted, 319 } 320 321 /// Durable outbox item and all current delivery evidence. 322 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 323 #[derive(Clone, Debug, Eq, PartialEq)] 324 pub struct OutboxRecord { 325 item_id: OutboxItemId, 326 operation_instance_id: OperationInstanceId, 327 plan_digest: DeliveryPlanDigest, 328 request: DeliveryRequest, 329 revision: OutboxRevision, 330 stage: OutboxStage, 331 lease: Option<OutboxLease>, 332 last_attempt: Option<DeliveryAttempt>, 333 evidence: Vec<TargetDeliveryEvidence>, 334 satisfaction: SatisfactionResult, 335 retry_not_before_unix_ms: Option<u64>, 336 created_at_unix_ms: u64, 337 updated_at_unix_ms: u64, 338 } 339 340 impl OutboxRecord { 341 fn from_enqueue(value: EnqueueOutboxItem) -> Self { 342 Self { 343 item_id: value.item_id, 344 operation_instance_id: value.operation_instance_id, 345 plan_digest: value.plan_digest, 346 request: value.request, 347 revision: OutboxRevision::INITIAL, 348 stage: OutboxStage::Pending, 349 lease: None, 350 last_attempt: None, 351 evidence: Vec::new(), 352 satisfaction: SatisfactionResult::Pending, 353 retry_not_before_unix_ms: None, 354 created_at_unix_ms: value.created_at_unix_ms, 355 updated_at_unix_ms: value.created_at_unix_ms, 356 } 357 } 358 359 /// Reconstructs and validates one record at a durable backend boundary. 360 #[allow(clippy::too_many_arguments)] 361 pub fn from_durable_parts( 362 enqueue: EnqueueOutboxItem, 363 revision: OutboxRevision, 364 stage: OutboxStage, 365 lease: Option<OutboxLease>, 366 last_attempt: Option<DeliveryAttempt>, 367 evidence: Vec<TargetDeliveryEvidence>, 368 satisfaction: SatisfactionResult, 369 retry_not_before_unix_ms: Option<u64>, 370 updated_at_unix_ms: u64, 371 ) -> Result<Self, Error> { 372 let invalid = [ 373 updated_at_unix_ms < enqueue.created_at_unix_ms, 374 matches!(stage, OutboxStage::Leased) != lease.is_some(), 375 matches!(stage, OutboxStage::Retryable) && last_attempt.is_none(), 376 matches!(stage, OutboxStage::Satisfied) 377 != matches!(satisfaction, SatisfactionResult::Satisfied), 378 matches!(stage, OutboxStage::Exhausted) 379 != matches!(satisfaction, SatisfactionResult::Exhausted), 380 matches!( 381 stage, 382 OutboxStage::Pending | OutboxStage::Leased | OutboxStage::Retryable 383 ) && !matches!(satisfaction, SatisfactionResult::Pending), 384 stage.is_terminal() && retry_not_before_unix_ms.is_some(), 385 ]; 386 if invalid.contains(&true) { 387 return Err(Error::CorruptOutboxRecord); 388 } 389 390 validate_evidence( 391 enqueue.request(), 392 last_attempt, 393 evidence.as_slice(), 394 satisfaction, 395 enqueue.created_at_unix_ms, 396 updated_at_unix_ms, 397 )?; 398 Ok(Self { 399 item_id: enqueue.item_id, 400 operation_instance_id: enqueue.operation_instance_id, 401 plan_digest: enqueue.plan_digest, 402 request: enqueue.request, 403 revision, 404 stage, 405 lease, 406 last_attempt, 407 evidence, 408 satisfaction, 409 retry_not_before_unix_ms, 410 created_at_unix_ms: enqueue.created_at_unix_ms, 411 updated_at_unix_ms, 412 }) 413 } 414 415 pub const fn item_id(&self) -> OutboxItemId { 416 self.item_id 417 } 418 pub const fn operation_instance_id(&self) -> OperationInstanceId { 419 self.operation_instance_id 420 } 421 pub const fn plan_digest(&self) -> DeliveryPlanDigest { 422 self.plan_digest 423 } 424 pub const fn request(&self) -> &DeliveryRequest { 425 &self.request 426 } 427 pub const fn revision(&self) -> OutboxRevision { 428 self.revision 429 } 430 pub const fn stage(&self) -> OutboxStage { 431 self.stage 432 } 433 pub const fn lease(&self) -> Option<&OutboxLease> { 434 self.lease.as_ref() 435 } 436 pub const fn last_attempt(&self) -> Option<DeliveryAttempt> { 437 self.last_attempt 438 } 439 pub fn evidence(&self) -> &[TargetDeliveryEvidence] { 440 self.evidence.as_slice() 441 } 442 /// Returns the latest evidence for one target without discarding history. 443 pub fn latest_target_evidence( 444 &self, 445 target: &TargetFingerprint, 446 ) -> Option<&TargetDeliveryEvidence> { 447 self.evidence 448 .iter() 449 .rev() 450 .find(|evidence| evidence.target() == target) 451 } 452 pub const fn satisfaction(&self) -> SatisfactionResult { 453 self.satisfaction 454 } 455 pub const fn retry_not_before_unix_ms(&self) -> Option<u64> { 456 self.retry_not_before_unix_ms 457 } 458 pub const fn created_at_unix_ms(&self) -> u64 { 459 self.created_at_unix_ms 460 } 461 pub const fn updated_at_unix_ms(&self) -> u64 { 462 self.updated_at_unix_ms 463 } 464 465 /// Claims this item if it is ready and has no active lease. 466 pub fn claim(&mut self, lease: OutboxLease) -> Result<(), Error> { 467 if self.stage.is_terminal() { 468 return Err(Error::OutboxItemTerminal); 469 } 470 if lease.acquired_at_unix_ms() < self.updated_at_unix_ms { 471 return Err(Error::InvalidOutboxTimestamp); 472 } 473 if matches!(self.retry_not_before_unix_ms, Some(not_before) if lease.acquired_at_unix_ms() < not_before) 474 { 475 return Err(Error::OutboxItemNotReady); 476 } 477 if self 478 .lease 479 .as_ref() 480 .is_some_and(|current| current.is_active_at(lease.acquired_at_unix_ms())) 481 { 482 return Err(Error::OutboxLeaseConflict); 483 } 484 self.revision = self.revision.next()?; 485 self.updated_at_unix_ms = lease.acquired_at_unix_ms(); 486 self.stage = OutboxStage::Leased; 487 self.lease = Some(lease); 488 Ok(()) 489 } 490 491 /// Applies one request-bound transport receipt under the active lease. 492 pub fn record_attempt(&mut self, value: DeliveryAttemptEvidence) -> Result<(), Error> { 493 self.validate_lease(value.lease_id, value.recorded_at_unix_ms)?; 494 if value.item_id != self.item_id || value.expected_revision != self.revision { 495 return Err(Error::OutboxRevisionConflict); 496 } 497 let expected_attempt = self 498 .last_attempt 499 .map_or(Ok(DeliveryAttempt::FIRST), DeliveryAttempt::next)?; 500 if value.attempt != expected_attempt { 501 return Err(Error::InvalidDeliveryAttempt); 502 } 503 value 504 .receipt 505 .validate_for_request(&self.request) 506 .map_err(|_| Error::InvalidDeliveryEvidence)?; 507 self.evidence 508 .extend( 509 value 510 .receipt 511 .target_receipts() 512 .iter() 513 .map(|receipt| TargetDeliveryEvidence { 514 target: receipt.target().fingerprint().clone(), 515 attempt: value.attempt, 516 attempted: receipt.was_attempted(), 517 outcome: receipt.outcome().clone(), 518 recorded_at_unix_ms: value.recorded_at_unix_ms, 519 }), 520 ); 521 self.last_attempt = Some(value.attempt); 522 self.satisfaction = evaluate_satisfaction(&self.request, &self.evidence); 523 self.stage = match self.satisfaction { 524 SatisfactionResult::Pending => OutboxStage::Retryable, 525 SatisfactionResult::Satisfied => OutboxStage::Satisfied, 526 SatisfactionResult::Exhausted => OutboxStage::Exhausted, 527 }; 528 self.lease = None; 529 self.retry_not_before_unix_ms = None; 530 self.updated_at_unix_ms = value.recorded_at_unix_ms; 531 self.revision = self.revision.next()?; 532 Ok(()) 533 } 534 535 /// Releases an active lease and optionally defers the next claim. 536 pub fn release( 537 &mut self, 538 lease_id: LeaseId, 539 expected_revision: OutboxRevision, 540 released_at_unix_ms: u64, 541 retry_not_before_unix_ms: Option<u64>, 542 ) -> Result<(), Error> { 543 self.validate_lease(lease_id, released_at_unix_ms)?; 544 if expected_revision != self.revision { 545 return Err(Error::OutboxRevisionConflict); 546 } 547 // `validate_lease` already proves this timestamp is within a lease 548 // whose acquisition timestamp is non-zero. 549 if matches!(retry_not_before_unix_ms, Some(value) if value <= released_at_unix_ms) { 550 return Err(Error::InvalidOutboxTimestamp); 551 } 552 self.lease = None; 553 self.stage = if self.last_attempt.is_some() { 554 OutboxStage::Retryable 555 } else { 556 OutboxStage::Pending 557 }; 558 self.retry_not_before_unix_ms = retry_not_before_unix_ms; 559 self.updated_at_unix_ms = released_at_unix_ms; 560 self.revision = self.revision.next()?; 561 Ok(()) 562 } 563 564 fn validate_lease(&self, lease_id: LeaseId, at_unix_ms: u64) -> Result<(), Error> { 565 let lease = self.lease.as_ref().ok_or(Error::OutboxLeaseConflict)?; 566 if lease.id() != lease_id { 567 return Err(Error::OutboxLeaseConflict); 568 } 569 if !lease.is_active_at(at_unix_ms) { 570 return Err(Error::OutboxLeaseExpired); 571 } 572 Ok(()) 573 } 574 } 575 576 /// Validated durable plan enqueue request. 577 #[derive(Clone, Debug, Eq, PartialEq)] 578 pub struct EnqueueOutboxItem { 579 item_id: OutboxItemId, 580 operation_instance_id: OperationInstanceId, 581 plan_digest: DeliveryPlanDigest, 582 request: DeliveryRequest, 583 created_at_unix_ms: u64, 584 } 585 586 impl EnqueueOutboxItem { 587 pub fn new( 588 item_id: OutboxItemId, 589 operation_instance_id: OperationInstanceId, 590 plan_digest: DeliveryPlanDigest, 591 request: DeliveryRequest, 592 created_at_unix_ms: u64, 593 ) -> Result<Self, Error> { 594 if created_at_unix_ms == 0 { 595 return Err(Error::InvalidOutboxTimestamp); 596 } 597 Ok(Self { 598 item_id, 599 operation_instance_id, 600 plan_digest, 601 request, 602 created_at_unix_ms, 603 }) 604 } 605 606 pub const fn item_id(&self) -> OutboxItemId { 607 self.item_id 608 } 609 pub const fn operation_instance_id(&self) -> OperationInstanceId { 610 self.operation_instance_id 611 } 612 pub const fn plan_digest(&self) -> DeliveryPlanDigest { 613 self.plan_digest 614 } 615 pub const fn request(&self) -> &DeliveryRequest { 616 &self.request 617 } 618 pub const fn created_at_unix_ms(&self) -> u64 { 619 self.created_at_unix_ms 620 } 621 pub fn into_record(self) -> OutboxRecord { 622 OutboxRecord::from_enqueue(self) 623 } 624 } 625 626 /// Idempotent enqueue result. 627 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 628 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] 629 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 630 pub enum EnqueueDisposition { 631 Created, 632 Replay, 633 } 634 635 #[derive(Clone, Debug, Eq, PartialEq)] 636 pub struct EnqueueReceipt { 637 disposition: EnqueueDisposition, 638 record: OutboxRecord, 639 } 640 641 impl EnqueueReceipt { 642 pub const fn new(disposition: EnqueueDisposition, record: OutboxRecord) -> Self { 643 Self { 644 disposition, 645 record, 646 } 647 } 648 pub const fn disposition(&self) -> EnqueueDisposition { 649 self.disposition 650 } 651 pub const fn record(&self) -> &OutboxRecord { 652 &self.record 653 } 654 } 655 656 /// Bounded lease acquisition request. 657 #[derive(Clone, Debug, Eq, PartialEq)] 658 pub struct ClaimOutboxItems { 659 owner: LeaseOwner, 660 lease_id_seed: LeaseId, 661 now_unix_ms: u64, 662 lease_expires_at_unix_ms: u64, 663 limit: u16, 664 } 665 666 impl ClaimOutboxItems { 667 pub fn new( 668 owner: LeaseOwner, 669 lease_id_seed: LeaseId, 670 now_unix_ms: u64, 671 lease_expires_at_unix_ms: u64, 672 limit: u16, 673 ) -> Result<Self, Error> { 674 if [now_unix_ms == 0, lease_expires_at_unix_ms <= now_unix_ms].contains(&true) { 675 return Err(Error::InvalidOutboxLease); 676 } 677 if [limit == 0, limit > OUTBOX_CLAIM_LIMIT_MAX].contains(&true) { 678 return Err(Error::InvalidOutboxClaimLimit); 679 } 680 Ok(Self { 681 owner, 682 lease_id_seed, 683 now_unix_ms, 684 lease_expires_at_unix_ms, 685 limit, 686 }) 687 } 688 pub const fn owner(&self) -> &LeaseOwner { 689 &self.owner 690 } 691 pub const fn lease_id_seed(&self) -> LeaseId { 692 self.lease_id_seed 693 } 694 /// Derives a stable, item-specific token from the caller's unique seed. 695 pub fn lease_id_for(&self, item_id: OutboxItemId) -> LeaseId { 696 let mut bytes = *self.lease_id_seed.as_bytes(); 697 for (byte, item_byte) in bytes.iter_mut().zip(item_id.as_bytes()) { 698 *byte ^= item_byte; 699 } 700 if bytes_are_zero(&bytes) { 701 bytes[0] = 1; 702 } 703 LeaseId(bytes) 704 } 705 pub const fn now_unix_ms(&self) -> u64 { 706 self.now_unix_ms 707 } 708 pub const fn lease_expires_at_unix_ms(&self) -> u64 { 709 self.lease_expires_at_unix_ms 710 } 711 pub const fn limit(&self) -> u16 { 712 self.limit 713 } 714 } 715 716 /// Claimed outbox item with exact lease authority. 717 #[derive(Clone, Debug, Eq, PartialEq)] 718 pub struct ClaimedOutboxItem { 719 record: OutboxRecord, 720 lease: OutboxLease, 721 } 722 723 impl ClaimedOutboxItem { 724 pub const fn new(record: OutboxRecord, lease: OutboxLease) -> Self { 725 Self { record, lease } 726 } 727 pub const fn record(&self) -> &OutboxRecord { 728 &self.record 729 } 730 pub const fn lease(&self) -> &OutboxLease { 731 &self.lease 732 } 733 } 734 735 /// Request-bound evidence for one complete adapter attempt. 736 #[derive(Clone, Debug, Eq, PartialEq)] 737 pub struct DeliveryAttemptEvidence { 738 item_id: OutboxItemId, 739 lease_id: LeaseId, 740 expected_revision: OutboxRevision, 741 attempt: DeliveryAttempt, 742 receipt: DeliveryReceipt, 743 recorded_at_unix_ms: u64, 744 } 745 746 impl DeliveryAttemptEvidence { 747 pub fn new( 748 item_id: OutboxItemId, 749 lease_id: LeaseId, 750 expected_revision: OutboxRevision, 751 attempt: DeliveryAttempt, 752 receipt: DeliveryReceipt, 753 recorded_at_unix_ms: u64, 754 ) -> Result<Self, Error> { 755 if recorded_at_unix_ms == 0 { 756 return Err(Error::InvalidOutboxTimestamp); 757 } 758 Ok(Self { 759 item_id, 760 lease_id, 761 expected_revision, 762 attempt, 763 receipt, 764 recorded_at_unix_ms, 765 }) 766 } 767 pub const fn item_id(&self) -> OutboxItemId { 768 self.item_id 769 } 770 pub const fn lease_id(&self) -> LeaseId { 771 self.lease_id 772 } 773 pub const fn expected_revision(&self) -> OutboxRevision { 774 self.expected_revision 775 } 776 pub const fn attempt(&self) -> DeliveryAttempt { 777 self.attempt 778 } 779 pub const fn receipt(&self) -> &DeliveryReceipt { 780 &self.receipt 781 } 782 pub const fn recorded_at_unix_ms(&self) -> u64 { 783 self.recorded_at_unix_ms 784 } 785 } 786 787 /// Passive outbox state summary. 788 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 789 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 790 pub struct OutboxStatus { 791 pub pending: u64, 792 pub leased: u64, 793 pub retryable: u64, 794 pub satisfied: u64, 795 pub exhausted: u64, 796 } 797 798 impl OutboxStatus { 799 pub fn total(self) -> Option<u64> { 800 self.pending 801 .checked_add(self.leased)? 802 .checked_add(self.retryable)? 803 .checked_add(self.satisfied)? 804 .checked_add(self.exhausted) 805 } 806 } 807 808 /// Backend-neutral durable delivery-plan SPI. 809 pub trait Outbox: Send + Sync { 810 fn enqueue(&self, item: EnqueueOutboxItem) -> BoxFuture<'_, Result<EnqueueReceipt, Error>>; 811 fn item(&self, item_id: OutboxItemId) -> BoxFuture<'_, Result<Option<OutboxRecord>, Error>>; 812 fn claim( 813 &self, 814 request: ClaimOutboxItems, 815 ) -> BoxFuture<'_, Result<Vec<ClaimedOutboxItem>, Error>>; 816 fn record_attempt( 817 &self, 818 evidence: DeliveryAttemptEvidence, 819 ) -> BoxFuture<'_, Result<OutboxRecord, Error>>; 820 fn release( 821 &self, 822 item_id: OutboxItemId, 823 lease_id: LeaseId, 824 expected_revision: OutboxRevision, 825 released_at_unix_ms: u64, 826 retry_not_before_unix_ms: Option<u64>, 827 ) -> BoxFuture<'_, Result<OutboxRecord, Error>>; 828 fn status(&self) -> BoxFuture<'_, Result<OutboxStatus, Error>>; 829 } 830 831 const fn bytes_are_zero(bytes: &[u8; 16]) -> bool { 832 let mut index = 0; 833 while index < bytes.len() { 834 if bytes[index] != 0 { 835 return false; 836 } 837 index += 1; 838 } 839 true 840 } 841 842 fn validate_evidence( 843 request: &DeliveryRequest, 844 last_attempt: Option<DeliveryAttempt>, 845 evidence: &[TargetDeliveryEvidence], 846 satisfaction: SatisfactionResult, 847 created_at_unix_ms: u64, 848 updated_at_unix_ms: u64, 849 ) -> Result<(), Error> { 850 let Some(last_attempt) = last_attempt else { 851 return if evidence.is_empty() & matches!(satisfaction, SatisfactionResult::Pending) { 852 Ok(()) 853 } else { 854 Err(Error::CorruptOutboxRecord) 855 }; 856 }; 857 let target_count = request.target_set().len(); 858 if evidence.len() != target_count.saturating_mul(last_attempt.get() as usize) { 859 return Err(Error::CorruptOutboxRecord); 860 } 861 if evidence.iter().any(|entry| { 862 [ 863 entry.recorded_at_unix_ms() < created_at_unix_ms, 864 entry.recorded_at_unix_ms() > updated_at_unix_ms, 865 !request 866 .target_set() 867 .targets() 868 .iter() 869 .any(|target| target.fingerprint() == entry.target()), 870 ] 871 .contains(&true) 872 }) { 873 return Err(Error::CorruptOutboxRecord); 874 } 875 876 let mut previous_recorded_at = created_at_unix_ms; 877 for attempt in 1..=last_attempt.get() { 878 let mut recorded_at = None; 879 let receipts = request 880 .target_set() 881 .targets() 882 .iter() 883 .map(|target| { 884 let mut matches = evidence.iter().filter(|entry| { 885 entry.attempt().get() == attempt && entry.target() == target.fingerprint() 886 }); 887 let entry = matches.next().ok_or(Error::CorruptOutboxRecord)?; 888 if matches.next().is_some() { 889 return Err(Error::CorruptOutboxRecord); 890 } 891 if recorded_at 892 .replace(entry.recorded_at_unix_ms()) 893 .is_some_and(|prior| prior != entry.recorded_at_unix_ms()) 894 { 895 return Err(Error::CorruptOutboxRecord); 896 } 897 if entry.was_attempted() { 898 Ok(DeliveryTargetReceipt::attempted( 899 target.clone(), 900 entry.outcome().clone(), 901 )) 902 } else { 903 DeliveryTargetReceipt::skipped(target.clone(), entry.outcome().clone()) 904 .map_err(|_| Error::CorruptOutboxRecord) 905 } 906 }) 907 .collect::<Result<Vec<_>, _>>()?; 908 let recorded_at = recorded_at.ok_or(Error::CorruptOutboxRecord)?; 909 if recorded_at < previous_recorded_at { 910 return Err(Error::CorruptOutboxRecord); 911 } 912 previous_recorded_at = recorded_at; 913 DeliveryReceipt::for_request(request, receipts).map_err(|_| Error::CorruptOutboxRecord)?; 914 } 915 let expected = evaluate_satisfaction(request, evidence); 916 if expected == satisfaction { 917 Ok(()) 918 } else { 919 Err(Error::CorruptOutboxRecord) 920 } 921 } 922 923 fn evaluate_satisfaction( 924 request: &DeliveryRequest, 925 evidence: &[TargetDeliveryEvidence], 926 ) -> SatisfactionResult { 927 match evaluate_transport_satisfaction( 928 request.satisfaction(), 929 request.target_set(), 930 evidence 931 .iter() 932 .map(|entry| (entry.target(), entry.outcome())), 933 ) { 934 Ok(SatisfactionState::Satisfied) => SatisfactionResult::Satisfied, 935 Ok(SatisfactionState::Pending) => SatisfactionResult::Pending, 936 Ok(SatisfactionState::Exhausted) | Err(_) => SatisfactionResult::Exhausted, 937 } 938 } 939 940 #[cfg(test)] 941 mod tests { 942 use super::*; 943 use radroots_event::{SignedEvent, wire::Nip01EventWire}; 944 use radroots_transport::sink::DeliveryPayload; 945 946 fn signed_event() -> SignedEvent { 947 let mut wire = Nip01EventWire { 948 id: "0".repeat(64), 949 pubkey: "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df".to_owned(), 950 created_at: 1_800_000_100, 951 kind: 0, 952 tags: vec![], 953 content: "{}".to_owned(), 954 sig: "42".repeat(64), 955 extra: Default::default(), 956 }; 957 wire.id = wire.computed_event_id().unwrap().to_hex(); 958 let raw = serde_json::json!({ 959 "id": wire.id, 960 "pubkey": wire.pubkey, 961 "created_at": wire.created_at, 962 "kind": wire.kind, 963 "tags": wire.tags, 964 "content": wire.content, 965 "sig": wire.sig, 966 }) 967 .to_string(); 968 SignedEvent::from_wire_verified_id(wire, raw).unwrap() 969 } 970 971 fn request() -> DeliveryRequest { 972 DeliveryRequest::new( 973 "storage-outbox-unit", 974 DeliveryPayload::new(signed_event()), 975 TargetSet::new(vec![Target::nostr_relay("wss://relay.example").unwrap()]).unwrap(), 976 SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::all()), 977 1_000, 978 ) 979 .unwrap() 980 } 981 982 fn enqueue() -> EnqueueOutboxItem { 983 EnqueueOutboxItem::new( 984 OutboxItemId::new([1; 16]).unwrap(), 985 OperationInstanceId::new([2; 16]).unwrap(), 986 DeliveryPlanDigest::new([3; 32]), 987 request(), 988 10, 989 ) 990 .unwrap() 991 } 992 993 fn lease(id: u8, acquired: u64, expires: u64) -> OutboxLease { 994 OutboxLease::new( 995 LeaseId::new([id; 16]).unwrap(), 996 LeaseOwner::parse("worker").unwrap(), 997 acquired, 998 expires, 999 ) 1000 .unwrap() 1001 } 1002 1003 fn receipt_for(request: &DeliveryRequest, outcome: DeliveryOutcome) -> DeliveryReceipt { 1004 DeliveryReceipt::for_request( 1005 request, 1006 vec![DeliveryTargetReceipt::attempted( 1007 request.target_set().targets()[0].clone(), 1008 outcome, 1009 )], 1010 ) 1011 .unwrap() 1012 } 1013 1014 #[test] 1015 fn scalar_types_and_lease_policy_cover_all_bounds() { 1016 assert_eq!(OutboxItemId::new([0; 16]), Err(Error::InvalidOutboxItemId)); 1017 assert_eq!(LeaseId::new([0; 16]), Err(Error::InvalidOutboxLease)); 1018 assert_eq!(OutboxRevision::new(0), Err(Error::InvalidOutboxRevision)); 1019 assert_eq!(DeliveryAttempt::new(0), Err(Error::InvalidDeliveryAttempt)); 1020 let item = OutboxItemId::new([1; 16]).unwrap(); 1021 let digest = DeliveryPlanDigest::new([2; 32]); 1022 assert_eq!(item.as_bytes(), &[1; 16]); 1023 assert_eq!(digest.as_bytes(), &[2; 32]); 1024 assert_eq!(OutboxRevision::INITIAL.get(), 1); 1025 assert_eq!(DeliveryAttempt::FIRST.get(), 1); 1026 assert_eq!( 1027 OutboxRevision(u64::MAX).next(), 1028 Err(Error::CorruptOutboxRecord) 1029 ); 1030 assert_eq!( 1031 DeliveryAttempt(u32::MAX).next(), 1032 Err(Error::CorruptOutboxRecord) 1033 ); 1034 1035 for invalid in ["", " worker", "worker ", "bad\nworker"] { 1036 assert_eq!( 1037 LeaseOwner::parse(invalid), 1038 Err(Error::InvalidOutboxLeaseOwner) 1039 ); 1040 } 1041 assert_eq!( 1042 LeaseOwner::parse("x".repeat(LEASE_OWNER_MAX_BYTES + 1)), 1043 Err(Error::InvalidOutboxLeaseOwner) 1044 ); 1045 let owner = LeaseOwner::parse("worker").unwrap(); 1046 assert_eq!(owner.as_str(), "worker"); 1047 let debug = format!("{owner:?}"); 1048 assert!(debug.contains("[REDACTED]")); 1049 assert!(!debug.contains("worker")); 1050 1051 let id = LeaseId::new([4; 16]).unwrap(); 1052 assert_eq!( 1053 OutboxLease::new(id, owner.clone(), 0, 2), 1054 Err(Error::InvalidOutboxLease) 1055 ); 1056 assert_eq!( 1057 OutboxLease::new(id, owner.clone(), 2, 2), 1058 Err(Error::InvalidOutboxLease) 1059 ); 1060 assert_eq!( 1061 OutboxLease::new(id, owner, 2, 1), 1062 Err(Error::InvalidOutboxLease) 1063 ); 1064 let value = lease(4, 2, 4); 1065 assert_eq!(value.id().as_bytes(), &[4; 16]); 1066 assert_eq!(value.owner().as_str(), "worker"); 1067 assert_eq!(value.acquired_at_unix_ms(), 2); 1068 assert_eq!(value.expires_at_unix_ms(), 4); 1069 assert!(!value.is_active_at(1)); 1070 assert!(value.is_active_at(2)); 1071 assert!(value.is_active_at(3)); 1072 assert!(!value.is_active_at(4)); 1073 assert!(!OutboxStage::Pending.is_terminal()); 1074 assert!(!OutboxStage::Leased.is_terminal()); 1075 assert!(!OutboxStage::Retryable.is_terminal()); 1076 assert!(OutboxStage::Satisfied.is_terminal()); 1077 assert!(OutboxStage::Exhausted.is_terminal()); 1078 } 1079 1080 #[test] 1081 fn enqueue_claim_and_evidence_models_cover_accessors_and_bounds() { 1082 assert_eq!( 1083 EnqueueOutboxItem::new( 1084 OutboxItemId::new([1; 16]).unwrap(), 1085 OperationInstanceId::new([2; 16]).unwrap(), 1086 DeliveryPlanDigest::new([3; 32]), 1087 request(), 1088 0, 1089 ), 1090 Err(Error::InvalidOutboxTimestamp) 1091 ); 1092 let value = enqueue(); 1093 assert_eq!(value.item_id().as_bytes(), &[1; 16]); 1094 assert_eq!(value.operation_instance_id().as_bytes(), &[2; 16]); 1095 assert_eq!(value.plan_digest().as_bytes(), &[3; 32]); 1096 assert_eq!(value.request().request_id().as_str(), "storage-outbox-unit"); 1097 assert_eq!(value.created_at_unix_ms(), 10); 1098 let record = value.into_record(); 1099 let receipt = EnqueueReceipt::new(EnqueueDisposition::Created, record.clone()); 1100 assert_eq!(receipt.disposition(), EnqueueDisposition::Created); 1101 assert_eq!(receipt.record(), &record); 1102 1103 for (now, expiry, limit, error) in [ 1104 (0, 2, 1, Error::InvalidOutboxLease), 1105 (2, 2, 1, Error::InvalidOutboxLease), 1106 (2, 3, 0, Error::InvalidOutboxClaimLimit), 1107 ( 1108 2, 1109 3, 1110 OUTBOX_CLAIM_LIMIT_MAX + 1, 1111 Error::InvalidOutboxClaimLimit, 1112 ), 1113 ] { 1114 assert_eq!( 1115 ClaimOutboxItems::new( 1116 LeaseOwner::parse("worker").unwrap(), 1117 LeaseId::new([1; 16]).unwrap(), 1118 now, 1119 expiry, 1120 limit, 1121 ), 1122 Err(error) 1123 ); 1124 } 1125 let claim = ClaimOutboxItems::new( 1126 LeaseOwner::parse("worker").unwrap(), 1127 LeaseId::new([1; 16]).unwrap(), 1128 2, 1129 3, 1130 1, 1131 ) 1132 .unwrap(); 1133 assert_eq!(claim.owner().as_str(), "worker"); 1134 assert_eq!(claim.lease_id_seed().as_bytes(), &[1; 16]); 1135 assert_eq!( 1136 claim 1137 .lease_id_for(OutboxItemId::new([1; 16]).unwrap()) 1138 .as_bytes()[0], 1139 1 1140 ); 1141 assert_eq!(claim.now_unix_ms(), 2); 1142 assert_eq!(claim.lease_expires_at_unix_ms(), 3); 1143 assert_eq!(claim.limit(), 1); 1144 1145 let claimed = ClaimedOutboxItem::new(record.clone(), lease(5, 10, 20)); 1146 assert_eq!(claimed.record(), &record); 1147 assert_eq!(claimed.lease().id().as_bytes(), &[5; 16]); 1148 assert_eq!( 1149 TargetDeliveryEvidence::new( 1150 record.request().target_set().targets()[0] 1151 .fingerprint() 1152 .clone(), 1153 DeliveryAttempt::FIRST, 1154 true, 1155 DeliveryOutcome::accepted(), 1156 0, 1157 ), 1158 Err(Error::InvalidDeliveryEvidence) 1159 ); 1160 let target_evidence = TargetDeliveryEvidence::new( 1161 record.request().target_set().targets()[0] 1162 .fingerprint() 1163 .clone(), 1164 DeliveryAttempt::FIRST, 1165 true, 1166 DeliveryOutcome::accepted(), 1167 12, 1168 ) 1169 .unwrap(); 1170 assert_eq!(target_evidence.attempt(), DeliveryAttempt::FIRST); 1171 assert!(target_evidence.was_attempted()); 1172 assert_eq!(target_evidence.recorded_at_unix_ms(), 12); 1173 assert!( 1174 target_evidence 1175 .outcome() 1176 .satisfies(SatisfactionClass::Accepted) 1177 ); 1178 } 1179 1180 #[test] 1181 fn durable_record_and_claim_reject_inconsistent_state() { 1182 let base = enqueue(); 1183 let valid = OutboxRecord::from_durable_parts( 1184 base.clone(), 1185 OutboxRevision::INITIAL, 1186 OutboxStage::Pending, 1187 None, 1188 None, 1189 vec![], 1190 SatisfactionResult::Pending, 1191 None, 1192 10, 1193 ) 1194 .unwrap(); 1195 assert_eq!(valid.item_id().as_bytes(), &[1; 16]); 1196 assert_eq!(valid.operation_instance_id().as_bytes(), &[2; 16]); 1197 assert_eq!(valid.plan_digest().as_bytes(), &[3; 32]); 1198 assert_eq!(valid.revision(), OutboxRevision::INITIAL); 1199 assert_eq!(valid.stage(), OutboxStage::Pending); 1200 assert!(valid.lease().is_none()); 1201 assert!(valid.last_attempt().is_none()); 1202 assert!(valid.evidence().is_empty()); 1203 assert!( 1204 valid 1205 .latest_target_evidence(valid.request().target_set().targets()[0].fingerprint()) 1206 .is_none() 1207 ); 1208 assert_eq!(valid.satisfaction(), SatisfactionResult::Pending); 1209 assert_eq!(valid.retry_not_before_unix_ms(), None); 1210 assert_eq!(valid.created_at_unix_ms(), 10); 1211 assert_eq!(valid.updated_at_unix_ms(), 10); 1212 1213 let cases = [ 1214 ( 1215 OutboxStage::Pending, 1216 None, 1217 None, 1218 SatisfactionResult::Pending, 1219 None, 1220 9, 1221 ), 1222 ( 1223 OutboxStage::Leased, 1224 None, 1225 None, 1226 SatisfactionResult::Pending, 1227 None, 1228 10, 1229 ), 1230 ( 1231 OutboxStage::Pending, 1232 Some(lease(1, 10, 20)), 1233 None, 1234 SatisfactionResult::Pending, 1235 None, 1236 10, 1237 ), 1238 ( 1239 OutboxStage::Retryable, 1240 None, 1241 None, 1242 SatisfactionResult::Pending, 1243 None, 1244 10, 1245 ), 1246 ( 1247 OutboxStage::Satisfied, 1248 None, 1249 None, 1250 SatisfactionResult::Pending, 1251 None, 1252 10, 1253 ), 1254 ( 1255 OutboxStage::Exhausted, 1256 None, 1257 None, 1258 SatisfactionResult::Pending, 1259 None, 1260 10, 1261 ), 1262 ( 1263 OutboxStage::Pending, 1264 None, 1265 None, 1266 SatisfactionResult::Satisfied, 1267 None, 1268 10, 1269 ), 1270 ( 1271 OutboxStage::Satisfied, 1272 None, 1273 None, 1274 SatisfactionResult::Satisfied, 1275 Some(20), 1276 10, 1277 ), 1278 ]; 1279 for (stage, lease, last_attempt, satisfaction, retry, updated) in cases { 1280 assert_eq!( 1281 OutboxRecord::from_durable_parts( 1282 base.clone(), 1283 OutboxRevision::INITIAL, 1284 stage, 1285 lease, 1286 last_attempt, 1287 vec![], 1288 satisfaction, 1289 retry, 1290 updated, 1291 ), 1292 Err(Error::CorruptOutboxRecord) 1293 ); 1294 } 1295 1296 let mut record = valid; 1297 assert_eq!( 1298 record.release(LeaseId::new([1; 16]).unwrap(), record.revision(), 11, None), 1299 Err(Error::OutboxLeaseConflict) 1300 ); 1301 record.claim(lease(1, 10, 20)).unwrap(); 1302 assert_eq!( 1303 record.claim(lease(2, 9, 20)), 1304 Err(Error::InvalidOutboxTimestamp) 1305 ); 1306 assert_eq!( 1307 record.claim(lease(2, 11, 20)), 1308 Err(Error::OutboxLeaseConflict) 1309 ); 1310 assert_eq!( 1311 record.release(LeaseId::new([2; 16]).unwrap(), record.revision(), 12, None), 1312 Err(Error::OutboxLeaseConflict) 1313 ); 1314 assert_eq!( 1315 record.release( 1316 LeaseId::new([1; 16]).unwrap(), 1317 OutboxRevision::INITIAL, 1318 12, 1319 None 1320 ), 1321 Err(Error::OutboxRevisionConflict) 1322 ); 1323 let revision = record.revision(); 1324 assert_eq!( 1325 record.release(LeaseId::new([1; 16]).unwrap(), revision, 12, Some(12)), 1326 Err(Error::InvalidOutboxTimestamp) 1327 ); 1328 record 1329 .release(LeaseId::new([1; 16]).unwrap(), revision, 12, Some(13)) 1330 .unwrap(); 1331 assert_eq!(record.stage(), OutboxStage::Pending); 1332 assert_eq!( 1333 record.claim(lease(3, 12, 20)), 1334 Err(Error::OutboxItemNotReady) 1335 ); 1336 record.claim(lease(3, 13, 20)).unwrap(); 1337 } 1338 1339 #[test] 1340 fn outbox_status_detects_each_overflow_position() { 1341 assert_eq!( 1342 OutboxStatus { 1343 pending: 1, 1344 leased: 2, 1345 retryable: 3, 1346 satisfied: 4, 1347 exhausted: 5 1348 } 1349 .total(), 1350 Some(15) 1351 ); 1352 for status in [ 1353 OutboxStatus { 1354 pending: u64::MAX, 1355 leased: 1, 1356 retryable: 0, 1357 satisfied: 0, 1358 exhausted: 0, 1359 }, 1360 OutboxStatus { 1361 pending: 0, 1362 leased: u64::MAX, 1363 retryable: 1, 1364 satisfied: 0, 1365 exhausted: 0, 1366 }, 1367 OutboxStatus { 1368 pending: 0, 1369 leased: 0, 1370 retryable: u64::MAX, 1371 satisfied: 1, 1372 exhausted: 0, 1373 }, 1374 OutboxStatus { 1375 pending: 0, 1376 leased: 0, 1377 retryable: 0, 1378 satisfied: u64::MAX, 1379 exhausted: 1, 1380 }, 1381 ] { 1382 assert_eq!(status.total(), None); 1383 } 1384 } 1385 1386 #[test] 1387 fn evidence_reconstruction_and_attempt_errors_are_fail_closed() { 1388 let enqueue = enqueue(); 1389 let target = enqueue.request().target_set().targets()[0] 1390 .fingerprint() 1391 .clone(); 1392 let accepted = TargetDeliveryEvidence::new( 1393 target.clone(), 1394 DeliveryAttempt::FIRST, 1395 true, 1396 DeliveryOutcome::accepted(), 1397 20, 1398 ) 1399 .unwrap(); 1400 assert_eq!( 1401 OutboxRecord::from_durable_parts( 1402 enqueue.clone(), 1403 OutboxRevision::new(2).unwrap(), 1404 OutboxStage::Satisfied, 1405 None, 1406 Some(DeliveryAttempt::FIRST), 1407 vec![accepted.clone()], 1408 SatisfactionResult::Satisfied, 1409 None, 1410 20, 1411 ) 1412 .unwrap() 1413 .latest_target_evidence(&target), 1414 Some(&accepted) 1415 ); 1416 let terminal = OutboxRecord::from_durable_parts( 1417 enqueue.clone(), 1418 OutboxRevision::new(2).unwrap(), 1419 OutboxStage::Satisfied, 1420 None, 1421 Some(DeliveryAttempt::FIRST), 1422 vec![accepted.clone()], 1423 SatisfactionResult::Satisfied, 1424 None, 1425 20, 1426 ) 1427 .unwrap(); 1428 let mut terminal = terminal; 1429 assert_eq!( 1430 terminal.claim(lease(2, 21, 30)), 1431 Err(Error::OutboxItemTerminal) 1432 ); 1433 1434 let retryable_evidence = TargetDeliveryEvidence::new( 1435 target.clone(), 1436 DeliveryAttempt::FIRST, 1437 true, 1438 DeliveryOutcome::unavailable(), 1439 20, 1440 ) 1441 .unwrap(); 1442 let mut retryable = OutboxRecord::from_durable_parts( 1443 enqueue.clone(), 1444 OutboxRevision::new(2).unwrap(), 1445 OutboxStage::Retryable, 1446 None, 1447 Some(DeliveryAttempt::FIRST), 1448 vec![retryable_evidence], 1449 SatisfactionResult::Pending, 1450 None, 1451 20, 1452 ) 1453 .unwrap(); 1454 retryable.claim(lease(3, 21, 30)).unwrap(); 1455 let revision = retryable.revision(); 1456 retryable 1457 .release(LeaseId::new([3; 16]).unwrap(), revision, 22, None) 1458 .unwrap(); 1459 assert_eq!(retryable.stage(), OutboxStage::Retryable); 1460 1461 let malformed = [ 1462 ( 1463 None, 1464 vec![accepted.clone()], 1465 SatisfactionResult::Pending, 1466 20, 1467 ), 1468 ( 1469 Some(DeliveryAttempt::FIRST), 1470 vec![], 1471 SatisfactionResult::Pending, 1472 20, 1473 ), 1474 ( 1475 Some(DeliveryAttempt::FIRST), 1476 vec![ 1477 TargetDeliveryEvidence::new( 1478 target.clone(), 1479 DeliveryAttempt::FIRST, 1480 true, 1481 DeliveryOutcome::accepted(), 1482 9, 1483 ) 1484 .unwrap(), 1485 ], 1486 SatisfactionResult::Satisfied, 1487 20, 1488 ), 1489 ( 1490 Some(DeliveryAttempt::FIRST), 1491 vec![ 1492 TargetDeliveryEvidence::new( 1493 target.clone(), 1494 DeliveryAttempt::FIRST, 1495 true, 1496 DeliveryOutcome::accepted(), 1497 21, 1498 ) 1499 .unwrap(), 1500 ], 1501 SatisfactionResult::Satisfied, 1502 20, 1503 ), 1504 ( 1505 Some(DeliveryAttempt::FIRST), 1506 vec![ 1507 TargetDeliveryEvidence::new( 1508 Target::nostr_relay("wss://foreign.example") 1509 .unwrap() 1510 .fingerprint() 1511 .clone(), 1512 DeliveryAttempt::FIRST, 1513 true, 1514 DeliveryOutcome::accepted(), 1515 20, 1516 ) 1517 .unwrap(), 1518 ], 1519 SatisfactionResult::Satisfied, 1520 20, 1521 ), 1522 ( 1523 Some(DeliveryAttempt::FIRST), 1524 vec![accepted.clone()], 1525 SatisfactionResult::Pending, 1526 20, 1527 ), 1528 ]; 1529 for (last_attempt, evidence, satisfaction, updated) in malformed { 1530 assert_eq!( 1531 OutboxRecord::from_durable_parts( 1532 enqueue.clone(), 1533 OutboxRevision::new(2).unwrap(), 1534 OutboxStage::Retryable, 1535 None, 1536 last_attempt, 1537 evidence, 1538 satisfaction, 1539 None, 1540 updated, 1541 ), 1542 Err(Error::CorruptOutboxRecord) 1543 ); 1544 } 1545 1546 let mut record = enqueue.into_record(); 1547 let active_lease = lease(1, 20, 40); 1548 record.claim(active_lease.clone()).unwrap(); 1549 let make_evidence = |item_id, revision, attempt, request: &DeliveryRequest| { 1550 DeliveryAttemptEvidence::new( 1551 item_id, 1552 active_lease.id(), 1553 revision, 1554 attempt, 1555 receipt_for(request, DeliveryOutcome::accepted()), 1556 30, 1557 ) 1558 .unwrap() 1559 }; 1560 assert_eq!( 1561 record.record_attempt(make_evidence( 1562 OutboxItemId::new([9; 16]).unwrap(), 1563 record.revision(), 1564 DeliveryAttempt::FIRST, 1565 record.request(), 1566 )), 1567 Err(Error::OutboxRevisionConflict) 1568 ); 1569 assert_eq!( 1570 record.record_attempt(make_evidence( 1571 record.item_id(), 1572 OutboxRevision::INITIAL, 1573 DeliveryAttempt::FIRST, 1574 record.request(), 1575 )), 1576 Err(Error::OutboxRevisionConflict) 1577 ); 1578 assert_eq!( 1579 record.record_attempt(make_evidence( 1580 record.item_id(), 1581 record.revision(), 1582 DeliveryAttempt::new(2).unwrap(), 1583 record.request(), 1584 )), 1585 Err(Error::InvalidDeliveryAttempt) 1586 ); 1587 let other = DeliveryRequest::new( 1588 "other-request", 1589 DeliveryPayload::new(signed_event()), 1590 TargetSet::new(vec![Target::nostr_relay("wss://other.example").unwrap()]).unwrap(), 1591 SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::all()), 1592 1_000, 1593 ) 1594 .unwrap(); 1595 assert_eq!( 1596 record.record_attempt(make_evidence( 1597 record.item_id(), 1598 record.revision(), 1599 DeliveryAttempt::FIRST, 1600 &other, 1601 )), 1602 Err(Error::InvalidDeliveryEvidence) 1603 ); 1604 let evidence = make_evidence( 1605 record.item_id(), 1606 record.revision(), 1607 DeliveryAttempt::FIRST, 1608 record.request(), 1609 ); 1610 assert_eq!(evidence.item_id(), record.item_id()); 1611 assert_eq!(evidence.lease_id(), active_lease.id()); 1612 assert_eq!(evidence.expected_revision(), record.revision()); 1613 assert_eq!(evidence.attempt(), DeliveryAttempt::FIRST); 1614 assert_eq!(evidence.recorded_at_unix_ms(), 30); 1615 assert_eq!( 1616 evidence.receipt().request_id(), 1617 record.request().request_id() 1618 ); 1619 record.record_attempt(evidence).unwrap(); 1620 assert_eq!(record.stage(), OutboxStage::Satisfied); 1621 } 1622 1623 #[test] 1624 fn evidence_validation_and_satisfaction_cover_multi_target_policy_edges() { 1625 let targets = vec![ 1626 Target::nostr_relay("wss://one.example").unwrap(), 1627 Target::nostr_relay("wss://two.example").unwrap(), 1628 ]; 1629 let request_with = |policy| { 1630 DeliveryRequest::new( 1631 "storage-outbox-policy-matrix", 1632 DeliveryPayload::new(signed_event()), 1633 TargetSet::new(targets.clone()).unwrap(), 1634 SatisfactionPolicy::new(SatisfactionClass::Accepted, policy), 1635 1_000, 1636 ) 1637 .unwrap() 1638 }; 1639 let evidence = |target: &Target, 1640 attempt: u32, 1641 was_attempted: bool, 1642 outcome: DeliveryOutcome, 1643 recorded_at_unix_ms| { 1644 TargetDeliveryEvidence::new( 1645 target.fingerprint().clone(), 1646 DeliveryAttempt::new(attempt).unwrap(), 1647 was_attempted, 1648 outcome, 1649 recorded_at_unix_ms, 1650 ) 1651 .unwrap() 1652 }; 1653 1654 let all_request = request_with(TargetPolicy::all()); 1655 let accepted = vec![ 1656 evidence(&targets[0], 1, true, DeliveryOutcome::accepted(), 20), 1657 evidence(&targets[1], 1, true, DeliveryOutcome::accepted(), 20), 1658 ]; 1659 assert_eq!( 1660 validate_evidence( 1661 &all_request, 1662 Some(DeliveryAttempt::FIRST), 1663 &accepted, 1664 SatisfactionResult::Satisfied, 1665 10, 1666 20, 1667 ), 1668 Ok(()) 1669 ); 1670 assert_eq!( 1671 validate_evidence( 1672 &all_request, 1673 None, 1674 &[], 1675 SatisfactionResult::Satisfied, 1676 10, 1677 20, 1678 ), 1679 Err(Error::CorruptOutboxRecord) 1680 ); 1681 1682 let duplicated = vec![accepted[0].clone(), accepted[0].clone()]; 1683 assert_eq!( 1684 validate_evidence( 1685 &all_request, 1686 Some(DeliveryAttempt::FIRST), 1687 &duplicated, 1688 SatisfactionResult::Satisfied, 1689 10, 1690 20, 1691 ), 1692 Err(Error::CorruptOutboxRecord) 1693 ); 1694 let mismatched_times = vec![ 1695 accepted[0].clone(), 1696 evidence(&targets[1], 1, true, DeliveryOutcome::accepted(), 21), 1697 ]; 1698 assert_eq!( 1699 validate_evidence( 1700 &all_request, 1701 Some(DeliveryAttempt::FIRST), 1702 &mismatched_times, 1703 SatisfactionResult::Satisfied, 1704 10, 1705 21, 1706 ), 1707 Err(Error::CorruptOutboxRecord) 1708 ); 1709 let skipped = vec![ 1710 evidence(&targets[0], 1, false, DeliveryOutcome::unavailable(), 20), 1711 evidence(&targets[1], 1, false, DeliveryOutcome::unavailable(), 20), 1712 ]; 1713 assert_eq!( 1714 validate_evidence( 1715 &all_request, 1716 Some(DeliveryAttempt::FIRST), 1717 &skipped, 1718 SatisfactionResult::Pending, 1719 10, 1720 20, 1721 ), 1722 Ok(()) 1723 ); 1724 let regressing = vec![ 1725 accepted[0].clone(), 1726 accepted[1].clone(), 1727 evidence(&targets[0], 2, true, DeliveryOutcome::accepted(), 19), 1728 evidence(&targets[1], 2, true, DeliveryOutcome::accepted(), 19), 1729 ]; 1730 assert_eq!( 1731 validate_evidence( 1732 &all_request, 1733 Some(DeliveryAttempt::new(2).unwrap()), 1734 ®ressing, 1735 SatisfactionResult::Satisfied, 1736 10, 1737 20, 1738 ), 1739 Err(Error::CorruptOutboxRecord) 1740 ); 1741 1742 let retryable = vec![evidence( 1743 &targets[0], 1744 1, 1745 true, 1746 DeliveryOutcome::unavailable(), 1747 20, 1748 )]; 1749 let one_accepted = vec![accepted[0].clone()]; 1750 let any_request = request_with(TargetPolicy::any()); 1751 assert_eq!( 1752 evaluate_satisfaction(&any_request, &[]), 1753 SatisfactionResult::Pending 1754 ); 1755 assert_eq!( 1756 evaluate_satisfaction(&any_request, &retryable), 1757 SatisfactionResult::Pending 1758 ); 1759 assert_eq!( 1760 evaluate_satisfaction(&any_request, &one_accepted), 1761 SatisfactionResult::Satisfied 1762 ); 1763 let quorum_request = request_with(TargetPolicy::quorum(2).unwrap()); 1764 assert_eq!( 1765 evaluate_satisfaction(&quorum_request, &one_accepted), 1766 SatisfactionResult::Pending 1767 ); 1768 let required_request = 1769 request_with(TargetPolicy::required(vec![targets[0].fingerprint().clone()]).unwrap()); 1770 assert_eq!( 1771 evaluate_satisfaction(&required_request, &one_accepted), 1772 SatisfactionResult::Satisfied 1773 ); 1774 assert_eq!( 1775 evaluate_satisfaction(&required_request, &retryable), 1776 SatisfactionResult::Pending 1777 ); 1778 1779 assert_eq!( 1780 DeliveryAttemptEvidence::new( 1781 OutboxItemId::new([1; 16]).unwrap(), 1782 LeaseId::new([2; 16]).unwrap(), 1783 OutboxRevision::INITIAL, 1784 DeliveryAttempt::FIRST, 1785 receipt_for(&request(), DeliveryOutcome::accepted()), 1786 0, 1787 ), 1788 Err(Error::InvalidOutboxTimestamp) 1789 ); 1790 } 1791 }