diff --git a/src/api/debug.rs b/src/api/debug.rs index 4c60097..fca5723 100644 --- a/src/api/debug.rs +++ b/src/api/debug.rs @@ -60,8 +60,6 @@ pub async fn handle_debug_count( .await .map_err(|_| StatusCode::BAD_REQUEST)?; - let prefix = keys::record_prefix_collection(&did, &req.collection); - let count = state .db .run(move |db| { @@ -69,6 +67,11 @@ pub async fn handle_debug_count( .keyspace_by_name("records") .expect("records keyspace exists in indexer mode"); + // lookup, not key_for: a debug read must not intern a collection + let prefix = keys::record_prefix_collection( + &did, + db.indexer.collections.lookup(&req.collection), + ); let start_key = prefix.clone(); let mut end_key = prefix.clone(); if let Some(msg) = end_key.last_mut() { diff --git a/src/backfill/sparse.rs b/src/backfill/sparse.rs index e31ece4..e7617ba 100644 --- a/src/backfill/sparse.rs +++ b/src/backfill/sparse.rs @@ -445,22 +445,13 @@ async fn persist_sparse_backfill( if !ephemeral { for guard in db.indexer.record_prefix(&prefix) { let (key, cid_bytes) = guard.into_inner().into_diagnostic()?; - let mut remaining = key[prefix.len()..].splitn(2, |b| keys::SEP.eq(b)); - let collection_raw = remaining - .next() - .ok_or_else(|| miette::miette!("invalid record key format: {key:?}"))?; - let rkey_raw = remaining - .next() - .ok_or_else(|| miette::miette!("invalid record key format: {key:?}"))?; - - let collection = std::str::from_utf8(collection_raw) - .map_err(|e| miette::miette!("invalid collection utf8: {e}"))?; + let (collection_key, rkey) = keys::parse_record_suffix(&key[prefix.len()..]) + .map_err(|e| miette::miette!("invalid record key '{key:?}': {e}"))?; + let collection = db.indexer.collections.text(collection_key)?; + let collection = collection.as_str(); if !filter.matches_collection(collection) { continue; } - - let rkey = keys::parse_rkey(rkey_raw) - .map_err(|e| miette::miette!("invalid rkey '{key:?}' for {did}: {e}"))?; let cid = cid::Cid::read_bytes(cid_bytes.as_ref()) .map_err(|e| miette::miette!("invalid cid '{cid_bytes:?}' for {did}: {e}"))? .to_smolstr(); diff --git a/src/backfill/worker/process.rs b/src/backfill/worker/process.rs index 177b969..5ae2d61 100644 --- a/src/backfill/worker/process.rs +++ b/src/backfill/worker/process.rs @@ -328,21 +328,12 @@ pub(super) async fn process_did( if !ephemeral { for guard in app_state.db.indexer.record_prefix(&prefix) { let (key, cid_bytes) = guard.into_inner().into_diagnostic()?; - // key is did|collection|rkey - // skip did| - let mut remaining = key[prefix.len()..].splitn(2, |b| keys::SEP.eq(b)); - let collection_raw = remaining - .next() - .ok_or_else(|| miette::miette!("invalid record key format: {key:?}"))?; - let rkey_raw = remaining - .next() - .ok_or_else(|| miette::miette!("invalid record key format: {key:?}"))?; - - let collection = std::str::from_utf8(collection_raw) - .map_err(|e| miette::miette!("invalid collection utf8: {e}"))?; - - let rkey = keys::parse_rkey(rkey_raw) - .map_err(|e| miette::miette!("invalid rkey '{key:?}' for {did}: {e}"))?; + // key is did|{collection}{rkey}; the did prefix is skipped by + // length and the collection segment carries its own terminator + let (collection_key, rkey) = keys::parse_record_suffix(&key[prefix.len()..]) + .map_err(|e| miette::miette!("invalid record key '{key:?}': {e}"))?; + let collection = app_state.db.indexer.collections.text(collection_key)?; + let collection = collection.as_str(); let cid = cid::Cid::read_bytes(cid_bytes.as_ref()) .map_err(|e| miette::miette!("invalid cid '{cid_bytes:?}' for {did}: {e}"))? diff --git a/src/backlinks/store.rs b/src/backlinks/store.rs index aa1e225..6c2a2fe 100644 --- a/src/backlinks/store.rs +++ b/src/backlinks/store.rs @@ -225,7 +225,8 @@ pub fn lookup_cid( ) -> Option { let did = Did::new(did).ok()?; let db_rkey = DbRkey::new(rkey); - let record_key = keys::record_key(&did, collection, &db_rkey); + // lookup, not key_for: resolving a backlink must not intern a collection + let record_key = keys::record_key(&did, indexer.collections.lookup(collection), &db_rkey); let cid_bytes = indexer.record(record_key).ok()??; Cid::new(&cid_bytes).ok().map(|c| c.as_str().to_string()) } diff --git a/src/config.rs b/src/config.rs index e14a32b..331ec75 100644 --- a/src/config.rs +++ b/src/config.rs @@ -279,6 +279,15 @@ pub struct Config { /// whether bloom filters are enabled for the records keyspace. /// set via `HYDRANT_DB_RECORDS_BLOOM_FILTERS`. pub db_records_bloom_filters: bool, + /// memory budget for the in-memory collection dictionary, in bytes. + /// + /// NSIDs are minted by whoever writes a record, so the dictionary is sized + /// by untrusted input and never evicted. past this budget new collections + /// keep their NSID text in keys instead of being interned, which costs the + /// pre-interning key size and nothing else. the network has ~26k + /// collections, using roughly 2 MB. + /// set via `HYDRANT_COLLECTION_DICT_MAX_MB`. + pub collection_dict_max_bytes: usize, /// replay batch size. /// @@ -373,6 +382,7 @@ impl Default for Config { db_events_memtable_size_mb: BASE_MEMTABLE_MB, db_records_memtable_size_mb: BASE_MEMTABLE_MB / 3 * 2, db_records_bloom_filters: false, + collection_dict_max_bytes: 32 * 1024 * 1024, stream_replay_chunk_size: 0, stream_replay_chunk_pause: Duration::ZERO, stream_pending_event_limit: 4096, @@ -506,6 +516,11 @@ impl fmt::Display for Config { format_args!("{} mb", self.db_records_memtable_size_mb) )?; config_line!(f, "db records bloom filters", self.db_records_bloom_filters)?; + config_line!( + f, + "collection dict budget", + format_args!("{} mb", self.collection_dict_max_bytes / (1024 * 1024)) + )?; let replay_chunk = if self.stream_replay_chunk_size == 0 { "auto".to_owned() } else { diff --git a/src/config/env.rs b/src/config/env.rs index b69ca93..bb347e5 100644 --- a/src/config/env.rs +++ b/src/config/env.rs @@ -195,6 +195,11 @@ impl Config { "DB_REPOS_MEMTABLE_SIZE_MB", defaults.db_repos_memtable_size_mb ); + let collection_dict_max_mb: usize = cfg!( + "COLLECTION_DICT_MAX_MB", + defaults.collection_dict_max_bytes / (1024 * 1024) + ); + let collection_dict_max_bytes = collection_dict_max_mb * 1024 * 1024; let stream_replay_chunk_size = cfg!( "STREAM_REPLAY_CHUNK_SIZE", defaults.stream_replay_chunk_size @@ -397,6 +402,7 @@ impl Config { db_events_memtable_size_mb, db_records_memtable_size_mb, db_records_bloom_filters, + collection_dict_max_bytes, stream_replay_chunk_size, stream_replay_chunk_pause, stream_pending_event_limit, diff --git a/src/control/repos/indexer.rs b/src/control/repos/indexer.rs index 2a307ae..b65c50a 100644 --- a/src/control/repos/indexer.rs +++ b/src/control/repos/indexer.rs @@ -164,21 +164,24 @@ impl ReposControl { ) -> Result>> { let dids: Vec> = dids.into_iter().map(|d| d.into_static()).collect(); - let queued = self.0.db.run(move |db| { - let mut txn = crate::db::Txn::new(db); - let mut queued: Vec> = Vec::new(); + let queued = self + .0 + .db + .run(move |db| { + let mut txn = crate::db::Txn::new(db); + let mut queued: Vec> = Vec::new(); - for did in dids { - if Self::_resync(db, &did, &mut txn)? { - queued.push(did); + for did in dids { + if Self::_resync(db, &did, &mut txn)? { + queued.push(did); + } } - } - txn.commit()?; - db.persist()?; - Ok(queued) - }) - .await?; + txn.commit()?; + db.persist()?; + Ok(queued) + }) + .await?; if !queued.is_empty() { self.0.notify_backfill(); @@ -198,48 +201,54 @@ impl ReposControl { ) -> Result>> { let dids: Vec> = dids.into_iter().map(|d| d.into_static()).collect(); - let queued = self.0.db.run(move |db| { - let mut txn = crate::db::Txn::new(db); - let mut queued: Vec> = Vec::new(); + let queued = self + .0 + .db + .run(move |db| { + let mut txn = crate::db::Txn::new(db); + let mut queued: Vec> = Vec::new(); - for did in dids { - let did_key = keys::repo_key(&did); - let metadata_key = keys::repo_metadata_key(&did); + for did in dids { + let did_key = keys::repo_key(&did); + let metadata_key = keys::repo_metadata_key(&did); - let metadata_bytes = db.repo_metadata.get(&metadata_key).into_diagnostic()?; - let existing_metadata = metadata_bytes - .map(|b| crate::db::deser_repo_meta(&b)) - .transpose()?; + let metadata_bytes = db.repo_metadata.get(&metadata_key).into_diagnostic()?; + let existing_metadata = metadata_bytes + .map(|b| crate::db::deser_repo_meta(&b)) + .transpose()?; - if let Some(metadata) = existing_metadata { - if !metadata.tracked && Self::_resync(db, &did, &mut txn)? { + if let Some(metadata) = existing_metadata { + if !metadata.tracked && Self::_resync(db, &did, &mut txn)? { + queued.push(did); + } + } else { + let repo_state = RepoState::backfilling(); + let metadata = RepoMetadata::backfilling(rand::random()); + txn.batch.insert( + &db.repos, + &did_key, + crate::db::ser_repo_state(&repo_state)?, + ); + txn.batch.insert( + &db.repo_metadata, + &metadata_key, + crate::db::ser_repo_meta(&metadata)?, + ); + txn.batch.insert( + &db.indexer.pending, + keys::pending_key(metadata.index_id), + &did_key, + ); + txn.counts.add_repos(1); + txn.transition_lifecycle(&did, GaugeState::Pending)?; queued.push(did); } - } else { - let repo_state = RepoState::backfilling(); - let metadata = RepoMetadata::backfilling(rand::random()); - txn.batch - .insert(&db.repos, &did_key, crate::db::ser_repo_state(&repo_state)?); - txn.batch.insert( - &db.repo_metadata, - &metadata_key, - crate::db::ser_repo_meta(&metadata)?, - ); - txn.batch.insert( - &db.indexer.pending, - keys::pending_key(metadata.index_id), - &did_key, - ); - txn.counts.add_repos(1); - txn.transition_lifecycle(&did, GaugeState::Pending)?; - queued.push(did); } - } - txn.commit()?; - db.persist()?; - Ok(queued) - }) - .await?; + txn.commit()?; + db.persist()?; + Ok(queued) + }) + .await?; self.0.notify_backfill(); Ok(queued) @@ -254,49 +263,53 @@ impl ReposControl { ) -> Result>> { let dids: Vec> = dids.into_iter().map(|d| d.into_static()).collect(); - let untracked = self.0.db.run(move |db| { - let mut txn = crate::db::Txn::new(db); - let mut untracked: Vec> = Vec::new(); - - for did in dids { - let did_key = keys::repo_key(&did); - let metadata_key = keys::repo_metadata_key(&did); - - let repo_bytes = db.repos.get(&did_key).into_diagnostic()?; - let existing = repo_bytes - .as_deref() - .map(db::deser_repo_state) - .transpose()?; - - if existing.is_some() { - let metadata_bytes = db.repo_metadata.get(&metadata_key).into_diagnostic()?; - let existing_metadata = metadata_bytes - .map(|b| crate::db::deser_repo_meta(&b)) + let untracked = self + .0 + .db + .run(move |db| { + let mut txn = crate::db::Txn::new(db); + let mut untracked: Vec> = Vec::new(); + + for did in dids { + let did_key = keys::repo_key(&did); + let metadata_key = keys::repo_metadata_key(&did); + + let repo_bytes = db.repos.get(&did_key).into_diagnostic()?; + let existing = repo_bytes + .as_deref() + .map(db::deser_repo_state) .transpose()?; - if let Some(mut metadata) = existing_metadata - && metadata.tracked - { - metadata.tracked = false; - txn.batch.insert( - &db.repo_metadata, - &metadata_key, - crate::db::ser_repo_meta(&metadata)?, - ); - txn.batch - .remove(&db.indexer.pending, keys::pending_key(metadata.index_id)); - txn.batch.remove(&db.indexer.resync, &did_key); - txn.transition_lifecycle(&did, GaugeState::Synced)?; - untracked.push(did); + if existing.is_some() { + let metadata_bytes = + db.repo_metadata.get(&metadata_key).into_diagnostic()?; + let existing_metadata = metadata_bytes + .map(|b| crate::db::deser_repo_meta(&b)) + .transpose()?; + + if let Some(mut metadata) = existing_metadata + && metadata.tracked + { + metadata.tracked = false; + txn.batch.insert( + &db.repo_metadata, + &metadata_key, + crate::db::ser_repo_meta(&metadata)?, + ); + txn.batch + .remove(&db.indexer.pending, keys::pending_key(metadata.index_id)); + txn.batch.remove(&db.indexer.resync, &did_key); + txn.transition_lifecycle(&did, GaugeState::Synced)?; + untracked.push(did); + } } } - } - txn.commit()?; - db.persist()?; - Ok(untracked) - }) - .await?; + txn.commit()?; + db.persist()?; + Ok(untracked) + }) + .await?; Ok(untracked) } @@ -310,35 +323,40 @@ impl<'i> RepoHandle<'i> { } let did = self.did.clone().into_static(); - let db_key = keys::record_key(&did, collection, &DbRkey::new(rkey)); - let collection = collection.to_smolstr(); - self.state.db.run(move |db| { - use miette::WrapErr; - - let cid_bytes = db.indexer.record(db_key).into_diagnostic()?; - let Some(cid_bytes) = cid_bytes else { - return Ok(None); - }; - - // lookup block using col|cid key - let block_key = keys::block_key(&collection, &cid_bytes); - let Some(block_bytes) = db.indexer.block(block_key).into_diagnostic()? else { - miette::bail!("block {cid_bytes:?} not found, this is a bug!!"); - }; - - let value = serde_ipld_dagcbor::from_slice::(&block_bytes) - .into_diagnostic() - .wrap_err("cant parse block")? - .into_static(); - let cid = Cid::new(&cid_bytes) - .into_diagnostic() - .wrap_err("cant parse block cid")?; - let cid = Cid::Str(cid.to_cowstr().into_static()); - - Ok(Some(Record { did, cid, value })) - }) - .await + let rkey = DbRkey::new(rkey); + self.state + .db + .run(move |db| { + use miette::WrapErr; + + // lookup, not key_for: a read must never intern a collection it has + // not seen, or a stream of misses would grow the dictionary + let db_key = + keys::record_key(&did, db.indexer.collections.lookup(&collection), &rkey); + let cid_bytes = db.indexer.record(db_key).into_diagnostic()?; + let Some(cid_bytes) = cid_bytes else { + return Ok(None); + }; + + // lookup block using col|cid key + let block_key = keys::block_key(&collection, &cid_bytes); + let Some(block_bytes) = db.indexer.block(block_key).into_diagnostic()? else { + miette::bail!("block {cid_bytes:?} not found, this is a bug!!"); + }; + + let value = serde_ipld_dagcbor::from_slice::(&block_bytes) + .into_diagnostic() + .wrap_err("cant parse block")? + .into_static(); + let cid = Cid::new(&cid_bytes) + .into_diagnostic() + .wrap_err("cant parse block cid")?; + let cid = Cid::Str(cid.to_cowstr().into_static()); + + Ok(Some(Record { did, cid, value })) + }) + .await } /// lists records from this repository. @@ -354,88 +372,92 @@ impl<'i> RepoHandle<'i> { } let did = self.did.clone().into_static(); - let prefix = keys::record_prefix_collection(&did, collection); let collection = collection.to_smolstr(); let cursor = cursor.map(|c| c.to_smolstr()); - self.state.db.run(move |db| { - let mut results = Vec::new(); - let mut next_cursor = None; + self.state + .db + .run(move |db| { + let prefix = keys::record_prefix_collection( + &did, + db.indexer.collections.lookup(&collection), + ); + let mut results = Vec::new(); + let mut next_cursor = None; - let iter: Box> = if !reverse { - let mut end_prefix = prefix.clone(); - if let Some(last) = end_prefix.last_mut() { - *last += 1; - } + let iter: Box> = if !reverse { + let mut end_prefix = prefix.clone(); + if let Some(last) = end_prefix.last_mut() { + *last += 1; + } - let end_key = if let Some(cursor) = &cursor { - let mut k = prefix.clone(); - let rkey = DbRkey::new(cursor); - keys::write_rkey(&mut k, &rkey); - k + let end_key = if let Some(cursor) = &cursor { + let mut k = prefix.clone(); + let rkey = DbRkey::new(cursor); + keys::write_rkey(&mut k, &rkey); + k + } else { + end_prefix + }; + + Box::new( + db.indexer + .record_range(prefix.as_slice()..end_key.as_slice()) + .rev(), + ) } else { - end_prefix + let start_key = if let Some(cursor) = &cursor { + let mut k = prefix.clone(); + let rkey = DbRkey::new(cursor); + keys::write_rkey(&mut k, &rkey); + k.push(0); + k + } else { + prefix.clone() + }; + + Box::new(db.indexer.record_range(start_key.as_slice()..)) }; - Box::new( - db - .indexer - .record_range(prefix.as_slice()..end_key.as_slice()) - .rev(), - ) - } else { - let start_key = if let Some(cursor) = &cursor { - let mut k = prefix.clone(); - let rkey = DbRkey::new(cursor); - keys::write_rkey(&mut k, &rkey); - k.push(0); - k - } else { - prefix.clone() - }; + for item in iter { + let (key, cid_bytes) = item.into_inner().into_diagnostic()?; - Box::new(db.indexer.record_range(start_key.as_slice()..)) - }; - - for item in iter { - let (key, cid_bytes) = item.into_inner().into_diagnostic()?; - - if !key.starts_with(prefix.as_slice()) { - break; - } + if !key.starts_with(prefix.as_slice()) { + break; + } - let rkey = keys::parse_rkey(&key[prefix.len()..])?; - if results.len() >= limit { - next_cursor = Some(rkey); - break; - } + let rkey = keys::parse_rkey(&key[prefix.len()..])?; + if results.len() >= limit { + next_cursor = Some(rkey); + break; + } - // look up using col|cid key built from collection and binary cid bytes - if let Ok(Some(block_bytes)) = db - .indexer - .block(keys::block_key(collection.as_str(), &cid_bytes)) - { - let value: Data = - serde_ipld_dagcbor::from_slice(&block_bytes).unwrap_or(Data::Null); - let cid = Cid::new(&cid_bytes).into_diagnostic()?; - let cid = Cid::Str(cid.to_cowstr().into_static()); - results.push(ListedRecord { - rkey: Rkey::new_cow(CowStr::Owned(rkey.to_smolstr())) - .expect("that rkey is validated"), - cid, - value: value.into_static(), - }); + // look up using col|cid key built from collection and binary cid bytes + if let Ok(Some(block_bytes)) = db + .indexer + .block(keys::block_key(collection.as_str(), &cid_bytes)) + { + let value: Data = + serde_ipld_dagcbor::from_slice(&block_bytes).unwrap_or(Data::Null); + let cid = Cid::new(&cid_bytes).into_diagnostic()?; + let cid = Cid::Str(cid.to_cowstr().into_static()); + results.push(ListedRecord { + rkey: Rkey::new_cow(CowStr::Owned(rkey.to_smolstr())) + .expect("that rkey is validated"), + cid, + value: value.into_static(), + }); + } } - } - Ok((results, next_cursor)) - }) - .await - .map(|(records, next_cursor)| RecordList { - records, - cursor: next_cursor.map(|rkey| { - Rkey::new_cow(CowStr::Owned(rkey.to_smolstr())).expect("that rkey is validated") - }), - }) + Ok((results, next_cursor)) + }) + .await + .map(|(records, next_cursor)| RecordList { + records, + cursor: next_cursor.map(|rkey| { + Rkey::new_cow(CowStr::Owned(rkey.to_smolstr())).expect("that rkey is validated") + }), + }) } /// generates a streaming CAR v1 response body for this repository. @@ -491,19 +513,9 @@ impl<'i> RepoHandle<'i> { for guard in app_state.db.indexer.record_prefix(&prefix) { let (key, cid_bytes) = guard.into_inner().into_diagnostic()?; - let rest = &key[prefix.len()..]; - let mut parts = rest.splitn(2, |b: &u8| *b == keys::SEP); - let collection_raw = parts - .next() - .ok_or_else(|| miette::miette!("missing collection in record key"))?; - let rkey_raw = parts - .next() - .ok_or_else(|| miette::miette!("missing rkey in record key"))?; - - let collection = std::str::from_utf8(collection_raw) - .into_diagnostic() - .wrap_err("collection is not valid utf8")?; - let rkey = keys::parse_rkey(rkey_raw)?; + let (collection_key, rkey) = keys::parse_record_suffix(&key[prefix.len()..])?; + let collection = app_state.db.indexer.collections.text(collection_key)?; + let collection = collection.as_str(); let mst_key = format!("{collection}/{rkey}"); let ipld_cid = cid::Cid::read_bytes(cid_bytes.as_ref()) @@ -583,7 +595,9 @@ impl<'i> RepoHandle<'i> { pub async fn count_records(&self, collection: &str) -> Result { let did = self.did.clone().into_static(); let collection = collection.to_string(); - self.state.db.run(move |db| db::get_record_count(db, &did, &collection)) + self.state + .db + .run(move |db| db::get_record_count(db, &did, &collection)) .await } } diff --git a/src/control/stats.rs b/src/control/stats.rs index 84ce371..0e43c92 100644 --- a/src/control/stats.rs +++ b/src/control/stats.rs @@ -64,14 +64,28 @@ impl Hydrant { state.db.jetstream.events.approximate_len() as u64, ); - let sizes = state.db.run(move |db| { - Ok(db - .all_keyspaces() - .into_iter() - .map(|(name, ks)| (name, ks.disk_space())) - .collect::>()) - }) - .await?; + #[cfg(feature = "indexer")] + { + // dictionary occupancy against its budget, so a run that hits the + // guard and falls back to NSID text in keys is visible rather than + // silent + counts.insert("collections", state.db.indexer.collections.len() as u64); + counts.insert( + "collections_bytes", + state.db.indexer.collections.used_bytes() as u64, + ); + } + + let sizes = state + .db + .run(move |db| { + Ok(db + .all_keyspaces() + .into_iter() + .map(|(name, ks)| (name, ks.disk_space())) + .collect::>()) + }) + .await?; Ok(StatsResponse { counts, sizes }) } diff --git a/src/db/collections.rs b/src/db/collections.rs new file mode 100644 index 0000000..3aa29f9 --- /dev/null +++ b/src/db/collections.rs @@ -0,0 +1,606 @@ +//! the collection NSID dictionary. +//! +//! holds both directions of the interning map fully in memory, loaded once at +//! [`Db::open`]. the network has ~26k distinct collections (UFOs census, +//! 2026-08-01), which costs roughly 2 MB, so nothing here is evicted and no +//! composite-key operation ever touches disk to resolve a collection. +//! +//! # the budget guard +//! +//! NSIDs are minted by whoever writes a record, so dictionary size is driven by +//! untrusted input and retained for the process lifetime. past +//! [`Config::collection_dict_max_mb`] the dictionary stops growing and new +//! collections are written as NSID text instead — [`CollectionKey::Text`], the +//! same form every key had before interning. that is a fallback to the *old key +//! layout*, not to a disk lookup: there is no cache miss path, so no read can be +//! slowed down by a full dictionary. the collections shut out are the long tail +//! that arrived last, while the collections carrying real traffic appear early +//! and keep their one-byte ids. +//! +//! the budget counts NSID bytes plus per-entry overhead rather than entry count, +//! so padding an NSID out to its 317-byte maximum buys no extra headroom. + +use std::collections::HashMap; + +use miette::{IntoDiagnostic, Result, WrapErr}; +use parking_lot::RwLock; +use smol_str::SmolStr; + +use super::collection_id::{self, CollectionId}; +use super::keys::SEP; +use super::schema::{self, Ks}; + +/// charged against the budget per entry on top of the NSID bytes, covering both +/// maps' slots and hashing overhead. +const ENTRY_OVERHEAD: usize = 64; + +/// how a collection is written into a composite key. +/// +/// both forms are load-bearing rather than transitional: the chunked migration +/// produces a mixed keyspace while it runs, and the budget guard can leave text +/// keys in place permanently. every read path therefore has to accept both, and +/// [`collection_id::is_interned`] is what tells them apart — an NSID always +/// starts with `[a-zA-Z]` and no allocatable id encodes to a leading ASCII +/// letter. +/// +/// # embedded versus terminal +/// +/// a collection sits in two different positions and needs a different encoding +/// in each, because the text form is only self-delimiting when a separator +/// follows it: +/// +/// - **embedded**, as in `records` (`{did}|{col}{rkey}`): [`Self::Id`] writes the +/// bare varint, which needs no separator because the encoding is prefix-free, +/// while [`Self::Text`] writes the NSID followed by [`SEP`] — exactly what +/// these keys held before interning. +/// - **terminal**, as in `counts` (`r|{did}|{col}`): nothing follows, so both +/// forms write bare bytes. adding a separator here would change the text +/// layout and force a rewrite of keys that do not otherwise need one. +/// +/// each reader consumes precisely what its writer produced, so both forms +/// coexist in one keyspace and a prefix built either way isolates exactly one +/// collection. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum CollectionKey<'a> { + Id(CollectionId), + Text(&'a str), +} + +impl<'a> CollectionKey<'a> { + /// bytes this segment occupies when something follows it, its separator + /// included. + pub fn embedded_len(&self) -> usize { + match self { + Self::Id(id) => id.encoded_len(), + Self::Text(nsid) => nsid.len() + 1, + } + } + + /// write this segment where more key follows it. + pub fn write_embedded(&self, buf: &mut Vec) { + match self { + Self::Id(id) => id.write_to_vec(buf), + Self::Text(nsid) => { + buf.extend_from_slice(nsid.as_bytes()); + buf.push(SEP); + } + } + } + + /// write this segment as the last part of a key. + pub fn write_terminal(&self, buf: &mut Vec) { + match self { + Self::Id(id) => id.write_to_vec(buf), + Self::Text(nsid) => buf.extend_from_slice(nsid.as_bytes()), + } + } + + /// decode an embedded segment from the front of `bytes`, with its remainder. + pub fn read_embedded(bytes: &[u8]) -> Result<(CollectionKey<'_>, &[u8])> { + if collection_id::is_interned(bytes) { + let (id, rest) = CollectionId::read_from(bytes)?; + return Ok((CollectionKey::Id(id), rest)); + } + let end = bytes + .iter() + .position(|&byte| byte == SEP) + .ok_or_else(|| miette::miette!("collection segment is missing its separator"))?; + let nsid = std::str::from_utf8(&bytes[..end]) + .into_diagnostic() + .wrap_err("collection segment is not valid utf8")?; + Ok((CollectionKey::Text(nsid), &bytes[end + 1..])) + } + + /// decode a segment that is the whole remainder of a key. + pub fn read_terminal(bytes: &[u8]) -> Result> { + if collection_id::is_interned(bytes) { + let (id, rest) = CollectionId::read_from(bytes)?; + if !rest.is_empty() { + miette::bail!("trailing bytes after a terminal collection id"); + } + return Ok(CollectionKey::Id(id)); + } + std::str::from_utf8(bytes) + .map(CollectionKey::Text) + .into_diagnostic() + .wrap_err("collection segment is not valid utf8") + } +} + +/// the in-memory dictionary. both directions, never evicted. +struct Inner { + forward: HashMap, + /// indexed by raw id. the ids skipped to keep encodings out of ASCII leave + /// at most 52 holes, all below 128, so a `Vec` stays dense; a hole holds the + /// empty string. + reverse: Vec, + /// the next id to hand out, `None` once the id space is exhausted. + next: Option, + used_bytes: usize, + /// whether the budget has already been reported as reached. + reported_full: bool, +} + +impl Inner { + fn record(&mut self, nsid: SmolStr, id: CollectionId) { + let slot = id.get() as usize; + if self.reverse.len() <= slot { + self.reverse.resize(slot + 1, SmolStr::default()); + } + self.reverse[slot] = nsid.clone(); + self.used_bytes += nsid.len() + ENTRY_OVERHEAD; + self.forward.insert(nsid, id); + } +} + +pub struct Interner { + ks: Ks, + budget_bytes: usize, + inner: RwLock, +} + +impl Interner { + /// load the whole dictionary from `ks`. + pub fn open(ks: Ks, budget_bytes: usize) -> Result { + let mut inner = Inner { + forward: HashMap::new(), + reverse: Vec::new(), + next: Some(CollectionId::first()), + used_bytes: 0, + reported_full: false, + }; + + let mut highest: Option = None; + for guard in ks.iter() { + let (key, value) = guard.into_inner().into_diagnostic()?; + let raw = u32::from_be_bytes( + key.as_ref() + .try_into() + .into_diagnostic() + .wrap_err("collection dictionary key expected to be 4 bytes")?, + ); + let id = CollectionId::from_raw(raw)?; + let nsid = std::str::from_utf8(&value) + .into_diagnostic() + .wrap_err("collection dictionary value is not valid utf8")?; + inner.record(SmolStr::new(nsid), id); + highest = Some(highest.map_or(id, |h: CollectionId| h.max(id))); + } + + // resume after the highest id ever handed out, never reusing one: a + // reused id would silently reassign every key already written under it + inner.next = match highest { + Some(id) => id.next(), + None => Some(CollectionId::first()), + }; + + if !inner.forward.is_empty() { + tracing::debug!( + "loaded {} interned collections ({} KiB)", + inner.forward.len(), + inner.used_bytes / 1024 + ); + } + + Ok(Self { + ks, + budget_bytes, + inner: RwLock::new(inner), + }) + } + + /// how a composite key should spell `nsid`, interning it if there is room. + /// + /// returns [`CollectionKey::Text`] when the budget is exhausted or the id + /// space is used up, which keeps writes working at the pre-interning key + /// size rather than failing. + pub fn key_for<'a>(&self, nsid: &'a str) -> Result> { + if let Some(id) = self.inner.read().forward.get(nsid) { + return Ok(CollectionKey::Id(*id)); + } + self.intern(nsid) + } + + /// how an already-written key spells `nsid`, without interning it. + /// + /// read paths must use this rather than [`Self::key_for`]: querying a + /// collection that was never written would otherwise mint an id for it, so a + /// stream of `getRecord` misses could grow the dictionary. falling back to + /// [`CollectionKey::Text`] is also the correct lookup for a collection the + /// budget guard shut out, since that is how its keys were written. + pub fn lookup<'a>(&self, nsid: &'a str) -> CollectionKey<'a> { + self.inner + .read() + .forward + .get(nsid) + .map(|id| CollectionKey::Id(*id)) + .unwrap_or(CollectionKey::Text(nsid)) + } + + /// the NSID behind `id`. + pub fn nsid(&self, id: CollectionId) -> Result { + self.inner + .read() + .reverse + .get(id.get() as usize) + .filter(|nsid| !nsid.is_empty()) + .cloned() + .ok_or_else(|| miette::miette!("collection id {} is not in the dictionary", id.get())) + } + + /// the NSID text for either form of a collection segment. + pub fn text(&self, key: CollectionKey<'_>) -> Result { + match key { + CollectionKey::Id(id) => self.nsid(id), + CollectionKey::Text(nsid) => Ok(SmolStr::new(nsid)), + } + } + + /// the underlying keyspace handle, for the registry's by-name tables. + pub fn keyspace(&self) -> fjall::Keyspace { + self.ks.keyspace() + } + + /// number of interned collections. + pub fn len(&self) -> usize { + self.inner.read().forward.len() + } + + /// bytes the dictionary is charged against its budget. + pub fn used_bytes(&self) -> usize { + self.inner.read().used_bytes + } + + fn intern<'a>(&self, nsid: &'a str) -> Result> { + let mut inner = self.inner.write(); + + // another writer may have interned it between the read and write locks + if let Some(id) = inner.forward.get(nsid) { + return Ok(CollectionKey::Id(*id)); + } + + if inner.used_bytes + nsid.len() + ENTRY_OVERHEAD > self.budget_bytes { + if !inner.reported_full { + inner.reported_full = true; + tracing::warn!( + interned = inner.forward.len(), + used_bytes = inner.used_bytes, + budget_bytes = self.budget_bytes, + "collection dictionary budget reached; further collections keep NSID text in keys" + ); + } + return Ok(CollectionKey::Text(nsid)); + } + + let Some(id) = inner.next else { + return Ok(CollectionKey::Text(nsid)); + }; + + // persist before memoizing, so a failed write leaves an unused id rather + // than an in-memory id whose dictionary entry never became durable. an + // id burned this way is never handed out again, which is the safe + // direction: reuse would remap keys already written under it. + self.ks + .insert(id.get().to_be_bytes(), nsid.as_bytes()) + .into_diagnostic() + .wrap_err_with(|| format!("failed to persist collection id for {nsid}"))?; + + inner.next = id.next(); + inner.record(SmolStr::new(nsid), id); + + tracing::debug!("interned collection {nsid} as id {}", id.get()); + Ok(CollectionKey::Id(id)) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::db::Db; + + fn open_db() -> Result<(tempfile::TempDir, Db)> { + let tmp = tempfile::tempdir().into_diagnostic()?; + let cfg = crate::config::Config { + database_path: tmp.path().to_path_buf(), + ..Default::default() + }; + let db = Db::open(&cfg)?; + Ok((tmp, db)) + } + + fn encode(key: CollectionKey<'_>) -> Vec { + let mut buf = Vec::new(); + key.write_embedded(&mut buf); + buf + } + + fn encode_terminal(key: CollectionKey<'_>) -> Vec { + let mut buf = Vec::new(); + key.write_terminal(&mut buf); + buf + } + + #[test] + fn both_key_forms_roundtrip_embedded_with_their_remainder() -> Result<()> { + let id = CollectionKey::Id(CollectionId::first()); + let text = CollectionKey::Text("app.bsky.feed.like"); + + for form in [id, text] { + let mut buf = encode(form); + buf.extend_from_slice(b"\x01rkeybytes"); + let (decoded, rest) = CollectionKey::read_embedded(&buf)?; + assert_eq!(decoded, form); + assert_eq!(rest, b"\x01rkeybytes"); + assert_eq!(encode(form).len(), form.embedded_len()); + } + Ok(()) + } + + #[test] + fn both_key_forms_roundtrip_terminal() -> Result<()> { + for form in [ + CollectionKey::Id(CollectionId::first()), + CollectionKey::Text("app.bsky.feed.like"), + ] { + let buf = encode_terminal(form); + assert_eq!(CollectionKey::read_terminal(&buf)?, form); + } + Ok(()) + } + + /// the terminal text form must stay byte-identical to the pre-interning + /// layout, or existing count keys would need a rewrite they do not deserve. + #[test] + fn terminal_text_is_the_bare_nsid() { + assert_eq!( + encode_terminal(CollectionKey::Text("app.bsky.feed.like")), + b"app.bsky.feed.like" + ); + } + + /// an interned id must never be mistaken for the text form or vice versa, + /// which is what lets both live in one keyspace. + #[test] + fn the_two_forms_are_never_confused() -> Result<()> { + let mut id = CollectionId::first(); + for _ in 0..500 { + let embedded = encode(CollectionKey::Id(id)); + let (decoded, _) = CollectionKey::read_embedded(&embedded)?; + assert_eq!(decoded, CollectionKey::Id(id), "id {}", id.get()); + assert_eq!( + CollectionKey::read_terminal(&encode_terminal(CollectionKey::Id(id)))?, + CollectionKey::Id(id) + ); + id = id.next().expect("id space"); + } + + // including an uppercase-leading NSID, which the live network has + for nsid in [ + "app.bsky.feed.like", + "C.d.e.f", + "a.b.c", + "sh.tangled.actor.profile", + ] { + let embedded = encode(CollectionKey::Text(nsid)); + let (decoded, _) = CollectionKey::read_embedded(&embedded)?; + assert_eq!(decoded, CollectionKey::Text(nsid)); + assert_eq!( + CollectionKey::read_terminal(&encode_terminal(CollectionKey::Text(nsid)))?, + CollectionKey::Text(nsid) + ); + } + Ok(()) + } + + #[test] + fn read_embedded_rejects_text_without_a_separator() { + assert!(CollectionKey::read_embedded(b"app.bsky.feed.like").is_err()); + } + + #[test] + fn read_terminal_rejects_trailing_bytes_after_an_id() { + let mut buf = encode_terminal(CollectionKey::Id(CollectionId::first())); + buf.extend_from_slice(b"extra"); + assert!(CollectionKey::read_terminal(&buf).is_err()); + } + + #[test] + fn interning_is_stable_and_reversible() -> Result<()> { + let (_tmp, db) = open_db()?; + let interner = &db.indexer.collections; + + let like = interner.key_for("app.bsky.feed.like")?; + let post = interner.key_for("app.bsky.feed.post")?; + assert_ne!(like, post); + // the same NSID always resolves to the same id + assert_eq!(interner.key_for("app.bsky.feed.like")?, like); + + let CollectionKey::Id(like_id) = like else { + panic!("expected an interned id"); + }; + assert_eq!(interner.nsid(like_id)?, "app.bsky.feed.like"); + assert_eq!(interner.text(like)?, "app.bsky.feed.like"); + assert_eq!(interner.len(), 2); + Ok(()) + } + + /// the hot collections must land in the one-byte range, which is the whole + /// point of the varint. + #[test] + fn the_first_collections_seen_get_one_byte_ids() -> Result<()> { + let (_tmp, db) = open_db()?; + let interner = &db.indexer.collections; + + for nsid in [ + "app.bsky.feed.like", + "app.bsky.feed.post", + "app.bsky.graph.follow", + "app.bsky.feed.repost", + "app.bsky.graph.block", + ] { + let key = interner.key_for(nsid)?; + assert_eq!(key.embedded_len(), 1, "{nsid} did not fit in one byte"); + } + Ok(()) + } + + #[test] + fn the_dictionary_survives_a_reopen() -> Result<()> { + let tmp = tempfile::tempdir().into_diagnostic()?; + let cfg = crate::config::Config { + database_path: tmp.path().to_path_buf(), + ..Default::default() + }; + + let (like, post, len) = { + let db = Db::open(&cfg)?; + let interner = &db.indexer.collections; + ( + interner.key_for("app.bsky.feed.like")?, + interner.key_for("app.bsky.feed.post")?, + interner.len(), + ) + }; + + let db = Db::open(&cfg)?; + let interner = &db.indexer.collections; + assert_eq!(interner.len(), len); + assert_eq!(interner.key_for("app.bsky.feed.like")?, like); + assert_eq!(interner.key_for("app.bsky.feed.post")?, post); + + // a collection interned after the reopen must not reuse an existing id + let follow = interner.key_for("app.bsky.graph.follow")?; + assert_ne!(follow, like); + assert_ne!(follow, post); + Ok(()) + } + + #[test] + fn a_reopen_never_reuses_an_id_after_a_gap() -> Result<()> { + let tmp = tempfile::tempdir().into_diagnostic()?; + let cfg = crate::config::Config { + database_path: tmp.path().to_path_buf(), + ..Default::default() + }; + + let highest = { + let db = Db::open(&cfg)?; + let interner = &db.indexer.collections; + let mut highest = CollectionId::first(); + for i in 0..300 { + let nsid = format!("test{i}.collection.name"); + let CollectionKey::Id(id) = interner.key_for(&nsid)? else { + panic!("expected an interned id"); + }; + highest = highest.max(id); + } + highest + }; + + let db = Db::open(&cfg)?; + let key = db.indexer.collections.key_for("fresh.collection.name")?; + let CollectionKey::Id(id) = key else { + panic!("expected an interned id"); + }; + assert!( + id > highest, + "reopened interner handed out {} which is not past {}", + id.get(), + highest.get() + ); + Ok(()) + } + + #[test] + fn the_budget_falls_back_to_text_instead_of_failing() -> Result<()> { + let tmp = tempfile::tempdir().into_diagnostic()?; + let cfg = crate::config::Config { + database_path: tmp.path().to_path_buf(), + // room for two entries at ~64 bytes of overhead each + collection_dict_max_bytes: 2 * (ENTRY_OVERHEAD + 20), + ..Default::default() + }; + let db = Db::open(&cfg)?; + let interner = &db.indexer.collections; + + assert!(matches!( + interner.key_for("app.bsky.feed.like")?, + CollectionKey::Id(_) + )); + assert!(matches!( + interner.key_for("app.bsky.feed.post")?, + CollectionKey::Id(_) + )); + + // past the budget, keys keep the pre-interning layout + let overflow = interner.key_for("app.bsky.graph.follow")?; + assert_eq!(overflow, CollectionKey::Text("app.bsky.graph.follow")); + + // and an already-interned collection still resolves to its id + assert!(matches!( + interner.key_for("app.bsky.feed.like")?, + CollectionKey::Id(_) + )); + + // the shut-out collection is not persisted, so a reopen agrees + drop(db); + let db = Db::open(&cfg)?; + assert_eq!(db.indexer.collections.len(), 2); + assert_eq!( + db.indexer.collections.key_for("app.bsky.graph.follow")?, + CollectionKey::Text("app.bsky.graph.follow") + ); + Ok(()) + } + + #[test] + fn the_budget_counts_nsid_bytes_so_padding_buys_nothing() -> Result<()> { + let tmp = tempfile::tempdir().into_diagnostic()?; + let cfg = crate::config::Config { + database_path: tmp.path().to_path_buf(), + collection_dict_max_bytes: 4 * ENTRY_OVERHEAD, + ..Default::default() + }; + let db = Db::open(&cfg)?; + let interner = &db.indexer.collections; + + // one padded NSID consumes the budget that several short ones would + let padded = format!("a{}.b.c", "-".repeat(200)); + assert!(matches!(interner.key_for(&padded)?, CollectionKey::Text(_))); + assert_eq!(interner.len(), 0); + + assert!(matches!( + interner.key_for("app.bsky.feed.like")?, + CollectionKey::Id(_) + )); + Ok(()) + } + + #[test] + fn nsid_rejects_an_id_that_was_never_interned() -> Result<()> { + let (_tmp, db) = open_db()?; + let unknown = CollectionId::from_raw(9_999)?; + assert!(db.indexer.collections.nsid(unknown).is_err()); + Ok(()) + } +} diff --git a/src/db/indexer.rs b/src/db/indexer.rs index 9c3c469..036178c 100644 --- a/src/db/indexer.rs +++ b/src/db/indexer.rs @@ -49,9 +49,10 @@ pub fn set_record_count( did: &Did<'_>, collection: &str, count: u64, -) { - let key = keys::count_collection_key(did, collection); +) -> Result<()> { + let key = keys::count_collection_key(did, db.indexer.collections.key_for(collection)?); batch.insert(&db.counts, key, count.to_be_bytes()); + Ok(()) } pub fn replace_record_counts<'a>( @@ -67,7 +68,7 @@ pub fn replace_record_counts<'a>( } for (collection, count) in counts { - set_record_count(batch, db, did, collection, count); + set_record_count(batch, db, did, collection, count)?; } Ok(()) @@ -86,14 +87,15 @@ pub fn replace_record_counts_matching<'a>( let collection = key .strip_prefix(prefix.as_slice()) .ok_or_else(|| miette::miette!("invalid collection count key: {key:?}")) - .and_then(|raw| std::str::from_utf8(raw).into_diagnostic())?; - if matches_collection(collection) { + .and_then(crate::db::collections::CollectionKey::read_terminal) + .and_then(|parsed| db.indexer.collections.text(parsed))?; + if matches_collection(&collection) { batch.remove(&db.counts, key); } } for (collection, count) in counts { - set_record_count(batch, db, did, collection, count); + set_record_count(batch, db, did, collection, count)?; } Ok(()) @@ -106,7 +108,7 @@ pub fn update_record_count( collection: &str, delta: i64, ) -> Result<()> { - let key = keys::count_collection_key(did, collection); + let key = keys::count_collection_key(did, db.indexer.collections.key_for(collection)?); let count = db .counts .get(&key) @@ -131,7 +133,8 @@ pub fn update_record_count( } pub fn get_record_count(db: &Db, did: &Did<'_>, collection: &str) -> Result { - let key = keys::count_collection_key(did, collection); + // a read must not intern: an unknown collection would otherwise get an id + let key = keys::count_collection_key(did, db.indexer.collections.lookup(collection)); let count = db .counts .get(&key) @@ -180,8 +183,8 @@ mod tests { let did = Did::new("did:plc:yk4q3id7id6p5z3bypvshc64").into_diagnostic()?; let mut batch = db.inner.batch(); - set_record_count(&mut batch, &db, &did, "app.bsky.feed.post", 3); - set_record_count(&mut batch, &db, &did, "app.bsky.feed.like", 2); + set_record_count(&mut batch, &db, &did, "app.bsky.feed.post", 3)?; + set_record_count(&mut batch, &db, &did, "app.bsky.feed.like", 2)?; batch.commit().into_diagnostic()?; let mut batch = db.inner.batch(); @@ -215,8 +218,8 @@ mod tests { let did = Did::new("did:plc:yk4q3id7id6p5z3bypvshc64").into_diagnostic()?; let mut batch = db.inner.batch(); - set_record_count(&mut batch, &db, &did, "sh.tangled.repo", 3); - set_record_count(&mut batch, &db, &did, "app.bsky.feed.post", 2); + set_record_count(&mut batch, &db, &did, "sh.tangled.repo", 3)?; + set_record_count(&mut batch, &db, &did, "app.bsky.feed.post", 2)?; batch.commit().into_diagnostic()?; let mut batch = db.inner.batch(); diff --git a/src/db/keys/indexer.rs b/src/db/keys/indexer.rs index 0add66f..a1366bf 100644 --- a/src/db/keys/indexer.rs +++ b/src/db/keys/indexer.rs @@ -2,6 +2,7 @@ use jacquard_common::types::string::Did; use smol_str::SmolStr; use super::SEP; +use crate::db::collections::CollectionKey; use crate::db::types::{DbRkey, DbTid, TrimmedDid}; #[cfg(feature = "indexer_stream")] @@ -28,29 +29,37 @@ pub fn record_prefix_did(did: &Did) -> Vec { prefix } -// prefix format: {DID}|{collection}| -pub fn record_prefix_collection(did: &Did, collection: &str) -> Vec { +// prefix format: {DID}|{collection} +// the collection segment carries its own terminator, so an interned id needs no +// separator while legacy NSID text keeps the one it always had +pub fn record_prefix_collection(did: &Did, collection: CollectionKey<'_>) -> Vec { let repo = TrimmedDid::from(did); - let mut prefix = Vec::with_capacity(repo.len() + 1 + collection.len() + 1); + let mut prefix = Vec::with_capacity(repo.len() + 1 + collection.embedded_len()); repo.write_to_vec(&mut prefix); prefix.push(SEP); - prefix.extend_from_slice(collection.as_bytes()); - prefix.push(SEP); + collection.write_embedded(&mut prefix); prefix } -// key format: {DID}|{collection}|{rkey} -pub fn record_key(did: &Did, collection: &str, rkey: &DbRkey) -> Vec { +// key format: {DID}|{collection}{rkey} +pub fn record_key(did: &Did, collection: CollectionKey<'_>, rkey: &DbRkey) -> Vec { let repo = TrimmedDid::from(did); - let mut key = Vec::with_capacity(repo.len() + 1 + collection.len() + 1 + rkey.len() + 1); + let mut key = Vec::with_capacity(repo.len() + 1 + collection.embedded_len() + rkey.len() + 1); repo.write_to_vec(&mut key); key.push(SEP); - key.extend_from_slice(collection.as_bytes()); - key.push(SEP); + collection.write_embedded(&mut key); write_rkey(&mut key, rkey); key } +/// split a record key's suffix, the part after `{DID}|`, into its collection and +/// rkey. +pub fn parse_record_suffix(suffix: &[u8]) -> miette::Result<(CollectionKey<'_>, DbRkey)> { + let (collection, rest) = CollectionKey::read_embedded(suffix)?; + let rkey = parse_rkey(rest)?; + Ok((collection, rkey)) +} + pub fn write_rkey(buf: &mut Vec, rkey: &DbRkey) { match rkey { DbRkey::Tid(tid) => { @@ -84,9 +93,11 @@ pub fn parse_rkey(raw: &[u8]) -> miette::Result { } // key format: r|{DID}|{collection} (DID trimmed) -pub fn count_collection_key(did: &Did, collection: &str) -> Vec { +// the collection is terminal here, so both forms write bare bytes and the text +// form stays byte-identical to what it was before interning +pub fn count_collection_key(did: &Did, collection: CollectionKey<'_>) -> Vec { let mut key = super::did_collection_prefix(did); - key.extend_from_slice(collection.as_bytes()); + collection.write_terminal(&mut key); key } @@ -197,3 +208,139 @@ pub fn block_key(collection: &str, cid: &[u8]) -> Vec { key.extend_from_slice(cid); key } + +#[cfg(test)] +mod tests { + use super::*; + use crate::db::collection_id::CollectionId; + use smol_str::SmolStr; + + fn did() -> Did<'static> { + Did::new_owned("did:plc:yk4q3id7id6p5z3bypvshc64").expect("valid did") + } + + /// a TID whose raw bytes contain `SEP`, the case that makes splitting a + /// record key on the separator unsound. + fn sep_bearing_tid() -> DbRkey { + let tid = DbTid::new_from_bytes([0x00, SEP, 0x11, SEP, 0x22, 0x33, SEP, 0x44]); + assert!(tid.as_bytes().contains(&SEP)); + DbRkey::Tid(tid) + } + + fn roundtrip(collection: CollectionKey<'_>, rkey: &DbRkey) -> miette::Result<()> { + let did = did(); + let key = record_key(&did, collection, rkey); + let prefix = record_prefix_did(&did); + assert!(key.starts_with(&prefix)); + + let (parsed_collection, parsed_rkey) = parse_record_suffix(&key[prefix.len()..])?; + assert_eq!(parsed_collection, collection); + assert_eq!(&parsed_rkey, rkey); + Ok(()) + } + + #[test] + fn record_keys_roundtrip_in_both_collection_forms() -> miette::Result<()> { + for rkey in [ + DbRkey::new("3jzfcijpj2z2a"), + DbRkey::Str(SmolStr::new("self")), + sep_bearing_tid(), + ] { + roundtrip(CollectionKey::Id(CollectionId::first()), &rkey)?; + roundtrip(CollectionKey::Text("app.bsky.feed.like"), &rkey)?; + } + Ok(()) + } + + /// a multi-byte id exercises the varint path rather than the single-byte + /// happy case. + #[test] + fn record_keys_roundtrip_with_a_multi_byte_collection_id() -> miette::Result<()> { + let mut id = CollectionId::first(); + for _ in 0..300 { + id = id.next().expect("id space"); + } + assert!(id.encoded_len() > 1); + roundtrip(CollectionKey::Id(id), &sep_bearing_tid()) + } + + /// the record prefix for one collection must not match another's records, + /// which is what makes listRecords correct. + #[test] + fn a_collection_prefix_matches_only_its_own_records() { + let did = did(); + let rkey = DbRkey::new("3jzfcijpj2z2a"); + + let first = CollectionId::first(); + let second = first.next().expect("id space"); + // a two-byte id whose leading byte differs from both single-byte ones + let mut wide = second; + for _ in 0..300 { + wide = wide.next().expect("id space"); + } + + for (a, b) in [ + (CollectionKey::Id(first), CollectionKey::Id(second)), + (CollectionKey::Id(first), CollectionKey::Id(wide)), + ( + CollectionKey::Text("app.bsky.feed.like"), + CollectionKey::Text("app.bsky.feed.post"), + ), + // a text collection that is a prefix of another must not match it + ( + CollectionKey::Text("app.bsky.feed"), + CollectionKey::Text("app.bsky.feed.post"), + ), + // the two forms must not match each other + ( + CollectionKey::Id(first), + CollectionKey::Text("app.bsky.feed.like"), + ), + ] { + let prefix = record_prefix_collection(&did, a); + let other = record_key(&did, b, &rkey); + assert!( + !other.starts_with(&prefix), + "prefix for {a:?} matched a record of {b:?}" + ); + assert!(record_key(&did, a, &rkey).starts_with(&prefix)); + } + } + + /// the exclusive upper bound used for reverse listRecords scans is built by + /// incrementing the prefix's last byte, which must never overflow. + #[test] + fn a_collection_prefix_never_ends_in_0xff() { + let did = did(); + let mut id = CollectionId::first(); + for _ in 0..5_000 { + let prefix = record_prefix_collection(&did, CollectionKey::Id(id)); + assert_ne!( + *prefix.last().expect("non-empty"), + 0xFF, + "id {} produced a prefix that cannot be incremented", + id.get() + ); + id = id.next().expect("id space"); + } + let text = record_prefix_collection(&did, CollectionKey::Text("app.bsky.feed.like")); + assert_eq!(*text.last().expect("non-empty"), SEP); + } + + #[test] + fn count_collection_keys_share_the_repo_prefix_and_roundtrip() -> miette::Result<()> { + let did = did(); + let prefix = super::super::did_collection_prefix(&did); + + for collection in [ + CollectionKey::Id(CollectionId::first()), + CollectionKey::Text("app.bsky.feed.like"), + ] { + let key = count_collection_key(&did, collection); + assert!(key.starts_with(&prefix)); + let parsed = CollectionKey::read_terminal(&key[prefix.len()..])?; + assert_eq!(parsed, collection); + } + Ok(()) + } +} diff --git a/src/db/keyspaces.rs b/src/db/keyspaces.rs index 190e586..05bdb39 100644 --- a/src/db/keyspaces.rs +++ b/src/db/keyspaces.rs @@ -51,6 +51,9 @@ impl OpenCx<'_> { #[cfg(feature = "indexer")] #[derive(Clone)] pub struct IndexerDb { + /// the collection NSID dictionary, held in memory and shared by every clone + /// of this handle so an id interned on one is visible on all of them + pub(crate) collections: Arc, /// maps `{DID}|{COL}|{RKey}` -> record CID pub(super) records: Ks, /// content-addressable storage of raw DAG-CBOR blocks @@ -69,6 +72,10 @@ pub struct IndexerDb { impl IndexerDb { pub(super) fn open(cx: &OpenCx) -> Result { Ok(Self { + collections: Arc::new(super::collections::Interner::open( + Ks::open(cx)?, + cx.cfg.collection_dict_max_bytes, + )?), records: Ks::open(cx)?, blocks: Ks::open(cx)?, pending: Ks::open(cx)?, diff --git a/src/db/mod.rs b/src/db/mod.rs index 21f08f4..eac53c0 100644 --- a/src/db/mod.rs +++ b/src/db/mod.rs @@ -11,10 +11,12 @@ use std::sync::atomic::AtomicU64; use std::sync::{Arc, Mutex}; use url::Url; -// consumed once composite keys carry ids instead of NSID text; until then the -// encoding invariants the key format depends on are held up by its own tests -#[allow(dead_code)] +// only indexer keys carry a collection segment, so the dictionary and its +// encoding are both indexer-only +#[cfg(feature = "indexer")] pub mod collection_id; +#[cfg(feature = "indexer")] +pub mod collections; pub mod compaction; pub mod counts; #[cfg(feature = "indexer")] diff --git a/src/db/registry.rs b/src/db/registry.rs index d147149..de93eb8 100644 --- a/src/db/registry.rs +++ b/src/db/registry.rs @@ -100,6 +100,8 @@ registry! { #[cfg(feature = "indexer")] indexer.blocks => schema::Blocks, #[cfg(feature = "indexer")] + indexer.collections => schema::Collections, + #[cfg(feature = "indexer")] indexer.pending => schema::Pending, #[cfg(feature = "indexer")] indexer.resync => schema::Resync, diff --git a/src/db/schema.rs b/src/db/schema.rs index e5bc458..f11bbb3 100644 --- a/src/db/schema.rs +++ b/src/db/schema.rs @@ -425,6 +425,48 @@ impl Schema for Blocks { } } +/// the collection dictionary, `{id}` (u32 BE) -> NSID text +/// +/// read in full at open and cached in memory, so it is only ever iterated. the +/// key is fixed-width rather than the varint used in composite keys, so the +/// highest id ever handed out is the last entry. +#[cfg(feature = "indexer")] +pub struct Collections; + +#[cfg(feature = "indexer")] +impl Schema for Collections { + const NAME: &'static str = "collections"; + + fn options(_cx: &OpenCx) -> KeyspaceCreateOptions { + KeyspaceCreateOptions::default() + // only iterated, at open, and appended to when a new collection + // shows up; the network has ~26k collections in total + .expect_point_read_hits(true) + .max_memtable_size(mb(1)) + .data_block_size_policy(BlockSizePolicy::all(kb(4))) + // nsids share long prefixes, but there are so few of them that the + // whole keyspace is a rounding error either way + .data_block_compression_policy(CompressionPolicy::disabled()) + .data_block_restart_interval_policy(RestartIntervalPolicy::all(8)) + } + + fn render_value(value: &[u8]) -> Value { + std::str::from_utf8(value) + .map(|nsid| Value::String(nsid.to_owned())) + .unwrap_or_else(|_| Value::String(hex::encode(value))) + } + + fn parse_debug_key(key: &str) -> Option> { + key.parse::().ok().map(|id| id.to_be_bytes().to_vec()) + } + + fn render_debug_key(key: &[u8]) -> String { + key.try_into() + .map(|arr| u32::from_be_bytes(arr).to_string()) + .unwrap_or_else(|_| "invalid_u32".to_string()) + } +} + /// backfill queue of `{ID}` (u64 BE) -> empty #[cfg(feature = "indexer")] pub struct Pending; diff --git a/src/db/txn.rs b/src/db/txn.rs index aded08f..3f85084 100644 --- a/src/db/txn.rs +++ b/src/db/txn.rs @@ -10,6 +10,10 @@ use jacquard_common::types::cid::IpldCid; use jacquard_common::types::did::Did; use miette::{IntoDiagnostic, Result}; +#[cfg(feature = "indexer")] +use crate::db::collection_id::CollectionId; +#[cfg(feature = "indexer")] +use crate::db::collections::CollectionKey; #[cfg(feature = "indexer")] use crate::db::types::{DbAction, DbRkey, DbTid}; use crate::db::{CountDeltas, Db}; @@ -21,6 +25,8 @@ use crate::ops::record_events::{EmitOp, RecordEmitter, RecordEventOrigin, Record use crate::state::AppState; #[cfg(feature = "indexer")] use crate::types::{GaugeState, RepoState}; +#[cfg(feature = "indexer")] +use smol_str::SmolStr; #[cfg(feature = "firehose-diagnostics")] type TxnInstant = std::time::Instant; @@ -125,6 +131,7 @@ impl<'db> Txn<'db> { records_delta: 0, blocks_count: 0, collection_deltas: HashMap::new(), + last_collection: None, } } @@ -210,10 +217,29 @@ pub(crate) struct RecordTxn<'txn, 'db, 'did, 'repo> { records_delta: i64, blocks_count: i64, collection_deltas: HashMap, + /// the last collection resolved to an id. + /// + /// records in a commit and in a backfill arrive grouped by collection, so + /// this turns one dictionary lookup per record into one per run. + last_collection: Option<(SmolStr, CollectionId)>, } #[cfg(feature = "indexer")] impl RecordTxn<'_, '_, '_, '_> { + /// how this collection should be spelled in a key, interning it if needed. + fn collection_key<'a>(&mut self, collection: &'a str) -> Result> { + if let Some((cached, id)) = &self.last_collection + && cached == collection + { + return Ok(CollectionKey::Id(*id)); + } + let key = self.txn.db.indexer.collections.key_for(collection)?; + if let CollectionKey::Id(id) = key { + self.last_collection = Some((SmolStr::new(collection), id)); + } + Ok(key) + } + pub(crate) fn put_record( &mut self, collection: &str, @@ -228,6 +254,7 @@ impl RecordTxn<'_, '_, '_, '_> { } if !self.ephemeral { + let collection_key = self.collection_key(collection)?; let cid_bytes = cid.to_bytes(); if !self.only_index_links { self.txn.batch.insert( @@ -238,7 +265,7 @@ impl RecordTxn<'_, '_, '_, '_> { } self.txn.batch.insert( &self.txn.db.indexer.records, - keys::record_key(self.did, collection, rkey), + keys::record_key(self.did, collection_key, rkey), cid_bytes, ); if action == DbAction::Create { @@ -277,9 +304,10 @@ impl RecordTxn<'_, '_, '_, '_> { } if !self.ephemeral { + let collection_key = self.collection_key(collection)?; self.txn.batch.remove( &self.txn.db.indexer.records, - keys::record_key(self.did, collection, rkey), + keys::record_key(self.did, collection_key, rkey), ); *self .collection_deltas diff --git a/tests/bench/collection_interning_ab.sh b/tests/bench/collection_interning_ab.sh new file mode 100755 index 0000000..053f40b --- /dev/null +++ b/tests/bench/collection_interning_ab.sh @@ -0,0 +1,215 @@ +#!/usr/bin/env bash +# collection interning storage A/B (hydrant-3kc.2 baseline + hydrant-3kc.6 arm). +# +# runs one arm per binary over the SAME uniform repo sample into a FRESH database, +# so the comparison needs no migration: only the interning read/write paths. +# +# both arms are driven from this one script, sequentially, never overlapping — +# the host is shared with production co-tenants and a single API port is reused. +# +# usage: +# ./collection_interning_ab.sh = [= ...] +# +# example: +# ./collection_interning_ab.sh /root/bench/ci-20260801 /root/bench/sample-30k.ndjson \ +# baseline=/root/bench/bin/hydrant-baseline interned=/root/bench/bin/hydrant-interned +set -euo pipefail + +RUN_DIR=${1:?run dir} +SAMPLE=${2:?sample ndjson} +shift 2 +[ $# -ge 1 ] || { echo "need at least one =" >&2; exit 1; } + +API_PORT=${API_PORT:-19310} +DEBUG_PORT=${DEBUG_PORT:-19311} +INTERVAL=${INTERVAL:-5} +# plc.klbr.net is dawn's unthrottled mirror; falls back if unreachable +PLC_URL=${PLC_URL:-https://plc.klbr.net} + +mkdir -p "$RUN_DIR" +echo "run dir: $RUN_DIR" +echo "sample: $SAMPLE ($(wc -l < "$SAMPLE") repos)" + +if ! curl -s --max-time 10 "$PLC_URL/did:plc:yk4q3id7id6p5z3bypvshc64" >/dev/null; then + echo "WARNING: $PLC_URL unreachable, falling back to plc.directory" >&2 + PLC_URL=https://plc.directory +fi +echo "plc: $PLC_URL" + +# a stray hydrant on our port would silently join the measurement +if curl -s --max-time 3 "http://127.0.0.1:$API_PORT/stats" >/dev/null 2>&1; then + echo "FATAL: something is already serving :$API_PORT" >&2 + exit 1 +fi + +run_arm() { + local arm=$1 bin=$2 + local dir="$RUN_DIR/$arm" + mkdir -p "$dir" + local db="$dir/hydrant.db" + rm -rf "$db" + + echo + echo "=== arm $arm ===" + "$bin" --version 2>/dev/null | head -1 || true + + # filtered mode over an explicit repo list: no crawler, no firehose. + # HYDRANT_RELAY_HOSTS and HYDRANT_CRAWLER_URLS must be set EMPTY rather than + # unset -- unset falls back to built-in defaults and quietly adds drift. + env \ + HYDRANT_DATABASE_PATH="$db" \ + HYDRANT_API_BIND="127.0.0.1:$API_PORT" \ + HYDRANT_ENABLE_DEBUG=true \ + HYDRANT_DEBUG_PORT="$DEBUG_PORT" \ + HYDRANT_ENABLE_FIREHOSE=false \ + HYDRANT_ENABLE_CRAWLER=false \ + HYDRANT_RELAY_HOSTS="" \ + HYDRANT_CRAWLER_URLS="" \ + HYDRANT_FULL_NETWORK=false \ + HYDRANT_BACKFILL_STRATEGY=auto \ + HYDRANT_PLC_URL="$PLC_URL" \ + HYDRANT_DATA_COMPRESSION=zstd \ + HYDRANT_JOURNAL_COMPRESSION=zstd \ + "$bin" > "$dir/hydrant.log" 2>&1 & + local pid=$! + echo "$pid" > "$dir/pid" + echo "pid $pid" + + # shellcheck disable=SC2064 + trap "kill $pid 2>/dev/null || true" EXIT + + for _ in $(seq 60); do + if curl -s --max-time 2 "http://127.0.0.1:$API_PORT/stats" >/dev/null 2>&1; then break; fi + if ! kill -0 "$pid" 2>/dev/null; then + echo "FATAL: $arm died on startup, see $dir/hydrant.log" >&2 + tail -20 "$dir/hydrant.log" >&2 + exit 1 + fi + sleep 1 + done + + date +%s > "$dir/start_epoch" + + # telemetry sampler: rss, db size, and journal size over time + ( + echo "epoch,rss_kb,vmhwm_kb,db_bytes,journal_bytes,records,pending,collections,collections_bytes" + while kill -0 "$pid" 2>/dev/null; do + local_rss=$(awk '/^VmRSS:/{print $2}' "/proc/$pid/status" 2>/dev/null || echo 0) + local_hwm=$(awk '/^VmHWM:/{print $2}' "/proc/$pid/status" 2>/dev/null || echo 0) + dbb=$(du -sb "$db" 2>/dev/null | cut -f1 || echo 0) + jrn=$(du -sb "$db"/journal* 2>/dev/null | awk '{s+=$1} END{print s+0}') + st=$(curl -s --max-time 4 "http://127.0.0.1:$API_PORT/stats" 2>/dev/null || echo '{}') + rec=$(echo "$st" | python3 -c 'import sys,json;d=json.load(sys.stdin);print(d.get("counts",{}).get("records",0))' 2>/dev/null || echo 0) + pend=$(echo "$st" | python3 -c 'import sys,json;d=json.load(sys.stdin);print(d.get("counts",{}).get("pending",0))' 2>/dev/null || echo 0) + col=$(echo "$st" | python3 -c 'import sys,json;d=json.load(sys.stdin);print(d.get("counts",{}).get("collections",0))' 2>/dev/null || echo 0) + colb=$(echo "$st" | python3 -c 'import sys,json;d=json.load(sys.stdin);print(d.get("counts",{}).get("collections_bytes",0))' 2>/dev/null || echo 0) + echo "$(date +%s),$local_rss,$local_hwm,$dbb,$jrn,$rec,$pend,$col,$colb" + sleep "$INTERVAL" + done + ) > "$dir/telemetry.csv" & + local sampler=$! + + echo "feeding $(wc -l < "$SAMPLE") repos..." + curl -s -X PUT "http://127.0.0.1:$API_PORT/repos" \ + -H 'content-type: application/x-ndjson' \ + --data-binary "@$SAMPLE" -o "$dir/put-repos.json" \ + -w 'put /repos -> %{http_code}\n' + + # drain: pending must hit zero and stay there, since a repo can be requeued + echo "draining..." + local zero=0 elapsed=0 + while [ "$zero" -lt 4 ]; do + sleep "$INTERVAL" + elapsed=$((elapsed + INTERVAL)) + local pend + pend=$(curl -s --max-time 5 "http://127.0.0.1:$API_PORT/stats" \ + | python3 -c 'import sys,json;d=json.load(sys.stdin);print(d.get("counts",{}).get("pending",0))' 2>/dev/null || echo 1) + if [ "$pend" = "0" ]; then zero=$((zero + 1)); else zero=0; fi + if [ $((elapsed % 60)) -lt "$INTERVAL" ]; then echo " ${elapsed}s pending=$pend"; fi + if [ "$elapsed" -gt 7200 ]; then echo " giving up after 2h with pending=$pend" >&2; break; fi + done + date +%s > "$dir/end_epoch" + + # sizes are only comparable after compaction; the sizing procedure notes + # pre-compact events are inflated + curl -s --max-time 60 "http://127.0.0.1:$API_PORT/stats" > "$dir/stats-precompact.json" + echo "compacting..." + curl -s -X POST --max-time 3600 "http://127.0.0.1:$API_PORT/db/compact" -o "$dir/compact.json" \ + -w 'compact -> %{http_code}\n' + curl -s --max-time 60 "http://127.0.0.1:$API_PORT/stats" > "$dir/stats-final.json" + + # per-collection distribution, for reading the bytes/record numbers + curl -s --max-time 300 "http://127.0.0.1:$DEBUG_PORT/debug/iter?partition=counts&limit=400000" \ + -o "$dir/counts-iter.json" || echo "(counts iter failed, non-fatal)" + + kill "$sampler" 2>/dev/null || true + kill "$pid" 2>/dev/null || true + wait "$pid" 2>/dev/null || true + trap - EXIT + + # apparent size after a clean shutdown, excluding the journal + du -sb "$db" > "$dir/db-du-bytes.txt" 2>/dev/null || true + echo "arm $arm done" +} + +for spec in "$@"; do + arm=${spec%%=*} + bin=${spec#*=} + [ -x "$bin" ] || { echo "FATAL: $bin is not executable" >&2; exit 1; } + run_arm "$arm" "$bin" +done + +echo +echo "=== summary ===" +python3 - "$RUN_DIR" "$@" <<'PY' +import json, sys, pathlib +run = pathlib.Path(sys.argv[1]) +arms = [s.split("=", 1)[0] for s in sys.argv[2:]] + +rows = {} +for arm in arms: + f = run / arm / "stats-final.json" + if not f.exists(): + print(f"{arm}: no stats-final.json") + continue + d = json.loads(f.read_text()) + counts, sizes = d.get("counts", {}), d.get("sizes", {}) + recs = counts.get("records", 0) or 0 + rows[arm] = { + "records": recs, + "collections": counts.get("collections", 0), + "collections_bytes": counts.get("collections_bytes", 0), + "records_bytes": sizes.get("records", 0), + "counts_bytes": sizes.get("counts", 0), + "blocks_bytes": sizes.get("blocks", 0), + "total_bytes": sum(sizes.values()), + "records_b_per_rec": (sizes.get("records", 0) / recs) if recs else 0, + "counts_b_per_rec": (sizes.get("counts", 0) / recs) if recs else 0, + "total_b_per_rec": (sum(sizes.values()) / recs) if recs else 0, + } + tel = run / arm / "telemetry.csv" + if tel.exists(): + peak = 0 + for line in tel.read_text().splitlines()[1:]: + p = line.split(",") + if len(p) > 2 and p[2].isdigit(): + peak = max(peak, int(p[2])) + rows[arm]["peak_rss_mb"] = peak / 1024 + +for arm, r in rows.items(): + print(f"\n[{arm}]") + for k, v in r.items(): + print(f" {k:22} {v:,.2f}" if isinstance(v, float) else f" {k:22} {v:,}") + +if len(rows) == 2: + a, b = list(rows) + print(f"\n=== {b} vs {a} ===") + for k in ("records_b_per_rec", "counts_b_per_rec", "total_b_per_rec", "peak_rss_mb"): + va, vb = rows[a].get(k, 0), rows[b].get(k, 0) + if va: + print(f" {k:22} {va:,.2f} -> {vb:,.2f} ({(vb-va)/va*100:+.2f}%)") + ra, rb = rows[a]["records"], rows[b]["records"] + if ra != rb: + print(f" !! record counts differ ({ra:,} vs {rb:,}); bytes/record still") + print(f" comparable but the arms did not index the same set") +PY