support.rs (11082B)
1 use std::sync::{ 2 Arc, Mutex, 3 atomic::{AtomicBool, Ordering}, 4 }; 5 6 use futures::future; 7 use radroots_transport::{ 8 BoxFuture, DeliveryReceipt, DeliveryRequest, Error, EventSink, EventSource, FetchPage, 9 FetchRequest, SinkFailure, SinkStatus, SourceStatus, TransportId, 10 capability::{Availability, Maturity, SinkCapabilities, SourceCapabilities}, 11 outcome::{DeliveryOutcome, FetchTargetOutcome, FetchTargetState}, 12 sink::DeliveryTargetReceipt, 13 source::NextPage, 14 }; 15 16 use crate::suite::{NOW_UNIX_MS, SinkConformanceHarness, SourceConformanceHarness}; 17 18 #[derive(Clone)] 19 enum Mode { 20 Success, 21 Fail(Error), 22 Pending, 23 } 24 25 #[derive(Default)] 26 pub(crate) struct SourceState { 27 request: Mutex<Option<FetchRequest>>, 28 pub(crate) published: AtomicBool, 29 pub(crate) cancelled_after_publish: AtomicBool, 30 } 31 32 impl SourceState { 33 pub(crate) fn request(&self) -> Option<FetchRequest> { 34 self.request.lock().expect("source request lock").clone() 35 } 36 } 37 38 #[derive(Default)] 39 pub(crate) struct SinkState { 40 request: Mutex<Option<DeliveryRequest>>, 41 pub(crate) published: AtomicBool, 42 pub(crate) cancelled_after_publish: AtomicBool, 43 } 44 45 impl SinkState { 46 pub(crate) fn request(&self) -> Option<DeliveryRequest> { 47 self.request.lock().expect("sink request lock").clone() 48 } 49 } 50 51 pub(crate) struct MockSource { 52 mode: Mode, 53 state: Arc<SourceState>, 54 } 55 56 impl MockSource { 57 pub(crate) fn successful() -> Self { 58 Self::new(Mode::Success) 59 } 60 61 pub(crate) fn failing(error: Error) -> Self { 62 Self::new(Mode::Fail(error)) 63 } 64 65 pub(crate) fn pending() -> Self { 66 Self::new(Mode::Pending) 67 } 68 69 fn new(mode: Mode) -> Self { 70 Self { 71 mode, 72 state: Arc::new(SourceState::default()), 73 } 74 } 75 } 76 77 impl EventSource for MockSource { 78 fn status(&self) -> BoxFuture<'_, Result<SourceStatus, Error>> { 79 Box::pin(async { Ok(source_status()) }) 80 } 81 82 fn fetch(&self, request: FetchRequest) -> BoxFuture<'_, Result<FetchPage, Error>> { 83 let state = Arc::clone(&self.state); 84 let mode = self.mode.clone(); 85 Box::pin(async move { 86 *state.request.lock().expect("source request lock") = Some(request.clone()); 87 if request.bounds().deadline_unix_ms() <= NOW_UNIX_MS { 88 return Err(Error::InvalidFetchDeadline); 89 } 90 match mode { 91 Mode::Fail(error) => Err(error), 92 Mode::Pending => { 93 state.published.store(true, Ordering::SeqCst); 94 let state_for_drop = Arc::clone(&state); 95 let _publish_guard = ScopeGuard::new(move || { 96 state_for_drop 97 .cancelled_after_publish 98 .store(true, Ordering::SeqCst); 99 }); 100 future::pending().await 101 } 102 Mode::Success => { 103 let outcomes = request 104 .target_set() 105 .targets() 106 .iter() 107 .enumerate() 108 .map(|(index, target)| { 109 FetchTargetOutcome::new( 110 target.fingerprint().clone(), 111 if index == 0 { 112 FetchTargetState::Complete 113 } else { 114 FetchTargetState::FailedRetryable 115 }, 116 ) 117 }) 118 .collect(); 119 FetchPage::for_request(&request, Vec::new(), outcomes, NextPage::Complete) 120 } 121 } 122 }) 123 } 124 } 125 126 impl SourceConformanceHarness for MockSource { 127 fn source(&self) -> &dyn EventSource { 128 self 129 } 130 131 fn expected_status(&self) -> SourceStatus { 132 source_status() 133 } 134 135 fn target_set(&self) -> radroots_transport::TargetSet { 136 target_set() 137 } 138 139 fn now_unix_ms(&self) -> u64 { 140 NOW_UNIX_MS 141 } 142 143 fn captured_request(&self) -> Option<FetchRequest> { 144 self.state.request() 145 } 146 147 fn published(&self) -> bool { 148 self.state.published.load(Ordering::SeqCst) 149 } 150 151 fn cancelled_after_publish(&self) -> bool { 152 self.state.cancelled_after_publish.load(Ordering::SeqCst) 153 } 154 } 155 156 pub(crate) struct MockSink { 157 mode: Mode, 158 state: Arc<SinkState>, 159 } 160 161 impl MockSink { 162 pub(crate) fn successful() -> Self { 163 Self::new(Mode::Success) 164 } 165 166 pub(crate) fn failing(error: Error) -> Self { 167 Self::new(Mode::Fail(error)) 168 } 169 170 pub(crate) fn pending() -> Self { 171 Self::new(Mode::Pending) 172 } 173 174 fn new(mode: Mode) -> Self { 175 Self { 176 mode, 177 state: Arc::new(SinkState::default()), 178 } 179 } 180 } 181 182 impl EventSink for MockSink { 183 fn status(&self) -> BoxFuture<'_, Result<SinkStatus, Error>> { 184 Box::pin(async { Ok(sink_status()) }) 185 } 186 187 fn deliver( 188 &self, 189 request: DeliveryRequest, 190 ) -> BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> { 191 let state = Arc::clone(&self.state); 192 let mode = self.mode.clone(); 193 Box::pin(async move { 194 *state.request.lock().expect("sink request lock") = Some(request.clone()); 195 if request.deadline_unix_ms() <= NOW_UNIX_MS { 196 return Err(SinkFailure::invalid_contract(&request)); 197 } 198 match mode { 199 Mode::Fail(_) => Err(SinkFailure::invalid_contract(&request)), 200 Mode::Pending => { 201 state.published.store(true, Ordering::SeqCst); 202 let state_for_drop = Arc::clone(&state); 203 let _publish_guard = ScopeGuard::new(move || { 204 state_for_drop 205 .cancelled_after_publish 206 .store(true, Ordering::SeqCst); 207 }); 208 future::pending().await 209 } 210 Mode::Success => { 211 let receipts = request 212 .target_set() 213 .targets() 214 .iter() 215 .enumerate() 216 .map(|(index, target)| { 217 DeliveryTargetReceipt::attempted( 218 target.clone(), 219 if index == 0 { 220 DeliveryOutcome::delivered() 221 } else { 222 DeliveryOutcome::unavailable() 223 }, 224 ) 225 }) 226 .collect(); 227 DeliveryReceipt::for_request(&request, receipts) 228 .map_err(|_| SinkFailure::invalid_contract(&request)) 229 } 230 } 231 }) 232 } 233 } 234 235 impl SinkConformanceHarness for MockSink { 236 fn sink(&self) -> &dyn EventSink { 237 self 238 } 239 240 fn expected_status(&self) -> SinkStatus { 241 sink_status() 242 } 243 244 fn target_set(&self) -> radroots_transport::TargetSet { 245 target_set() 246 } 247 248 fn now_unix_ms(&self) -> u64 { 249 NOW_UNIX_MS 250 } 251 252 fn captured_request(&self) -> Option<DeliveryRequest> { 253 self.state.request() 254 } 255 256 fn published(&self) -> bool { 257 self.state.published.load(Ordering::SeqCst) 258 } 259 260 fn cancelled_after_publish(&self) -> bool { 261 self.state.cancelled_after_publish.load(Ordering::SeqCst) 262 } 263 } 264 265 struct ScopeGuard<F: FnOnce()>(Option<F>); 266 267 impl<F: FnOnce()> ScopeGuard<F> { 268 fn new(callback: F) -> Self { 269 Self(Some(callback)) 270 } 271 } 272 273 impl<F: FnOnce()> Drop for ScopeGuard<F> { 274 fn drop(&mut self) { 275 if let Some(callback) = self.0.take() { 276 callback(); 277 } 278 } 279 } 280 281 pub(crate) struct CombinedAdapter { 282 source: MockSource, 283 sink: MockSink, 284 } 285 286 impl CombinedAdapter { 287 pub(crate) fn successful() -> Self { 288 Self { 289 source: MockSource::successful(), 290 sink: MockSink::successful(), 291 } 292 } 293 } 294 295 impl EventSource for CombinedAdapter { 296 fn status(&self) -> BoxFuture<'_, Result<SourceStatus, Error>> { 297 EventSource::status(&self.source) 298 } 299 300 fn fetch(&self, request: FetchRequest) -> BoxFuture<'_, Result<FetchPage, Error>> { 301 self.source.fetch(request) 302 } 303 } 304 305 impl EventSink for CombinedAdapter { 306 fn status(&self) -> BoxFuture<'_, Result<SinkStatus, Error>> { 307 EventSink::status(&self.sink) 308 } 309 310 fn deliver( 311 &self, 312 request: DeliveryRequest, 313 ) -> BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> { 314 self.sink.deliver(request) 315 } 316 } 317 318 impl SourceConformanceHarness for CombinedAdapter { 319 fn source(&self) -> &dyn EventSource { 320 self 321 } 322 323 fn expected_status(&self) -> SourceStatus { 324 source_status() 325 } 326 327 fn target_set(&self) -> radroots_transport::TargetSet { 328 target_set() 329 } 330 331 fn now_unix_ms(&self) -> u64 { 332 NOW_UNIX_MS 333 } 334 335 fn captured_request(&self) -> Option<FetchRequest> { 336 self.source.state.request() 337 } 338 339 fn published(&self) -> bool { 340 self.source.state.published.load(Ordering::SeqCst) 341 } 342 343 fn cancelled_after_publish(&self) -> bool { 344 self.source 345 .state 346 .cancelled_after_publish 347 .load(Ordering::SeqCst) 348 } 349 } 350 351 impl SinkConformanceHarness for CombinedAdapter { 352 fn sink(&self) -> &dyn EventSink { 353 self 354 } 355 356 fn expected_status(&self) -> SinkStatus { 357 sink_status() 358 } 359 360 fn target_set(&self) -> radroots_transport::TargetSet { 361 target_set() 362 } 363 364 fn now_unix_ms(&self) -> u64 { 365 NOW_UNIX_MS 366 } 367 368 fn captured_request(&self) -> Option<DeliveryRequest> { 369 self.sink.state.request() 370 } 371 372 fn published(&self) -> bool { 373 self.sink.state.published.load(Ordering::SeqCst) 374 } 375 376 fn cancelled_after_publish(&self) -> bool { 377 self.sink 378 .state 379 .cancelled_after_publish 380 .load(Ordering::SeqCst) 381 } 382 } 383 384 fn target_set() -> radroots_transport::TargetSet { 385 radroots_transport::TargetSet::new(vec![ 386 radroots_transport::Target::local("local:conformance-a").expect("first target"), 387 radroots_transport::Target::local("local:conformance-b").expect("second target"), 388 ]) 389 .expect("target set") 390 } 391 392 fn source_status() -> SourceStatus { 393 SourceStatus::new( 394 TransportId::LOCAL, 395 true, 396 Maturity::Stable, 397 Availability::Available, 398 SourceCapabilities::FETCH, 399 "mock source ready", 400 ) 401 } 402 403 fn sink_status() -> SinkStatus { 404 SinkStatus::new( 405 TransportId::LOCAL, 406 true, 407 Maturity::Stable, 408 Availability::Available, 409 SinkCapabilities::DELIVER, 410 "mock sink ready", 411 ) 412 }