pub mod backfill_progress; pub mod collection_index; pub mod error; pub mod firehose_cursor; pub mod list_hosts_cursor; pub mod meta; pub mod pds_host; pub mod repo; pub mod resync_buffer; pub mod resync_queue; pub(crate) use error::{StorageError, StorageResult}; pub(crate) use meta::StatsRef; pub(crate) use repo::Account; // --------------------------------------------------------------------------- // key prefixes // // could be a single byte each really, but making them slightly human-readable // is nice. fjall's key compression probably means we could make these full // words, but they feel better short. // --------------------------------------------------------------------------- /// Fixed-length (3 byte) key prefix per data type type KeyPrefix = [u8; 3]; /// Main collection index (collection → did). See [`collection_index`]. pub(super) const PREFIX_RBC: KeyPrefix = *b"rbc"; /// Reversed collection index (did → collection). See [`collection_index`]. pub(super) const PREFIX_CBR: KeyPrefix = *b"cbr"; /// Per-repo state and account status. See [`repo`]. pub(super) const PREFIX_REPO: KeyPrefix = *b"rep"; /// Per-repo transient sync state (rev + prevData CID). See [`repo`]. pub(super) const PREFIX_REPO_PREV: KeyPrefix = *b"rev"; /// Firehose subscription cursor (per relay host). See [`firehose_cursor`]. pub(super) const PREFIX_SUBSCRIBE_REPOS: KeyPrefix = *b"sub"; /// listRepos backfill walk progress (per relay host). See [`backfill_progress`]. pub(super) const PREFIX_LIST_REPOS: KeyPrefix = *b"lsr"; /// Timestamp-ordered resync work queue. See [`resync_queue`]. pub(super) const PREFIX_RESYNC_QUEUE: KeyPrefix = *b"rsq"; /// Per-repo buffered firehose events during resync. See [`resync_buffer`]. pub(super) const PREFIX_RESYNC_BUFFER: KeyPrefix = *b"rsb"; /// Per-PDS host state (sync1.1 mode, trust, listRepos cursor/done). See [`pds_host`]. pub(super) const PREFIX_PDS_HOST: KeyPrefix = *b"pdh"; /// listHosts walk cursor (per upstream relay host). See [`list_hosts_cursor`]. pub(super) const PREFIX_LIST_HOSTS: KeyPrefix = *b"lhs"; /// Persistent system stats and cardinality sketches. See [`meta`]. pub(super) const PREFIX_META: KeyPrefix = *b"met"; use std::path::Path; use std::sync::Arc; /// Shared handle to the fjall database and its per-concern keyspaces. /// /// In fjall 3.x, `Database` is the top-level multi-keyspace container and /// `Keyspace` is an individual column-family (the old `PartitionHandle`). pub struct Db { pub(crate) database: fjall::Database, /// General-purpose keyspace: repo state, queues, cursors, etc. pub(crate) ks: fjall::Keyspace, /// Collection index keyspace: rbc + cbr ranges. /// /// Tuned for scan-heavy access: 64 KiB blocks (amortises per-block overhead /// across sequential reads) and Lz4 compression at all levels (higher /// on-disk density means more data fits in the block cache). pub(crate) index_ks: fjall::Keyspace, /// Persistent system stats and cardinality sketches, loaded on open. pub(crate) stats: StatsRef, } /// Point-in-time snapshot of fjall storage stats. pub struct StorageStats { /// Total disk space used by the database (journals + all keyspaces). pub disk_bytes: u64, /// Disk space used by the default keyspace (repo state, queues, cursors). pub default_ks_disk_bytes: u64, /// Disk space used by the index keyspace (rbc + cbr). pub index_ks_disk_bytes: u64, /// Number of journal files on disk. pub journal_count: usize, /// Number of compactions currently running. pub active_compactions: usize, /// Total compactions completed since the database was opened. pub compactions_completed: usize, /// Total time spent compacting since the database was opened. pub time_compacting: std::time::Duration, } impl Db { /// Flush the journal write buffer to the OS page cache. /// /// With `manual_journal_persist` enabled, individual writes skip the flush; /// call this periodically to batch many writes into a single flush. /// Does NOT fsync — crash recovery may lose up to one flush interval of /// writes, which is acceptable since all data can be re-fetched. pub fn persist_journal(&self) -> StorageResult<()> { self.database.persist(fjall::PersistMode::Buffer)?; Ok(()) } /// Collect a snapshot of fjall storage stats. pub fn storage_stats(&self) -> StorageStats { StorageStats { disk_bytes: self.database.disk_space().unwrap_or(0), default_ks_disk_bytes: self.ks.disk_space(), index_ks_disk_bytes: self.index_ks.disk_space(), journal_count: self.database.journal_count(), active_compactions: self.database.active_compactions(), compactions_completed: self.database.compactions_completed(), time_compacting: self.database.time_compacting(), } } } /// Cheaply-cloneable reference to the shared database. pub type DbRef = Arc; /// Open (or create) the fjall database at `path` and return a shared handle. /// /// `worker_threads`: number of fjall background threads for flush + compaction. /// `None` uses fjall's own default (`min(cores, 4)`). pub fn open(path: &Path, cache_mb: u64, worker_threads: Option) -> StorageResult { open_inner( path, DbConfig::ForReal { cache_mb, worker_threads, }, ) } enum DbConfig { /// temporary db for tests #[allow(dead_code)] Testing, /// bumpable cache for prod ForReal { cache_mb: u64, worker_threads: Option, }, } /// Open a temporary database that deletes itself on drop. For tests only. #[cfg(test)] pub(crate) fn open_temporary() -> StorageResult { use std::sync::atomic::{AtomicU64, Ordering}; static COUNTER: AtomicU64 = AtomicU64::new(0); let n = COUNTER.fetch_add(1, Ordering::Relaxed); let path = std::env::temp_dir().join(format!("lightrail-test-{}-{}", std::process::id(), n)); open_inner(&path, DbConfig::Testing) } fn open_inner(path: &Path, config: DbConfig) -> StorageResult { let builder = fjall::Database::builder(path).manual_journal_persist(true); let builder = match config { DbConfig::Testing => builder.temporary(true), DbConfig::ForReal { cache_mb, worker_threads, } => { let b = builder.cache_size(cache_mb * 2_u64.pow(20)); if let Some(n) = worker_threads { b.worker_threads(n) } else { b } } }; let database = builder.open()?; let ks = database.keyspace("default", fjall::KeyspaceCreateOptions::default)?; let index_ks = database.keyspace("index", || { fjall::KeyspaceCreateOptions::default() .data_block_size_policy(fjall::config::BlockSizePolicy::all(64 * 1_024)) .data_block_compression_policy(fjall::config::CompressionPolicy::all( fjall::CompressionType::Lz4, )) })?; let stats = meta::load(&ks)?; Ok(Arc::new(Db { database, ks, index_ks, stats, })) }