lib

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

projection.rs (16670B)


      1 use std::sync::{
      2     Arc, Mutex,
      3     atomic::{AtomicU64, Ordering},
      4 };
      5 
      6 use futures_executor::block_on;
      7 use radroots_event::{
      8     SignedEvent,
      9     admission::{AdmissionPolicy, RawEvent, SignatureVerifier, VisibilityPolicy},
     10     draft::SignedEventParts,
     11     envelope::EventEnvelope,
     12     wire::compute_canonical_nip01_event_id,
     13 };
     14 use radroots_storage::{
     15     EventStore, ProjectionStore,
     16     event::{EventAdmission, SourceGeneration, StoredVisibleEvent},
     17     memory::MemoryStorage,
     18     projection::{
     19         InvalidationReason, ProjectionGeneration, ProjectionHealth, ProjectionId,
     20         ProjectionInvalidation, RawSourceDigest, RebuildFailure, RebuildTicketId,
     21     },
     22 };
     23 use radroots_sync::{
     24     Engine,
     25     policy::{Clock, DeadlinePolicy, Error, IdSource, OperationKind, SyncId, SyncStorage},
     26     projection::{Reducer, ReducerError, RefreshKind, RefreshRequest, RefreshState},
     27 };
     28 use radroots_transport::{
     29     Error as TransportError, EventSource, FetchPage, FetchRequest, SourceStatus, Target,
     30     TransportId,
     31     source::{EventProvenance, ObservedEvent},
     32 };
     33 
     34 const PUBKEY: &str = "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df";
     35 const CONTENT: &str = "{\"display_name\":\"Moss Street Farm\",\"bot\":false}";
     36 
     37 struct MockSource;
     38 
     39 impl EventSource for MockSource {
     40     fn status(&self) -> radroots_transport::BoxFuture<'_, Result<SourceStatus, TransportError>> {
     41         Box::pin(async { unreachable!("projection refresh does not inspect source") })
     42     }
     43 
     44     fn fetch(
     45         &self,
     46         _request: FetchRequest,
     47     ) -> radroots_transport::BoxFuture<'_, Result<FetchPage, TransportError>> {
     48         Box::pin(async { unreachable!("projection refresh does not fetch") })
     49     }
     50 }
     51 
     52 struct TestClock(AtomicU64);
     53 
     54 impl Clock for TestClock {
     55     fn now_unix_ms(&self) -> Result<u64, Error> {
     56         Ok(self.0.fetch_add(1, Ordering::Relaxed))
     57     }
     58 }
     59 
     60 struct TestIds(Mutex<u8>);
     61 
     62 impl IdSource for TestIds {
     63     fn next_id(&self, operation: OperationKind) -> Result<SyncId, Error> {
     64         assert_eq!(operation, OperationKind::Projection);
     65         let mut value = self.0.lock().expect("ids");
     66         let current = *value;
     67         *value += 1;
     68         SyncId::new([current; 16])
     69     }
     70 }
     71 
     72 struct Allow;
     73 
     74 impl SignatureVerifier for Allow {
     75     fn verify_signature(&self, _event: &EventEnvelope) -> Result<(), radroots_event::Error> {
     76         Ok(())
     77     }
     78 }
     79 
     80 impl AdmissionPolicy for Allow {
     81     type Error = core::convert::Infallible;
     82     fn policy_id(&self) -> &'static str {
     83         "test.projection.admission.v1"
     84     }
     85     fn admit(
     86         &self,
     87         _event: &radroots_event::admission::ContractValidatedEvent,
     88     ) -> Result<(), Self::Error> {
     89         Ok(())
     90     }
     91 }
     92 
     93 impl VisibilityPolicy for Allow {
     94     type Error = core::convert::Infallible;
     95     fn policy_id(&self) -> &'static str {
     96         "test.projection.visibility.v1"
     97     }
     98     fn make_visible(
     99         &self,
    100         _event: &radroots_event::admission::AdmittedEvent,
    101     ) -> Result<(), Self::Error> {
    102         Ok(())
    103     }
    104 }
    105 
    106 struct CountingReducer {
    107     projection_id: ProjectionId,
    108     generation: ProjectionGeneration,
    109     fail: bool,
    110     regress: bool,
    111 }
    112 
    113 impl Reducer for CountingReducer {
    114     fn projection_id(&self) -> &ProjectionId {
    115         &self.projection_id
    116     }
    117     fn generation(&self) -> ProjectionGeneration {
    118         self.generation
    119     }
    120     fn begin_rebuild(
    121         &self,
    122         _ticket_id: RebuildTicketId,
    123         _source_generation: SourceGeneration,
    124         _source_digest: RawSourceDigest,
    125     ) -> Result<(), ReducerError> {
    126         if self.fail { Err(ReducerError) } else { Ok(()) }
    127     }
    128     fn reduce(
    129         &self,
    130         events: &[StoredVisibleEvent],
    131         prior_projected_rows: u64,
    132         _rebuild_ticket: Option<RebuildTicketId>,
    133     ) -> Result<u64, ReducerError> {
    134         if self.fail {
    135             return Err(ReducerError);
    136         }
    137         if self.regress {
    138             return Ok(prior_projected_rows.saturating_sub(1));
    139         }
    140         prior_projected_rows
    141             .checked_add(u64::try_from(events.len()).expect("event count"))
    142             .ok_or(ReducerError)
    143     }
    144     fn abort_rebuild(
    145         &self,
    146         _ticket_id: RebuildTicketId,
    147         _failure: RebuildFailure,
    148     ) -> Result<(), ReducerError> {
    149         Ok(())
    150     }
    151 }
    152 
    153 fn setup() -> (Engine, Arc<MemoryStorage>, ProjectionId) {
    154     let storage = Arc::new(MemoryStorage::new(
    155         SourceGeneration::new([9; 32]).expect("generation"),
    156     ));
    157     let storage_capability: Arc<dyn SyncStorage> = storage.clone();
    158     let engine = Engine::builder(
    159         storage_capability,
    160         Arc::new(TestClock(AtomicU64::new(1_000))),
    161         Arc::new(TestIds(Mutex::new(1))),
    162         DeadlinePolicy::new(100, 100, 100).expect("deadlines"),
    163     )
    164     .source(Arc::new(MockSource))
    165     .build()
    166     .expect("engine");
    167     (
    168         engine,
    169         storage,
    170         ProjectionId::parse("test.projection").expect("projection id"),
    171     )
    172 }
    173 
    174 fn signed_event(created_at: u64) -> SignedEvent {
    175     let tags: Vec<Vec<String>> = vec![];
    176     let id = compute_canonical_nip01_event_id(PUBKEY, created_at, 1, &tags, CONTENT)
    177         .expect("event id")
    178         .to_hex();
    179     let signature = "42".repeat(64);
    180     let raw_json = format!(
    181         "{{\"id\":\"{id}\",\"pubkey\":\"{PUBKEY}\",\"created_at\":{created_at},\"kind\":1,\"tags\":[],\"content\":{content:?},\"sig\":\"{signature}\"}}",
    182         content = CONTENT,
    183     );
    184     SignedEvent::new(SignedEventParts {
    185         id,
    186         pubkey: PUBKEY.to_owned(),
    187         created_at,
    188         kind: 1,
    189         tags,
    190         content: CONTENT.to_owned(),
    191         sig: signature,
    192         raw_json,
    193     })
    194     .expect("signed event")
    195 }
    196 
    197 fn seed(storage: &MemoryStorage, count: u64) {
    198     for offset in 0..count {
    199         let event = signed_event(1_800_000_100 + offset);
    200         let visible = RawEvent::new(event.envelope().clone())
    201             .verify_id()
    202             .expect("id")
    203             .verify_signature(&Allow)
    204             .expect("signature")
    205             .validate_contract()
    206             .expect("contract")
    207             .admit_with(&Allow)
    208             .expect("admission")
    209             .make_visible_with(&Allow)
    210             .expect("visibility");
    211         let target = Target::new(TransportId::NOSTR, "wss://relay.example").expect("target");
    212         let provenance = EventProvenance::new(
    213             TransportId::NOSTR,
    214             target.fingerprint().clone(),
    215             1_900_000_000_000 + offset,
    216         )
    217         .expect("provenance");
    218         block_on(
    219             storage.admit(
    220                 EventAdmission::visible(ObservedEvent::new(event, provenance), visible)
    221                     .expect("visible admission"),
    222             ),
    223         )
    224         .expect("seed event");
    225     }
    226 }
    227 
    228 fn reducer(id: &ProjectionId, generation: u8, fail: bool) -> CountingReducer {
    229     CountingReducer {
    230         projection_id: id.clone(),
    231         generation: ProjectionGeneration::new([generation; 32]).expect("generation"),
    232         fail,
    233         regress: false,
    234     }
    235 }
    236 
    237 #[test]
    238 fn incremental_refresh_checkpoints_visible_events() {
    239     let (engine, storage, id) = setup();
    240     seed(&storage, 1);
    241     let reducer = reducer(&id, 1, false);
    242     let request = RefreshRequest::new(id.clone(), reducer.generation(), 10, 1).expect("request");
    243     assert_eq!(request.projection_id(), &id);
    244     assert_eq!(request.generation(), reducer.generation());
    245     assert_eq!(request.batch_limit(), 10);
    246     assert_eq!(request.max_batches(), 1);
    247     let receipt = block_on(engine.refresh_projection(request, &reducer)).expect("refresh");
    248     assert_eq!(receipt.kind(), RefreshKind::Incremental);
    249     assert_eq!(receipt.state(), RefreshState::Complete);
    250     assert_eq!(receipt.events_reduced(), 1);
    251     assert_eq!(receipt.batches(), 1);
    252     assert!(receipt.rebuild_ticket().is_none());
    253     assert_eq!(
    254         receipt.checkpoint().expect("checkpoint").projected_rows(),
    255         1
    256     );
    257     assert_eq!(
    258         block_on(ProjectionStore::status(&*storage, id))
    259             .expect("status")
    260             .expect("projection")
    261             .health(),
    262         ProjectionHealth::Ready
    263     );
    264 }
    265 
    266 #[test]
    267 fn generation_change_rebuilds_and_reducer_failure_is_durable() {
    268     let (engine, storage, id) = setup();
    269     seed(&storage, 1);
    270     let first = reducer(&id, 1, false);
    271     block_on(engine.refresh_projection(
    272         RefreshRequest::new(id.clone(), first.generation(), 10, 1).expect("request"),
    273         &first,
    274     ))
    275     .expect("initial refresh");
    276 
    277     let replacement = reducer(&id, 2, false);
    278     let rebuilt = block_on(engine.refresh_projection(
    279         RefreshRequest::new(id.clone(), replacement.generation(), 10, 1).expect("request"),
    280         &replacement,
    281     ))
    282     .expect("rebuild");
    283     assert_eq!(rebuilt.kind(), RefreshKind::Rebuild);
    284     assert_eq!(rebuilt.state(), RefreshState::Complete);
    285 
    286     let failing = reducer(&id, 3, true);
    287     let failed = block_on(engine.refresh_projection(
    288         RefreshRequest::new(id.clone(), failing.generation(), 10, 1).expect("request"),
    289         &failing,
    290     ))
    291     .expect("normalized failure");
    292     assert_eq!(failed.state(), RefreshState::Failed);
    293     assert_eq!(
    294         block_on(ProjectionStore::status(&*storage, id))
    295             .expect("status")
    296             .expect("projection")
    297             .health(),
    298         ProjectionHealth::Ready
    299     );
    300     let retried_failure = block_on(
    301         engine.refresh_projection(
    302             RefreshRequest::new(failing.projection_id().clone(), failing.generation(), 10, 1)
    303                 .expect("failed generation request"),
    304             &failing,
    305         ),
    306     )
    307     .expect("retry failed generation");
    308     assert_eq!(retried_failure.state(), RefreshState::Failed);
    309 }
    310 
    311 #[test]
    312 fn obsolete_generation_request_fails_closed_after_invalidation() {
    313     let (engine, storage, id) = setup();
    314     seed(&storage, 1);
    315     let active = reducer(&id, 1, false);
    316     block_on(engine.refresh_projection(
    317         RefreshRequest::new(id.clone(), active.generation(), 10, 1).expect("request"),
    318         &active,
    319     ))
    320     .expect("initial refresh");
    321 
    322     let replacement = ProjectionGeneration::new([2; 32]).expect("replacement generation");
    323     block_on(ProjectionStore::invalidate(
    324         &*storage,
    325         ProjectionInvalidation::new(
    326             id.clone(),
    327             active.generation(),
    328             replacement,
    329             InvalidationReason::ProjectionGenerationChanged,
    330             2_000,
    331         )
    332         .expect("invalidation"),
    333     ))
    334     .expect("invalidate active generation");
    335 
    336     assert_eq!(
    337         block_on(engine.refresh_projection(
    338             RefreshRequest::new(id, active.generation(), 10, 1).expect("obsolete request"),
    339             &active,
    340         )),
    341         Err(Error::StorageFailed)
    342     );
    343 }
    344 
    345 #[test]
    346 fn partial_rebuild_resumes_and_rejects_concurrent_generation() {
    347     let (engine, storage, id) = setup();
    348     seed(&storage, 2);
    349     let first = reducer(&id, 1, false);
    350     block_on(engine.refresh_projection(
    351         RefreshRequest::new(id.clone(), first.generation(), 10, 1).expect("request"),
    352         &first,
    353     ))
    354     .expect("initial refresh");
    355 
    356     let replacement = reducer(&id, 2, false);
    357     let partial = block_on(engine.refresh_projection(
    358         RefreshRequest::new(id.clone(), replacement.generation(), 1, 1).expect("request"),
    359         &replacement,
    360     ))
    361     .expect("partial rebuild");
    362     assert_eq!(partial.state(), RefreshState::Partial);
    363     assert!(partial.rebuild_ticket().is_some());
    364     let visible_status = block_on(ProjectionStore::status(&*storage, id.clone()))
    365         .expect("status")
    366         .expect("projection");
    367     assert_eq!(visible_status.generation(), first.generation());
    368     assert_eq!(visible_status.health(), ProjectionHealth::Rebuilding);
    369 
    370     let concurrent = reducer(&id, 3, false);
    371     assert_eq!(
    372         block_on(engine.refresh_projection(
    373             RefreshRequest::new(id.clone(), concurrent.generation(), 1, 1).expect("request"),
    374             &concurrent,
    375         )),
    376         Err(Error::StorageConflict)
    377     );
    378 
    379     let second = block_on(engine.refresh_projection(
    380         RefreshRequest::new(id.clone(), replacement.generation(), 1, 1).expect("request"),
    381         &replacement,
    382     ))
    383     .expect("second batch");
    384     assert_eq!(second.state(), RefreshState::Partial);
    385     let complete = block_on(engine.refresh_projection(
    386         RefreshRequest::new(id.clone(), replacement.generation(), 1, 1).expect("request"),
    387         &replacement,
    388     ))
    389     .expect("complete rebuild");
    390     assert_eq!(complete.state(), RefreshState::Complete);
    391     for invalid in [
    392         RefreshRequest::new(id.clone(), replacement.generation(), 0, 1),
    393         RefreshRequest::new(
    394             id.clone(),
    395             replacement.generation(),
    396             radroots_storage::event::EVENT_QUERY_LIMIT_MAX + 1,
    397             1,
    398         ),
    399         RefreshRequest::new(id.clone(), replacement.generation(), 1, 0),
    400         RefreshRequest::new(
    401             id,
    402             replacement.generation(),
    403             1,
    404             radroots_sync::projection::PROJECTION_REFRESH_MAX_BATCHES + 1,
    405         ),
    406     ] {
    407         assert_eq!(invalid, Err(Error::InvalidProjectionRequest));
    408     }
    409 }
    410 
    411 #[test]
    412 fn source_change_fails_rebuild_and_preserves_prior_generation() {
    413     let (engine, storage, id) = setup();
    414     seed(&storage, 2);
    415     let active = reducer(&id, 1, false);
    416     block_on(engine.refresh_projection(
    417         RefreshRequest::new(id.clone(), active.generation(), 10, 1).expect("request"),
    418         &active,
    419     ))
    420     .expect("initial refresh");
    421 
    422     let replacement = reducer(&id, 2, false);
    423     let partial = block_on(engine.refresh_projection(
    424         RefreshRequest::new(id.clone(), replacement.generation(), 1, 1).expect("request"),
    425         &replacement,
    426     ))
    427     .expect("partial rebuild");
    428     let ticket_id = partial.rebuild_ticket().expect("ticket");
    429     seed(&storage, 3);
    430 
    431     let failed = block_on(engine.refresh_projection(
    432         RefreshRequest::new(id.clone(), replacement.generation(), 1, 1).expect("request"),
    433         &replacement,
    434     ))
    435     .expect("source change is normalized");
    436     assert_eq!(failed.state(), RefreshState::Failed);
    437     let status = block_on(ProjectionStore::status(&*storage, id))
    438         .expect("status")
    439         .expect("projection");
    440     assert_eq!(status.generation(), active.generation());
    441     assert_eq!(status.health(), ProjectionHealth::Ready);
    442     let ticket = block_on(storage.rebuild(ticket_id))
    443         .expect("ticket lookup")
    444         .expect("durable ticket");
    445     assert_eq!(ticket.failure(), Some(RebuildFailure::SourceChanged));
    446 }
    447 
    448 #[test]
    449 fn reducer_identity_progress_and_multi_batch_boundaries_fail_closed() {
    450     let (engine, storage, id) = setup();
    451     seed(&storage, 3);
    452     let active_reducer = reducer(&id, 1, false);
    453     let request =
    454         RefreshRequest::new(id.clone(), active_reducer.generation(), 1, 2).expect("request");
    455     let wrong_id = reducer(
    456         &ProjectionId::parse("different-projection").expect("projection id"),
    457         1,
    458         false,
    459     );
    460     assert_eq!(
    461         block_on(engine.refresh_projection(request.clone(), &wrong_id)),
    462         Err(Error::InvalidProjectionRequest)
    463     );
    464     let wrong_generation = reducer(&id, 2, false);
    465     assert_eq!(
    466         block_on(engine.refresh_projection(request.clone(), &wrong_generation)),
    467         Err(Error::InvalidProjectionRequest)
    468     );
    469     let partial =
    470         block_on(engine.refresh_projection(request, &active_reducer)).expect("two batches");
    471     assert_eq!(partial.state(), RefreshState::Partial);
    472     assert_eq!(partial.batches(), 2);
    473 
    474     let (engine, storage, id) = setup();
    475     seed(&storage, 1);
    476     let failing = reducer(&id, 1, true);
    477     let failed = block_on(engine.refresh_projection(
    478         RefreshRequest::new(id.clone(), failing.generation(), 1, 1).expect("request"),
    479         &failing,
    480     ))
    481     .expect("normalized incremental failure");
    482     assert_eq!(failed.state(), RefreshState::Failed);
    483     assert!(failed.rebuild_ticket().is_none());
    484 
    485     let (engine, storage, id) = setup();
    486     seed(&storage, 1);
    487     let initial = reducer(&id, 1, false);
    488     block_on(engine.refresh_projection(
    489         RefreshRequest::new(id.clone(), initial.generation(), 1, 1).expect("request"),
    490         &initial,
    491     ))
    492     .expect("initial projection");
    493     seed(&storage, 2);
    494     let regressing = CountingReducer {
    495         projection_id: id.clone(),
    496         generation: initial.generation(),
    497         fail: false,
    498         regress: true,
    499     };
    500     assert_eq!(
    501         block_on(engine.refresh_projection(
    502             RefreshRequest::new(id, regressing.generation(), 1, 1).expect("request"),
    503             &regressing,
    504         )),
    505         Err(Error::InvalidReducerOutput)
    506     );
    507 }