pull.rs (14277B)
1 use std::{ 2 collections::VecDeque, 3 sync::{Arc, Mutex}, 4 }; 5 6 #[path = "pull/summary.rs"] 7 mod summary; 8 9 use futures_executor::block_on; 10 use radroots_event::{SignedEvent, draft::SignedEventParts}; 11 use radroots_storage::{event::SourceGeneration, memory::MemoryStorage}; 12 use radroots_sync::{ 13 Engine, PullRequest, 14 ingest::RegistryPolicy, 15 policy::{Clock, DeadlinePolicy, Error, IdSource, OperationKind, SyncId, SyncStorage}, 16 pull::{PULL_MAX_PAGES, PullTermination}, 17 }; 18 use radroots_transport::{ 19 Error as TransportError, EventSource, FetchPage, FetchRequest, SourceStatus, Target, TargetSet, 20 TransportId, 21 outcome::{FetchTargetOutcome, FetchTargetState}, 22 source::{EventProvenance, FetchCursor, NextPage, ObservedEvent}, 23 }; 24 25 const EVENT_ID: &str = "762bee187e9e645b81ec26ade05a69b5e8398caf527be8de0d9a45311ed0c7a0"; 26 const PUBKEY: &str = "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df"; 27 const SIGNATURE: &str = "4290da0bb6422986647bc8cd5f63bd52d49f41e7b665d3b47105b8109183e8d596f322c531d4061df53e1d2b70fda12d5d1c14f3720d7a56d9d0a03746af5109"; 28 const CONTENT: &str = "{\"display_name\":\"Moss Street Farm\",\"bot\":false,\"website\":\"https://mossstreet.example\",\"picture\":42}"; 29 30 enum Response { 31 Page { 32 events: Vec<ObservedEvent>, 33 state: FetchTargetState, 34 next: NextPage, 35 }, 36 Failure, 37 } 38 39 #[derive(Clone, Debug, Eq, PartialEq)] 40 struct RequestEvidence { 41 cursor: Option<String>, 42 limit: u16, 43 deadline_unix_ms: u64, 44 kinds: Vec<u32>, 45 } 46 47 struct ScriptedSource { 48 responses: Mutex<VecDeque<Response>>, 49 requests: Mutex<Vec<RequestEvidence>>, 50 } 51 52 impl ScriptedSource { 53 fn new(responses: Vec<Response>) -> Self { 54 Self { 55 responses: Mutex::new(responses.into()), 56 requests: Mutex::new(Vec::new()), 57 } 58 } 59 60 fn requests(&self) -> Vec<RequestEvidence> { 61 self.requests.lock().expect("requests").clone() 62 } 63 } 64 65 impl EventSource for ScriptedSource { 66 fn status(&self) -> radroots_transport::BoxFuture<'_, Result<SourceStatus, TransportError>> { 67 Box::pin(async { unreachable!("pull does not inspect source status") }) 68 } 69 70 fn fetch( 71 &self, 72 request: FetchRequest, 73 ) -> radroots_transport::BoxFuture<'_, Result<FetchPage, TransportError>> { 74 Box::pin(async move { 75 self.requests 76 .lock() 77 .expect("requests") 78 .push(RequestEvidence { 79 cursor: request.cursor().map(|cursor| cursor.as_str().to_owned()), 80 limit: request.bounds().limit(), 81 deadline_unix_ms: request.bounds().deadline_unix_ms(), 82 kinds: request.selector().kinds().to_vec(), 83 }); 84 match self 85 .responses 86 .lock() 87 .expect("responses") 88 .pop_front() 89 .expect("scripted response") 90 { 91 Response::Failure => Err(TransportError::UnsupportedOperation), 92 Response::Page { 93 events, 94 state, 95 next, 96 } => { 97 let target = request.target_set().targets()[0].fingerprint().clone(); 98 FetchPage::for_request( 99 &request, 100 events, 101 vec![FetchTargetOutcome::new(target, state)], 102 next, 103 ) 104 } 105 } 106 }) 107 } 108 } 109 110 struct FixedClock(u64); 111 112 impl Clock for FixedClock { 113 fn now_unix_ms(&self) -> Result<u64, Error> { 114 Ok(self.0) 115 } 116 } 117 118 struct DeadlineClock(Mutex<VecDeque<u64>>); 119 120 impl Clock for DeadlineClock { 121 fn now_unix_ms(&self) -> Result<u64, Error> { 122 self.0 123 .lock() 124 .expect("clock") 125 .pop_front() 126 .ok_or(Error::ClockUnavailable) 127 } 128 } 129 130 struct SequenceIds(Mutex<u8>); 131 132 impl IdSource for SequenceIds { 133 fn next_id(&self, _operation: OperationKind) -> Result<SyncId, Error> { 134 let mut next = self.0.lock().expect("ids"); 135 let value = *next; 136 *next = next.checked_add(1).ok_or(Error::InvalidSyncId)?; 137 SyncId::new([value; 16]) 138 } 139 } 140 141 fn target() -> Target { 142 Target::new(TransportId::NOSTR, "wss://relay.example").expect("target") 143 } 144 145 fn targets() -> TargetSet { 146 TargetSet::new(vec![target()]).expect("target set") 147 } 148 149 fn signed_event(signature: &str) -> SignedEvent { 150 let raw_json = format!( 151 "{{\"id\":\"{EVENT_ID}\",\"pubkey\":\"{PUBKEY}\",\"created_at\":1800000100,\"kind\":0,\"tags\":[],\"content\":{content:?},\"sig\":\"{signature}\"}}", 152 content = CONTENT, 153 ); 154 SignedEvent::new(SignedEventParts { 155 id: EVENT_ID.to_owned(), 156 pubkey: PUBKEY.to_owned(), 157 created_at: 1_800_000_100, 158 kind: 0, 159 tags: vec![], 160 content: CONTENT.to_owned(), 161 sig: signature.to_owned(), 162 raw_json, 163 }) 164 .expect("ID-valid event") 165 } 166 167 fn observed(signature: &str, observed_at: u64) -> ObservedEvent { 168 let target = target(); 169 ObservedEvent::new( 170 signed_event(signature), 171 EventProvenance::new( 172 TransportId::NOSTR, 173 target.fingerprint().clone(), 174 observed_at, 175 ) 176 .expect("provenance"), 177 ) 178 } 179 180 fn engine(source: Arc<dyn EventSource>, clock: Arc<dyn Clock>, timeout_ms: u64) -> Engine { 181 let storage: Arc<dyn SyncStorage> = Arc::new(MemoryStorage::new( 182 SourceGeneration::new([8; 32]).expect("generation"), 183 )); 184 Engine::builder( 185 storage, 186 clock, 187 Arc::new(SequenceIds(Mutex::new(1))), 188 DeadlinePolicy::new(timeout_ms, 10, 10).expect("deadlines"), 189 ) 190 .source(source) 191 .build() 192 .expect("engine") 193 } 194 195 #[test] 196 fn single_and_multiple_pages_propagate_cursor_deadline_and_ingest_results() { 197 let single_source = Arc::new(ScriptedSource::new(vec![Response::Page { 198 events: vec![observed(SIGNATURE, 1)], 199 state: FetchTargetState::Complete, 200 next: NextPage::Complete, 201 }])); 202 let single = engine(single_source.clone(), Arc::new(FixedClock(100)), 50); 203 let request = PullRequest::new(targets(), 20, 1).expect("request"); 204 assert_eq!(request.targets().len(), 1); 205 assert_eq!(request.page_limit(), 20); 206 assert_eq!(request.max_pages(), 1); 207 assert!(request.cursor().is_none()); 208 assert!(request.selector().kinds().is_empty()); 209 let receipt = block_on(single.pull(request, &RegistryPolicy::visible())).expect("pull"); 210 assert_eq!(receipt.termination(), PullTermination::Complete); 211 assert_eq!(receipt.pages_fetched(), 1); 212 assert_eq!(receipt.events_observed(), 1); 213 assert_ne!(receipt.sync_id().as_bytes(), &[0; 16]); 214 assert_eq!(receipt.deadline_unix_ms(), 150); 215 assert!(receipt.ingest_outcomes()[0].is_ok()); 216 assert_eq!(single_source.requests()[0].deadline_unix_ms, 150); 217 218 let next = FetchCursor::parse("page-2").expect("cursor"); 219 let multiple_source = Arc::new(ScriptedSource::new(vec![ 220 Response::Page { 221 events: vec![], 222 state: FetchTargetState::Partial, 223 next: NextPage::Cursor(next.clone()), 224 }, 225 Response::Page { 226 events: vec![observed(SIGNATURE, 2)], 227 state: FetchTargetState::Complete, 228 next: NextPage::Complete, 229 }, 230 ])); 231 let multiple = engine(multiple_source.clone(), Arc::new(FixedClock(200)), 50); 232 let receipt = block_on( 233 multiple.pull( 234 PullRequest::new(targets(), 10, 2) 235 .expect("request") 236 .with_cursor(FetchCursor::parse("starting").expect("initial cursor")), 237 &RegistryPolicy::visible(), 238 ), 239 ) 240 .expect("pull"); 241 assert_eq!(receipt.pages_fetched(), 2); 242 assert_eq!(receipt.termination(), PullTermination::Complete); 243 assert_eq!( 244 receipt.target_outcomes()[0].state(), 245 FetchTargetState::Complete 246 ); 247 let requests = multiple_source.requests(); 248 assert_eq!(requests[0].cursor.as_deref(), Some("starting")); 249 assert_eq!(requests[1].cursor.as_deref(), Some(next.as_str())); 250 assert_eq!(requests[0].deadline_unix_ms, requests[1].deadline_unix_ms); 251 } 252 253 #[cfg(feature = "serde")] 254 #[test] 255 fn later_complete_page_retains_earlier_incomplete_evidence() { 256 let source = Arc::new(ScriptedSource::new(vec![ 257 Response::Page { 258 events: vec![], 259 state: FetchTargetState::Partial, 260 next: NextPage::Cursor(FetchCursor::parse("next").expect("cursor")), 261 }, 262 Response::Page { 263 events: vec![], 264 state: FetchTargetState::Complete, 265 next: NextPage::Complete, 266 }, 267 ])); 268 let pull = engine(source, Arc::new(FixedClock(100)), 50); 269 let receipt = block_on(pull.pull( 270 PullRequest::new(targets(), 10, 2).expect("request"), 271 &RegistryPolicy::visible(), 272 )) 273 .expect("pull"); 274 assert_eq!(receipt.termination(), PullTermination::Complete); 275 assert_eq!( 276 receipt.target_outcomes()[0].state(), 277 FetchTargetState::Complete 278 ); 279 let wire = serde_json::to_value(receipt).expect("receipt JSON"); 280 assert_eq!(wire["target_summaries"][0]["incomplete_pages"], 1); 281 assert_eq!(wire["target_summaries"][0]["last_incomplete"], "partial"); 282 } 283 284 #[test] 285 fn pull_propagates_the_exact_selector_to_every_page() { 286 let next = FetchCursor::parse("page-2").expect("cursor"); 287 let source = Arc::new(ScriptedSource::new(vec![ 288 Response::Page { 289 events: vec![], 290 state: FetchTargetState::Partial, 291 next: NextPage::Cursor(next), 292 }, 293 Response::Page { 294 events: vec![], 295 state: FetchTargetState::Complete, 296 next: NextPage::Complete, 297 }, 298 ])); 299 let pull = engine(source.clone(), Arc::new(FixedClock(100)), 50); 300 let selector = radroots_transport::source::FetchSelector::all() 301 .with_kinds(vec![0, 1, 5, 1111, 30402, 31922, 31923]) 302 .expect("selector"); 303 let request = PullRequest::new(targets(), 20, 2) 304 .expect("request") 305 .with_selector(selector.clone()); 306 assert_eq!(request.selector(), &selector); 307 308 let receipt = block_on(pull.pull(request, &RegistryPolicy::visible())).expect("pull"); 309 assert_eq!(receipt.pages_fetched(), 2); 310 assert_eq!( 311 source 312 .requests() 313 .into_iter() 314 .map(|request| request.kinds) 315 .collect::<Vec<_>>(), 316 vec![selector.kinds().to_vec(), selector.kinds().to_vec()] 317 ); 318 } 319 320 #[test] 321 fn source_failure_and_cancelled_page_return_resumable_partial_receipts() { 322 let cursor = FetchCursor::parse("resume").expect("cursor"); 323 let source = Arc::new(ScriptedSource::new(vec![ 324 Response::Page { 325 events: vec![observed(SIGNATURE, 1)], 326 state: FetchTargetState::Partial, 327 next: NextPage::Cursor(cursor.clone()), 328 }, 329 Response::Failure, 330 ])); 331 let pull = engine(source, Arc::new(FixedClock(100)), 50); 332 let receipt = block_on(pull.pull( 333 PullRequest::new(targets(), 10, 3).expect("request"), 334 &RegistryPolicy::visible(), 335 )) 336 .expect("partial receipt"); 337 assert_eq!(receipt.termination(), PullTermination::SourceFailed); 338 assert_eq!(receipt.pages_fetched(), 1); 339 assert_eq!( 340 receipt.resume_from().map(FetchCursor::as_str), 341 Some("resume") 342 ); 343 assert!(receipt.ingest_outcomes()[0].is_ok()); 344 345 let cancelled_from = FetchCursor::parse("cancelled-at").expect("cursor"); 346 let source = Arc::new(ScriptedSource::new(vec![Response::Page { 347 events: vec![], 348 state: FetchTargetState::Cancelled, 349 next: NextPage::Cancelled { 350 resume_from: Some(cancelled_from.clone()), 351 }, 352 }])); 353 let pull = engine(source, Arc::new(FixedClock(100)), 50); 354 let receipt = block_on(pull.pull( 355 PullRequest::new(targets(), 10, 3).expect("request"), 356 &RegistryPolicy::visible(), 357 )) 358 .expect("cancelled receipt"); 359 assert_eq!(receipt.termination(), PullTermination::Cancelled); 360 assert_eq!( 361 receipt.resume_from().map(FetchCursor::as_str), 362 Some(cancelled_from.as_str()) 363 ); 364 } 365 366 #[test] 367 fn page_and_deadline_limits_stop_without_hidden_fetches() { 368 assert_eq!( 369 PullRequest::new(targets(), 0, 1), 370 Err(Error::InvalidPullRequest) 371 ); 372 assert_eq!( 373 PullRequest::new( 374 targets(), 375 radroots_transport::source::FETCH_PAGE_MAX_EVENTS + 1, 376 1 377 ), 378 Err(Error::InvalidPullRequest) 379 ); 380 assert_eq!( 381 PullRequest::new(targets(), 1, 0), 382 Err(Error::InvalidPullRequest) 383 ); 384 assert_eq!( 385 PullRequest::new(targets(), 1, PULL_MAX_PAGES + 1), 386 Err(Error::InvalidPullRequest) 387 ); 388 389 let cursor = FetchCursor::parse("more").expect("cursor"); 390 let page_limited_source = Arc::new(ScriptedSource::new(vec![Response::Page { 391 events: vec![], 392 state: FetchTargetState::Partial, 393 next: NextPage::Cursor(cursor.clone()), 394 }])); 395 let pull = engine(page_limited_source.clone(), Arc::new(FixedClock(100)), 10); 396 let receipt = block_on(pull.pull( 397 PullRequest::new(targets(), 1, 1).expect("request"), 398 &RegistryPolicy::visible(), 399 )) 400 .expect("limited receipt"); 401 assert_eq!(receipt.termination(), PullTermination::PageLimit); 402 assert_eq!(page_limited_source.requests().len(), 1); 403 404 let deadline_source = Arc::new(ScriptedSource::new(vec![Response::Page { 405 events: vec![], 406 state: FetchTargetState::Partial, 407 next: NextPage::Cursor(cursor), 408 }])); 409 let clock = Arc::new(DeadlineClock(Mutex::new(VecDeque::from([100, 110])))); 410 let pull = engine(deadline_source.clone(), clock, 10); 411 let receipt = block_on(pull.pull( 412 PullRequest::new(targets(), 1, 2).expect("request"), 413 &RegistryPolicy::visible(), 414 )) 415 .expect("deadline receipt"); 416 assert_eq!(receipt.termination(), PullTermination::Deadline); 417 assert_eq!(receipt.deadline_unix_ms(), 110); 418 assert_eq!(deadline_source.requests().len(), 1); 419 }