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 }