delivery_selection.rs (11481B)
1 use super::*; 2 use futures::channel::oneshot; 3 use radroots_storage::authored_delivery::DeliveryAttemptOutcome; 4 use radroots_transport::policy::SatisfactionState; 5 6 type Pending = ( 7 DeliveryRequest, 8 TargetSet, 9 oneshot::Sender<Result<DeliveryReceipt, SinkFailure>>, 10 ); 11 12 #[derive(Default)] 13 struct SelectedSink { 14 pending: Mutex<VecDeque<Pending>>, 15 calls: AtomicUsize, 16 } 17 18 impl EventSink for SelectedSink { 19 fn status(&self) -> radroots_transport::BoxFuture<'_, Result<SinkStatus, TransportError>> { 20 Box::pin(async { panic!("no status probe") }) 21 } 22 fn deliver( 23 &self, 24 _: DeliveryRequest, 25 ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> { 26 Box::pin(async { panic!("selected attempts must not fall back") }) 27 } 28 fn deliver_selected( 29 &self, 30 request: DeliveryRequest, 31 selected: TargetSet, 32 ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> { 33 Box::pin(async move { 34 request.validate_target_selection(&selected).unwrap(); 35 self.calls.fetch_add(1, Ordering::SeqCst); 36 let (send, receive) = oneshot::channel(); 37 self.pending 38 .lock() 39 .unwrap() 40 .push_back((request, selected, send)); 41 receive.await.unwrap() 42 }) 43 } 44 } 45 46 fn selected_receipt(request: &DeliveryRequest, selected: &TargetSet) -> DeliveryReceipt { 47 DeliveryReceipt::for_request( 48 request, 49 request 50 .target_set() 51 .targets() 52 .iter() 53 .map(|target| { 54 if selected.targets().contains(target) { 55 DeliveryTargetReceipt::attempted(target.clone(), DeliveryOutcome::accepted()) 56 } else { 57 DeliveryTargetReceipt::skipped(target.clone(), DeliveryOutcome::unavailable()) 58 .unwrap() 59 } 60 }) 61 .collect(), 62 ) 63 .unwrap() 64 } 65 66 fn poll_pending(future: &mut (impl std::future::Future + Unpin)) { 67 assert!( 68 future 69 .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref())) 70 .is_pending() 71 ); 72 } 73 74 fn engine(storage: Arc<dyn SyncStorage>, clock: Arc<TestClock>, sink: Arc<SelectedSink>) -> Engine { 75 Engine::builder( 76 storage, 77 clock, 78 Arc::new(TestIds(AtomicU64::new(10))), 79 DeadlinePolicy::new(10_000, 10_000, 10_000).unwrap(), 80 ) 81 .sink(sink) 82 .signer(Arc::new(MockSigner::new(SignBehavior::Success { 83 completed_at_unix_ms: 1_800_000_200_500, 84 }))) 85 .build() 86 .unwrap() 87 } 88 89 fn push() -> PushRequest { 90 request_with_policy( 91 211, 92 &["wss://one.example", "wss://two.example"], 93 SatisfactionClass::Accepted, 94 TargetPolicy::all(), 95 ) 96 } 97 98 #[test] 99 fn invalid_selection_cannot_claim_and_out_of_selection_results_fail_closed() { 100 for failure in [false, true] { 101 let storage = Arc::new(MemoryStorage::new( 102 SourceGeneration::new([211; 32]).unwrap(), 103 )); 104 let clock = Arc::new(TestClock(AtomicU64::new(1_800_000_200_000))); 105 let sink = Arc::new(SelectedSink::default()); 106 let engine = engine(storage, clock, sink.clone()); 107 let push = push(); 108 execute_to_admitted(&engine, &push); 109 let before = block_on(engine.push_status(push.operation_id())) 110 .unwrap() 111 .unwrap(); 112 let original = before.delivery_plan().request().unwrap(); 113 let foreign = 114 TargetSet::new(vec![Target::nostr_relay("wss://foreign.example").unwrap()]).unwrap(); 115 assert!(matches!( 116 block_on(engine.deliver_push_selected(push.operation_id(), foreign)), 117 Err(Error::InvalidDeliveryRequest) 118 )); 119 let unchanged = block_on(engine.push_status(push.operation_id())) 120 .unwrap() 121 .unwrap(); 122 assert_eq!(unchanged.delivery_plan(), before.delivery_plan()); 123 assert_eq!(sink.calls.load(Ordering::SeqCst), 0); 124 let selected = TargetSet::new(vec![original.target_set().targets()[0].clone()]).unwrap(); 125 let mut future = Box::pin(engine.deliver_push_selected(push.operation_id(), selected)); 126 poll_pending(&mut future); 127 let (request, _, sender) = sink.pending.lock().unwrap().pop_front().unwrap(); 128 assert_eq!(&request, original); 129 let forbidden = DeliveryTargetReceipt::attempted( 130 request.target_set().targets()[1].clone(), 131 DeliveryOutcome::accepted(), 132 ); 133 let result = if failure { 134 Err(SinkFailure::for_request( 135 &request, 136 "upstream_failure", 137 Retryability::Retryable, 138 None, 139 None, 140 vec![forbidden], 141 ) 142 .unwrap()) 143 } else { 144 Ok(receipt(&request, vec![DeliveryOutcome::accepted(); 2]).unwrap()) 145 }; 146 sender.send(result).unwrap(); 147 let result = block_on(future).unwrap(); 148 assert_eq!(result.plan().request(), Some(original)); 149 let DeliveryAttemptOutcome::SinkFailure(failure) = 150 result.plan().delivery_facts()[0].outcome() 151 else { 152 panic!("invalid adapter evidence is never acceptance") 153 }; 154 assert_eq!(failure.code(), "invalid_transport_contract"); 155 assert!(failure.partial_evidence().is_empty()); 156 assert_ne!( 157 result.plan().delivery_satisfaction().unwrap(), 158 SatisfactionState::Satisfied 159 ); 160 } 161 } 162 163 #[test] 164 fn selected_late_acceptance_survives_stop_without_scheduling_new_targets() { 165 let storage = Arc::new(MemoryStorage::new( 166 SourceGeneration::new([212; 32]).unwrap(), 167 )); 168 let clock = Arc::new(TestClock(AtomicU64::new(1_800_000_200_000))); 169 let sink = Arc::new(SelectedSink::default()); 170 let engine = engine(storage.clone(), clock.clone(), sink.clone()); 171 let push = push(); 172 execute_to_admitted(&engine, &push); 173 let before = block_on(engine.push_status(push.operation_id())) 174 .unwrap() 175 .unwrap(); 176 let request = before.delivery_plan().request().unwrap(); 177 let selected = TargetSet::new(vec![request.target_set().targets()[0].clone()]).unwrap(); 178 let mut future = Box::pin(engine.deliver_push_selected(push.operation_id(), selected.clone())); 179 poll_pending(&mut future); 180 block_on(engine.cancel_push(push.operation_id())).unwrap(); 181 let (captured, captured_selection, sender) = sink.pending.lock().unwrap().pop_front().unwrap(); 182 assert_eq!(&captured, request); 183 assert_eq!(captured_selection, selected); 184 let expected = selected_receipt(&captured, &selected); 185 sender.send(Ok(expected.clone())).unwrap(); 186 let late = block_on(future).unwrap(); 187 assert_eq!(late.plan().state(), AuthoredDeliveryState::Cancelled); 188 assert_eq!(late.plan().attempt_count(), 0); 189 assert_eq!( 190 late.plan().delivery_facts()[0].outcome(), 191 &DeliveryAttemptOutcome::Receipt(expected) 192 ); 193 let recovery = delivery_evidence::source_only(storage, clock); 194 let replay = block_on(recovery.deliver_push_selected(push.operation_id(), selected)).unwrap(); 195 assert!(replay.is_replay()); 196 assert_eq!(replay.plan(), late.plan()); 197 assert_eq!(sink.calls.load(Ordering::SeqCst), 1); 198 } 199 200 #[tokio::test] 201 async fn sqlite_selected_facts_survive_reopen_and_next_target_finishes_same_request() { 202 let directory = tempfile::tempdir().unwrap(); 203 let paths = Paths::from_directory(directory.path()).unwrap(); 204 let storage = Arc::new( 205 SqliteStorage::open( 206 OpenOptions::new(paths.clone(), OpenMode::Create) 207 .with_source_generation(SourceGeneration::new([213; 32]).unwrap(), 1) 208 .unwrap(), 209 ) 210 .await 211 .unwrap(), 212 ); 213 let clock = Arc::new(TestClock(AtomicU64::new(1_800_000_200_000))); 214 let sink = Arc::new(SelectedSink::default()); 215 let first_engine = engine(storage.clone(), clock.clone(), sink.clone()); 216 let push = push(); 217 first_engine.sign_prepared(push.clone()).await.unwrap(); 218 first_engine 219 .admit_signed(push.operation_id()) 220 .await 221 .unwrap(); 222 let before = first_engine 223 .push_status(push.operation_id()) 224 .await 225 .unwrap() 226 .unwrap(); 227 let original = before.delivery_plan().request().unwrap().clone(); 228 let selected_a = TargetSet::new(vec![original.target_set().targets()[0].clone()]).unwrap(); 229 let first = { 230 let mut future = 231 Box::pin(first_engine.deliver_push_selected(push.operation_id(), selected_a.clone())); 232 // SQLite storage needs an executor turn before the sink is reached. 233 let responder = async { 234 let pending = loop { 235 if let Some(pending) = sink.pending.lock().unwrap().pop_front() { 236 break pending; 237 } 238 tokio::task::yield_now().await; 239 }; 240 assert_eq!(pending.0, original); 241 assert_eq!(pending.1, selected_a); 242 pending 243 .2 244 .send(Ok(selected_receipt(&pending.0, &pending.1))) 245 .unwrap(); 246 }; 247 tokio::time::timeout(std::time::Duration::from_secs(10), async { 248 let (result, ()) = tokio::join!(&mut future, responder); 249 result.unwrap() 250 }) 251 .await 252 .unwrap() 253 }; 254 assert_ne!(first.plan().state(), AuthoredDeliveryState::Satisfied); 255 assert_eq!(first.plan().attempt_count(), 1); 256 clock.0.store( 257 first.plan().retry().unwrap().not_before_unix_ms(), 258 Ordering::SeqCst, 259 ); 260 drop(first_engine); 261 storage.close().await.unwrap(); 262 let storage = Arc::new( 263 SqliteStorage::open(OpenOptions::new(paths, OpenMode::ReadWriteExisting)) 264 .await 265 .unwrap(), 266 ); 267 let second_engine = engine(storage.clone(), clock.clone(), sink.clone()); 268 let recovered = second_engine 269 .push_status(push.operation_id()) 270 .await 271 .unwrap() 272 .unwrap(); 273 assert_eq!( 274 recovered.delivery_plan().delivery_facts(), 275 first.plan().delivery_facts() 276 ); 277 let selected_b = TargetSet::new(vec![original.target_set().targets()[1].clone()]).unwrap(); 278 let responder = async { 279 let pending = loop { 280 if let Some(pending) = sink.pending.lock().unwrap().pop_front() { 281 break pending; 282 } 283 tokio::task::yield_now().await; 284 }; 285 assert_eq!(pending.0, original); 286 assert_eq!(pending.1, selected_b); 287 pending 288 .2 289 .send(Ok(selected_receipt(&pending.0, &pending.1))) 290 .unwrap(); 291 }; 292 let finished = tokio::time::timeout(std::time::Duration::from_secs(10), async { 293 let (result, ()) = tokio::join!( 294 second_engine.deliver_push_selected(push.operation_id(), selected_b.clone()), 295 responder 296 ); 297 result.unwrap() 298 }) 299 .await 300 .unwrap(); 301 assert_eq!(finished.plan().request(), Some(&original)); 302 assert_eq!(finished.plan().state(), AuthoredDeliveryState::Satisfied); 303 assert_eq!(finished.plan().attempt_count(), 2); 304 assert_eq!(finished.plan().delivery_facts().len(), 2); 305 assert_eq!(sink.calls.load(Ordering::SeqCst), 2); 306 drop(second_engine); 307 storage.close().await.unwrap(); 308 }