lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

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 }