diff --git a/AGENTS.md b/AGENTS.md index 1be86ae..105aa8c 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -74,7 +74,7 @@ To minimize latency in `apply_commit` and the backfill worker, events are stored Hydrant uses multiple `fjall` keyspaces: - `repos`: Maps `{DID}` -> `RepoState` (MessagePack). -- `records`: Partitioned by collection. Maps `{DID}|{RKey}` -> `{CID}` (Binary). +- `records`: Maps `{DID}|{COL}|{RKey}` -> `{CID}` (Binary). - `blocks`: Maps `{CID}` -> `Block Data` (Raw CBOR). - `events`: Maps `{ID}` (u64) -> `StoredEvent` (MessagePack). This is the source for the JSON stream API. - `cursors`: Maps `firehose_cursor` or `crawler_cursor` -> `Value` (u64/i64 BE Bytes). diff --git a/src/api/debug.rs b/src/api/debug.rs index 701bb8e..19dc1f8 100644 --- a/src/api/debug.rs +++ b/src/api/debug.rs @@ -1,6 +1,6 @@ +use crate::api::AppState; use crate::db::keys; use crate::types::{RepoState, ResyncState, StoredEvent}; -use crate::{api::AppState, db::types::TrimmedDid}; use axum::{ Json, extract::{Query, State}, @@ -42,14 +42,10 @@ pub async fn handle_debug_count( .map_err(|_| StatusCode::BAD_REQUEST)?; let db = &state.db; - let ks = db - .record_partition(&req.collection) - .map_err(|_| StatusCode::NOT_FOUND)?; + let ks = db.records.clone(); - // {TrimmedDid}\x00 - let mut prefix = Vec::new(); - TrimmedDid::from(&did).write_to_vec(&mut prefix); - prefix.push(keys::SEP); + // {TrimmedDid}|{collection}| + let prefix = keys::record_prefix_collection(&did, &req.collection); let count = tokio::task::spawn_blocking(move || { let start_key = prefix.clone(); @@ -255,12 +251,7 @@ fn get_keyspace_by_name(db: &crate::db::Db, name: &str) -> Result Ok(db.resync.clone()), "events" => Ok(db.events.clone()), "counts" => Ok(db.counts.clone()), - _ => { - if let Some(col) = name.strip_prefix(crate::db::RECORDS_PARTITION_PREFIX) { - db.record_partition(col).map_err(|_| StatusCode::NOT_FOUND) - } else { - Err(StatusCode::BAD_REQUEST) - } - } + "records" => Ok(db.records.clone()), + _ => Err(StatusCode::BAD_REQUEST), } } diff --git a/src/api/xrpc.rs b/src/api/xrpc.rs index 354cc86..c69b406 100644 --- a/src/api/xrpc.rs +++ b/src/api/xrpc.rs @@ -78,13 +78,13 @@ pub async fn handle_get_record( .await .map_err(|e| bad_request(GetRecord::NSID, e))?; - let partition = db - .record_partition(req.collection.as_str()) - .map_err(|e| internal_error(GetRecord::NSID, e))?; - - let db_key = keys::record_key(&did, &DbRkey::new(req.rkey.0.as_str())); + let db_key = keys::record_key( + &did, + req.collection.as_str(), + &DbRkey::new(req.rkey.0.as_str()), + ); - let cid_bytes = Db::get(partition, db_key) + let cid_bytes = Db::get(db.records.clone(), db_key) .await .map_err(|e| internal_error(GetRecord::NSID, e))?; @@ -132,11 +132,9 @@ pub async fn handle_list_records( .await .map_err(|e| bad_request(ListRecords::NSID, e))?; - let ks = db - .record_partition(req.collection.as_str()) - .map_err(|e| internal_error(ListRecords::NSID, e))?; + let ks = db.records.clone(); - let prefix = keys::record_prefix(&did); + let prefix = keys::record_prefix_collection(&did, req.collection.as_str()); let limit = req.limit.unwrap_or(50).min(100) as usize; let reverse = req.reverse.unwrap_or(false); diff --git a/src/backfill/mod.rs b/src/backfill/mod.rs index a4ac62d..4660980 100644 --- a/src/backfill/mod.rs +++ b/src/backfill/mod.rs @@ -29,6 +29,17 @@ pub mod manager; use crate::ingest::{BufferTx, IngestMessage}; +trait SliceSplitExt { + fn split_once<'a>(&'a self, delimiter: impl Fn(&u8) -> bool) -> Option<(&'a [u8], &'a [u8])>; +} + +impl SliceSplitExt for [u8] { + fn split_once<'a>(&'a self, delimiter: impl Fn(&u8) -> bool) -> Option<(&'a [u8], &'a [u8])> { + let idx = self.iter().position(delimiter)?; + Some((&self[..idx], &self[idx + 1..])) + } +} + struct AdaptiveLimiter { current_limit: usize, max_limit: usize, @@ -608,26 +619,29 @@ async fn process_did<'i>( let mut batch = app_state.db.inner.batch(); let store = mst.storage(); - let prefix = keys::record_prefix(&did); + let prefix = keys::record_prefix_did(&did); let mut existing_cids: HashMap<(SmolStr, DbRkey), SmolStr> = HashMap::new(); - let mut partitions = Vec::new(); - app_state.db.record_partitions.iter_sync(|col, ks| { - partitions.push((col.clone(), ks.clone())); - true - }); - - for (col_name, ks) in partitions { - for guard in ks.prefix(&prefix) { - let (key, cid_bytes) = guard.into_inner().into_diagnostic()?; - let rkey = keys::parse_rkey(&key[prefix.len()..]) - .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(); - - existing_cids.insert((col_name.as_str().into(), rkey), cid); - } + for guard in app_state.db.records.prefix(&prefix) { + let (key, cid_bytes) = guard.into_inner().into_diagnostic()?; + // key is did|collection|rkey + // skip did| + let remaining = &key[prefix.len()..]; + let (collection_bytes, rkey_bytes) = + SliceSplitExt::split_once(remaining, |b| *b == keys::SEP) + .ok_or_else(|| miette::miette!("invalid record key format: {key:?}"))?; + + let collection = std::str::from_utf8(collection_bytes) + .map_err(|e| miette::miette!("invalid collection utf8: {e}"))?; + + let rkey = keys::parse_rkey(rkey_bytes) + .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(); + + existing_cids.insert((collection.into(), rkey), cid); } for (key, cid) in leaves { @@ -640,7 +654,6 @@ async fn process_did<'i>( let rkey = DbRkey::new(rkey); let path = (collection.to_smolstr(), rkey.clone()); let cid_obj = Cid::ipld(cid); - let partition = app_state.db.record_partition(collection)?; // check if this record already exists with same CID let (action, is_new) = if let Some(existing_cid) = existing_cids.remove(&path) { @@ -654,11 +667,11 @@ async fn process_did<'i>( }; trace!("{action} {did}/{collection}/{rkey} ({cid})"); - // Key is just did|rkey - let db_key = keys::record_key(&did, &rkey); + // Key is did|collection|rkey + let db_key = keys::record_key(&did, collection, &rkey); batch.insert(&app_state.db.blocks, cid.to_bytes(), val.as_ref()); - batch.insert(&partition, db_key, cid.to_bytes()); + batch.insert(&app_state.db.records, db_key, cid.to_bytes()); added_blocks += 1; if is_new { @@ -686,9 +699,11 @@ async fn process_did<'i>( // remove any remaining existing records (they weren't in the new MST) for ((collection, rkey), cid) in existing_cids { trace!("remove {did}/{collection}/{rkey} ({cid})"); - let partition = app_state.db.record_partition(collection.as_str())?; - batch.remove(&partition, keys::record_key(&did, &rkey)); + batch.remove( + &app_state.db.records, + keys::record_key(&did, &collection, &rkey), + ); let event_id = app_state.db.next_event_id.fetch_add(1, Ordering::SeqCst); let evt = StoredEvent { diff --git a/src/config.rs b/src/config.rs index 6a29326..35e5575 100644 --- a/src/config.rs +++ b/src/config.rs @@ -1,4 +1,4 @@ -use miette::{IntoDiagnostic, Result}; +use miette::Result; use smol_str::SmolStr; use std::fmt; use std::path::PathBuf; @@ -61,8 +61,7 @@ pub struct Config { pub db_blocks_memtable_size_mb: u64, pub db_repos_memtable_size_mb: u64, pub db_events_memtable_size_mb: u64, - pub db_records_default_memtable_size_mb: u64, - pub db_records_partition_overrides: Vec<(glob::Pattern, u64)>, + pub db_records_memtable_size_mb: u64, pub crawler_max_pending_repos: usize, pub crawler_resume_pending_repos: usize, } @@ -126,11 +125,9 @@ impl Config { default_db_worker_threads, default_db_max_journaling_size_mb, default_db_memtable_size_mb, - default_records_memtable_size_mb, - default_partition_overrides, - ): (usize, u64, u64, u64, &str) = full_network - .then_some((8usize, 1024u64, 192u64, 8u64, "app.bsky.*=64")) - .unwrap_or((4usize, 512u64, 64u64, 16u64, "")); + ): (usize, u64, u64) = full_network + .then_some((8usize, 1024u64, 192u64)) + .unwrap_or((4usize, 512u64, 64u64)); let db_worker_threads = cfg!("DB_WORKER_THREADS", default_db_worker_threads); let db_max_journaling_size_mb = cfg!( @@ -145,33 +142,12 @@ impl Config { cfg!("DB_REPOS_MEMTABLE_SIZE_MB", default_db_memtable_size_mb); let db_events_memtable_size_mb = cfg!("DB_EVENTS_MEMTABLE_SIZE_MB", default_db_memtable_size_mb); - let db_records_default_memtable_size_mb = cfg!( - "DB_RECORDS_DEFAULT_MEMTABLE_SIZE_MB", - default_records_memtable_size_mb - ); + let db_records_memtable_size_mb = + cfg!("DB_RECORDS_MEMTABLE_SIZE_MB", default_db_memtable_size_mb); let crawler_max_pending_repos = cfg!("CRAWLER_MAX_PENDING_REPOS", 2000usize); let crawler_resume_pending_repos = cfg!("CRAWLER_RESUME_PENDING_REPOS", 1000usize); - let db_records_partition_overrides: Vec<(glob::Pattern, u64)> = - std::env::var("HYDRANT_DB_RECORDS_PARTITION_OVERRIDES") - .unwrap_or_else(|_| default_partition_overrides.to_string()) - .split(',') - .filter(|s| !s.is_empty()) - .map(|s| { - let mut parts = s.split('='); - let pattern = parts - .next() - .ok_or_else(|| miette::miette!("invalid partition override format"))?; - let size = parts - .next() - .ok_or_else(|| miette::miette!("invalid partition override format"))? - .parse::() - .into_diagnostic()?; - Ok((glob::Pattern::new(pattern).into_diagnostic()?, size)) - }) - .collect::>>()?; - Ok(Self { database_path, relay_host, @@ -197,8 +173,7 @@ impl Config { db_blocks_memtable_size_mb, db_repos_memtable_size_mb, db_events_memtable_size_mb, - db_records_default_memtable_size_mb, - db_records_partition_overrides: db_records_partition_overrides, + db_records_memtable_size_mb, crawler_max_pending_repos, crawler_resume_pending_repos, }) @@ -274,8 +249,8 @@ impl fmt::Display for Config { )?; writeln!( f, - " db records def memtable: {} mb", - self.db_records_default_memtable_size_mb + " db records memtable: {} mb", + self.db_records_memtable_size_mb )?; writeln!( diff --git a/src/db/keys.rs b/src/db/keys.rs index 916efef..01a84a6 100644 --- a/src/db/keys.rs +++ b/src/db/keys.rs @@ -15,8 +15,8 @@ pub fn repo_key<'a>(did: &'a Did) -> Vec { vec } -// prefix format: {DID}\x00 -pub fn record_prefix(did: &Did) -> Vec { +// prefix format: {DID}| (DID trimmed) +pub fn record_prefix_did(did: &Did) -> Vec { let repo = TrimmedDid::from(did); let mut prefix = Vec::with_capacity(repo.len() + 1); repo.write_to_vec(&mut prefix); @@ -24,12 +24,25 @@ pub fn record_prefix(did: &Did) -> Vec { prefix } -// key format: {DID}\x00{rkey} -pub fn record_key(did: &Did, rkey: &DbRkey) -> Vec { +// prefix format: {DID}|{collection}| +pub fn record_prefix_collection(did: &Did, collection: &str) -> Vec { + let repo = TrimmedDid::from(did); + let mut prefix = Vec::with_capacity(repo.len() + 1 + collection.len() + 1); + repo.write_to_vec(&mut prefix); + prefix.push(SEP); + prefix.extend_from_slice(collection.as_bytes()); + prefix.push(SEP); + prefix +} + +// key format: {DID}|{collection}|{rkey} +pub fn record_key(did: &Did, collection: &str, rkey: &DbRkey) -> Vec { let repo = TrimmedDid::from(did); - let mut key = Vec::with_capacity(repo.len() + rkey.len() + 1); + let mut key = Vec::with_capacity(repo.len() + 1 + collection.len() + 1 + rkey.len() + 1); repo.write_to_vec(&mut key); key.push(SEP); + key.extend_from_slice(collection.as_bytes()); + key.push(SEP); write_rkey(&mut key, rkey); key } diff --git a/src/db/mod.rs b/src/db/mod.rs index eeb0847..422749f 100644 --- a/src/db/mod.rs +++ b/src/db/mod.rs @@ -4,7 +4,7 @@ use jacquard::IntoStatic; use jacquard_common::types::string::Did; use miette::{Context, IntoDiagnostic, Result}; use scc::HashMap; -use smol_str::{SmolStr, format_smolstr}; +use smol_str::SmolStr; use std::sync::Arc; @@ -15,8 +15,6 @@ use std::sync::atomic::AtomicU64; use tokio::sync::broadcast; use tracing::error; -pub const RECORDS_PARTITION_PREFIX: &str = "r:"; - fn default_opts() -> KeyspaceCreateOptions { KeyspaceCreateOptions::default() } @@ -24,7 +22,7 @@ fn default_opts() -> KeyspaceCreateOptions { pub struct Db { pub inner: Arc, pub repos: Keyspace, - pub record_partitions: HashMap, + pub records: Keyspace, pub blocks: Keyspace, pub cursors: Keyspace, pub pending: Keyspace, @@ -35,8 +33,6 @@ pub struct Db { pub event_tx: broadcast::Sender, pub next_event_id: Arc, pub counts_map: HashMap, - pub record_partition_overrides: Vec<(glob::Pattern, u64)>, - pub record_partition_default_size: u64, } impl Db { @@ -86,18 +82,10 @@ impl Db { )?; let counts = open_ks("counts", opts().expect_point_read_hits(true))?; - let record_partitions = HashMap::new(); - { - let names = db.list_keyspace_names(); - for name in names { - let name_str: &str = name.as_ref(); - if let Some(collection) = name_str.strip_prefix(RECORDS_PARTITION_PREFIX) { - let opts = Self::get_record_partition_opts(cfg, collection); - let ks = db.keyspace(name_str, move || opts).into_diagnostic()?; - let _ = record_partitions.insert_sync(collection.to_string(), ks); - } - } - } + let records = open_ks( + "records", + opts().max_memtable_size(cfg.db_records_memtable_size_mb * 1024 * 1024), + )?; let mut last_id = 0; if let Some(guard) = events.iter().next_back() { @@ -128,7 +116,7 @@ impl Db { Ok(Self { inner: db, repos, - record_partitions, + records, blocks, cursors, pending, @@ -139,49 +127,9 @@ impl Db { event_tx, counts_map, next_event_id: Arc::new(AtomicU64::new(last_id + 1)), - record_partition_overrides: cfg.db_records_partition_overrides.clone(), - record_partition_default_size: cfg.db_records_default_memtable_size_mb, }) } - fn get_record_partition_opts( - cfg: &crate::config::Config, - collection: &str, - ) -> KeyspaceCreateOptions { - let size = cfg - .db_records_partition_overrides - .iter() - .find(|(p, _)| p.matches(collection)) - .map(|(_, s)| *s) - .unwrap_or(cfg.db_records_default_memtable_size_mb); - - default_opts().max_memtable_size(size * 1024 * 1024) - } - - pub fn record_partition(&self, collection: &str) -> Result { - use scc::hash_map::Entry; - match self.record_partitions.entry_sync(collection.to_string()) { - Entry::Occupied(o) => Ok(o.get().clone()), - Entry::Vacant(v) => { - let name = format_smolstr!("{}{}", RECORDS_PARTITION_PREFIX, collection); - let size = self - .record_partition_overrides - .iter() - .find(|(p, _)| p.matches(collection)) - .map(|(_, s)| *s) - .unwrap_or(self.record_partition_default_size); - - let ks = self - .inner - .keyspace(&name, move || { - default_opts().max_memtable_size(size * 1024 * 1024) - }) - .into_diagnostic()?; - Ok(v.insert_entry(ks).get().clone()) - } - } - } - pub fn persist(&self) -> Result<()> { self.inner.persist(PersistMode::SyncAll).into_diagnostic()?; Ok(()) diff --git a/src/ops.rs b/src/ops.rs index 81a11aa..17cc36b 100644 --- a/src/ops.rs +++ b/src/ops.rs @@ -80,19 +80,11 @@ pub fn delete_repo<'batch>( batch.remove(&db.resync_buffer, k); } - // 2. delete from records (all partitions) - let mut partitions = Vec::new(); - db.record_partitions.iter_sync(|_, v| { - partitions.push(v.clone()); - true - }); - - let records_prefix = keys::record_prefix(did); - for ks in partitions { - for guard in ks.prefix(&records_prefix) { - let k = guard.key().into_diagnostic()?; - batch.remove(&ks, k); - } + // 2. delete from records + let records_prefix = keys::record_prefix_did(did); + for guard in db.records.prefix(&records_prefix) { + let k = guard.key().into_diagnostic()?; + batch.remove(&db.records, k); } // 3. reset collection counts @@ -249,8 +241,7 @@ pub fn apply_commit<'batch, 'db, 'commit, 's>( for op in &commit.ops { let (collection, rkey) = parse_path(&op.path)?; let rkey = DbRkey::new(rkey); - let partition = db.record_partition(collection)?; - let db_key = keys::record_key(did, &rkey); + let db_key = keys::record_key(did, collection, &rkey); let event_id = db.next_event_id.fetch_add(1, Ordering::SeqCst); @@ -261,7 +252,7 @@ pub fn apply_commit<'batch, 'db, 'commit, 's>( continue; }; batch.insert( - &partition, + &db.records, db_key.clone(), cid.to_ipld() .into_diagnostic() @@ -276,7 +267,7 @@ pub fn apply_commit<'batch, 'db, 'commit, 's>( } } DbAction::Delete => { - batch.remove(&partition, db_key); + batch.remove(&db.records, db_key); // accumulate counts records_delta -= 1;