mod.rs (42954B)
1 use crate::backend::map_backend; 2 use radroots_event::SignedEvent; 3 use radroots_event_codec::Codec; 4 use radroots_storage::{ 5 Error, EventStore, 6 event::{ 7 AdmissionDisposition, AdmissionReceipt, AdmissionStage, BoxFuture, EventAdmission, 8 EventCursor, EventId, EventPage, EventPosition, EventQuery, EventQueryBounds, 9 EventSequence, SourceGeneration, StoredEventProvenance, StoredRawEvent, 10 StoredVerifiedEvent, StoredVisibleEvent, VisibilityEvaluation, VisibilityInput, 11 VisibilitySnapshot, evaluate_visibility, 12 }, 13 status::{EventStoreHealth, EventStoreMode, EventStoreStatus}, 14 }; 15 use sqlx::{QueryBuilder, Row, Sqlite, SqlitePool}; 16 use std::{ 17 path::PathBuf, 18 sync::{Arc, Mutex}, 19 }; 20 21 use crate::lock::WriterLock; 22 use crate::status::StorageLifecycle; 23 24 #[derive(Clone)] 25 pub struct SqliteStorage { 26 pub(crate) pool: SqlitePool, 27 pub(crate) private_pool: SqlitePool, 28 pub(crate) generation: SourceGeneration, 29 pub(crate) mode: EventStoreMode, 30 pub(crate) lifecycle: Arc<StorageLifecycle>, 31 pub(crate) backup_root: Option<Arc<PathBuf>>, 32 pub(crate) paths: Option<Arc<crate::Paths>>, 33 pub(crate) reliability: Arc<Mutex<crate::backup::ReliabilityState>>, 34 } 35 36 struct StoredEventRow { 37 position: EventPosition, 38 raw_json: String, 39 event: SignedEvent, 40 stage: AdmissionStage, 41 } 42 43 impl SqliteStorage { 44 #[allow(dead_code)] // Single-pool in-memory scaffold retained for focused backend tests. 45 pub(crate) fn new( 46 pool: SqlitePool, 47 generation: SourceGeneration, 48 mode: EventStoreMode, 49 ) -> Self { 50 Self { 51 private_pool: pool.clone(), 52 pool, 53 generation, 54 mode, 55 lifecycle: Arc::new(StorageLifecycle::scaffold(mode)), 56 backup_root: None, 57 paths: None, 58 reliability: Arc::new(Mutex::new(crate::backup::ReliabilityState::default())), 59 } 60 } 61 62 #[allow(dead_code)] // Two-pool in-memory scaffold retained for focused backend tests. 63 pub(crate) fn with_private_pool( 64 pool: SqlitePool, 65 private_pool: SqlitePool, 66 generation: SourceGeneration, 67 mode: EventStoreMode, 68 ) -> Self { 69 Self { 70 pool, 71 private_pool, 72 generation, 73 mode, 74 lifecycle: Arc::new(StorageLifecycle::scaffold(mode)), 75 backup_root: None, 76 paths: None, 77 reliability: Arc::new(Mutex::new(crate::backup::ReliabilityState::default())), 78 } 79 } 80 81 pub(crate) fn from_opened( 82 pool: SqlitePool, 83 private_pool: SqlitePool, 84 generation: SourceGeneration, 85 options: &crate::OpenOptions, 86 writer_lock: Option<WriterLock>, 87 ) -> Self { 88 let mode = if options.mode().is_writable() { 89 EventStoreMode::ReadWrite 90 } else { 91 EventStoreMode::ReadOnly 92 }; 93 Self { 94 pool, 95 private_pool, 96 generation, 97 mode, 98 lifecycle: Arc::new(StorageLifecycle::new( 99 options.mode(), 100 options.busy_timeout(), 101 writer_lock, 102 )), 103 backup_root: options 104 .backup_root() 105 .map(|path| Arc::new(path.to_path_buf())), 106 paths: Some(Arc::new(options.paths().clone())), 107 reliability: Arc::new(Mutex::new(crate::backup::ReliabilityState::default())), 108 } 109 } 110 111 pub(crate) const fn pool(&self) -> &SqlitePool { 112 &self.pool 113 } 114 115 pub(crate) const fn private_pool(&self) -> &SqlitePool { 116 &self.private_pool 117 } 118 119 pub(crate) const fn event_mode(&self) -> EventStoreMode { 120 self.mode 121 } 122 123 async fn selected( 124 &self, 125 query: &EventQuery, 126 minimum_stage: AdmissionStage, 127 ) -> Result<(Vec<StoredEventRow>, Option<EventCursor>), Error> { 128 self.validate_cursor(query)?; 129 let after = query 130 .bounds() 131 .cursor() 132 .map_or(0, |cursor| cursor.sequence().get()); 133 let fetch_limit = u64::from(query.bounds().limit()) + 1; 134 let mut builder = QueryBuilder::<Sqlite>::new( 135 "SELECT source_generation, source_sequence, signed_event, admission_stage, \ 136 admitted_contract_id, admitted_registry_version \ 137 FROM radroots_runtime_events WHERE source_generation = ", 138 ); 139 builder.push_bind(self.generation.as_bytes().as_slice()); 140 builder.push(" AND source_sequence > "); 141 builder.push_bind(i64_from_u64(after)?); 142 if minimum_stage == AdmissionStage::Verified { 143 builder.push(" AND admission_stage IN ('verified', 'visible')"); 144 } else if minimum_stage == AdmissionStage::Visible { 145 builder.push(" AND admission_stage = 'visible'"); 146 } 147 if !query.event_ids().is_empty() { 148 builder.push(" AND event_id IN ("); 149 let mut separated = builder.separated(", "); 150 for event_id in query.event_ids() { 151 separated.push_bind(event_id.as_bytes().to_vec()); 152 } 153 separated.push_unseparated(")"); 154 } 155 builder.push(" ORDER BY source_sequence LIMIT "); 156 builder.push_bind(i64_from_u64(fetch_limit)?); 157 158 let rows = builder 159 .build() 160 .fetch_all(&self.pool) 161 .await 162 .map_err(map_backend)?; 163 let mut decoded = rows 164 .iter() 165 .map(|row| self.decode_event_row(row)) 166 .collect::<Result<Vec<_>, _>>()?; 167 let next = if decoded.len() > usize::from(query.bounds().limit()) { 168 decoded.truncate(usize::from(query.bounds().limit())); 169 decoded.last().map(|row| row.position) 170 } else { 171 None 172 }; 173 Ok((decoded, next)) 174 } 175 176 async fn current_event_rows(&self) -> Result<Vec<StoredEventRow>, Error> { 177 let rows = sqlx::query( 178 "SELECT source_generation, source_sequence, signed_event, admission_stage, 179 admitted_contract_id, admitted_registry_version 180 FROM radroots_runtime_events 181 WHERE source_generation = ? 182 ORDER BY source_sequence", 183 ) 184 .bind(self.generation.as_bytes().as_slice()) 185 .fetch_all(&self.pool) 186 .await 187 .map_err(map_backend)?; 188 rows.iter().map(|row| self.decode_event_row(row)).collect() 189 } 190 191 fn visibility_for_rows(&self, rows: &[StoredEventRow]) -> Result<VisibilityEvaluation, Error> { 192 evaluate_visibility( 193 self.generation, 194 rows.iter() 195 .map(|row| VisibilityInput::new(row.position, &row.event, row.stage)), 196 ) 197 } 198 199 fn validate_cursor(&self, query: &EventQuery) -> Result<(), Error> { 200 if query 201 .bounds() 202 .cursor() 203 .is_some_and(|cursor| cursor.generation() != self.generation) 204 { 205 return Err(Error::SourceGenerationChanged); 206 } 207 Ok(()) 208 } 209 210 fn decode_event_row(&self, row: &sqlx::sqlite::SqliteRow) -> Result<StoredEventRow, Error> { 211 let generation = source_generation(row.try_get("source_generation").map_err(map_corrupt)?)?; 212 if generation != self.generation { 213 return Err(Error::CorruptStoredEvent); 214 } 215 let sequence = event_sequence(row.try_get("source_sequence").map_err(map_corrupt)?)?; 216 let raw_json = String::from_utf8(row.try_get("signed_event").map_err(map_corrupt)?) 217 .map_err(|_| Error::CorruptStoredEvent)?; 218 let event = 219 Codec::decode_signed_event(raw_json.as_str()).map_err(|_| Error::CorruptStoredEvent)?; 220 let stage = admission_stage(row.try_get("admission_stage").map_err(map_corrupt)?)?; 221 let contract_id = row 222 .try_get::<Option<String>, _>("admitted_contract_id") 223 .map_err(map_corrupt)?; 224 let registry_version = row 225 .try_get::<Option<i64>, _>("admitted_registry_version") 226 .map_err(map_corrupt)? 227 .map(|value| u32::try_from(value).map_err(|_| Error::CorruptStoredEvent)) 228 .transpose()?; 229 if contract_id.is_some() != registry_version.is_some() 230 || registry_version.is_some_and(|version| { 231 version != radroots_event::contract::RegistryVersion::CURRENT.get() 232 }) 233 || contract_id.as_deref().is_some_and(|stored| { 234 radroots_event::contract::registry_v7::validate_event_contract_for_admission( 235 event.envelope(), 236 stored, 237 ) 238 .is_err() 239 }) 240 { 241 return Err(Error::CorruptStoredEvent); 242 } 243 Ok(StoredEventRow { 244 position: EventPosition::new(generation, sequence), 245 raw_json, 246 event, 247 stage, 248 }) 249 } 250 251 async fn store_provenance( 252 transaction: &mut sqlx::Transaction<'_, Sqlite>, 253 admission: &EventAdmission, 254 ) -> Result<(), Error> { 255 let provenance = admission.provenance(); 256 let cursor = provenance.cursor().map_or("", |cursor| cursor.as_str()); 257 sqlx::query( 258 "INSERT OR IGNORE INTO radroots_runtime_event_provenance ( 259 event_id, transport_id, target_fingerprint, observed_at_unix_ms, cursor 260 ) VALUES (?, ?, ?, ?, ?)", 261 ) 262 .bind(admission.event_id().as_bytes().as_slice()) 263 .bind(provenance.transport_id().as_str()) 264 .bind(provenance.target().as_str()) 265 .bind(i64_from_u64(provenance.observed_at_unix_ms())?) 266 .bind(cursor) 267 .execute(&mut **transaction) 268 .await 269 .map_err(map_backend)?; 270 Ok(()) 271 } 272 273 pub(crate) async fn admit_transaction( 274 &self, 275 transaction: &mut sqlx::Transaction<'_, Sqlite>, 276 admission: EventAdmission, 277 ) -> Result<AdmissionReceipt, Error> { 278 let existing = sqlx::query( 279 "SELECT source_generation, source_sequence, signed_event, admission_stage, 280 admitted_contract_id, admitted_registry_version 281 FROM radroots_runtime_events WHERE event_id = ?", 282 ) 283 .bind(admission.event_id().as_bytes().as_slice()) 284 .fetch_optional(&mut **transaction) 285 .await 286 .map_err(map_backend)?; 287 288 let (position, disposition) = if let Some(row) = existing { 289 let stored = self.decode_event_row(&row)?; 290 let metadata = admission_contract_metadata(&admission); 291 if stored.raw_json.as_bytes() != admission.event().raw_json().as_bytes() { 292 return Err(Error::EventConflict); 293 } 294 if admission.stage() < stored.stage { 295 return Err(Error::AdmissionRegression); 296 } 297 let disposition = if admission.stage() == stored.stage { 298 AdmissionDisposition::Duplicate 299 } else { 300 sqlx::query( 301 "UPDATE radroots_runtime_events 302 SET admission_stage = ?, updated_at_unix_ms = MAX(updated_at_unix_ms, ?), 303 admitted_contract_id = COALESCE(admitted_contract_id, ?), 304 admitted_registry_version = COALESCE(admitted_registry_version, ?) 305 WHERE event_id = ?", 306 ) 307 .bind(stage_name(admission.stage())) 308 .bind(i64_from_u64(admission.provenance().observed_at_unix_ms())?) 309 .bind(metadata.map(|value| value.0)) 310 .bind(metadata.map(|value| i64::from(value.1))) 311 .bind(admission.event_id().as_bytes().as_slice()) 312 .execute(&mut **transaction) 313 .await 314 .map_err(map_backend)?; 315 AdmissionDisposition::Advanced 316 }; 317 (stored.position, disposition) 318 } else { 319 let metadata = admission_contract_metadata(&admission); 320 let next = sqlx::query_scalar::<_, i64>( 321 "UPDATE radroots_runtime_source_generations 322 SET sequence_head = sequence_head + 1 323 WHERE generation = ? AND state = 'active' 324 RETURNING sequence_head", 325 ) 326 .bind(self.generation.as_bytes().as_slice()) 327 .fetch_optional(&mut **transaction) 328 .await 329 .map_err(map_backend)? 330 .ok_or(Error::SourceGenerationChanged)?; 331 let sequence = event_sequence(next)?; 332 let observed_at = i64_from_u64(admission.provenance().observed_at_unix_ms())?; 333 sqlx::query( 334 "INSERT INTO radroots_runtime_events ( 335 source_generation, source_sequence, event_id, admission_stage, 336 signed_event, admitted_at_unix_ms, updated_at_unix_ms, 337 admitted_contract_id, admitted_registry_version 338 ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)", 339 ) 340 .bind(self.generation.as_bytes().as_slice()) 341 .bind(next) 342 .bind(admission.event_id().as_bytes().as_slice()) 343 .bind(stage_name(admission.stage())) 344 .bind(admission.event().raw_json().as_bytes()) 345 .bind(observed_at) 346 .bind(observed_at) 347 .bind(metadata.map(|value| value.0)) 348 .bind(metadata.map(|value| i64::from(value.1))) 349 .execute(&mut **transaction) 350 .await 351 .map_err(map_backend)?; 352 ( 353 EventPosition::new(self.generation, sequence), 354 AdmissionDisposition::Inserted, 355 ) 356 }; 357 Self::store_provenance(transaction, &admission).await?; 358 Ok(AdmissionReceipt::new( 359 *admission.event_id(), 360 position, 361 admission.stage(), 362 disposition, 363 )) 364 } 365 } 366 367 fn admission_contract_metadata(admission: &EventAdmission) -> Option<(&'static str, u32)> { 368 admission.visible_event().map(|event| { 369 ( 370 event.admitted_event().validated_event().contract_id(), 371 radroots_event::contract::RegistryVersion::CURRENT.get(), 372 ) 373 }) 374 } 375 376 #[cfg_attr(coverage_nightly, coverage(off))] 377 impl EventStore for SqliteStorage { 378 fn status(&self) -> BoxFuture<'_, Result<EventStoreStatus, Error>> { 379 Box::pin(async move { 380 let row = sqlx::query( 381 "SELECT 382 COUNT(*) AS raw_events, 383 COALESCE(SUM(CASE WHEN admission_stage IN ('verified', 'visible') THEN 1 ELSE 0 END), 0) 384 AS verified_events 385 FROM radroots_runtime_events WHERE source_generation = ?", 386 ) 387 .bind(self.generation.as_bytes().as_slice()) 388 .fetch_one(&self.pool) 389 .await 390 .map_err(map_backend)?; 391 let visibility_rows = self.current_event_rows().await?; 392 let visible_events = u64::try_from( 393 self.visibility_for_rows(&visibility_rows)? 394 .snapshot() 395 .visible_event_ids() 396 .len(), 397 ) 398 .map_err(|_| Error::CorruptStoredEvent)?; 399 EventStoreStatus::new( 400 self.generation, 401 self.mode, 402 EventStoreHealth::Available, 403 u64_from_i64(row.try_get("raw_events").map_err(map_corrupt)?)?, 404 u64_from_i64(row.try_get("verified_events").map_err(map_corrupt)?)?, 405 visible_events, 406 ) 407 }) 408 } 409 410 fn admit(&self, admission: EventAdmission) -> BoxFuture<'_, Result<AdmissionReceipt, Error>> { 411 Box::pin(async move { 412 if self.mode == EventStoreMode::ReadOnly { 413 return Err(Error::BackendUnavailable); 414 } 415 let mut transaction = self 416 .pool 417 .begin_with("BEGIN IMMEDIATE") 418 .await 419 .map_err(map_backend)?; 420 let receipt = self.admit_transaction(&mut transaction, admission).await?; 421 transaction.commit().await.map_err(map_backend)?; 422 Ok(receipt) 423 }) 424 } 425 426 fn query_raw( 427 &self, 428 query: EventQuery, 429 ) -> BoxFuture<'_, Result<EventPage<StoredRawEvent>, Error>> { 430 Box::pin(async move { 431 let (rows, next) = self.selected(&query, AdmissionStage::Raw).await?; 432 let items = rows 433 .into_iter() 434 .map(|row| StoredRawEvent::new(row.position, row.event, row.stage)) 435 .collect(); 436 EventPage::new(self.generation, items, next, query.bounds()) 437 }) 438 } 439 440 fn query_verified( 441 &self, 442 query: EventQuery, 443 ) -> BoxFuture<'_, Result<EventPage<StoredVerifiedEvent>, Error>> { 444 Box::pin(async move { 445 let (rows, next) = self.selected(&query, AdmissionStage::Verified).await?; 446 let items = rows 447 .into_iter() 448 .map(|row| StoredVerifiedEvent::new(row.position, row.event)) 449 .collect(); 450 EventPage::new(self.generation, items, next, query.bounds()) 451 }) 452 } 453 454 fn query_visible( 455 &self, 456 query: EventQuery, 457 ) -> BoxFuture<'_, Result<EventPage<StoredVisibleEvent>, Error>> { 458 Box::pin(async move { 459 self.validate_cursor(&query)?; 460 let after = query 461 .bounds() 462 .cursor() 463 .map_or(0, |cursor| cursor.sequence().get()); 464 let rows = self.current_event_rows().await?; 465 let visibility = self.visibility_for_rows(&rows)?; 466 let mut selected = rows 467 .into_iter() 468 .filter(|row| { 469 row.position.sequence().get() > after 470 && query.selects(row.event.id()) 471 && visibility.is_visible(row.event.id()) 472 }) 473 .take(usize::from(query.bounds().limit()) + 1) 474 .collect::<Vec<_>>(); 475 let next = if selected.len() > usize::from(query.bounds().limit()) { 476 selected.truncate(usize::from(query.bounds().limit())); 477 selected.last().map(|row| row.position) 478 } else { 479 None 480 }; 481 let items = selected 482 .into_iter() 483 .map(|row| StoredVisibleEvent::new(row.position, row.event)) 484 .collect(); 485 EventPage::new(self.generation, items, next, query.bounds()) 486 }) 487 } 488 489 fn rebuild_visibility(&self) -> BoxFuture<'_, Result<VisibilitySnapshot, Error>> { 490 Box::pin(async move { 491 let rows = self.current_event_rows().await?; 492 Ok(self.visibility_for_rows(&rows)?.into_snapshot()) 493 }) 494 } 495 496 fn query_provenance( 497 &self, 498 event_id: EventId, 499 bounds: EventQueryBounds, 500 ) -> BoxFuture<'_, Result<EventPage<StoredEventProvenance>, Error>> { 501 Box::pin(async move { 502 if bounds 503 .cursor() 504 .is_some_and(|cursor| cursor.generation() != self.generation) 505 { 506 return Err(Error::SourceGenerationChanged); 507 } 508 let event_row = sqlx::query( 509 "SELECT source_generation, source_sequence 510 FROM radroots_runtime_events WHERE event_id = ?", 511 ) 512 .bind(event_id.as_bytes().as_slice()) 513 .fetch_optional(&self.pool) 514 .await 515 .map_err(map_backend)? 516 .ok_or(Error::EventNotFound)?; 517 let generation = source_generation( 518 event_row 519 .try_get("source_generation") 520 .map_err(map_corrupt)?, 521 )?; 522 if generation != self.generation { 523 return Err(Error::CorruptStoredEvent); 524 } 525 let position = EventPosition::new( 526 generation, 527 event_sequence(event_row.try_get("source_sequence").map_err(map_corrupt)?)?, 528 ); 529 let after = bounds.cursor().map_or(0, |cursor| cursor.sequence().get()); 530 let items = if position.sequence().get() <= after { 531 Vec::new() 532 } else { 533 let rows = sqlx::query( 534 "SELECT transport_id, target_fingerprint, observed_at_unix_ms, cursor 535 FROM radroots_runtime_event_provenance 536 WHERE event_id = ? 537 ORDER BY observed_at_unix_ms, transport_id, target_fingerprint, cursor 538 LIMIT ?", 539 ) 540 .bind(event_id.as_bytes().as_slice()) 541 .bind(i64::from(bounds.limit())) 542 .fetch_all(&self.pool) 543 .await 544 .map_err(map_backend)?; 545 rows.iter() 546 .map(|row| { 547 let cursor = row.try_get::<String, _>("cursor").map_err(map_corrupt)?; 548 StoredEventProvenance::from_stored_parts( 549 position, 550 row.try_get::<String, _>("transport_id") 551 .map_err(map_corrupt)? 552 .as_str(), 553 row.try_get::<String, _>("target_fingerprint") 554 .map_err(map_corrupt)? 555 .as_str(), 556 u64_from_i64(row.try_get("observed_at_unix_ms").map_err(map_corrupt)?)?, 557 (!cursor.is_empty()).then_some(cursor.as_str()), 558 ) 559 }) 560 .collect::<Result<Vec<_>, Error>>()? 561 }; 562 EventPage::new(self.generation, items, None, bounds) 563 }) 564 } 565 } 566 567 const fn stage_name(stage: AdmissionStage) -> &'static str { 568 match stage { 569 AdmissionStage::Raw => "raw", 570 AdmissionStage::Verified => "verified", 571 AdmissionStage::Visible => "visible", 572 } 573 } 574 575 fn admission_stage(value: String) -> Result<AdmissionStage, Error> { 576 match value.as_str() { 577 "raw" => Ok(AdmissionStage::Raw), 578 "verified" => Ok(AdmissionStage::Verified), 579 "visible" => Ok(AdmissionStage::Visible), 580 _ => Err(Error::CorruptStoredEvent), 581 } 582 } 583 584 fn source_generation(value: Vec<u8>) -> Result<SourceGeneration, Error> { 585 SourceGeneration::new(value.try_into().map_err(|_| Error::CorruptStoredEvent)?) 586 .map_err(|_| Error::CorruptStoredEvent) 587 } 588 589 fn event_sequence(value: i64) -> Result<EventSequence, Error> { 590 EventSequence::new(u64_from_i64(value)?).map_err(|_| Error::CorruptStoredEvent) 591 } 592 593 fn i64_from_u64(value: u64) -> Result<i64, Error> { 594 i64::try_from(value).map_err(|_| Error::CorruptStoredEvent) 595 } 596 597 fn u64_from_i64(value: i64) -> Result<u64, Error> { 598 u64::try_from(value).map_err(|_| Error::CorruptStoredEvent) 599 } 600 601 fn map_corrupt(_: sqlx::Error) -> Error { 602 Error::CorruptStoredEvent 603 } 604 605 #[cfg(test)] 606 #[cfg_attr(coverage_nightly, coverage(off))] 607 mod tests { 608 use super::*; 609 use crate::migration::runtime::{MIGRATIONS, migration_sql}; 610 use radroots_event::{ 611 SignedEvent, 612 admission::{AdmissionPolicy, RawEvent, VisibilityPolicy, VisibleEvent}, 613 wire::Nip01EventWire, 614 }; 615 use radroots_storage::event::EventQueryBounds; 616 use radroots_transport::{ 617 Target, TransportId, 618 source::{EventProvenance, FetchCursor, ObservedEvent}, 619 }; 620 use sqlx::sqlite::SqlitePoolOptions; 621 622 struct Allow; 623 624 impl radroots_event::admission::SignatureVerifier for Allow { 625 fn verify_signature( 626 &self, 627 _event: &radroots_event::Event, 628 ) -> Result<(), radroots_event::Error> { 629 Ok(()) 630 } 631 } 632 633 impl AdmissionPolicy for Allow { 634 type Error = core::convert::Infallible; 635 636 fn policy_id(&self) -> &'static str { 637 "test.storage-sqlite.admission.v1" 638 } 639 640 fn admit( 641 &self, 642 _event: &radroots_event::admission::ContractValidatedEvent, 643 ) -> Result<(), Self::Error> { 644 Ok(()) 645 } 646 } 647 648 impl VisibilityPolicy for Allow { 649 type Error = core::convert::Infallible; 650 651 fn policy_id(&self) -> &'static str { 652 "test.storage-sqlite.visibility.v1" 653 } 654 655 fn make_visible( 656 &self, 657 _event: &radroots_event::admission::AdmittedEvent, 658 ) -> Result<(), Self::Error> { 659 Ok(()) 660 } 661 } 662 663 async fn store(generation: SourceGeneration) -> SqliteStorage { 664 let pool = SqlitePoolOptions::new() 665 .max_connections(1) 666 .connect("sqlite::memory:") 667 .await 668 .expect("memory SQLite"); 669 sqlx::query("PRAGMA foreign_keys = ON") 670 .execute(&pool) 671 .await 672 .expect("foreign keys"); 673 for migration in MIGRATIONS { 674 sqlx::raw_sql(migration_sql(migration.version()).expect("registered SQL")) 675 .execute(&pool) 676 .await 677 .expect("runtime migration"); 678 } 679 sqlx::query( 680 "INSERT INTO radroots_runtime_source_generations ( 681 generation, sequence_head, state, created_at_unix_ms, retired_at_unix_ms 682 ) VALUES (?, 0, 'active', 1, NULL)", 683 ) 684 .bind(generation.as_bytes().as_slice()) 685 .execute(&pool) 686 .await 687 .expect("source generation"); 688 SqliteStorage::new(pool, generation, EventStoreMode::ReadWrite) 689 } 690 691 fn signed_event(content: &str, pretty: bool) -> SignedEvent { 692 signed_event_with(content, pretty, 1_800_000_100, 0, vec![]) 693 } 694 695 fn signed_event_with( 696 content: &str, 697 pretty: bool, 698 created_at: u64, 699 kind: u32, 700 tags: Vec<Vec<String>>, 701 ) -> SignedEvent { 702 let mut wire = Nip01EventWire { 703 id: "0".repeat(64), 704 pubkey: "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df".to_owned(), 705 created_at, 706 kind, 707 tags, 708 content: content.to_owned(), 709 sig: "42".repeat(64), 710 extra: Default::default(), 711 }; 712 wire.id = wire 713 .computed_event_id() 714 .expect("canonical event id") 715 .to_hex(); 716 let value = serde_json::json!({ 717 "id": &wire.id, 718 "pubkey": &wire.pubkey, 719 "created_at": wire.created_at, 720 "kind": wire.kind, 721 "tags": &wire.tags, 722 "content": &wire.content, 723 "sig": &wire.sig, 724 }); 725 let raw_json = if pretty { 726 serde_json::to_string_pretty(&value).expect("pretty event JSON") 727 } else { 728 value.to_string() 729 }; 730 SignedEvent::from_wire_verified_id(wire, raw_json).expect("signed event") 731 } 732 733 fn observed(event: SignedEvent, at: u64, cursor: Option<&str>) -> ObservedEvent { 734 let target = Target::new(TransportId::NOSTR, "wss://relay.example").expect("target"); 735 let mut provenance = 736 EventProvenance::new(TransportId::NOSTR, target.fingerprint().clone(), at) 737 .expect("provenance"); 738 if let Some(cursor) = cursor { 739 provenance = provenance.with_cursor(FetchCursor::parse(cursor).expect("cursor")); 740 } 741 ObservedEvent::new(event, provenance) 742 } 743 744 fn verified(event: &SignedEvent) -> radroots_event::VerifiedEvent { 745 RawEvent::new(event.envelope().clone()) 746 .verify_id() 747 .expect("event id") 748 .verify_signature(&Allow) 749 .expect("signature") 750 } 751 752 fn visible(event: &SignedEvent) -> VisibleEvent { 753 let verified = verified(event); 754 let validated = if event.envelope().kind_u32() == 5 { 755 verified 756 .validate_contract_for_admission("radroots.social.deletion_request.v1") 757 .expect("admission-selected contract") 758 } else { 759 verified.validate_contract().expect("contract") 760 }; 761 validated 762 .admit_with(&Allow) 763 .expect("admission") 764 .make_visible_with(&Allow) 765 .expect("visibility") 766 } 767 768 #[tokio::test] 769 async fn admission_is_idempotent_monotonic_and_conflict_safe() { 770 let generation = SourceGeneration::new([7; 32]).expect("generation"); 771 let store = store(generation).await; 772 let event = signed_event( 773 "{\"display_name\":\"Moss Street Farm\",\"bot\":false}", 774 false, 775 ); 776 777 let inserted = store 778 .admit(EventAdmission::raw(observed(event.clone(), 10, None))) 779 .await 780 .expect("insert raw"); 781 assert_eq!(inserted.disposition(), AdmissionDisposition::Inserted); 782 assert_eq!(inserted.position().sequence().get(), 1); 783 784 let advanced = store 785 .admit( 786 EventAdmission::visible( 787 observed(event.clone(), 11, Some("relay-page-1")), 788 visible(&event), 789 ) 790 .expect("visible admission"), 791 ) 792 .await 793 .expect("advance visible"); 794 assert_eq!(advanced.disposition(), AdmissionDisposition::Advanced); 795 assert_eq!(advanced.position(), inserted.position()); 796 let metadata = sqlx::query( 797 "SELECT admitted_contract_id, admitted_registry_version 798 FROM radroots_runtime_events WHERE event_id = ?", 799 ) 800 .bind(event.id().as_bytes().as_slice()) 801 .fetch_one(&store.pool) 802 .await 803 .expect("event contract metadata"); 804 assert_eq!( 805 metadata.get::<String, _>("admitted_contract_id"), 806 "radroots.profile.metadata.v1" 807 ); 808 assert_eq!(metadata.get::<i64, _>("admitted_registry_version"), 7); 809 810 let duplicate = store 811 .admit( 812 EventAdmission::visible( 813 observed(event.clone(), 11, Some("relay-page-1")), 814 visible(&event), 815 ) 816 .expect("visible admission"), 817 ) 818 .await 819 .expect("duplicate visible"); 820 assert_eq!(duplicate.disposition(), AdmissionDisposition::Duplicate); 821 assert_eq!( 822 store 823 .admit(EventAdmission::raw(observed(event.clone(), 12, None))) 824 .await, 825 Err(Error::AdmissionRegression) 826 ); 827 828 let same_id_different_bytes = signed_event( 829 "{\"display_name\":\"Moss Street Farm\",\"bot\":false}", 830 true, 831 ); 832 assert_eq!(same_id_different_bytes.id(), event.id()); 833 assert_eq!( 834 store 835 .admit(EventAdmission::raw(observed( 836 same_id_different_bytes, 837 13, 838 None, 839 ))) 840 .await, 841 Err(Error::EventConflict) 842 ); 843 } 844 845 #[tokio::test] 846 async fn queries_preserve_stage_bounds_cursors_and_exact_provenance() { 847 let generation = SourceGeneration::new([8; 32]).expect("generation"); 848 let store = store(generation).await; 849 let empty_status = store.status().await.expect("empty status"); 850 assert_eq!(empty_status.raw_events(), 0); 851 assert_eq!(empty_status.verified_events(), 0); 852 assert_eq!(empty_status.visible_events(), 0); 853 let raw_event = signed_event("raw", false); 854 let visible_event = 855 signed_event("{\"display_name\":\"Visible Farm\",\"bot\":false}", false); 856 store 857 .admit(EventAdmission::raw(observed(raw_event.clone(), 20, None))) 858 .await 859 .expect("raw event"); 860 store 861 .admit( 862 EventAdmission::visible( 863 observed(visible_event.clone(), 21, Some("relay-page-2")), 864 visible(&visible_event), 865 ) 866 .expect("visible admission"), 867 ) 868 .await 869 .expect("visible event"); 870 871 let first = store 872 .query_raw(EventQuery::all(EventQueryBounds::first(1).expect("bounds"))) 873 .await 874 .expect("first page"); 875 assert_eq!(first.items().len(), 1); 876 assert_eq!(first.items()[0].event(), &raw_event); 877 let next = first.next_cursor().expect("continuation cursor"); 878 let second = store 879 .query_raw(EventQuery::all( 880 EventQueryBounds::first(1).expect("bounds").after(next), 881 )) 882 .await 883 .expect("second page"); 884 assert_eq!(second.items()[0].event(), &visible_event); 885 assert!(second.next_cursor().is_none()); 886 887 let verified_page = store 888 .query_verified(EventQuery::all( 889 EventQueryBounds::first(10).expect("bounds"), 890 )) 891 .await 892 .expect("verified page"); 893 assert_eq!(verified_page.items().len(), 1); 894 assert_eq!(verified_page.items()[0].event(), &visible_event); 895 let visible_page = store 896 .query_visible( 897 EventQuery::for_ids( 898 EventQueryBounds::first(10).expect("bounds"), 899 vec![*visible_event.id()], 900 ) 901 .expect("id query"), 902 ) 903 .await 904 .expect("visible page"); 905 assert_eq!(visible_page.items().len(), 1); 906 907 let provenance = store 908 .query_provenance( 909 *visible_event.id(), 910 EventQueryBounds::first(10).expect("bounds"), 911 ) 912 .await 913 .expect("provenance"); 914 assert_eq!(provenance.items().len(), 1); 915 assert_eq!(provenance.items()[0].provenance().observed_at_unix_ms(), 21); 916 assert_eq!( 917 provenance.items()[0] 918 .provenance() 919 .cursor() 920 .expect("cursor") 921 .as_str(), 922 "relay-page-2" 923 ); 924 925 let status = store.status().await.expect("status"); 926 assert_eq!(status.raw_events(), 2); 927 assert_eq!(status.verified_events(), 1); 928 assert_eq!(status.visible_events(), 1); 929 let foreign_cursor = EventPosition::new( 930 SourceGeneration::new([9; 32]).expect("foreign generation"), 931 EventSequence::new(1).expect("sequence"), 932 ); 933 assert_eq!( 934 store 935 .query_raw(EventQuery::all( 936 EventQueryBounds::first(1) 937 .expect("bounds") 938 .after(foreign_cursor), 939 )) 940 .await, 941 Err(Error::SourceGenerationChanged) 942 ); 943 } 944 945 #[tokio::test] 946 async fn verified_replacement_and_same_count_advance_match_memory_after_reopen() { 947 let generation = SourceGeneration::new([19; 32]).unwrap(); 948 let store = store(generation).await; 949 let memory = radroots_storage::memory::MemoryStorage::new(generation); 950 let old = signed_event_with( 951 r#"{"display_name":"Old Farm","bot":false}"#, 952 false, 953 10, 954 0, 955 vec![], 956 ); 957 let newer = signed_event_with("malformed profile", false, 20, 0, vec![]); 958 let admissions = [ 959 EventAdmission::visible(observed(old.clone(), 10, None), visible(&old)).unwrap(), 960 EventAdmission::raw(observed(newer.clone(), 20, None)), 961 EventAdmission::verified(observed(newer.clone(), 20, None), verified(&newer)).unwrap(), 962 ]; 963 let mut previous_digest = None; 964 for (index, admission) in admissions.into_iter().enumerate() { 965 assert_eq!( 966 store.admit(admission.clone()).await.unwrap(), 967 memory.admit(admission).await.unwrap() 968 ); 969 let reopened = 970 SqliteStorage::new(store.pool.clone(), generation, EventStoreMode::ReadWrite); 971 let snapshot = reopened.rebuild_visibility().await.unwrap(); 972 assert_eq!(snapshot, memory.rebuild_visibility().await.unwrap()); 973 let query = EventQuery::all(EventQueryBounds::first(10).unwrap()); 974 assert_eq!( 975 reopened.query_visible(query.clone()).await.unwrap(), 976 memory.query_visible(query).await.unwrap() 977 ); 978 if index == 2 { 979 assert_ne!(previous_digest, Some(snapshot.digest())); 980 assert!(snapshot.visible_event_ids().is_empty()); 981 assert_eq!(snapshot.current_heads()[0].event_id, *newer.id()); 982 assert_eq!(reopened.status().await.unwrap().raw_events(), 2); 983 } else { 984 assert_eq!(snapshot.visible_event_ids(), &[*old.id()]); 985 } 986 previous_digest = Some(snapshot.digest()); 987 } 988 } 989 990 #[tokio::test] 991 async fn visibility_rebuild_survives_reopen_and_matches_current_head_deletion_queries() { 992 let generation = SourceGeneration::new([18; 32]).expect("generation"); 993 let store = store(generation).await; 994 let old = signed_event_with( 995 r#"{"display_name":"Old Farm","bot":false}"#, 996 false, 997 1_800_000_100, 998 0, 999 vec![], 1000 ); 1001 let current = signed_event_with( 1002 r#"{"display_name":"Current Farm","bot":false}"#, 1003 false, 1004 1_800_000_200, 1005 0, 1006 vec![], 1007 ); 1008 let deletion = signed_event_with( 1009 "retired profile", 1010 false, 1011 1_800_000_300, 1012 5, 1013 vec![vec!["e".to_owned(), current.id().to_hex()]], 1014 ); 1015 for (event, observed_at) in [ 1016 (old.clone(), 100), 1017 (current.clone(), 200), 1018 (deletion.clone(), 300), 1019 ] { 1020 store 1021 .admit( 1022 EventAdmission::visible( 1023 observed(event.clone(), observed_at, None), 1024 visible(&event), 1025 ) 1026 .expect("visible admission"), 1027 ) 1028 .await 1029 .expect("admit visible event"); 1030 } 1031 1032 let before = store 1033 .rebuild_visibility() 1034 .await 1035 .expect("visibility rebuild"); 1036 let reopened = 1037 SqliteStorage::new(store.pool.clone(), generation, EventStoreMode::ReadWrite); 1038 let after = reopened 1039 .rebuild_visibility() 1040 .await 1041 .expect("reopened visibility rebuild"); 1042 assert_eq!(before, after); 1043 assert_eq!(after.current_heads()[0].event_id, *current.id()); 1044 assert_eq!(after.visible_event_ids(), &[*deletion.id()]); 1045 assert_eq!(after.suppressed_event_ids(), &[*current.id()]); 1046 assert_eq!(after.superseded_event_ids(), &[*old.id()]); 1047 let page = reopened 1048 .query_visible(EventQuery::all( 1049 EventQueryBounds::first(10).expect("bounds"), 1050 )) 1051 .await 1052 .expect("visible page"); 1053 assert_eq!(page.items().len(), 1); 1054 assert_eq!(page.items()[0].event().id(), deletion.id()); 1055 assert_eq!(reopened.status().await.expect("status").visible_events(), 1); 1056 } 1057 1058 #[tokio::test] 1059 async fn corrupt_rows_fail_closed_and_source_history_is_immutable() { 1060 let generation = SourceGeneration::new([10; 32]).expect("generation"); 1061 let store = store(generation).await; 1062 let event = signed_event( 1063 "{\"display_name\":\"Corruption Probe\",\"bot\":false}", 1064 false, 1065 ); 1066 store 1067 .admit(EventAdmission::raw(observed(event.clone(), 30, None))) 1068 .await 1069 .expect("event"); 1070 1071 let wrong_generation = SourceGeneration::new([11; 32]).expect("generation"); 1072 let mismatched_store = SqliteStorage::new( 1073 store.pool.clone(), 1074 wrong_generation, 1075 EventStoreMode::ReadWrite, 1076 ); 1077 assert_eq!( 1078 mismatched_store 1079 .admit(EventAdmission::raw(observed(event, 31, None))) 1080 .await, 1081 Err(Error::CorruptStoredEvent) 1082 ); 1083 1084 assert!( 1085 sqlx::query("DELETE FROM radroots_runtime_events") 1086 .execute(&store.pool) 1087 .await 1088 .is_err() 1089 ); 1090 assert!( 1091 sqlx::query("DELETE FROM radroots_runtime_source_generations") 1092 .execute(&store.pool) 1093 .await 1094 .is_err() 1095 ); 1096 sqlx::query("DROP TRIGGER radroots_runtime_events_contract_metadata_guard") 1097 .execute(&store.pool) 1098 .await 1099 .expect("drop metadata guard for corruption probe"); 1100 sqlx::query( 1101 "UPDATE radroots_runtime_events 1102 SET admitted_contract_id = 'radroots.social.geochat.v1'", 1103 ) 1104 .execute(&store.pool) 1105 .await 1106 .expect("forge contract metadata"); 1107 assert_eq!( 1108 store 1109 .query_raw(EventQuery::all(EventQueryBounds::first(1).expect("bounds"),)) 1110 .await, 1111 Err(Error::CorruptStoredEvent) 1112 ); 1113 sqlx::query( 1114 "UPDATE radroots_runtime_events 1115 SET admitted_contract_id = 'radroots.social.geochat.v1', 1116 admitted_registry_version = 7", 1117 ) 1118 .execute(&store.pool) 1119 .await 1120 .expect("forge mismatched selected contract"); 1121 assert_eq!( 1122 store 1123 .query_raw(EventQuery::all(EventQueryBounds::first(1).expect("bounds"),)) 1124 .await, 1125 Err(Error::CorruptStoredEvent) 1126 ); 1127 sqlx::query( 1128 "UPDATE radroots_runtime_events 1129 SET admitted_registry_version = 6", 1130 ) 1131 .execute(&store.pool) 1132 .await 1133 .expect("forge registry version"); 1134 assert_eq!( 1135 store 1136 .query_raw(EventQuery::all(EventQueryBounds::first(1).expect("bounds"),)) 1137 .await, 1138 Err(Error::CorruptStoredEvent) 1139 ); 1140 sqlx::query("PRAGMA ignore_check_constraints = ON") 1141 .execute(&store.pool) 1142 .await 1143 .expect("disable checks for corruption probe"); 1144 sqlx::query("UPDATE radroots_runtime_events SET admission_stage = 'corrupt'") 1145 .execute(&store.pool) 1146 .await 1147 .expect("forge corrupt stage"); 1148 assert_eq!( 1149 store 1150 .query_raw(EventQuery::all(EventQueryBounds::first(1).expect("bounds"),)) 1151 .await, 1152 Err(Error::CorruptStoredEvent) 1153 ); 1154 } 1155 }