summary.rs (13533B)
1 use super::*; 2 use radroots_sync::PullReceipt; 3 use radroots_transport::target::TARGET_SET_MAX_ITEMS; 4 5 struct Source { 6 pages: Mutex<VecDeque<Option<Vec<Option<FetchTargetState>>>>>, 7 finish: NextPage, 8 } 9 10 impl EventSource for Source { 11 fn status(&self) -> radroots_transport::BoxFuture<'_, Result<SourceStatus, TransportError>> { 12 Box::pin(async { unreachable!("explicit pull only") }) 13 } 14 15 fn fetch( 16 &self, 17 request: FetchRequest, 18 ) -> radroots_transport::BoxFuture<'_, Result<FetchPage, TransportError>> { 19 Box::pin(async move { 20 let mut pages = self.pages.lock().expect("pages"); 21 let states = pages 22 .pop_front() 23 .expect("bounded scripted fetch") 24 .ok_or(TransportError::UnsupportedOperation)?; 25 assert_eq!(states.len(), request.target_set().len()); 26 let outcomes = request 27 .target_set() 28 .targets() 29 .iter() 30 .zip(states) 31 .filter_map(|(target, state)| { 32 state.map(|state| FetchTargetOutcome::new(target.fingerprint().clone(), state)) 33 }) 34 .collect(); 35 let next = if pages.is_empty() { 36 self.finish.clone() 37 } else { 38 NextPage::Cursor( 39 FetchCursor::parse(format!("remaining-{}", pages.len())).expect("cursor"), 40 ) 41 }; 42 FetchPage::for_request(&request, vec![], outcomes, next) 43 }) 44 } 45 } 46 47 fn target_set(count: usize) -> TargetSet { 48 TargetSet::new( 49 (0..count) 50 .rev() 51 .map(|index| { 52 Target::nostr_relay(format!("wss://relay-{index}.example")).expect("target") 53 }) 54 .collect(), 55 ) 56 .expect("targets") 57 } 58 59 fn pull( 60 pages: Vec<Option<Vec<Option<FetchTargetState>>>>, 61 targets: TargetSet, 62 max_pages: u16, 63 finish: NextPage, 64 clock: Arc<dyn Clock>, 65 ) -> PullReceipt { 66 let source = Arc::new(Source { 67 pages: Mutex::new(pages.into()), 68 finish, 69 }); 70 block_on(engine(source, clock, 50).pull( 71 PullRequest::new(targets, 1, max_pages).expect("request"), 72 &RegistryPolicy::visible(), 73 )) 74 .expect("receipt") 75 } 76 77 fn ordered(states: &[Option<FetchTargetState>]) -> PullReceipt { 78 pull( 79 states.iter().map(|state| Some(vec![*state])).collect(), 80 target_set(1), 81 states.len() as u16, 82 NextPage::Complete, 83 Arc::new(FixedClock(100)), 84 ) 85 } 86 87 #[test] 88 fn every_incomplete_state_survives_later_success_and_missing_outcomes() { 89 use FetchTargetState::*; 90 for state in [ 91 Partial, 92 Unavailable, 93 FailedRetryable, 94 FailedTerminal, 95 Cancelled, 96 ] { 97 for states in [ 98 vec![Some(state), Some(Complete)], 99 vec![Some(Complete), Some(state)], 100 vec![Some(state), None, Some(Complete)], 101 vec![Some(state), Some(Complete), None], 102 ] { 103 let receipt = ordered(&states); 104 let summaries = receipt.target_summaries().expect("measured"); 105 assert_eq!(summaries.len(), 1); 106 let summary = &summaries[0]; 107 assert_eq!(summary.target(), target_set(1).targets()[0].fingerprint()); 108 assert_eq!(summary.pages_observed(), states.len() as u16); 109 assert_eq!(summary.incomplete_pages(), 1); 110 assert_eq!( 111 summary.missing_outcome_pages(), 112 states.iter().filter(|value| value.is_none()).count() as u16 113 ); 114 assert_eq!(summary.last_incomplete(), Some(state)); 115 assert!(!summary.all_pages_complete()); 116 assert_eq!( 117 receipt.target_outcomes()[0].state(), 118 states.iter().rev().flatten().next().copied().unwrap() 119 ); 120 #[cfg(feature = "serde")] 121 assert_eq!( 122 serde_json::from_str::<PullReceipt>(&serde_json::to_string(&receipt).unwrap()) 123 .unwrap(), 124 receipt 125 ); 126 } 127 } 128 let receipt = ordered(&[Some(Partial), Some(FailedTerminal), Some(Complete)]); 129 let summary = &receipt.target_summaries().unwrap()[0]; 130 assert_eq!(summary.incomplete_pages(), 2); 131 assert_eq!(summary.last_incomplete(), Some(FailedTerminal)); 132 } 133 134 #[test] 135 fn omitted_outcomes_never_become_positive_evidence() { 136 use FetchTargetState::Complete; 137 for states in [ 138 vec![None], 139 vec![None, Some(Complete)], 140 vec![Some(Complete), None], 141 ] { 142 let receipt = ordered(&states); 143 assert_eq!(receipt.termination(), PullTermination::Complete); 144 let summary = &receipt.target_summaries().unwrap()[0]; 145 assert_eq!(summary.missing_outcome_pages(), 1); 146 assert_eq!(summary.incomplete_pages(), 0); 147 assert_eq!(summary.last_incomplete(), None); 148 assert!(!summary.all_pages_complete()); 149 #[cfg(feature = "serde")] 150 assert_eq!( 151 serde_json::from_value::<PullReceipt>(serde_json::to_value(&receipt).unwrap()).unwrap(), 152 receipt 153 ); 154 } 155 } 156 157 #[test] 158 fn maximum_inventory_preserves_request_order_and_counts_without_page_history() { 159 let targets = target_set(TARGET_SET_MAX_ITEMS); 160 let receipt = pull( 161 vec![ 162 Some(vec![Some(FetchTargetState::Complete); TARGET_SET_MAX_ITEMS]); 163 usize::from(PULL_MAX_PAGES) 164 ], 165 targets.clone(), 166 PULL_MAX_PAGES, 167 NextPage::Complete, 168 Arc::new(FixedClock(100)), 169 ); 170 assert_eq!(receipt.pages_fetched(), PULL_MAX_PAGES); 171 let summaries = receipt.target_summaries().unwrap(); 172 assert_eq!(summaries.len(), TARGET_SET_MAX_ITEMS); 173 for (summary, target) in summaries.iter().zip(targets.targets()) { 174 assert_eq!(summary.target(), target.fingerprint()); 175 assert_eq!(summary.pages_observed(), PULL_MAX_PAGES); 176 assert!(summary.all_pages_complete()); 177 } 178 #[cfg(feature = "serde")] 179 { 180 let encoded = serde_json::to_string(&receipt).unwrap(); 181 assert!( 182 encoded.len() < 32_768, 183 "bounded evidence excludes per-page history" 184 ); 185 assert_eq!( 186 serde_json::from_str::<PullReceipt>(&encoded).unwrap(), 187 receipt 188 ); 189 } 190 } 191 192 #[test] 193 fn multiple_targets_retain_independent_evidence() { 194 use FetchTargetState::*; 195 let receipt = pull( 196 vec![ 197 Some(vec![Some(Partial), Some(Complete), None]), 198 Some(vec![Some(Complete), Some(FailedRetryable), Some(Complete)]), 199 ], 200 target_set(3), 201 2, 202 NextPage::Complete, 203 Arc::new(FixedClock(100)), 204 ); 205 let summaries = receipt.target_summaries().unwrap(); 206 assert_eq!( 207 summaries 208 .iter() 209 .map(|s| s.incomplete_pages()) 210 .collect::<Vec<_>>(), 211 [1, 1, 0] 212 ); 213 assert_eq!( 214 summaries 215 .iter() 216 .map(|s| s.missing_outcome_pages()) 217 .collect::<Vec<_>>(), 218 [0, 0, 1] 219 ); 220 assert_eq!( 221 summaries 222 .iter() 223 .map(|s| s.last_incomplete()) 224 .collect::<Vec<_>>(), 225 [Some(Partial), Some(FailedRetryable), None] 226 ); 227 assert!(summaries.iter().all(|s| !s.all_pages_complete())); 228 } 229 230 #[test] 231 fn termination_and_zero_returned_pages_remain_explicit() { 232 use FetchTargetState::*; 233 let failed = pull( 234 vec![None], 235 target_set(1), 236 1, 237 NextPage::Complete, 238 Arc::new(FixedClock(100)), 239 ); 240 assert_eq!(failed.termination(), PullTermination::SourceFailed); 241 assert_eq!(failed.target_summaries().unwrap()[0].pages_observed(), 0); 242 assert!(!failed.target_summaries().unwrap()[0].all_pages_complete()); 243 let later_failure = pull( 244 vec![Some(vec![Some(Complete)]), None], 245 target_set(1), 246 2, 247 NextPage::Complete, 248 Arc::new(FixedClock(100)), 249 ); 250 assert_eq!(later_failure.termination(), PullTermination::SourceFailed); 251 assert!(later_failure.target_summaries().unwrap()[0].all_pages_complete()); 252 assert_eq!( 253 later_failure.target_summaries().unwrap()[0].pages_observed(), 254 1 255 ); 256 let limited = pull( 257 vec![Some(vec![Some(Partial)]); 2], 258 target_set(1), 259 1, 260 NextPage::Complete, 261 Arc::new(FixedClock(100)), 262 ); 263 assert_eq!(limited.termination(), PullTermination::PageLimit); 264 let deadline = pull( 265 vec![Some(vec![Some(Partial)]); 2], 266 target_set(1), 267 2, 268 NextPage::Complete, 269 Arc::new(DeadlineClock(Mutex::new(VecDeque::from([100, 150])))), 270 ); 271 assert_eq!(deadline.termination(), PullTermination::Deadline); 272 let cancelled = pull( 273 vec![Some(vec![Some(Cancelled)])], 274 target_set(1), 275 1, 276 NextPage::Cancelled { resume_from: None }, 277 Arc::new(FixedClock(100)), 278 ); 279 assert_eq!(cancelled.termination(), PullTermination::Cancelled); 280 for receipt in [failed, later_failure, limited, deadline, cancelled] { 281 #[cfg(feature = "serde")] 282 assert_eq!( 283 serde_json::from_value::<PullReceipt>(serde_json::to_value(&receipt).unwrap()).unwrap(), 284 receipt 285 ); 286 assert_eq!( 287 receipt.target_summaries().unwrap()[0].pages_observed(), 288 receipt.pages_fetched() 289 ); 290 } 291 } 292 293 #[cfg(feature = "serde")] 294 #[test] 295 fn legacy_receipts_remain_unknown_and_new_receipts_reject_inconsistent_evidence() { 296 use serde_json::json; 297 let receipt = ordered(&[Some(FetchTargetState::Complete)]); 298 let original = serde_json::to_value(&receipt).unwrap(); 299 let mut legacy = original.clone(); 300 legacy.as_object_mut().unwrap().remove("target_summaries"); 301 assert!( 302 serde_json::from_value::<PullReceipt>(legacy.clone()) 303 .unwrap() 304 .target_summaries() 305 .is_none() 306 ); 307 legacy["target_summaries"] = json!(null); 308 assert!( 309 serde_json::from_value::<PullReceipt>(legacy) 310 .unwrap() 311 .target_summaries() 312 .is_none() 313 ); 314 for replacement in [json!([]), json!({}), json!(42)] { 315 let mut changed = original.clone(); 316 changed["target_summaries"] = replacement; 317 assert!(serde_json::from_value::<PullReceipt>(changed).is_err()); 318 } 319 let mut duplicate = original.clone(); 320 duplicate["target_summaries"] 321 .as_array_mut() 322 .unwrap() 323 .push(original["target_summaries"][0].clone()); 324 assert!(serde_json::from_value::<PullReceipt>(duplicate).is_err()); 325 for (field, value) in [ 326 ("pages_observed", json!(1001)), 327 ("pages_observed", json!(2)), 328 ("incomplete_pages", json!(1)), 329 ("incomplete_pages", json!(65535)), 330 ("missing_outcome_pages", json!(2)), 331 ("missing_outcome_pages", json!(1)), 332 ("last_incomplete", json!("partial")), 333 ("target", json!("bad")), 334 ("extra", json!(true)), 335 ] { 336 let mut changed = original.clone(); 337 changed["target_summaries"][0][field] = value; 338 assert!( 339 serde_json::from_value::<PullReceipt>(changed).is_err(), 340 "{field}" 341 ); 342 } 343 let mut complete_as_incomplete = original.clone(); 344 complete_as_incomplete["target_summaries"][0]["incomplete_pages"] = json!(1); 345 complete_as_incomplete["target_summaries"][0]["last_incomplete"] = json!("complete"); 346 assert!(serde_json::from_value::<PullReceipt>(complete_as_incomplete).is_err()); 347 for change in 0..5 { 348 let mut changed = original.clone(); 349 match change { 350 0 => { 351 changed["target_outcomes"] = json!([]); 352 } 353 1 => { 354 changed["target_outcomes"] 355 .as_array_mut() 356 .unwrap() 357 .push(original["target_outcomes"][0].clone()); 358 } 359 2 => { 360 changed["target_outcomes"][0]["target"] = 361 json!(target_set(2).targets()[0].fingerprint()); 362 } 363 3 => { 364 changed["target_outcomes"][0]["state"] = json!("partial"); 365 } 366 _ => { 367 changed["target_summaries"][0]["incomplete_pages"] = json!(1); 368 changed["target_summaries"][0]["last_incomplete"] = json!("partial"); 369 } 370 } 371 assert!( 372 serde_json::from_value::<PullReceipt>(changed).is_err(), 373 "case {change}" 374 ); 375 } 376 let mut all_missing = serde_json::to_value(ordered(&[None])).unwrap(); 377 all_missing["target_summaries"][0]["missing_outcome_pages"] = json!(0); 378 assert!(serde_json::from_value::<PullReceipt>(all_missing).is_err()); 379 let maximum = pull( 380 vec![Some(vec![ 381 Some(FetchTargetState::Complete); 382 TARGET_SET_MAX_ITEMS 383 ])], 384 target_set(TARGET_SET_MAX_ITEMS), 385 1, 386 NextPage::Complete, 387 Arc::new(FixedClock(100)), 388 ); 389 let mut too_many = serde_json::to_value(maximum).unwrap(); 390 too_many["target_summaries"] 391 .as_array_mut() 392 .unwrap() 393 .push(json!({"unparsed_extra": [1, 2, 3]})); 394 let error = serde_json::from_str::<PullReceipt>(&serde_json::to_string(&too_many).unwrap()) 395 .unwrap_err(); 396 assert!(error.to_string().contains("too many pull target summaries")); 397 }