lib

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

outbox.rs (13948B)


      1 use futures_executor::block_on;
      2 use radroots_event::{SignedEvent, wire::Nip01EventWire};
      3 use radroots_storage::{
      4     Error, Outbox,
      5     journal::OperationInstanceId,
      6     outbox::{
      7         ClaimOutboxItems, ClaimedOutboxItem, DeliveryAttempt, DeliveryAttemptEvidence,
      8         DeliveryPlanDigest, EnqueueDisposition, EnqueueOutboxItem, EnqueueReceipt, LeaseId,
      9         LeaseOwner, OutboxItemId, OutboxLease, OutboxRecord, OutboxStage, OutboxStatus,
     10         SatisfactionResult,
     11     },
     12 };
     13 use radroots_transport::{
     14     BoxFuture, DeliveryReceipt, DeliveryRequest, Target, TargetSet,
     15     outcome::DeliveryOutcome,
     16     policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy},
     17     sink::{DeliveryPayload, DeliveryTargetReceipt},
     18 };
     19 use std::{collections::BTreeMap, sync::Mutex};
     20 
     21 struct MemoryOutbox {
     22     records: Mutex<BTreeMap<OutboxItemId, OutboxRecord>>,
     23 }
     24 
     25 impl MemoryOutbox {
     26     fn new() -> Self {
     27         Self {
     28             records: Mutex::new(BTreeMap::new()),
     29         }
     30     }
     31 }
     32 
     33 impl Outbox for MemoryOutbox {
     34     fn enqueue(&self, item: EnqueueOutboxItem) -> BoxFuture<'_, Result<EnqueueReceipt, Error>> {
     35         Box::pin(async move {
     36             let mut records = self.records.lock().expect("test outbox lock");
     37             if let Some(existing) = records.get(&item.item_id()) {
     38                 let exact = existing.operation_instance_id() == item.operation_instance_id()
     39                     && existing.plan_digest() == item.plan_digest()
     40                     && existing.request() == item.request();
     41                 return if exact {
     42                     Ok(EnqueueReceipt::new(
     43                         EnqueueDisposition::Replay,
     44                         existing.clone(),
     45                     ))
     46                 } else {
     47                     Err(Error::OutboxPlanConflict)
     48                 };
     49             }
     50             let record = item.into_record();
     51             records.insert(record.item_id(), record.clone());
     52             Ok(EnqueueReceipt::new(EnqueueDisposition::Created, record))
     53         })
     54     }
     55 
     56     fn item(&self, item_id: OutboxItemId) -> BoxFuture<'_, Result<Option<OutboxRecord>, Error>> {
     57         Box::pin(async move {
     58             Ok(self
     59                 .records
     60                 .lock()
     61                 .expect("test outbox lock")
     62                 .get(&item_id)
     63                 .cloned())
     64         })
     65     }
     66 
     67     fn claim(
     68         &self,
     69         request: ClaimOutboxItems,
     70     ) -> BoxFuture<'_, Result<Vec<ClaimedOutboxItem>, Error>> {
     71         Box::pin(async move {
     72             let mut records = self.records.lock().expect("test outbox lock");
     73             let mut claimed = Vec::new();
     74             for record in records.values_mut() {
     75                 if claimed.len() >= usize::from(request.limit()) || record.stage().is_terminal() {
     76                     continue;
     77                 }
     78                 if matches!(record.retry_not_before_unix_ms(), Some(value) if value > request.now_unix_ms())
     79                 {
     80                     continue;
     81                 }
     82                 if record
     83                     .lease()
     84                     .is_some_and(|lease| lease.is_active_at(request.now_unix_ms()))
     85                 {
     86                     continue;
     87                 }
     88                 let lease = OutboxLease::new(
     89                     request.lease_id_for(record.item_id()),
     90                     request.owner().clone(),
     91                     request.now_unix_ms(),
     92                     request.lease_expires_at_unix_ms(),
     93                 )?;
     94                 record.claim(lease.clone())?;
     95                 claimed.push(ClaimedOutboxItem::new(record.clone(), lease));
     96             }
     97             Ok(claimed)
     98         })
     99     }
    100 
    101     fn record_attempt(
    102         &self,
    103         evidence: DeliveryAttemptEvidence,
    104     ) -> BoxFuture<'_, Result<OutboxRecord, Error>> {
    105         Box::pin(async move {
    106             let mut records = self.records.lock().expect("test outbox lock");
    107             let record = records
    108                 .get_mut(&evidence.item_id())
    109                 .ok_or(Error::OutboxItemNotFound)?;
    110             record.record_attempt(evidence)?;
    111             Ok(record.clone())
    112         })
    113     }
    114 
    115     fn release(
    116         &self,
    117         item_id: OutboxItemId,
    118         lease_id: LeaseId,
    119         expected_revision: radroots_storage::outbox::OutboxRevision,
    120         released_at_unix_ms: u64,
    121         retry_not_before_unix_ms: Option<u64>,
    122     ) -> BoxFuture<'_, Result<OutboxRecord, Error>> {
    123         Box::pin(async move {
    124             let mut records = self.records.lock().expect("test outbox lock");
    125             let record = records.get_mut(&item_id).ok_or(Error::OutboxItemNotFound)?;
    126             record.release(
    127                 lease_id,
    128                 expected_revision,
    129                 released_at_unix_ms,
    130                 retry_not_before_unix_ms,
    131             )?;
    132             Ok(record.clone())
    133         })
    134     }
    135 
    136     fn status(&self) -> BoxFuture<'_, Result<OutboxStatus, Error>> {
    137         Box::pin(async move {
    138             let mut status = OutboxStatus {
    139                 pending: 0,
    140                 leased: 0,
    141                 retryable: 0,
    142                 satisfied: 0,
    143                 exhausted: 0,
    144             };
    145             for record in self.records.lock().expect("test outbox lock").values() {
    146                 match record.stage() {
    147                     OutboxStage::Pending => status.pending += 1,
    148                     OutboxStage::Leased => status.leased += 1,
    149                     OutboxStage::Retryable => status.retryable += 1,
    150                     OutboxStage::Satisfied => status.satisfied += 1,
    151                     OutboxStage::Exhausted => status.exhausted += 1,
    152                 }
    153             }
    154             Ok(status)
    155         })
    156     }
    157 }
    158 
    159 fn signed_event() -> SignedEvent {
    160     let mut wire = Nip01EventWire {
    161         id: "0".repeat(64),
    162         pubkey: "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df".to_owned(),
    163         created_at: 1_800_000_100,
    164         kind: 0,
    165         tags: vec![],
    166         content: "{\"display_name\":\"Moss Street Farm\",\"bot\":false}".to_owned(),
    167         sig: "42".repeat(64),
    168         extra: Default::default(),
    169     };
    170     wire.id = wire
    171         .computed_event_id()
    172         .expect("canonical event id")
    173         .to_hex();
    174     let raw_json = serde_json::json!({
    175         "id": &wire.id,
    176         "pubkey": &wire.pubkey,
    177         "created_at": wire.created_at,
    178         "kind": wire.kind,
    179         "tags": &wire.tags,
    180         "content": &wire.content,
    181         "sig": &wire.sig,
    182     })
    183     .to_string();
    184     SignedEvent::from_wire_verified_id(wire, raw_json).expect("signed event")
    185 }
    186 
    187 fn targets() -> Vec<Target> {
    188     vec![
    189         Target::nostr_relay("wss://one.example").expect("first target"),
    190         Target::nostr_relay("wss://two.example").expect("second target"),
    191     ]
    192 }
    193 
    194 fn request() -> DeliveryRequest {
    195     DeliveryRequest::new(
    196         "outbox-test-request",
    197         DeliveryPayload::new(signed_event()),
    198         TargetSet::new(targets()).expect("target set"),
    199         SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::all()),
    200         10_000,
    201     )
    202     .expect("delivery request")
    203 }
    204 
    205 fn enqueue(item_byte: u8, digest_byte: u8) -> EnqueueOutboxItem {
    206     EnqueueOutboxItem::new(
    207         OutboxItemId::new([item_byte; 16]).expect("item id"),
    208         OperationInstanceId::new([9; 16]).expect("operation instance"),
    209         DeliveryPlanDigest::new([digest_byte; 32]),
    210         request(),
    211         10,
    212     )
    213     .expect("enqueue request")
    214 }
    215 
    216 fn claim(store: &dyn Outbox, now: u64, expiry: u64, seed: u8) -> ClaimedOutboxItem {
    217     let request = ClaimOutboxItems::new(
    218         LeaseOwner::parse("worker-a").expect("owner"),
    219         LeaseId::new([seed; 16]).expect("lease seed"),
    220         now,
    221         expiry,
    222         1,
    223     )
    224     .expect("claim request");
    225     block_on(store.claim(request))
    226         .expect("claim")
    227         .pop()
    228         .expect("claimed item")
    229 }
    230 
    231 fn receipt(request: &DeliveryRequest, outcomes: [DeliveryOutcome; 2]) -> DeliveryReceipt {
    232     let receipts = request
    233         .target_set()
    234         .targets()
    235         .iter()
    236         .cloned()
    237         .zip(outcomes)
    238         .map(|(target, outcome)| DeliveryTargetReceipt::attempted(target, outcome))
    239         .collect();
    240     DeliveryReceipt::for_request(request, receipts).expect("delivery receipt")
    241 }
    242 
    243 #[test]
    244 fn enqueue_is_idempotent_and_rejects_a_conflicting_plan() {
    245     let store = MemoryOutbox::new();
    246     let first = enqueue(1, 2);
    247     let created = block_on(store.enqueue(first.clone())).expect("created");
    248     assert_eq!(created.disposition(), EnqueueDisposition::Created);
    249     let replay = block_on(store.enqueue(first)).expect("exact replay");
    250     assert_eq!(replay.disposition(), EnqueueDisposition::Replay);
    251 
    252     let conflict = enqueue(1, 3);
    253     assert_eq!(
    254         block_on(store.enqueue(conflict)),
    255         Err(Error::OutboxPlanConflict)
    256     );
    257 }
    258 
    259 #[test]
    260 fn leases_exclude_concurrent_claims_expire_and_defer_retries() {
    261     let store = MemoryOutbox::new();
    262     block_on(store.enqueue(enqueue(1, 2))).expect("enqueue");
    263     let first = claim(&store, 100, 200, 3);
    264 
    265     let concurrent = ClaimOutboxItems::new(
    266         LeaseOwner::parse("worker-b").expect("owner"),
    267         LeaseId::new([4; 16]).expect("seed"),
    268         150,
    269         250,
    270         1,
    271     )
    272     .expect("claim request");
    273     assert!(block_on(store.claim(concurrent)).expect("claim").is_empty());
    274 
    275     let stale = DeliveryAttemptEvidence::new(
    276         first.record().item_id(),
    277         first.lease().id(),
    278         first.record().revision(),
    279         DeliveryAttempt::FIRST,
    280         receipt(
    281             first.record().request(),
    282             [DeliveryOutcome::accepted(), DeliveryOutcome::unavailable()],
    283         ),
    284         200,
    285     )
    286     .expect("evidence");
    287     assert_eq!(
    288         block_on(store.record_attempt(stale)),
    289         Err(Error::OutboxLeaseExpired)
    290     );
    291 
    292     let reclaimed = claim(&store, 200, 300, 5);
    293     let released = block_on(store.release(
    294         reclaimed.record().item_id(),
    295         reclaimed.lease().id(),
    296         reclaimed.record().revision(),
    297         210,
    298         Some(250),
    299     ))
    300     .expect("release");
    301     assert_eq!(released.stage(), OutboxStage::Pending);
    302     assert_eq!(released.retry_not_before_unix_ms(), Some(250));
    303 }
    304 
    305 #[test]
    306 fn partial_retryable_evidence_advances_to_satisfaction() {
    307     let store = MemoryOutbox::new();
    308     block_on(store.enqueue(enqueue(1, 2))).expect("enqueue");
    309     let first = claim(&store, 100, 200, 3);
    310     let first_result = block_on(
    311         store.record_attempt(
    312             DeliveryAttemptEvidence::new(
    313                 first.record().item_id(),
    314                 first.lease().id(),
    315                 first.record().revision(),
    316                 DeliveryAttempt::FIRST,
    317                 receipt(
    318                     first.record().request(),
    319                     [DeliveryOutcome::accepted(), DeliveryOutcome::unavailable()],
    320                 ),
    321                 150,
    322             )
    323             .expect("attempt evidence"),
    324         ),
    325     )
    326     .expect("record partial attempt");
    327     assert_eq!(first_result.stage(), OutboxStage::Retryable);
    328     assert_eq!(first_result.satisfaction(), SatisfactionResult::Pending);
    329     assert_eq!(first_result.evidence().len(), 2);
    330     assert!(first_result.evidence()[1].outcome().is_retryable());
    331 
    332     let second = claim(&store, 250, 350, 4);
    333     let second_result = block_on(
    334         store.record_attempt(
    335             DeliveryAttemptEvidence::new(
    336                 second.record().item_id(),
    337                 second.lease().id(),
    338                 second.record().revision(),
    339                 DeliveryAttempt::new(2).expect("second attempt"),
    340                 receipt(
    341                     second.record().request(),
    342                     [DeliveryOutcome::accepted(), DeliveryOutcome::delivered()],
    343                 ),
    344                 300,
    345             )
    346             .expect("attempt evidence"),
    347         ),
    348     )
    349     .expect("record successful attempt");
    350     assert_eq!(second_result.stage(), OutboxStage::Satisfied);
    351     assert_eq!(second_result.satisfaction(), SatisfactionResult::Satisfied);
    352     assert_eq!(second_result.evidence().len(), 4);
    353     let second_target = &second_result.request().target_set().targets()[1];
    354     assert_eq!(
    355         second_result
    356             .latest_target_evidence(second_target.fingerprint())
    357             .expect("latest target evidence")
    358             .attempt()
    359             .get(),
    360         2
    361     );
    362     assert_eq!(block_on(store.status()).expect("status").satisfied, 1);
    363 }
    364 
    365 #[test]
    366 fn terminal_outcomes_exhaust_the_plan_and_models_are_bounded() {
    367     let store = MemoryOutbox::new();
    368     block_on(store.enqueue(enqueue(1, 2))).expect("enqueue");
    369     let claimed = claim(&store, 100, 200, 3);
    370     let exhausted = block_on(
    371         store.record_attempt(
    372             DeliveryAttemptEvidence::new(
    373                 claimed.record().item_id(),
    374                 claimed.lease().id(),
    375                 claimed.record().revision(),
    376                 DeliveryAttempt::FIRST,
    377                 receipt(
    378                     claimed.record().request(),
    379                     [DeliveryOutcome::rejected(), DeliveryOutcome::rejected()],
    380                 ),
    381                 150,
    382             )
    383             .expect("attempt evidence"),
    384         ),
    385     )
    386     .expect("record terminal attempt");
    387     assert_eq!(exhausted.stage(), OutboxStage::Exhausted);
    388     assert_eq!(exhausted.satisfaction(), SatisfactionResult::Exhausted);
    389     assert_eq!(block_on(store.status()).expect("status").total(), Some(1));
    390 
    391     assert_eq!(OutboxItemId::new([0; 16]), Err(Error::InvalidOutboxItemId));
    392     assert_eq!(DeliveryAttempt::new(0), Err(Error::InvalidDeliveryAttempt));
    393     assert_eq!(
    394         ClaimOutboxItems::new(
    395             LeaseOwner::parse("worker").expect("owner"),
    396             LeaseId::new([1; 16]).expect("seed"),
    397             1,
    398             2,
    399             0,
    400         ),
    401         Err(Error::InvalidOutboxClaimLimit)
    402     );
    403     assert!(
    404         format!("{:?}", LeaseOwner::parse("secret-worker").expect("owner")).contains("[REDACTED]")
    405     );
    406 }