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 }