lib

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

ingest.rs (17648B)


      1 use std::sync::{
      2     Arc,
      3     atomic::{AtomicU8, Ordering},
      4 };
      5 
      6 use futures_executor::block_on;
      7 use radroots_event::{
      8     SignedEvent,
      9     draft::SignedEventParts,
     10     food::availability::{
     11         FoodAvailabilityDetails, FoodAvailabilityDetailsParts, FoodAvailabilityStatus, FoodContent,
     12         FoodCurrency, FoodIdentifier, FoodPrice, FoodPublishedAt, FoodText, FoodUnit,
     13     },
     14 };
     15 use radroots_event_codec::authoring::AuthoredEventPlan;
     16 use radroots_storage::{
     17     EventStore,
     18     event::{AdmissionDisposition, AdmissionStage, EventQuery, EventQueryBounds, SourceGeneration},
     19     memory::MemoryStorage,
     20 };
     21 use radroots_sync::{
     22     Engine,
     23     ingest::{AdmissionDecision, AdmissionPolicy, RegistryPolicy},
     24     policy::{Clock, DeadlinePolicy, Error, IdSource, OperationKind, SyncId, SyncStorage},
     25 };
     26 use radroots_transport::{
     27     Error as TransportError, EventSource, FetchPage, FetchRequest, SourceStatus, Target,
     28     TransportId,
     29     source::{EventProvenance, ObservedEvent},
     30 };
     31 use secp256k1::{Keypair, Message, Secp256k1, SecretKey};
     32 
     33 const EVENT_ID: &str = "762bee187e9e645b81ec26ade05a69b5e8398caf527be8de0d9a45311ed0c7a0";
     34 const PUBKEY: &str = "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df";
     35 const SIGNATURE: &str = "4290da0bb6422986647bc8cd5f63bd52d49f41e7b665d3b47105b8109183e8d596f322c531d4061df53e1d2b70fda12d5d1c14f3720d7a56d9d0a03746af5109";
     36 const CONTENT: &str = "{\"display_name\":\"Moss Street Farm\",\"bot\":false,\"website\":\"https://mossstreet.example\",\"picture\":42}";
     37 
     38 struct MockSource;
     39 
     40 impl EventSource for MockSource {
     41     fn status(&self) -> radroots_transport::BoxFuture<'_, Result<SourceStatus, TransportError>> {
     42         Box::pin(async { unreachable!("ingest does not inspect source status") })
     43     }
     44 
     45     fn fetch(
     46         &self,
     47         _request: FetchRequest,
     48     ) -> radroots_transport::BoxFuture<'_, Result<FetchPage, TransportError>> {
     49         Box::pin(async { unreachable!("ingest does not fetch") })
     50     }
     51 }
     52 
     53 struct FixedClock;
     54 
     55 impl Clock for FixedClock {
     56     fn now_unix_ms(&self) -> Result<u64, Error> {
     57         Ok(1_800_000_200_000)
     58     }
     59 }
     60 
     61 struct SequenceIds(AtomicU8);
     62 
     63 impl IdSource for SequenceIds {
     64     fn next_id(&self, operation: OperationKind) -> Result<SyncId, Error> {
     65         assert_eq!(operation, OperationKind::Ingest);
     66         let byte = self.0.fetch_add(1, Ordering::Relaxed);
     67         SyncId::new([byte; 16])
     68     }
     69 }
     70 
     71 struct ConstantIds(u8);
     72 
     73 impl IdSource for ConstantIds {
     74     fn next_id(&self, operation: OperationKind) -> Result<SyncId, Error> {
     75         assert_eq!(operation, OperationKind::Ingest);
     76         SyncId::new([self.0; 16])
     77     }
     78 }
     79 
     80 struct Reject;
     81 
     82 impl AdmissionPolicy for Reject {
     83     fn policy_id(&self) -> &'static str {
     84         "test.reject.v1"
     85     }
     86 
     87     fn decide(
     88         &self,
     89         _event: &radroots_event::admission::ContractValidatedEvent,
     90     ) -> AdmissionDecision {
     91         AdmissionDecision::Reject
     92     }
     93 }
     94 
     95 struct FoodAvailabilityPolicy;
     96 
     97 impl AdmissionPolicy for FoodAvailabilityPolicy {
     98     fn policy_id(&self) -> &'static str {
     99         "test.food-availability.v1"
    100     }
    101 
    102     fn select_contract(
    103         &self,
    104         _event: &radroots_event::admission::SignatureVerifiedEvent,
    105     ) -> Option<&'static str> {
    106         Some("radroots.food.availability.v1")
    107     }
    108 
    109     fn decide(
    110         &self,
    111         _event: &radroots_event::admission::ContractValidatedEvent,
    112     ) -> AdmissionDecision {
    113         AdmissionDecision::Visible
    114     }
    115 }
    116 
    117 fn setup_engine(first_id: u8) -> (Engine, Arc<MemoryStorage>) {
    118     setup_engine_with_ids(Arc::new(SequenceIds(AtomicU8::new(first_id))))
    119 }
    120 
    121 fn setup_engine_with_ids(ids: Arc<dyn IdSource>) -> (Engine, Arc<MemoryStorage>) {
    122     let storage = Arc::new(MemoryStorage::new(
    123         SourceGeneration::new([7; 32]).expect("source generation"),
    124     ));
    125     let storage_capability: Arc<dyn SyncStorage> = storage.clone();
    126     let engine = Engine::builder(
    127         storage_capability,
    128         Arc::new(FixedClock),
    129         ids,
    130         DeadlinePolicy::new(10_000, 10_000, 10_000).expect("deadlines"),
    131     )
    132     .source(Arc::new(MockSource))
    133     .build()
    134     .expect("engine");
    135     (engine, storage)
    136 }
    137 
    138 fn signed_event(signature: &str) -> SignedEvent {
    139     let raw_json = format!(
    140         "{{\"id\":\"{EVENT_ID}\",\"pubkey\":\"{PUBKEY}\",\"created_at\":1800000100,\"kind\":0,\"tags\":[],\"content\":{content:?},\"sig\":\"{signature}\"}}",
    141         content = CONTENT,
    142     );
    143     SignedEvent::new(SignedEventParts {
    144         id: EVENT_ID.to_owned(),
    145         pubkey: PUBKEY.to_owned(),
    146         created_at: 1_800_000_100,
    147         kind: 0,
    148         tags: vec![],
    149         content: CONTENT.to_owned(),
    150         sig: signature.to_owned(),
    151         raw_json,
    152     })
    153     .expect("ID-valid signed event")
    154 }
    155 
    156 fn observed(signature: &str, observed_at: u64) -> ObservedEvent {
    157     let target = Target::new(TransportId::NOSTR, "wss://relay.example").expect("target");
    158     let provenance = EventProvenance::new(
    159         TransportId::NOSTR,
    160         target.fingerprint().clone(),
    161         observed_at,
    162     )
    163     .expect("provenance");
    164     ObservedEvent::new(signed_event(signature), provenance)
    165 }
    166 
    167 fn observed_food(observed_at: u64) -> ObservedEvent {
    168     let created_at = 1_800_000_100;
    169     let keypair = Keypair::from_secret_key(
    170         &Secp256k1::new(),
    171         &SecretKey::from_slice(&[1; 32]).expect("food fixture secret"),
    172     );
    173     let public_key = keypair.x_only_public_key().0.to_string();
    174     let details = FoodAvailabilityDetails::new(FoodAvailabilityDetailsParts {
    175         content: FoodContent::new("Carrots available this week.").expect("content"),
    176         identifier: FoodIdentifier::parse("nantes-carrots").expect("identifier"),
    177         title: FoodText::new("Nantes Carrots").expect("title"),
    178         summary: FoodText::new("Fresh bunches").expect("summary"),
    179         published_at: FoodPublishedAt::new(created_at).expect("published at"),
    180         location: FoodText::new("Central Saanich, BC").expect("location"),
    181         price: FoodPrice::new(
    182             "3",
    183             FoodCurrency::parse("CAD").expect("currency"),
    184             FoodUnit::Pound,
    185         )
    186         .expect("price"),
    187         quantity: None,
    188         status: FoodAvailabilityStatus::Active,
    189         images: Vec::new(),
    190     })
    191     .expect("food availability");
    192     let plan = AuthoredEventPlan::from_food_availability(&details, created_at, &public_key)
    193         .expect("food plan");
    194     let id = plan.expected_event_id().to_hex();
    195     let signature = Secp256k1::new()
    196         .sign_schnorr_no_aux_rand(
    197             &Message::from_digest(*plan.expected_event_id().as_bytes()),
    198             &keypair,
    199         )
    200         .to_string();
    201     let raw_json = format!(
    202         "{{\"id\":\"{id}\",\"pubkey\":\"{public_key}\",\"created_at\":{created_at},\"kind\":{},\"tags\":{:?},\"content\":{content:?},\"sig\":\"{signature}\"}}",
    203         plan.body().kind(),
    204         plan.body().tags(),
    205         content = plan.body().content(),
    206     );
    207     let event = SignedEvent::new(SignedEventParts {
    208         id,
    209         pubkey: public_key,
    210         created_at,
    211         kind: plan.body().kind(),
    212         tags: plan.body().tags().to_vec(),
    213         content: plan.body().content().to_owned(),
    214         sig: signature,
    215         raw_json,
    216     })
    217     .expect("signed food event");
    218     let target = Target::new(TransportId::NOSTR, "wss://relay.example").expect("target");
    219     let provenance = EventProvenance::new(
    220         TransportId::NOSTR,
    221         target.fingerprint().clone(),
    222         observed_at,
    223     )
    224     .expect("provenance");
    225     ObservedEvent::new(event, provenance)
    226 }
    227 
    228 #[test]
    229 fn valid_visible_ingest_is_atomic_and_preserves_provenance() {
    230     let (engine, storage) = setup_engine(1);
    231     let receipt = block_on(engine.ingest(
    232         observed(SIGNATURE, 1_800_000_100_000),
    233         &RegistryPolicy::visible(),
    234     ))
    235     .expect("visible ingest");
    236     assert_eq!(receipt.sync_id().as_bytes(), &[1; 16]);
    237     assert_eq!(
    238         receipt.commit_disposition(),
    239         radroots_storage::atomic::AtomicCommitDisposition::Committed
    240     );
    241     assert_eq!(receipt.committed_at_unix_ms(), 1_800_000_200_000);
    242     assert_eq!(receipt.admission().stage(), AdmissionStage::Visible);
    243     assert_eq!(
    244         receipt.admission().disposition(),
    245         AdmissionDisposition::Inserted
    246     );
    247 
    248     let bounds = EventQueryBounds::first(10).expect("bounds");
    249     let visible = block_on(storage.query_visible(EventQuery::all(bounds))).expect("visible query");
    250     assert_eq!(visible.items().len(), 1);
    251     let provenance = block_on(storage.query_provenance(*receipt.admission().event_id(), bounds))
    252         .expect("provenance query");
    253     assert_eq!(provenance.items().len(), 1);
    254     assert_eq!(
    255         provenance.items()[0].provenance().observed_at_unix_ms(),
    256         1_800_000_100_000
    257     );
    258 }
    259 
    260 struct RetainInvalid {
    261     calls: AtomicU8,
    262     explicit_contract: bool,
    263 }
    264 
    265 impl AdmissionPolicy for RetainInvalid {
    266     fn policy_id(&self) -> &'static str {
    267         "test.retain-invalid.v1"
    268     }
    269     fn select_contract(
    270         &self,
    271         _: &radroots_event::admission::SignatureVerifiedEvent,
    272     ) -> Option<&'static str> {
    273         self.explicit_contract
    274             .then_some("radroots.food.availability.v1")
    275     }
    276     fn contract_failure(
    277         &self,
    278         _: &radroots_event::admission::SignatureVerifiedEvent,
    279     ) -> radroots_sync::ingest::ContractFailureDecision {
    280         self.calls.fetch_add(1, Ordering::Relaxed);
    281         radroots_sync::ingest::ContractFailureDecision::Verified
    282     }
    283     fn decide(&self, _: &radroots_event::admission::ContractValidatedEvent) -> AdmissionDecision {
    284         AdmissionDecision::Visible
    285     }
    286 }
    287 
    288 fn observed_profile(created_at: u64, content: &str, valid_signature: bool) -> ObservedEvent {
    289     let pair = Keypair::from_secret_key(
    290         &Secp256k1::new(),
    291         &SecretKey::from_slice(&[1; 32]).expect("fixture key"),
    292     );
    293     let mut wire = radroots_event::wire::Nip01EventWire {
    294         id: "0".repeat(64),
    295         pubkey: pair.x_only_public_key().0.to_string(),
    296         created_at,
    297         kind: 0,
    298         tags: vec![],
    299         content: content.to_owned(),
    300         sig: "42".repeat(64),
    301         extra: Default::default(),
    302     };
    303     let id = wire.computed_event_id().expect("id");
    304     wire.id = id.to_hex();
    305     if valid_signature {
    306         wire.sig = Secp256k1::new()
    307             .sign_schnorr_no_aux_rand(&Message::from_digest(*id.as_bytes()), &pair)
    308             .to_string();
    309     }
    310     let raw = serde_json::json!({"id":wire.id,"pubkey":wire.pubkey,"created_at":wire.created_at,"kind":wire.kind,"tags":wire.tags,"content":wire.content,"sig":wire.sig}).to_string();
    311     let event = SignedEvent::from_wire_verified_id(wire, raw).expect("ID-valid signed profile");
    312     let target = Target::new(TransportId::NOSTR, "wss://relay.example").expect("target");
    313     ObservedEvent::new(
    314         event,
    315         EventProvenance::new(
    316             TransportId::NOSTR,
    317             target.fingerprint().clone(),
    318             created_at * 1000,
    319         )
    320         .expect("provenance"),
    321     )
    322 }
    323 
    324 #[test]
    325 fn contract_failure_defaults_to_reject_and_opt_in_retains_only_verified_heads() {
    326     let (engine, storage) = setup_engine(1);
    327     let old = observed_profile(10, r#"{"display_name":"Old Farm","bot":false}"#, true);
    328     let newer = observed_profile(20, "not JSON", true);
    329     let forged = observed_profile(30, "not JSON", false);
    330     block_on(engine.ingest(old.clone(), &RegistryPolicy::visible())).expect("old visible");
    331     assert_eq!(
    332         block_on(engine.ingest(newer.clone(), &RegistryPolicy::visible())).unwrap_err(),
    333         Error::VerificationFailed
    334     );
    335     assert_eq!(block_on(storage.status()).expect("status").raw_events(), 1);
    336     let policy = RetainInvalid {
    337         calls: AtomicU8::new(0),
    338         explicit_contract: false,
    339     };
    340     let batch = block_on(engine.ingest_batch(vec![forged, newer.clone(), old], &policy));
    341     assert_eq!(batch.accepted(), 2);
    342     assert_eq!(batch.rejected(), 1);
    343     assert_eq!(batch.outcomes()[0], Err(Error::VerificationFailed));
    344     assert_eq!(
    345         batch.outcomes()[1].as_ref().unwrap().admission().stage(),
    346         AdmissionStage::Verified
    347     );
    348     assert_eq!(policy.calls.load(Ordering::Relaxed), 1);
    349     let status = block_on(storage.status()).expect("status");
    350     assert_eq!(status.raw_events(), 2);
    351     assert_eq!(status.verified_events(), 2);
    352     assert_eq!(status.visible_events(), 0);
    353     assert!(
    354         block_on(storage.query_visible(EventQuery::all(EventQueryBounds::first(10).unwrap())))
    355             .unwrap()
    356             .items()
    357             .is_empty()
    358     );
    359     let snapshot = block_on(storage.rebuild_visibility()).unwrap();
    360     assert_eq!(snapshot.current_heads()[0].event_id, *newer.event().id());
    361 }
    362 
    363 #[test]
    364 fn explicit_contract_failure_retention_advances_to_visible_without_new_raw_record() {
    365     let (engine, storage) = setup_engine(1);
    366     let profile = observed_profile(10, r#"{"display_name":"Farm","bot":false}"#, true);
    367     let policy = RetainInvalid {
    368         calls: AtomicU8::new(0),
    369         explicit_contract: true,
    370     };
    371     let retained = block_on(engine.ingest(profile.clone(), &policy)).expect("retained");
    372     assert_eq!(retained.admission().stage(), AdmissionStage::Verified);
    373     let before = block_on(storage.rebuild_visibility()).unwrap();
    374     assert!(before.visible_event_ids().is_empty());
    375     let advanced = block_on(engine.ingest(profile, &RegistryPolicy::visible())).expect("advanced");
    376     assert_eq!(
    377         advanced.admission().disposition(),
    378         AdmissionDisposition::Advanced
    379     );
    380     assert_eq!(block_on(storage.status()).unwrap().raw_events(), 1);
    381     let after = block_on(storage.rebuild_visibility()).unwrap();
    382     assert_ne!(before.digest(), after.digest());
    383     assert_eq!(
    384         after.visible_event_ids(),
    385         &[*retained.admission().event_id()]
    386     );
    387     assert_eq!(policy.calls.load(Ordering::Relaxed), 1);
    388 }
    389 
    390 #[cfg(feature = "serde")]
    391 #[test]
    392 fn contract_failure_wire_decision_cannot_authorize_visibility() {
    393     use radroots_sync::ingest::ContractFailureDecision;
    394     for value in [
    395         ContractFailureDecision::Reject,
    396         ContractFailureDecision::Verified,
    397     ] {
    398         let encoded = serde_json::to_string(&value).unwrap();
    399         assert_eq!(
    400             serde_json::from_str::<ContractFailureDecision>(&encoded).unwrap(),
    401             value
    402         );
    403     }
    404     assert!(serde_json::from_str::<ContractFailureDecision>(r#""visible""#).is_err());
    405 }
    406 
    407 #[test]
    408 fn admission_policy_selects_and_fully_validates_admission_only_wire_profiles() {
    409     let (engine, storage) = setup_engine(1);
    410     assert_eq!(
    411         block_on(engine.ingest(observed_food(1), &RegistryPolicy::visible())),
    412         Err(Error::VerificationFailed)
    413     );
    414     let receipt = block_on(engine.ingest(observed_food(2), &FoodAvailabilityPolicy))
    415         .expect("policy-selected food admission");
    416     assert_eq!(receipt.admission().stage(), AdmissionStage::Visible);
    417     let visible = block_on(storage.query_visible(EventQuery::all(
    418         EventQueryBounds::first(10).expect("bounds"),
    419     )))
    420     .expect("visible query");
    421     assert_eq!(visible.items().len(), 1);
    422     assert_eq!(visible.items()[0].event().kind(), 30_402);
    423 }
    424 
    425 #[test]
    426 fn invalid_policy_rejected_and_verified_only_inputs_fail_closed() {
    427     let (engine, storage) = setup_engine(1);
    428     let invalid_signature = format!("0{}", &SIGNATURE[1..]);
    429     assert_eq!(
    430         block_on(engine.ingest(observed(&invalid_signature, 1), &RegistryPolicy::visible())),
    431         Err(Error::VerificationFailed)
    432     );
    433     assert_eq!(
    434         block_on(engine.ingest(observed(SIGNATURE, 2), &Reject)),
    435         Err(Error::PolicyRejected)
    436     );
    437 
    438     let receipt = block_on(engine.ingest(observed(SIGNATURE, 3), &RegistryPolicy::verified()))
    439         .expect("verified ingest");
    440     assert_eq!(receipt.admission().stage(), AdmissionStage::Verified);
    441     let page = block_on(storage.query_visible(EventQuery::all(
    442         EventQueryBounds::first(10).expect("bounds"),
    443     )))
    444     .expect("visible query");
    445     assert!(page.items().is_empty());
    446 }
    447 
    448 #[test]
    449 fn duplicate_conflict_and_partial_batch_outcomes_are_normalized() {
    450     let (engine, storage) = setup_engine(1);
    451     let inserted = block_on(engine.ingest(observed(SIGNATURE, 10), &RegistryPolicy::visible()))
    452         .expect("insert");
    453     let duplicate = block_on(engine.ingest(observed(SIGNATURE, 11), &RegistryPolicy::visible()))
    454         .expect("duplicate");
    455     assert_eq!(
    456         inserted.admission().position(),
    457         duplicate.admission().position()
    458     );
    459     assert_eq!(
    460         duplicate.admission().disposition(),
    461         AdmissionDisposition::Duplicate
    462     );
    463     let provenance = block_on(storage.query_provenance(
    464         *inserted.admission().event_id(),
    465         EventQueryBounds::first(10).expect("bounds"),
    466     ))
    467     .expect("provenance query");
    468     assert_eq!(provenance.items().len(), 2);
    469 
    470     let (collision_engine, _) = setup_engine_with_ids(Arc::new(ConstantIds(9)));
    471     block_on(collision_engine.ingest(observed(SIGNATURE, 20), &RegistryPolicy::visible()))
    472         .expect("first identity use");
    473     assert_eq!(
    474         block_on(collision_engine.ingest(observed(SIGNATURE, 21), &RegistryPolicy::visible())),
    475         Err(Error::StorageConflict)
    476     );
    477 
    478     let invalid_signature = format!("0{}", &SIGNATURE[1..]);
    479     let batch = block_on(engine.ingest_batch(
    480         vec![
    481             observed(SIGNATURE, 30),
    482             observed(&invalid_signature, 31),
    483             observed(SIGNATURE, 32),
    484         ],
    485         &RegistryPolicy::visible(),
    486     ));
    487     assert_eq!(batch.accepted(), 2);
    488     assert_eq!(batch.rejected(), 1);
    489     assert_eq!(batch.outcomes()[1], Err(Error::VerificationFailed));
    490 }