lib

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

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 }