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 ®ressing, 504 )), 505 Err(Error::InvalidReducerOutput) 506 ); 507 }