From 725b546893827bf78f0dabae4f2d142d6a88f1df Mon Sep 17 00:00:00 2001 From: dawn <90008@klbr.net> Date: Sat, 11 Jul 2026 18:35:27 +0300 Subject: [PATCH] [db] group mode-specific keyspaces into per-mode structs Db now embeds IndexerDb/StreamDb/JetstreamDb/RelayDb (declared in db/keyspaces.rs, the single place mode keyspaces live) behind one cfg each, instead of 17 cfg'd fields sprinkled through the struct, open, and compact. each group owns its keyspace options, post-migration init (sequence/watermark seeding), and compaction targets. ~200 callsites renamed to the grouped paths (db.indexer.records, db.stream.events, db.jetstream.tx, db.relay.next_seq, ...). --- src/api/debug.rs | 24 +- src/backfill/manager.rs | 16 +- src/backfill/sparse.rs | 16 +- src/backfill/worker.rs | 2 +- src/backfill/worker/process.rs | 32 +-- src/backfill/worker/task.rs | 12 +- src/backlinks/mod.rs | 2 +- src/control/hydrant/run.rs | 4 +- src/control/indexer.rs | 4 +- src/control/repos/indexer.rs | 32 +-- src/control/repos/mod.rs | 2 +- src/control/stats.rs | 22 +- src/control/stream/indexer.rs | 10 +- src/control/stream/jetstream.rs | 14 +- src/control/stream/relay.rs | 8 +- src/crawler/worker.rs | 2 +- src/db/ephemeral.rs | 76 +++---- src/db/keyspaces.rs | 384 ++++++++++++++++++++++++++++++++ src/db/lifecycle_counts.rs | 22 +- src/db/migration/v8.rs | 10 +- src/db/mod.rs | 87 +++----- src/db/open.rs | 266 ++-------------------- src/db/train.rs | 6 +- src/ingest/indexer/shard.rs | 20 +- src/ingest/relay/sink/relay.rs | 12 +- src/jetstream.rs | 8 +- src/ops.rs | 48 ++-- 27 files changed, 640 insertions(+), 501 deletions(-) create mode 100644 src/db/keyspaces.rs diff --git a/src/api/debug.rs b/src/api/debug.rs index 34d0f6c..d07517f 100644 --- a/src/api/debug.rs +++ b/src/api/debug.rs @@ -63,7 +63,7 @@ pub async fn handle_debug_count( .map_err(|_| StatusCode::BAD_REQUEST)?; let db = &state.db; - let ks = db.records.clone(); + let ks = db.indexer.records.clone(); // {TrimmedDid}|{collection}| let prefix = keys::record_prefix_collection(&did, &req.collection); @@ -287,19 +287,19 @@ fn get_keyspace_by_name(db: &crate::db::Db, name: &str) -> Result Ok(db.counts.clone()), "cursors" => Ok(db.cursors.clone()), #[cfg(feature = "indexer")] - "blocks" => Ok(db.blocks.clone()), + "blocks" => Ok(db.indexer.blocks.clone()), #[cfg(feature = "indexer")] - "pending" => Ok(db.pending.clone()), + "pending" => Ok(db.indexer.pending.clone()), #[cfg(feature = "indexer")] - "resync" => Ok(db.resync.clone()), + "resync" => Ok(db.indexer.resync.clone()), #[cfg(feature = "indexer_stream")] - "events" => Ok(db.events.clone()), + "events" => Ok(db.stream.events.clone()), #[cfg(feature = "jetstream")] - "jetstream_events" => Ok(db.jetstream_events.clone()), + "jetstream_events" => Ok(db.jetstream.events.clone()), #[cfg(feature = "relay")] - "relay_events" => Ok(db.relay_events.clone()), + "relay_events" => Ok(db.relay.events.clone()), #[cfg(feature = "indexer")] - "records" => Ok(db.records.clone()), + "records" => Ok(db.indexer.records.clone()), _ => Err(StatusCode::BAD_REQUEST), } } @@ -413,9 +413,9 @@ pub async fn handle_debug_seed_events( for _ in 0..req.count { let seq = state .db - .next_event_id + .stream.next_event_id .fetch_add(1, std::sync::atomic::Ordering::SeqCst); - batch.insert(&state.db.events, crate::db::keys::event_key(seq), b"dummy"); + batch.insert(&state.db.stream.events, crate::db::keys::event_key(seq), b"dummy"); } } } else if req.partition == "relay_events" { @@ -424,10 +424,10 @@ pub async fn handle_debug_seed_events( for _ in 0..req.count { let seq = state .db - .next_relay_seq + .relay.next_seq .fetch_add(1, std::sync::atomic::Ordering::SeqCst); batch.insert( - &state.db.relay_events, + &state.db.relay.events, crate::db::keys::relay_event_key(seq), b"dummy", ); diff --git a/src/backfill/manager.rs b/src/backfill/manager.rs index 8f531a6..c1f9f82 100644 --- a/src/backfill/manager.rs +++ b/src/backfill/manager.rs @@ -15,7 +15,7 @@ pub fn queue_gone_backfills(state: &Arc) -> Result<()> { let mut batch = state.db.inner.batch(); let mut lifecycle_counts = state.db.lifecycle_counts(); - for guard in state.db.resync.iter() { + for guard in state.db.indexer.resync.iter() { let (key, val) = guard.into_inner().into_diagnostic()?; let did = match TrimmedDid::try_from(key.as_ref()) { Ok(did) => did.to_did(), @@ -48,12 +48,12 @@ pub fn queue_gone_backfills(state: &Arc) -> Result<()> { let mut metadata = crate::db::deser_repo_meta(&metadata_bytes)?; // move from resync back into pending - batch.remove(&state.db.resync, key.clone()); + batch.remove(&state.db.indexer.resync, key.clone()); let old_pending = keys::pending_key(metadata.index_id); - batch.remove(&state.db.pending, old_pending); + batch.remove(&state.db.indexer.pending, old_pending); metadata.index_id = rand::random::(); batch.insert( - &state.db.pending, + &state.db.indexer.pending, keys::pending_key(metadata.index_id), key.clone(), ); @@ -95,7 +95,7 @@ pub fn retry_worker(state: Arc) { let mut batch = state.db.inner.batch(); let mut lifecycle_counts = state.db.lifecycle_counts(); - for guard in db.resync.iter() { + for guard in db.indexer.resync.iter() { let (key, value) = match guard.into_inner() { Ok(t) => t, Err(e) => { @@ -186,9 +186,9 @@ pub fn retry_worker(state: Arc) { } // move from resync back into pending - batch.remove(&state.db.resync, key.clone()); - batch.remove(&state.db.pending, old_pending); - batch.insert(&state.db.pending, new_pending, key.clone()); + batch.remove(&state.db.indexer.resync, key.clone()); + batch.remove(&state.db.indexer.pending, old_pending); + batch.insert(&state.db.indexer.pending, new_pending, key.clone()); batch.insert(&state.db.repo_metadata, &metadata_key, serialized_metadata); transitions += 1; } diff --git a/src/backfill/sparse.rs b/src/backfill/sparse.rs index 99ceecb..b357540 100644 --- a/src/backfill/sparse.rs +++ b/src/backfill/sparse.rs @@ -436,7 +436,7 @@ async fn persist_sparse_backfill( let mut existing_cids: HashMap<(SmolStr, DbRkey), SmolStr> = HashMap::new(); if !ephemeral { - for guard in app_state.db.records.prefix(&prefix) { + for guard in app_state.db.indexer.records.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 @@ -503,9 +503,9 @@ async fn persist_sparse_backfill( let block_key = Slice::from(keys::block_key(collection, &cid_raw)); if !ephemeral { if !only_index_links { - batch.insert(&app_state.db.blocks, block_key.clone(), val.as_ref()); + batch.insert(&app_state.db.indexer.blocks, block_key.clone(), val.as_ref()); } - batch.insert(&app_state.db.records, db_key, cid_raw); + batch.insert(&app_state.db.indexer.records, db_key, cid_raw); #[cfg(feature = "backlinks")] if let Ok(value) = serde_ipld_dagcbor::from_slice::(val.as_ref()) @@ -528,7 +528,7 @@ async fn persist_sparse_backfill( #[cfg(feature = "indexer_stream")] { - let event_id = app_state.db.next_event_id.fetch_add(1, Ordering::SeqCst); + let event_id = app_state.db.stream.next_event_id.fetch_add(1, Ordering::SeqCst); let evt = StoredEvent { live: false, did: TrimmedDid::from(&did), @@ -545,7 +545,7 @@ async fn persist_sparse_backfill( }, }; let bytes = rmp_serde::to_vec(&evt).into_diagnostic()?; - batch.insert(&app_state.db.events, keys::event_key(event_id), bytes); + batch.insert(&app_state.db.stream.events, keys::event_key(event_id), bytes); #[cfg(feature = "jetstream")] { @@ -566,7 +566,7 @@ async fn persist_sparse_backfill( trace!(collection = %collection, rkey = %rkey, cid = %cid, "remove sparse-stale record"); batch.remove( - &app_state.db.records, + &app_state.db.indexer.records, keys::record_key(&did, &collection, &rkey), ); #[cfg(feature = "backlinks")] @@ -580,7 +580,7 @@ async fn persist_sparse_backfill( #[cfg(feature = "indexer_stream")] { - let event_id = app_state.db.next_event_id.fetch_add(1, Ordering::SeqCst); + let event_id = app_state.db.stream.next_event_id.fetch_add(1, Ordering::SeqCst); let evt = StoredEvent { live: false, did: TrimmedDid::from(&did), @@ -591,7 +591,7 @@ async fn persist_sparse_backfill( data: StoredData::Nothing, }; let bytes = rmp_serde::to_vec(&evt).into_diagnostic()?; - batch.insert(&app_state.db.events, keys::event_key(event_id), bytes); + batch.insert(&app_state.db.stream.events, keys::event_key(event_id), bytes); #[cfg(feature = "jetstream")] { diff --git a/src/backfill/worker.rs b/src/backfill/worker.rs index 6671286..5e2cfa5 100644 --- a/src/backfill/worker.rs +++ b/src/backfill/worker.rs @@ -91,7 +91,7 @@ impl BackfillWorker { self.enabled.wait_enabled("backfill").await; let mut spawned = 0; - for guard in self.state.db.pending.iter() { + for guard in self.state.db.indexer.pending.iter() { let (key, value) = match guard.into_inner() { Ok(kv) => kv, Err(e) => { diff --git a/src/backfill/worker/process.rs b/src/backfill/worker/process.rs index f18c48f..c5154a9 100644 --- a/src/backfill/worker/process.rs +++ b/src/backfill/worker/process.rs @@ -91,7 +91,7 @@ pub(crate) async fn process_did( .then_some(None) .unwrap_or_else(|| Some(status.into())), }; - let _ = app_state.db.event_tx.send(ops::make_account_event(db, evt)); + let _ = app_state.db.stream.event_tx.send(ops::make_account_event(db, evt)); }; if strategy != BackfillStrategy::Full { @@ -151,7 +151,7 @@ pub(crate) async fn process_did( pending_key.as_ref(), GaugeState::Synced, )?; - batch.remove(&app_state_clone.db.pending, pending_key.clone()); + batch.remove(&app_state_clone.db.indexer.pending, pending_key.clone()); if applied { batch.remove(&app_state_clone.db.repos, &did_key); batch.remove(&app_state_clone.db.repo_metadata, &metadata_key); @@ -234,7 +234,7 @@ pub(crate) async fn process_did( pending_key.as_ref(), GaugeState::Synced, )?; - batch.remove(&db.pending, pending_key.clone()); + batch.remove(&db.indexer.pending, pending_key.clone()); if applied { if let Err(e) = crate::ops::delete_repo(&mut batch, db, did, &state) { error!(err = %e, "failed to wipe repo during backfill"); @@ -278,7 +278,7 @@ pub(crate) async fn process_did( pending_key.as_ref(), GaugeState::Resync(None), )?; - batch.remove(&db.pending, pending_key.clone()); + batch.remove(&db.indexer.pending, pending_key.clone()); if applied { Db::update_repo_state( &mut batch, @@ -287,7 +287,7 @@ pub(crate) async fn process_did( move |state, (key, batch)| { state.active = false; state.status = status; - batch.insert(&db.resync, key, resync_bytes); + batch.insert(&db.indexer.resync, key, resync_bytes); Ok((true, ())) }, )?; @@ -409,7 +409,7 @@ pub(crate) async fn process_did( let mut existing_cids: HashMap<(SmolStr, DbRkey), SmolStr> = HashMap::new(); if !ephemeral { - for guard in app_state.db.records.prefix(&prefix) { + for guard in app_state.db.indexer.records.prefix(&prefix) { let (key, cid_bytes) = guard.into_inner().into_diagnostic()?; // key is did|collection|rkey // skip did| @@ -480,9 +480,9 @@ pub(crate) async fn process_did( let block_key = Slice::from(keys::block_key(collection, &cid_raw)); if !ephemeral { if !only_index_links { - batch.insert(&app_state.db.blocks, block_key.clone(), val.as_ref()); + batch.insert(&app_state.db.indexer.blocks, block_key.clone(), val.as_ref()); } - batch.insert(&app_state.db.records, db_key, cid_raw); + batch.insert(&app_state.db.indexer.records, db_key, cid_raw); #[cfg(feature = "backlinks")] if let Ok(value) = serde_ipld_dagcbor::from_slice::(val.as_ref()) { crate::backlinks::store::index_record( @@ -503,7 +503,7 @@ pub(crate) async fn process_did( #[cfg(feature = "indexer_stream")] { - let event_id = app_state.db.next_event_id.fetch_add(1, Ordering::SeqCst); + let event_id = app_state.db.stream.next_event_id.fetch_add(1, Ordering::SeqCst); let evt = StoredEvent { live: false, did: TrimmedDid::from(&did), @@ -520,7 +520,7 @@ pub(crate) async fn process_did( }, }; let bytes = rmp_serde::to_vec(&evt).into_diagnostic()?; - batch.insert(&app_state.db.events, keys::event_key(event_id), bytes); + batch.insert(&app_state.db.stream.events, keys::event_key(event_id), bytes); #[cfg(feature = "jetstream")] { @@ -549,7 +549,7 @@ pub(crate) async fn process_did( // we dont have to put if ephemeral around here since // existing_cids will be empty anyway batch.remove( - &app_state.db.records, + &app_state.db.indexer.records, keys::record_key(&did, &collection, &rkey), ); #[cfg(feature = "backlinks")] @@ -563,7 +563,7 @@ pub(crate) async fn process_did( #[cfg(feature = "indexer_stream")] { - let event_id = app_state.db.next_event_id.fetch_add(1, Ordering::SeqCst); + let event_id = app_state.db.stream.next_event_id.fetch_add(1, Ordering::SeqCst); let evt = StoredEvent { live: false, did: TrimmedDid::from(&did), @@ -574,7 +574,7 @@ pub(crate) async fn process_did( data: StoredData::Nothing, }; let bytes = rmp_serde::to_vec(&evt).into_diagnostic()?; - batch.insert(&app_state.db.events, keys::event_key(event_id), bytes); + batch.insert(&app_state.db.stream.events, keys::event_key(event_id), bytes); #[cfg(feature = "jetstream")] { @@ -667,7 +667,7 @@ pub(crate) async fn process_did( backfill_pending_key.as_ref(), GaugeState::Synced, )?; - batch.remove(&app_state.db.pending, backfill_pending_key.clone()); + batch.remove(&app_state.db.indexer.pending, backfill_pending_key.clone()); if applied { batch.remove(&app_state.db.repos, &did_key); batch.remove(&app_state.db.repo_metadata, &metadata_key); @@ -693,8 +693,8 @@ pub(crate) async fn process_did( ); #[cfg(feature = "indexer_stream")] - let _ = db.event_tx.send(BroadcastEvent::Persisted( - db.next_event_id.load(Ordering::SeqCst) - 1, + let _ = db.stream.event_tx.send(BroadcastEvent::Persisted( + db.stream.next_event_id.load(Ordering::SeqCst) - 1, )); trace!("complete"); diff --git a/src/backfill/worker/task.rs b/src/backfill/worker/task.rs index 1bc9a12..de6648e 100644 --- a/src/backfill/worker/task.rs +++ b/src/backfill/worker/task.rs @@ -53,9 +53,9 @@ pub(crate) async fn did_task( pending_key.as_ref(), GaugeState::Synced, )?; - batch.remove(&db.pending, pending_key.clone()); + batch.remove(&db.indexer.pending, pending_key.clone()); if applied { - batch.remove(&db.resync, &did_key); + batch.remove(&db.indexer.resync, &did_key); } let lifecycle_reservation = lifecycle_counts.stage(&mut batch); batch.commit().into_diagnostic()?; @@ -106,7 +106,7 @@ pub(crate) async fn did_task( pending_key.as_ref(), GaugeState::Synced, )?; - batch.remove(&db.pending, pending_key); + batch.remove(&db.indexer.pending, pending_key); let lifecycle_reservation = lifecycle_counts.stage(&mut batch); batch.commit().into_diagnostic()?; db.apply_lifecycle_counts(lifecycle_reservation); @@ -143,7 +143,7 @@ pub(crate) async fn did_task( let did_key = keys::repo_key(did); // 1. get current retry count - let existing_state = Db::get(db.resync.clone(), &did_key).await.and_then(|b| { + let existing_state = Db::get(db.indexer.resync.clone(), &did_key).await.and_then(|b| { b.map(|b| rmp_serde::from_slice::(&b).into_diagnostic()) .transpose() })?; @@ -200,9 +200,9 @@ pub(crate) async fn did_task( pending_key.as_ref(), GaugeState::Resync(Some(error_kind)), )?; - batch.remove(&state.db.pending, pending_key.clone()); + batch.remove(&state.db.indexer.pending, pending_key.clone()); if applied { - batch.insert(&state.db.resync, &did_key, serialized_resync_state); + batch.insert(&state.db.indexer.resync, &did_key, serialized_resync_state); if let Some(state_bytes) = serialized_repo_state { batch.insert(&state.db.repos, &did_key, state_bytes); } diff --git a/src/backlinks/mod.rs b/src/backlinks/mod.rs index 1753cff..6f7cca0 100644 --- a/src/backlinks/mod.rs +++ b/src/backlinks/mod.rs @@ -175,7 +175,7 @@ impl BacklinksFetch { continue; } - let Some(cid) = store::lookup_cid_from_ks(&db.records, did, col, rkey) else { + let Some(cid) = store::lookup_cid_from_ks(&db.indexer.records, did, col, rkey) else { // record deleted after backlink was written continue; }; diff --git a/src/control/hydrant/run.rs b/src/control/hydrant/run.rs index 254abbe..ae73849 100644 --- a/src/control/hydrant/run.rs +++ b/src/control/hydrant/run.rs @@ -178,9 +178,9 @@ impl Hydrant { let state = state.clone(); let get_id = |state: &AppState| { #[cfg(feature = "indexer_stream")] - let id = state.db.next_event_id.load(Ordering::Relaxed); + let id = state.db.stream.next_event_id.load(Ordering::Relaxed); #[cfg(feature = "relay")] - let id = state.db.next_relay_seq.load(Ordering::Relaxed); + let id = state.db.relay.next_seq.load(Ordering::Relaxed); id }; let mut last_id = get_id(&state); diff --git a/src/control/indexer.rs b/src/control/indexer.rs index b06d620..d0e2110 100644 --- a/src/control/indexer.rs +++ b/src/control/indexer.rs @@ -115,7 +115,7 @@ impl Hydrant { let collection = CowStr::Borrowed("app.bsky.feed.post").into_static(); for i in 0..count { - let event_id = db.next_event_id.fetch_add(1, Ordering::SeqCst); + let event_id = db.stream.next_event_id.fetch_add(1, Ordering::SeqCst); let rkey = DbRkey::Str(smol_str::format_smolstr!("r{i}")); let evt = StoredEvent { live: false, @@ -128,7 +128,7 @@ impl Hydrant { }; let bytes = rmp_serde::to_vec(&evt).expect("msgpack serialization cannot fail"); total_bytes += bytes.len(); - batch.insert(&db.events, keys::event_key(event_id), bytes); + batch.insert(&db.stream.events, keys::event_key(event_id), bytes); } batch.commit().expect("failed to commit events batch"); diff --git a/src/control/repos/indexer.rs b/src/control/repos/indexer.rs index dc3efa4..23a7656 100644 --- a/src/control/repos/indexer.rs +++ b/src/control/repos/indexer.rs @@ -21,7 +21,7 @@ impl ReposControl { let state = self.0.clone(); self.0 .db - .pending + .indexer.pending .range((start_bound, std::ops::Bound::Unbounded)) .map(move |g| { let (id_raw, did_key) = g.into_inner().into_diagnostic()?; @@ -71,7 +71,7 @@ impl ReposControl { let state = self.0.clone(); self.0 .db - .resync + .indexer.resync .range((start_bound, std::ops::Bound::Unbounded)) .map(move |g| { let did_key = g.key().into_diagnostic()?; @@ -124,7 +124,7 @@ impl ReposControl { // skip if already in pending queue let is_pending = db - .pending + .indexer.pending .get(keys::pending_key(metadata.index_id)) .into_diagnostic()? .is_some(); @@ -132,10 +132,10 @@ impl ReposControl { metadata.tracked = true; // insert into pending with new index_id let old_pending = keys::pending_key(metadata.index_id); - batch.remove(&db.pending, old_pending); + batch.remove(&db.indexer.pending, old_pending); metadata.index_id = rand::Rng::next_u64(&mut rand::rng()); - batch.insert(&db.pending, keys::pending_key(metadata.index_id), &did_key); - batch.remove(&db.resync, &did_key); + batch.insert(&db.indexer.pending, keys::pending_key(metadata.index_id), &did_key); + batch.remove(&db.indexer.resync, &did_key); batch.insert( &db.repo_metadata, &metadata_key, @@ -234,7 +234,7 @@ impl ReposControl { &metadata_key, crate::db::ser_repo_meta(&metadata)?, ); - batch.insert(&db.pending, keys::pending_key(metadata.index_id), &did_key); + batch.insert(&db.indexer.pending, keys::pending_key(metadata.index_id), &did_key); count_deltas.add_repos(1); lifecycle_counts.transition(&mut batch, &did, GaugeState::Pending)?; queued.push(did); @@ -297,8 +297,8 @@ impl ReposControl { &metadata_key, crate::db::ser_repo_meta(&metadata)?, ); - batch.remove(&db.pending, keys::pending_key(metadata.index_id)); - batch.remove(&db.resync, &did_key); + batch.remove(&db.indexer.pending, keys::pending_key(metadata.index_id)); + batch.remove(&db.indexer.resync, &did_key); lifecycle_counts.transition(&mut batch, &did, GaugeState::Synced)?; untracked.push(did); } @@ -332,14 +332,14 @@ impl<'i> RepoHandle<'i> { tokio::task::spawn_blocking(move || { use miette::WrapErr; - let cid_bytes = state.db.records.get(db_key).into_diagnostic()?; + let cid_bytes = state.db.indexer.records.get(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) = state.db.blocks.get(block_key).into_diagnostic()? else { + let Some(block_bytes) = state.db.indexer.blocks.get(block_key).into_diagnostic()? else { miette::bail!("block {cid_bytes:?} not found, this is a bug!!"); }; @@ -398,7 +398,7 @@ impl<'i> RepoHandle<'i> { Box::new( state .db - .records + .indexer.records .range(prefix.as_slice()..end_key.as_slice()) .rev(), ) @@ -413,7 +413,7 @@ impl<'i> RepoHandle<'i> { prefix.clone() }; - Box::new(state.db.records.range(start_key.as_slice()..)) + Box::new(state.db.indexer.records.range(start_key.as_slice()..)) }; for item in iter { @@ -432,7 +432,7 @@ impl<'i> RepoHandle<'i> { // look up using col|cid key built from collection and binary cid bytes if let Ok(Some(block_bytes)) = state .db - .blocks + .indexer.blocks .get(keys::block_key(collection.as_str(), &cid_bytes)) { let value: Data = @@ -509,7 +509,7 @@ impl<'i> RepoHandle<'i> { let mut mst = mst; let prefix = keys::record_prefix_did(&did); - for guard in app_state.db.records.prefix(&prefix) { + for guard in app_state.db.indexer.records.prefix(&prefix) { let (key, cid_bytes) = guard.into_inner().into_diagnostic()?; let rest = &key[prefix.len()..]; @@ -534,7 +534,7 @@ impl<'i> RepoHandle<'i> { let block_key = keys::block_key(collection, cid_bytes.as_ref()); let block_bytes = app_state .db - .blocks + .indexer.blocks .get(&block_key) .into_diagnostic()? .ok_or_else(|| miette::miette!("block missing for record {mst_key}"))?; diff --git a/src/control/repos/mod.rs b/src/control/repos/mod.rs index 71c3760..9c68268 100644 --- a/src/control/repos/mod.rs +++ b/src/control/repos/mod.rs @@ -359,7 +359,7 @@ impl<'i> RepoHandle<'i> { let metadata = crate::db::deser_repo_meta(metadata_bytes.as_ref())?; Ok(app_state .db - .pending + .indexer.pending .get(crate::db::keys::pending_key(metadata.index_id)) .into_diagnostic()? .is_some()) diff --git a/src/control/stats.rs b/src/control/stats.rs index 4926873..2ee3b73 100644 --- a/src/control/stats.rs +++ b/src/control/stats.rs @@ -51,17 +51,17 @@ impl Hydrant { .collect(); #[cfg(feature = "indexer_stream")] - counts.insert("events", state.db.events.approximate_len() as u64); + counts.insert("events", state.db.stream.events.approximate_len() as u64); #[cfg(feature = "relay")] counts.insert( "relay_events", - state.db.relay_events.approximate_len() as u64, + state.db.relay.events.approximate_len() as u64, ); #[cfg(feature = "jetstream")] counts.insert( "jetstream_events", - state.db.jetstream_events.approximate_len() as u64, + state.db.jetstream.events.approximate_len() as u64, ); let sizes = tokio::task::spawn_blocking(move || { @@ -74,19 +74,19 @@ impl Hydrant { #[cfg(feature = "indexer")] { - s.insert("records", state.db.records.disk_space()); - s.insert("blocks", state.db.blocks.disk_space()); - s.insert("pending", state.db.pending.disk_space()); - s.insert("resync", state.db.resync.disk_space()); - s.insert("resync_buffer", state.db.resync_buffer.disk_space()); + s.insert("records", state.db.indexer.records.disk_space()); + s.insert("blocks", state.db.indexer.blocks.disk_space()); + s.insert("pending", state.db.indexer.pending.disk_space()); + s.insert("resync", state.db.indexer.resync.disk_space()); + s.insert("resync_buffer", state.db.indexer.resync_buffer.disk_space()); } #[cfg(feature = "indexer_stream")] - s.insert("events", state.db.events.disk_space()); + s.insert("events", state.db.stream.events.disk_space()); #[cfg(feature = "relay")] - s.insert("relay_events", state.db.relay_events.disk_space()); + s.insert("relay_events", state.db.relay.events.disk_space()); #[cfg(feature = "jetstream")] - s.insert("jetstream_events", state.db.jetstream_events.disk_space()); + s.insert("jetstream_events", state.db.jetstream.events.disk_space()); #[cfg(feature = "backlinks")] s.insert("backlinks", state.db.backlinks.disk_space()); diff --git a/src/control/stream/indexer.rs b/src/control/stream/indexer.rs index a024788..281c03d 100644 --- a/src/control/stream/indexer.rs +++ b/src/control/stream/indexer.rs @@ -23,14 +23,14 @@ pub(crate) fn event_stream_thread( opts: StreamOptions, ) { let db = &state.db; - let event_rx = db.event_tx.subscribe(); - let ks = db.events.clone(); + let event_rx = db.stream.event_tx.subscribe(); + let ks = db.stream.events.clone(); let current_id = match cursor { Some(c) => c.checked_sub(1), - None => db.next_event_id.load(Ordering::SeqCst).checked_sub(1), + None => db.stream.next_event_id.load(Ordering::SeqCst).checked_sub(1), }; let catch_up_target = cursor - .and_then(|_| db.next_event_id.load(Ordering::SeqCst).checked_sub(1)) + .and_then(|_| db.stream.next_event_id.load(Ordering::SeqCst).checked_sub(1)) .filter(|target| stream_seq_after(*target, current_id)); let replay_state = state.clone(); @@ -157,7 +157,7 @@ pub(crate) fn stored_to_event( } else { let block = state .db - .blocks + .indexer.blocks .get(keys::block_key(collection.as_str(), &cid.to_bytes())); match block { Ok(Some(bytes)) => { diff --git a/src/control/stream/jetstream.rs b/src/control/stream/jetstream.rs index 8e5f772..e30540a 100644 --- a/src/control/stream/jetstream.rs +++ b/src/control/stream/jetstream.rs @@ -29,7 +29,7 @@ pub(crate) fn jetstream_stream_thread( filter: JetstreamFilter, opts: StreamOptions, ) { - let mut event_rx = state.db.jetstream_tx.subscribe(); + let mut event_rx = state.db.jetstream.tx.subscribe(); let head = latest_jetstream_head(&state); let replay = cursor @@ -93,7 +93,7 @@ pub(crate) fn jetstream_stream_thread( } fn latest_jetstream_head(state: &AppState) -> Option<(u64, u64)> { - let guard = state.db.jetstream_events.iter().next_back()?; + let guard = state.db.jetstream.events.iter().next_back()?; let key = match guard.key() { Ok(key) => key, Err(e) => { @@ -211,7 +211,7 @@ fn read_jetstream_replay_chunk( let end_key = keys::jetstream_event_key(target_time_us, u64::MAX); let mut events = Vec::with_capacity(chunk_size); let mut exhausted = false; - let mut iter = state.db.jetstream_events.range(start_key..=&end_key[..]); + let mut iter = state.db.jetstream.events.range(start_key..=&end_key[..]); while events.len() < chunk_size { let Some(item) = iter.next() else { @@ -330,7 +330,7 @@ fn jetstream_event_to_bytes( match &event.event { #[cfg(feature = "indexer_stream")] StoredJetstreamEvent::Commit { event_id, live, .. } => { - let bytes = state.db.events.get(keys::event_key(*event_id)).ok()??; + let bytes = state.db.stream.events.get(keys::event_key(*event_id)).ok()??; let stored: StoredEvent = rmp_serde::from_slice(&bytes).ok()?; let evt = stored_to_event(state, *event_id, stored, None)?; let rec = evt.record?; @@ -363,7 +363,7 @@ fn jetstream_event_to_bytes( } => { let frame = state .db - .relay_events + .relay.events .get(keys::relay_event_key(*relay_seq)) .ok()??; let SubscribeReposMessage::Commit(commit) = decode_frame(frame.as_ref()).ok()? else { @@ -411,7 +411,7 @@ fn jetstream_event_to_bytes( StoredJetstreamEvent::RelayAccount { relay_seq, .. } => { let frame = state .db - .relay_events + .relay.events .get(keys::relay_event_key(*relay_seq)) .ok()??; let SubscribeReposMessage::Account(account) = decode_frame(frame.as_ref()).ok()? else { @@ -442,7 +442,7 @@ fn jetstream_event_to_bytes( StoredJetstreamEvent::RelayIdentity { relay_seq, .. } => { let frame = state .db - .relay_events + .relay.events .get(keys::relay_event_key(*relay_seq)) .ok()??; let SubscribeReposMessage::Identity(identity) = decode_frame(frame.as_ref()).ok()? diff --git a/src/control/stream/relay.rs b/src/control/stream/relay.rs index ffce9c2..639c040 100644 --- a/src/control/stream/relay.rs +++ b/src/control/stream/relay.rs @@ -16,14 +16,14 @@ pub(crate) fn relay_stream_thread( cursor: Option, opts: StreamOptions, ) { - let relay_rx = state.db.relay_broadcast_tx.subscribe(); - let ks = state.db.relay_events.clone(); + let relay_rx = state.db.relay.broadcast_tx.subscribe(); + let ks = state.db.relay.events.clone(); let current_seq = match cursor { Some(c) => Some(c.saturating_sub(1)), None => Some( state .db - .next_relay_seq + .relay.next_seq .load(Ordering::SeqCst) .saturating_sub(1), ), @@ -32,7 +32,7 @@ pub(crate) fn relay_stream_thread( .and_then(|_| { state .db - .next_relay_seq + .relay.next_seq .load(Ordering::SeqCst) .checked_sub(1) }) diff --git a/src/crawler/worker.rs b/src/crawler/worker.rs index 5cc6525..37d1dd8 100644 --- a/src/crawler/worker.rs +++ b/src/crawler/worker.rs @@ -167,7 +167,7 @@ impl CrawlerWorker { ); #[cfg(feature = "indexer")] batch.insert( - &app_state.db.pending, + &app_state.db.indexer.pending, keys::pending_key(metadata.index_id), &did_key, ); diff --git a/src/db/ephemeral.rs b/src/db/ephemeral.rs index e002922..c7d9951 100644 --- a/src/db/ephemeral.rs +++ b/src/db/ephemeral.rs @@ -41,26 +41,26 @@ pub fn relay_events_ttl_worker(state: Arc) { #[cfg(feature = "indexer_stream")] pub fn ephemeral_ttl_tick(db: &Db, ttl: &Duration) -> miette::Result<()> { - let current_seq = db.next_event_id.load(Ordering::SeqCst); + let current_seq = db.stream.next_event_id.load(Ordering::SeqCst); ttl_tick_inner( db, ttl, keys::EVENT_WATERMARK_PREFIX, keys::event_watermark_key, - &db.events, + &db.stream.events, current_seq, ) } #[cfg(feature = "relay")] pub fn relay_events_ttl_tick(db: &Db, ttl: &Duration) -> miette::Result<()> { - let current_seq = db.next_relay_seq.load(Ordering::SeqCst); + let current_seq = db.relay.next_seq.load(Ordering::SeqCst); ttl_tick_inner( db, ttl, keys::RELAY_EVENT_WATERMARK_PREFIX, keys::relay_event_watermark_key, - &db.relay_events, + &db.relay.events, current_seq, ) } @@ -71,19 +71,19 @@ pub fn jetstream_events_ttl_tick(db: &Db, ttl: &Duration) -> miette::Result<()> let cutoff_ts = now.saturating_sub(ttl.as_secs()); let cutoff_us = cutoff_ts.saturating_mul(1_000_000); - db.jetstream_events + db.jetstream.events .rotate_memtable_and_wait() .into_diagnostic() .wrap_err("failed to rotate memtable before Jetstream TTL range drop")?; - let before_space = db.jetstream_events.disk_space(); - let before_tables = db.jetstream_events.table_count(); - db.jetstream_events + let before_space = db.jetstream.events.disk_space(); + let before_tables = db.jetstream.events.table_count(); + db.jetstream.events .drop_range(..keys::jetstream_event_key(cutoff_us, 0)) .into_diagnostic() .wrap_err("failed Jetstream TTL range drop for old events")?; - let after_space = db.jetstream_events.disk_space(); - let after_tables = db.jetstream_events.table_count(); + let after_space = db.jetstream.events.disk_space(); + let after_tables = db.jetstream.events.table_count(); info!( cutoff_us, @@ -231,7 +231,7 @@ mod tests { } fn first_relay_seq(db: &crate::db::Db) -> miette::Result { - let Some(guard) = db.relay_events.iter().next() else { + let Some(guard) = db.relay.events.iter().next() else { miette::bail!("expected at least one relay event"); }; let key = guard.key().into_diagnostic()?; @@ -262,10 +262,10 @@ mod tests { ) -> miette::Result<()> { let mut batch = db.inner.batch(); for seq in start_seq..start_seq + count { - batch.insert(&db.relay_events, keys::relay_event_key(seq), payload); + batch.insert(&db.relay.events, keys::relay_event_key(seq), payload); } batch.commit().into_diagnostic()?; - db.next_relay_seq.store(start_seq + count, Ordering::SeqCst); + db.relay.next_seq.store(start_seq + count, Ordering::SeqCst); Ok(()) } @@ -276,7 +276,7 @@ mod tests { payload: &[u8], ) -> miette::Result<()> { insert_relay_events(db, start_seq, count, payload)?; - db.relay_events + db.relay.events .rotate_memtable_and_wait() .into_diagnostic()?; Ok(()) @@ -288,7 +288,7 @@ mod tests { payload: &[u8], ) -> miette::Result<()> { seed_relay_event_table(db, 0, count, payload)?; - db.relay_events.major_compact().into_diagnostic()?; + db.relay.events.major_compact().into_diagnostic()?; Ok(()) } @@ -303,14 +303,14 @@ mod tests { } fn compact_relay_events_once(db: &crate::db::Db) -> miette::Result<()> { - db.relay_events + db.relay.events .compact(Arc::new(fjall::compaction::Leveled::default())) .into_diagnostic() } fn wait_for_relay_table_count(db: &crate::db::Db, min_tables: usize) -> miette::Result<()> { for _ in 0..200 { - if db.relay_events.table_count() >= min_tables { + if db.relay.events.table_count() >= min_tables { return Ok(()); } std::thread::sleep(Duration::from_millis(10)); @@ -318,7 +318,7 @@ mod tests { miette::bail!( "timed out waiting for at least {min_tables} relay event tables, saw {}", - db.relay_events.table_count() + db.relay.events.table_count() ); } @@ -328,10 +328,10 @@ mod tests { ) -> miette::Result<()> { let mut batch = db.inner.batch(); for seq in 0..cutoff_seq { - batch.remove(&db.relay_events, keys::relay_event_key(seq)); + batch.remove(&db.relay.events, keys::relay_event_key(seq)); } batch.commit().into_diagnostic()?; - db.relay_events + db.relay.events .rotate_memtable_and_wait() .into_diagnostic()?; Ok(()) @@ -353,17 +353,17 @@ mod tests { start_seq += events_per_table; } - let before_prune = db.relay_events.disk_space(); - let before_tables = db.relay_events.table_count(); + let before_prune = db.relay.events.disk_space(); + let before_tables = db.relay.events.table_count(); seed_past_relay_watermark(&db, old_count)?; relay_events_ttl_tick(&db, &Duration::from_secs(60 * 60))?; - let after_prune = db.relay_events.disk_space(); - let after_tables = db.relay_events.table_count(); + let after_prune = db.relay.events.disk_space(); + let after_tables = db.relay.events.table_count(); assert_eq!( retained_count as usize, - db.relay_events.iter().count(), + db.relay.events.iter().count(), "TTL should keep the tables at or after the cutoff sequence" ); assert!( @@ -394,12 +394,12 @@ mod tests { start_seq += events_per_table; } - let before_tables = db.relay_events.table_count(); + let before_tables = db.relay.events.table_count(); seed_past_relay_watermark(&db, cutoff_seq)?; relay_events_ttl_tick(&db, &Duration::from_secs(60 * 60))?; - let after_tables = db.relay_events.table_count(); + let after_tables = db.relay.events.table_count(); assert!( after_tables <= before_tables.saturating_sub(fully_expired_tables as usize), "drop_range should drop tables fully below the cutoff; before={before_tables}, after={after_tables}" @@ -411,7 +411,7 @@ mod tests { ); assert_eq!( ((table_count - fully_expired_tables) * events_per_table) as usize, - db.relay_events.iter().count(), + db.relay.events.iter().count(), "only the boundary table and newer tables should remain" ); @@ -433,8 +433,8 @@ mod tests { } let cutoff_seq = events_per_second * ttl_seconds; - let before_space = db.relay_events.disk_space(); - let before_tables = db.relay_events.table_count(); + let before_space = db.relay.events.disk_space(); + let before_tables = db.relay.events.table_count(); assert!( before_tables >= total_seconds as usize, "sustained writes should have produced multiple SSTs; tables={before_tables}" @@ -443,10 +443,10 @@ mod tests { seed_past_relay_watermark(&db, cutoff_seq)?; relay_events_ttl_tick(&db, &Duration::from_secs(60 * 60))?; - let after_space = db.relay_events.disk_space(); - let after_tables = db.relay_events.table_count(); + let after_space = db.relay.events.disk_space(); + let after_tables = db.relay.events.table_count(); let first_seq = first_relay_seq(&db)?; - let remaining_events = db.relay_events.iter().count() as u64; + let remaining_events = db.relay.events.iter().count() as u64; assert!( after_tables < before_tables, @@ -480,18 +480,18 @@ mod tests { let payload = vec![0x42; 8 * 1024]; seed_compacted_relay_events(&db, old_count + retained_count, &payload)?; - let before_prune = db.relay_events.disk_space(); + let before_prune = db.relay.events.disk_space(); prune_prefix_with_per_key_tombstones(&db, old_count)?; - let after_delete = db.relay_events.disk_space(); + let after_delete = db.relay.events.disk_space(); for _ in 0..16 { - db.relay_events + db.relay.events .compact(Arc::new(fjall::compaction::Leveled::default())) .into_diagnostic()?; } - let after_leveled = db.relay_events.disk_space(); - assert_eq!(retained_count as usize, db.relay_events.iter().count()); + let after_leveled = db.relay.events.disk_space(); + assert_eq!(retained_count as usize, db.relay.events.iter().count()); assert!( after_leveled >= before_prune * 9 / 10, "leveled compaction should not be expected to reclaim prefix tombstones quickly; before={before_prune}, after_delete={after_delete}, after_leveled={after_leveled}" diff --git a/src/db/keyspaces.rs b/src/db/keyspaces.rs new file mode 100644 index 0000000..3f2c7db --- /dev/null +++ b/src/db/keyspaces.rs @@ -0,0 +1,384 @@ +//! per-mode keyspace groups. this is the single place where mode-specific +//! database state is declared: each feature contributes one group struct, +//! opened and initialized here, and `Db` embeds each group behind one cfg. + +use fjall::{CompressionType, Database, Keyspace, KeyspaceCreateOptions}; +use miette::{IntoDiagnostic, Result}; +use std::sync::Arc; + +use crate::config::Config; + +#[cfg(any(feature = "indexer_stream", feature = "jetstream", feature = "relay"))] +use miette::Context; +#[cfg(any(feature = "indexer_stream", feature = "jetstream", feature = "relay"))] +use std::sync::atomic::AtomicU64; + +#[cfg_attr( + not(any( + feature = "indexer", + feature = "indexer_stream", + feature = "jetstream", + feature = "relay" + )), + allow(dead_code) +)] +const fn kb(v: u32) -> u32 { + v * 1024 +} +pub(super) const fn mb(v: u64) -> u64 { + v * 1024 * 1024 +} + +/// everything a mode group needs to open its keyspaces. +#[cfg_attr( + not(any( + feature = "indexer", + feature = "indexer_stream", + feature = "jetstream", + feature = "relay" + )), + allow(dead_code) +)] +pub(super) struct OpenCx<'a> { + pub(super) db: &'a Arc, + pub(super) cfg: &'a Config, + pub(super) compression: &'a dyn Fn(&str, i32) -> CompressionType, +} + +impl OpenCx<'_> { + pub(super) fn open_ks(&self, name: &str, opts: KeyspaceCreateOptions) -> Result { + self.db.keyspace(name, move || opts).into_diagnostic() + } +} + +#[cfg(feature = "indexer")] +pub struct IndexerDb { + /// maps `{DID}|{COL}|{RKey}` -> record CID + pub records: Keyspace, + /// content-addressable storage of raw DAG-CBOR blocks + pub blocks: Keyspace, + /// backfill queue of `{ID}` -> empty + pub pending: Keyspace, + /// per-repo resync/retry state + pub resync: Keyspace, + /// live events buffered during backfill + pub resync_buffer: Keyspace, + /// serializes lifecycle count rebuilds + pub(crate) lifecycle_count_lock: Arc>, +} + +#[cfg(feature = "indexer")] +impl IndexerDb { + pub(super) fn open(cx: &OpenCx) -> Result { + use fjall::config::{BlockSizePolicy, CompressionPolicy, RestartIntervalPolicy}; + + let opts = KeyspaceCreateOptions::default; + let pending = cx.open_ks( + "pending", + opts() + // iterated over as a queue, no point reads are used so bloom filters are disabled + .expect_point_read_hits(true) + .max_memtable_size(mb(8)) + // its just index of id (int) -> did, and dids arent compressable (especially with the ids being random) + .data_block_size_policy(BlockSizePolicy::all(kb(8))) + // and we'll transition from pending to synced anyway, no point trying to compress + .data_block_compression_policy(CompressionPolicy::disabled()) + // ids are sequential and share prefix so we can use large interval to save space + .data_block_restart_interval_policy(RestartIntervalPolicy::all(64)), + )?; + let resync = cx.open_ks( + "resync", + opts() + // we only point read in backfill when we check for existing resync state + // ...and also in repos api. so we can disable bloom filters + .expect_point_read_hits(true) + .max_memtable_size(mb(8)) + // did -> error state, so its gonna be basically random, cant compress well + .data_block_size_policy(BlockSizePolicy::all(kb(4))) + // and we arent going to have many of these anyway, no point trying + .data_block_compression_policy(CompressionPolicy::disabled()) + .data_block_restart_interval_policy(RestartIntervalPolicy::all(4)), + )?; + // this is used in non-ephemeral mode + let blocks = cx.open_ks( + "blocks", + opts() + // point reads are used a lot by stream, we know the blocks exist though + .expect_point_read_hits(true) + .max_memtable_size(mb(cx.cfg.db_blocks_memtable_size_mb)) + // 16 - 128 kb, as the newer blocks will be in the first level (or memtable) + // and any consumers will probably be streaming the newer events... + // and blocks are pretty big-ish like around 5kb... + // replaying will hit later levels so it will be slower but thats honestly + // an acceptable tradeoff to save more space... + // todo: we can probably decrease these when we get zstd dict compression? + .data_block_size_policy(BlockSizePolicy::new([kb(16), kb(64), kb(128)])) + // lets not compress first level so the reads for new blocks are faster + // since we will be streaming them to consumers + .data_block_compression_policy(CompressionPolicy::new([ + CompressionType::None, + (cx.compression)("blocks", 3), + (cx.compression)("blocks", 3), + (cx.compression)("blocks", 5), + ])) + .data_block_restart_interval_policy(RestartIntervalPolicy::new([8, 16, 32])), + )?; + let records = cx.open_ks( + "records", + opts() + // point reads might miss when using getRecord + // but we assume thats not going to happen often... + // since this keyspace is big, turning off bloom filters will help a lot with memory/disk space, + // but leaves point reads vulnerable to disk I/O misses under public query load + .expect_point_read_hits(!cx.cfg.db_records_bloom_filters) + .max_memtable_size(mb(cx.cfg.db_records_memtable_size_mb)) + // its just did|col|rkey -> cid, very small (84 bytes for bsky post) + .data_block_size_policy(BlockSizePolicy::new([kb(8), kb(16)])) + // cids arent compressable, most rkeys are TIDs so they will get compressed + // by prefix truncation anyway + .data_block_compression_policy(CompressionPolicy::disabled()) + .data_block_restart_interval_policy(RestartIntervalPolicy::new([16, 32])), + )?; + let resync_buffer = cx.open_ks( + "resync_buffer", + opts() + // iterated during backfill, no point reads + .expect_point_read_hits(true) + .max_memtable_size(mb(16)) + .data_block_size_policy(BlockSizePolicy::all(kb(32))) + // dont have to compress here since resync buffer will be emptied at some point anyway + .data_block_compression_policy(CompressionPolicy::disabled()) + .data_block_restart_interval_policy(RestartIntervalPolicy::all(16)), + )?; + + Ok(Self { + records, + blocks, + pending, + resync, + resync_buffer, + lifecycle_count_lock: Arc::new(std::sync::Mutex::new(())), + }) + } + + pub(super) fn keyspaces(&self) -> impl Iterator { + [ + self.records.clone(), + self.blocks.clone(), + self.pending.clone(), + self.resync.clone(), + self.resync_buffer.clone(), + ] + .into_iter() + } +} + +#[cfg(feature = "indexer_stream")] +pub struct StreamDb { + /// maps `{ID}` (u64 BE) -> `StoredEvent`, the source for the json stream api + pub events: Keyspace, + pub(crate) event_tx: tokio::sync::broadcast::Sender, + pub next_event_id: Arc, +} + +#[cfg(feature = "indexer_stream")] +impl StreamDb { + pub(super) fn open(cx: &OpenCx) -> Result { + use fjall::config::{BlockSizePolicy, CompressionPolicy, RestartIntervalPolicy}; + + let events = cx.open_ks( + "events", + KeyspaceCreateOptions::default() + // only iterators are used here, no point reads + .expect_point_read_hits(true) + .max_memtable_size(mb(cx.cfg.db_events_memtable_size_mb)) + // the compression here wont be quite as good since events are quite random + // eg. by many different repos and different records etc. + // since its sequential we should still go with bigger block size though + // backfills will be sequential though... + .data_block_size_policy( + cx.cfg + .ephemeral + .then(|| BlockSizePolicy::new([kb(64), kb(128), kb(256)])) + .unwrap_or_else(|| BlockSizePolicy::new([kb(16), kb(64)])), + ) + // we are streaming the new events to consumers so we dont want to compress them + .data_block_compression_policy( + cx.cfg + .ephemeral + .then(|| { + CompressionPolicy::new([ + CompressionType::None, + (cx.compression)("events", 3), + ]) + }) + .unwrap_or_else(|| { + CompressionPolicy::new([ + CompressionType::None, + (cx.compression)("events", 3), + (cx.compression)("events", 3), + (cx.compression)("events", 5), + ]) + }), + ) + // ids are int, we can prefix truncate a lot + .data_block_restart_interval_policy(RestartIntervalPolicy::new([64, 128])), + )?; + + let (event_tx, _) = tokio::sync::broadcast::channel(512); + + Ok(Self { + events, + event_tx, + next_event_id: Arc::new(AtomicU64::new(0)), + }) + } + + /// resume event ids after the last stored event. + pub(super) fn init(&self) -> Result<()> { + let mut last_id = 0; + if let Some(guard) = self.events.iter().next_back() { + let k = guard.key().into_diagnostic()?; + last_id = u64::from_be_bytes( + k.as_ref() + .try_into() + .into_diagnostic() + .wrap_err("expected to be id (8 bytes)")?, + ); + } + // relaxed is fine since we are just initializing the db + self.next_event_id + .store(last_id + 1, std::sync::atomic::Ordering::Relaxed); + Ok(()) + } + + pub(super) fn keyspaces(&self) -> impl Iterator { + [self.events.clone()].into_iter() + } +} + +#[cfg(feature = "jetstream")] +pub(crate) struct JetstreamDb { + /// maps `{time_us}|{ID}` (16 bytes) -> jetstream event data + pub(crate) events: Keyspace, + pub(crate) tx: tokio::sync::broadcast::Sender, + pub(crate) next_id: Arc, + pub(crate) last_time_us: Arc, + /// serializes jetstream staging with batch commit in relay mode + #[cfg(feature = "relay")] + pub(crate) lock: Arc>, +} + +#[cfg(feature = "jetstream")] +impl JetstreamDb { + pub(super) fn open(cx: &OpenCx) -> Result { + use fjall::config::{BlockSizePolicy, CompressionPolicy, RestartIntervalPolicy}; + + let events = cx.open_ks( + "jetstream_events", + KeyspaceCreateOptions::default() + // time-ordered append-only stream metadata, only iterated for replay. + .expect_point_read_hits(true) + .max_memtable_size(mb(cx.cfg.db_events_memtable_size_mb)) + .data_block_size_policy(BlockSizePolicy::new([kb(16), kb(64), kb(128)])) + .data_block_compression_policy(CompressionPolicy::new([ + CompressionType::None, + (cx.compression)("jetstream_events", 3), + (cx.compression)("jetstream_events", 5), + ])) + .data_block_restart_interval_policy(RestartIntervalPolicy::new([64, 128])), + )?; + + let (tx, _) = tokio::sync::broadcast::channel(512); + + Ok(Self { + events, + tx, + next_id: Arc::new(AtomicU64::new(0)), + last_time_us: Arc::new(std::sync::atomic::AtomicI64::new(0)), + #[cfg(feature = "relay")] + lock: Arc::new(parking_lot::Mutex::new(())), + }) + } + + /// resume jetstream ids and time watermark after the last stored event. + pub(super) fn init(&self) -> Result<()> { + let mut last_id = 0; + let mut last_time_us = 0; + if let Some(guard) = self.events.iter().next_back() { + let k = guard.key().into_diagnostic()?; + let (time_us, id) = super::keys::parse_jetstream_event_key(&k)?; + last_id = id; + last_time_us = time_us; + } + self.next_id + .store(last_id + 1, std::sync::atomic::Ordering::Relaxed); + self.last_time_us + .store(last_time_us as i64, std::sync::atomic::Ordering::Relaxed); + Ok(()) + } + + pub(super) fn keyspaces(&self) -> impl Iterator { + [self.events.clone()].into_iter() + } +} + +#[cfg(feature = "relay")] +pub(crate) struct RelayDb { + /// maps `{SEQ}` (u64 BE) -> re-encoded relay frame + pub(crate) events: Keyspace, + pub(crate) next_seq: Arc, + pub(crate) broadcast_tx: tokio::sync::broadcast::Sender, +} + +#[cfg(feature = "relay")] +impl RelayDb { + pub(super) fn open(cx: &OpenCx) -> Result { + use fjall::config::{BlockSizePolicy, CompressionPolicy, RestartIntervalPolicy}; + + let events = cx.open_ks( + "relay_events", + KeyspaceCreateOptions::default() + // only iterated for cursor replay + .expect_point_read_hits(true) + .max_memtable_size(mb(cx.cfg.db_events_memtable_size_mb)) + .data_block_size_policy(BlockSizePolicy::new([kb(64), kb(128), kb(256)])) + .data_block_compression_policy(CompressionPolicy::new([ + CompressionType::None, + (cx.compression)("events", 3), + (cx.compression)("events", 3), + (cx.compression)("events", 5), + ])) + .data_block_restart_interval_policy(RestartIntervalPolicy::new([64, 128])), + )?; + + let (broadcast_tx, _) = tokio::sync::broadcast::channel(512); + + Ok(Self { + events, + next_seq: Arc::new(AtomicU64::new(0)), + broadcast_tx, + }) + } + + /// resume relay sequence numbers after the last stored frame. + pub(super) fn init(&self) -> Result<()> { + let mut last_relay_seq = 0u64; + if let Some(guard) = self.events.iter().next_back() { + let k = guard.key().into_diagnostic()?; + last_relay_seq = u64::from_be_bytes( + k.as_ref() + .try_into() + .into_diagnostic() + .wrap_err("relay_events: invalid key length")?, + ); + } + self.next_seq + .store(last_relay_seq + 1, std::sync::atomic::Ordering::Relaxed); + Ok(()) + } + + pub(super) fn keyspaces(&self) -> impl Iterator { + [self.events.clone()].into_iter() + } +} diff --git a/src/db/lifecycle_counts.rs b/src/db/lifecycle_counts.rs index 9251d23..3fc8ff3 100644 --- a/src/db/lifecycle_counts.rs +++ b/src/db/lifecycle_counts.rs @@ -27,7 +27,7 @@ impl Db { LifecycleCountBatch { db: self, _lock: self - .lifecycle_count_lock + .indexer.lifecycle_count_lock .lock() .expect("lifecycle count lock poisoned"), deltas: CountDeltas::default(), @@ -112,7 +112,7 @@ impl<'a> LifecycleCountBatch<'a> { if let Some(index_id) = index_id { if self .db - .pending + .indexer.pending .get(keys::pending_key(index_id)) .into_diagnostic()? .is_some() @@ -123,7 +123,7 @@ impl<'a> LifecycleCountBatch<'a> { let gauge = if is_pending { GaugeState::Pending - } else if let Some(resync_bytes) = self.db.resync.get(did_key).into_diagnostic()? { + } else if let Some(resync_bytes) = self.db.indexer.resync.get(did_key).into_diagnostic()? { gauge_from_resync(&resync_bytes) } else { GaugeState::Synced @@ -141,7 +141,7 @@ impl<'a> LifecycleCountBatch<'a> { }; Ok(current.as_slice() == pending_key - && self.db.pending.get(current).into_diagnostic()?.is_some()) + && self.db.indexer.pending.get(current).into_diagnostic()?.is_some()) } fn stage_membership( @@ -232,7 +232,7 @@ mod tests { let mut batch = db.inner.batch(); batch.insert(&db.repos, &did_key, ser_repo_state(&state)?); batch.insert(&db.repo_metadata, metadata_key, ser_repo_meta(&metadata)?); - batch.insert(&db.pending, pending_key, did_key); + batch.insert(&db.indexer.pending, pending_key, did_key); let mut lifecycle = db.lifecycle_counts(); lifecycle.transition(&mut batch, did, GaugeState::Pending)?; let lifecycle_reservation = lifecycle.stage(&mut batch); @@ -251,7 +251,7 @@ mod tests { GaugeState::Synced, )?; if applied { - batch.remove(&db.pending, pending_key); + batch.remove(&db.indexer.pending, pending_key); } let lifecycle_reservation = lifecycle.stage(&mut batch); batch.commit().into_diagnostic()?; @@ -289,8 +289,8 @@ mod tests { let mut metadata = RepoMetadata::backfilling(2); metadata.tracked = true; let mut batch = db.inner.batch(); - batch.remove(&db.pending, stale_pending); - batch.insert(&db.pending, current_pending, &did_key); + batch.remove(&db.indexer.pending, stale_pending); + batch.insert(&db.indexer.pending, current_pending, &did_key); batch.insert(&db.repo_metadata, metadata_key, ser_repo_meta(&metadata)?); let mut lifecycle = db.lifecycle_counts(); lifecycle.transition(&mut batch, &did, GaugeState::Pending)?; @@ -301,7 +301,7 @@ mod tests { assert_eq!(db.get_count_sync("pending"), 1); assert!(!complete_pending(&db, &did, stale_pending)?); assert_eq!(db.get_count_sync("pending"), 1); - assert!(db.pending.get(current_pending).into_diagnostic()?.is_some()); + assert!(db.indexer.pending.get(current_pending).into_diagnostic()?.is_some()); Ok(()) } @@ -326,7 +326,7 @@ mod tests { GaugeState::Synced, )?; assert!(applied); - batch.remove(&db.pending, pending_key); + batch.remove(&db.indexer.pending, pending_key); let lifecycle_reservation = lifecycle.stage(&mut batch); batch.commit().into_diagnostic()?; db.persist()?; @@ -337,7 +337,7 @@ mod tests { let db = Db::open(&cfg)?; assert_eq!(db.get_count_sync("pending"), 0); assert!( - db.pending + db.indexer.pending .get(keys::pending_key(1)) .into_diagnostic()? .is_none() diff --git a/src/db/migration/v8.rs b/src/db/migration/v8.rs index 0d2e537..018765d 100644 --- a/src/db/migration/v8.rs +++ b/src/db/migration/v8.rs @@ -80,7 +80,7 @@ fn primary_lifecycle_gauge(db: &Db, repo_key: &[u8]) -> Result { let metadata = deser_repo_meta(metadata_bytes.as_ref()) .wrap_err("invalid repo metadata during lifecycle count rebuild")?; if db - .pending + .indexer.pending .get(keys::pending_key(metadata.index_id)) .into_diagnostic()? .is_some() @@ -89,7 +89,7 @@ fn primary_lifecycle_gauge(db: &Db, repo_key: &[u8]) -> Result { } } - db.resync + db.indexer.resync .get(repo_key) .into_diagnostic()? .map(|bytes| crate::db::lifecycle_counts::gauge_from_resync(bytes.as_ref())) @@ -187,12 +187,12 @@ mod tests { insert_repo(&db, &mut batch, &synced_did, 40)?; batch.insert( - &db.pending, + &db.indexer.pending, keys::pending_key(10), keys::repo_key(&pending_did), ); batch.insert( - &db.resync, + &db.indexer.resync, keys::repo_key(&ratelimited_did), rmp_serde::to_vec(&ResyncState::Error { kind: ResyncErrorKind::Ratelimited, @@ -202,7 +202,7 @@ mod tests { .into_diagnostic()?, ); batch.insert( - &db.resync, + &db.indexer.resync, keys::repo_key(&gone_did), rmp_serde::to_vec(&ResyncState::Gone { status: RepoStatus::Deactivated, diff --git a/src/db/mod.rs b/src/db/mod.rs index 2ee4ddb..76d64b5 100644 --- a/src/db/mod.rs +++ b/src/db/mod.rs @@ -1,20 +1,12 @@ use crate::types::{RepoMetadata, RepoState}; -#[cfg(feature = "indexer_stream")] -use crate::types::BroadcastEvent; -#[cfg(feature = "jetstream")] -use crate::types::JetstreamBroadcast; -#[cfg(feature = "relay")] -use crate::types::RelayBroadcast; - use fjall::{Database, Keyspace, PersistMode, Slice}; use miette::{Context, IntoDiagnostic, Result}; use scc::HashMap; use smol_str::SmolStr; use std::collections::BTreeSet; -#[cfg(feature = "jetstream")] -use std::sync::atomic::AtomicI64; + use std::sync::atomic::AtomicU64; use std::sync::{Arc, Mutex}; use url::Url; @@ -33,16 +25,25 @@ pub mod migration; pub mod pds_meta; pub mod types; +pub mod keyspaces; mod open; mod train; +#[cfg(feature = "indexer")] +pub use keyspaces::IndexerDb; +#[cfg(feature = "jetstream")] +pub(crate) use keyspaces::JetstreamDb; +#[cfg(feature = "relay")] +pub(crate) use keyspaces::RelayDb; +#[cfg(feature = "indexer_stream")] +pub use keyspaces::StreamDb; + #[cfg(feature = "indexer")] pub(crate) use lifecycle_counts::LifecycleCountBatch; use tracing::error; -#[cfg(any(feature = "indexer_stream", feature = "relay"))] -use tokio::sync::broadcast; + pub struct Db { pub inner: Arc, @@ -54,46 +55,20 @@ pub struct Db { pub filter: Keyspace, pub crawler: Keyspace, #[cfg(feature = "indexer")] - pub records: Keyspace, - #[cfg(feature = "indexer")] - pub blocks: Keyspace, - #[cfg(feature = "indexer")] - pub pending: Keyspace, - #[cfg(feature = "indexer")] - pub resync: Keyspace, - #[cfg(feature = "indexer")] - pub resync_buffer: Keyspace, + pub indexer: IndexerDb, #[cfg(feature = "indexer_stream")] - pub events: Keyspace, + pub stream: StreamDb, #[cfg(feature = "jetstream")] - pub(crate) jetstream_events: Keyspace, + pub(crate) jetstream: JetstreamDb, + #[cfg(feature = "relay")] + pub(crate) relay: RelayDb, #[cfg(feature = "backlinks")] pub backlinks: Keyspace, - #[cfg(feature = "indexer_stream")] - pub(crate) event_tx: broadcast::Sender, - #[cfg(feature = "indexer_stream")] - pub next_event_id: Arc, - #[cfg(feature = "jetstream")] - pub(crate) jetstream_tx: broadcast::Sender, - #[cfg(feature = "jetstream")] - pub(crate) next_jetstream_id: Arc, - #[cfg(feature = "jetstream")] - pub(crate) last_jetstream_time_us: Arc, - #[cfg(all(feature = "jetstream", feature = "relay"))] - pub(crate) jetstream_lock: Arc>, - #[cfg(feature = "relay")] - pub(crate) relay_events: Keyspace, - #[cfg(feature = "relay")] - pub(crate) next_relay_seq: Arc, - #[cfg(feature = "relay")] - pub(crate) relay_broadcast_tx: broadcast::Sender, pub counts_map: HashMap, next_count_delta_id: Arc, count_delta_checkpoint_watermark: Arc, count_delta_gc_watermark: Arc, count_delta_in_flight: Arc>>, - #[cfg(feature = "indexer")] - lifecycle_count_lock: Arc>, pub(crate) compaction_running: Arc, } @@ -130,7 +105,16 @@ impl Db { .into_diagnostic()? }; - #[cfg_attr(not(any(feature = "indexer", feature = "indexer_stream", feature = "jetstream", feature = "relay", feature = "backlinks")), allow(unused_mut))] + #[cfg_attr( + not(any( + feature = "indexer", + feature = "indexer_stream", + feature = "jetstream", + feature = "relay", + feature = "backlinks" + )), + allow(unused_mut) + )] let mut tasks = vec![ compact(self.repos.clone()), compact(self.cursors.clone()), @@ -141,22 +125,13 @@ impl Db { ]; #[cfg(feature = "indexer")] - { - tasks.push(compact(self.records.clone())); - tasks.push(compact(self.blocks.clone())); - tasks.push(compact(self.pending.clone())); - tasks.push(compact(self.resync.clone())); - tasks.push(compact(self.resync_buffer.clone())); - } + tasks.extend(self.indexer.keyspaces().map(&compact)); #[cfg(feature = "indexer_stream")] - tasks.push(compact(self.events.clone())); - + tasks.extend(self.stream.keyspaces().map(&compact)); #[cfg(feature = "jetstream")] - tasks.push(compact(self.jetstream_events.clone())); - + tasks.extend(self.jetstream.keyspaces().map(&compact)); #[cfg(feature = "relay")] - tasks.push(compact(self.relay_events.clone())); - + tasks.extend(self.relay.keyspaces().map(&compact)); #[cfg(feature = "backlinks")] tasks.push(compact(self.backlinks.clone())); diff --git a/src/db/open.rs b/src/db/open.rs index de850ab..9ba8852 100644 --- a/src/db/open.rs +++ b/src/db/open.rs @@ -7,26 +7,20 @@ use miette::{Context, IntoDiagnostic, Result}; use scc::HashMap; use smol_str::SmolStr; use std::collections::BTreeSet; -#[cfg(feature = "jetstream")] -use std::sync::atomic::AtomicI64; use std::sync::atomic::{AtomicU64, Ordering}; use std::sync::{Arc, Mutex}; -#[cfg(any(feature = "indexer_stream", feature = "relay", feature = "jetstream"))] -use tokio::sync::broadcast; use crate::config::{Compression, Config}; use super::compaction::CountsGcFilterFactory; use super::counts::{load_count_delta_watermark, read_u64_counter, replay_count_deltas}; use super::keys; +use super::keyspaces::{OpenCx, mb}; use super::{Db, migration}; const fn kb(v: u32) -> u32 { v * 1024 } -const fn mb(v: u64) -> u64 { - v * 1024 * 1024 -} fn default_opts() -> KeyspaceCreateOptions { KeyspaceCreateOptions::default() } @@ -61,9 +55,6 @@ impl Db { let db = Arc::new(db); let opts = default_opts; - let open_ks = |name: &str, opts: KeyspaceCreateOptions| { - db.keyspace(name, move || opts).into_diagnostic() - }; let load_dict = |name: &str| -> Option> { let path = cfg.database_path.join(format!("dict_{name}.bin")); @@ -98,6 +89,12 @@ impl Db { .unwrap_or_else(|| CompressionType::Zstd { level }), Compression::None => CompressionType::None, }; + let cx = OpenCx { + db: &db, + cfg, + compression: &get_compression, + }; + let open_ks = |name: &str, opts: KeyspaceCreateOptions| cx.open_ks(name, opts); let repos = open_ks( "repos", @@ -136,75 +133,9 @@ impl Db { .data_block_restart_interval_policy(RestartIntervalPolicy::new([2, 4])), )?; #[cfg(feature = "indexer")] - let pending = open_ks( - "pending", - opts() - // iterated over as a queue, no point reads are used so bloom filters are disabled - .expect_point_read_hits(true) - .max_memtable_size(mb(8)) - // its just index of id (int) -> did, and dids arent compressable (especially with the ids being random) - .data_block_size_policy(BlockSizePolicy::all(kb(8))) - // and we'll transition from pending to synced anyway, no point trying to compress - .data_block_compression_policy(CompressionPolicy::disabled()) - // ids are sequential and share prefix so we can use large interval to save space - .data_block_restart_interval_policy(RestartIntervalPolicy::all(64)), - )?; - #[cfg(feature = "indexer")] - let resync = open_ks( - "resync", - opts() - // we only point read in backfill when we check for existing resync state - // ...and also in repos api. so we can disable bloom filters - .expect_point_read_hits(true) - .max_memtable_size(mb(8)) - // did -> error state, so its gonna be basically random, cant compress well - .data_block_size_policy(BlockSizePolicy::all(kb(4))) - // and we arent going to have many of these anyway, no point trying - .data_block_compression_policy(CompressionPolicy::disabled()) - .data_block_restart_interval_policy(RestartIntervalPolicy::all(4)), - )?; - // this is used in non-ephemeral mode - #[cfg(feature = "indexer")] - let blocks = open_ks( - "blocks", - opts() - // point reads are used a lot by stream, we know the blocks exist though - .expect_point_read_hits(true) - .max_memtable_size(mb(cfg.db_blocks_memtable_size_mb)) - // 16 - 128 kb, as the newer blocks will be in the first level (or memtable) - // and any consumers will probably be streaming the newer events... - // and blocks are pretty big-ish like around 5kb... - // replaying will hit later levels so it will be slower but thats honestly - // an acceptable tradeoff to save more space... - // todo: we can probably decrease these when we get zstd dict compression? - .data_block_size_policy(BlockSizePolicy::new([kb(16), kb(64), kb(128)])) - // lets not compress first level so the reads for new blocks are faster - // since we will be streaming them to consumers - .data_block_compression_policy(CompressionPolicy::new([ - CompressionType::None, - get_compression("blocks", 3), - get_compression("blocks", 3), - get_compression("blocks", 5), - ])) - .data_block_restart_interval_policy(RestartIntervalPolicy::new([8, 16, 32])), - )?; - #[cfg(feature = "indexer")] - let records = open_ks( - "records", - opts() - // point reads might miss when using getRecord - // but we assume thats not going to happen often... - // since this keyspace is big, turning off bloom filters will help a lot with memory/disk space, - // but leaves point reads vulnerable to disk I/O misses under public query load - .expect_point_read_hits(!cfg.db_records_bloom_filters) - .max_memtable_size(mb(cfg.db_records_memtable_size_mb)) - // its just did|col|rkey -> cid, very small (84 bytes for bsky post) - .data_block_size_policy(BlockSizePolicy::new([kb(8), kb(16)])) - // cids arent compressable, most rkeys are TIDs so they will get compressed - // by prefix truncation anyway - .data_block_compression_policy(CompressionPolicy::disabled()) - .data_block_restart_interval_policy(RestartIntervalPolicy::new([16, 32])), - )?; + let indexer = super::keyspaces::IndexerDb::open(&cx)?; + #[cfg(feature = "indexer_stream")] + let stream = super::keyspaces::StreamDb::open(&cx)?; let cursors = open_ks( "cursors", opts() @@ -216,71 +147,9 @@ impl Db { .data_block_compression_policy(CompressionPolicy::disabled()) .data_block_restart_interval_policy(RestartIntervalPolicy::all(1)), )?; - #[cfg(feature = "indexer")] - let resync_buffer = open_ks( - "resync_buffer", - opts() - // iterated during backfill, no point reads - .expect_point_read_hits(true) - .max_memtable_size(mb(16)) - .data_block_size_policy(BlockSizePolicy::all(kb(32))) - // dont have to compress here since resync buffer will be emptied at some point anyway - .data_block_compression_policy(CompressionPolicy::disabled()) - .data_block_restart_interval_policy(RestartIntervalPolicy::all(16)), - )?; - #[cfg(feature = "indexer_stream")] - let events = open_ks( - "events", - opts() - // only iterators are used here, no point reads - .expect_point_read_hits(true) - .max_memtable_size(mb(cfg.db_events_memtable_size_mb)) - // the compression here wont be quite as good since events are quite random - // eg. by many different repos and different records etc. - // since its sequential we should still go with bigger block size though - // backfills will be sequential though... - .data_block_size_policy( - cfg.ephemeral - .then(|| BlockSizePolicy::new([kb(64), kb(128), kb(256)])) - .unwrap_or_else(|| BlockSizePolicy::new([kb(16), kb(64)])), - ) - // we are streaming the new events to consumers so we dont want to compress them - .data_block_compression_policy( - cfg.ephemeral - .then(|| { - CompressionPolicy::new([ - CompressionType::None, - get_compression("events", 3), - ]) - }) - .unwrap_or_else(|| { - CompressionPolicy::new([ - CompressionType::None, - get_compression("events", 3), - get_compression("events", 3), - get_compression("events", 5), - ]) - }), - ) - // ids are int, we can prefix truncate a lot - .data_block_restart_interval_policy(RestartIntervalPolicy::new([64, 128])), - )?; #[cfg(feature = "jetstream")] - let jetstream_events = open_ks( - "jetstream_events", - opts() - // time-ordered append-only stream metadata, only iterated for replay. - .expect_point_read_hits(true) - .max_memtable_size(mb(cfg.db_events_memtable_size_mb)) - .data_block_size_policy(BlockSizePolicy::new([kb(16), kb(64), kb(128)])) - .data_block_compression_policy(CompressionPolicy::new([ - CompressionType::None, - get_compression("jetstream_events", 3), - get_compression("jetstream_events", 5), - ])) - .data_block_restart_interval_policy(RestartIntervalPolicy::new([64, 128])), - )?; + let jetstream = super::keyspaces::JetstreamDb::open(&cx)?; let counts = open_ks( "counts", opts() @@ -318,21 +187,7 @@ impl Db { )?; #[cfg(feature = "relay")] - let relay_events = open_ks( - "relay_events", - opts() - // only iterated for cursor replay - .expect_point_read_hits(true) - .max_memtable_size(mb(cfg.db_events_memtable_size_mb)) - .data_block_size_policy(BlockSizePolicy::new([kb(64), kb(128), kb(256)])) - .data_block_compression_policy(CompressionPolicy::new([ - CompressionType::None, - get_compression("events", 3), - get_compression("events", 3), - get_compression("events", 5), - ])) - .data_block_restart_interval_policy(RestartIntervalPolicy::new([64, 128])), - )?; + let relay = super::keyspaces::RelayDb::open(&cx)?; #[cfg(feature = "backlinks")] let backlinks = open_ks( @@ -354,118 +209,43 @@ impl Db { // when adding new keyspaces, make sure to add them to the /stats endpoint // and also update any relevant /debug/* endpoints - #[cfg(feature = "indexer_stream")] - let (event_tx, _) = broadcast::channel(512); - #[cfg(feature = "relay")] - let (relay_broadcast_tx, _) = broadcast::channel(512); - - #[cfg(feature = "jetstream")] - let (jetstream_tx, _) = broadcast::channel(512); let this = Self { inner: db, path: cfg.database_path.clone(), repos, repo_metadata, - #[cfg(feature = "indexer")] - records, - #[cfg(feature = "indexer")] - blocks, cursors, - #[cfg(feature = "indexer")] - pending, - #[cfg(feature = "indexer")] - resync, - #[cfg(feature = "indexer")] - resync_buffer, - #[cfg(feature = "indexer_stream")] - events, - #[cfg(feature = "jetstream")] - jetstream_events, counts, filter, crawler, - #[cfg(feature = "backlinks")] - backlinks, - #[cfg(feature = "indexer_stream")] - event_tx, - counts_map: HashMap::new(), + #[cfg(feature = "indexer")] + indexer, #[cfg(feature = "indexer_stream")] - next_event_id: Arc::new(AtomicU64::new(0)), + stream, #[cfg(feature = "jetstream")] - jetstream_tx, - #[cfg(feature = "jetstream")] - next_jetstream_id: Arc::new(AtomicU64::new(0)), - #[cfg(feature = "jetstream")] - last_jetstream_time_us: Arc::new(AtomicI64::new(0)), - #[cfg(all(feature = "jetstream", feature = "relay"))] - jetstream_lock: Arc::new(parking_lot::Mutex::new(())), - #[cfg(feature = "relay")] - relay_events, + jetstream, #[cfg(feature = "relay")] - next_relay_seq: Arc::new(AtomicU64::new(0)), - #[cfg(feature = "relay")] - relay_broadcast_tx, + relay, + #[cfg(feature = "backlinks")] + backlinks, + counts_map: HashMap::new(), next_count_delta_id: Arc::new(AtomicU64::new(0)), count_delta_checkpoint_watermark: Arc::new(AtomicU64::new(0)), count_delta_gc_watermark, count_delta_in_flight: Arc::new(Mutex::new(BTreeSet::new())), - #[cfg(feature = "indexer")] - lifecycle_count_lock: Arc::new(Mutex::new(())), compaction_running: Arc::new(std::sync::atomic::AtomicBool::new(false)), }; migration::run(&this)?; #[cfg(feature = "relay")] - { - let mut last_relay_seq = 0u64; - if let Some(guard) = this.relay_events.iter().next_back() { - let k = guard.key().into_diagnostic()?; - last_relay_seq = u64::from_be_bytes( - k.as_ref() - .try_into() - .into_diagnostic() - .wrap_err("relay_events: invalid key length")?, - ); - } - this.next_relay_seq - .store(last_relay_seq + 1, std::sync::atomic::Ordering::Relaxed); - } - + this.relay.init()?; #[cfg(feature = "indexer_stream")] - { - let mut last_id = 0; - if let Some(guard) = this.events.iter().next_back() { - let k = guard.key().into_diagnostic()?; - last_id = u64::from_be_bytes( - k.as_ref() - .try_into() - .into_diagnostic() - .wrap_err("expected to be id (8 bytes)")?, - ); - } - // relaxed is fine since we are just initializing the db - this.next_event_id - .store(last_id + 1, std::sync::atomic::Ordering::Relaxed); - } - + this.stream.init()?; #[cfg(feature = "jetstream")] - { - let mut last_id = 0; - let mut last_time_us = 0; - if let Some(guard) = this.jetstream_events.iter().next_back() { - let k = guard.key().into_diagnostic()?; - let (time_us, id) = keys::parse_jetstream_event_key(&k)?; - last_id = id; - last_time_us = time_us; - } - this.next_jetstream_id - .store(last_id + 1, std::sync::atomic::Ordering::Relaxed); - this.last_jetstream_time_us - .store(last_time_us as i64, std::sync::atomic::Ordering::Relaxed); - } + this.jetstream.init()?; // load counts into memory for guard in this.counts.prefix(keys::COUNT_KS_PREFIX) { diff --git a/src/db/train.rs b/src/db/train.rs index 059986f..4da0243 100644 --- a/src/db/train.rs +++ b/src/db/train.rs @@ -14,11 +14,11 @@ impl Db { pub fn train_dict(&self, ks_name: &str) -> Result<()> { let ks = match ks_name { #[cfg(feature = "indexer")] - "blocks" => &self.blocks, + "blocks" => &self.indexer.blocks, #[cfg(feature = "indexer_stream")] - "events" => &self.events, + "events" => &self.stream.events, #[cfg(feature = "jetstream")] - "jetstream_events" => &self.jetstream_events, + "jetstream_events" => &self.jetstream.events, "repos" => &self.repos, #[cfg(feature = "backlinks")] "backlinks" => &self.backlinks, diff --git a/src/ingest/indexer/shard.rs b/src/ingest/indexer/shard.rs index 812da76..a674c97 100644 --- a/src/ingest/indexer/shard.rs +++ b/src/ingest/indexer/shard.rs @@ -287,11 +287,11 @@ impl FirehoseWorker { drop(reservation); #[cfg(feature = "indexer_stream")] for evt in broadcast_events.drain(..) { - let _ = state.db.event_tx.send(evt); + let _ = state.db.stream.event_tx.send(evt); } #[cfg(feature = "jetstream")] for evt in jetstream_events.drain(..) { - let _ = state.db.jetstream_tx.send(evt); + let _ = state.db.jetstream.tx.send(evt); } // state.db.inner.persist(fjall::PersistMode::Buffer).ok(); @@ -337,7 +337,7 @@ impl FirehoseWorker { let metadata_bytes = db.repo_metadata.get(&metadata_key).into_diagnostic()?; let is_backfilling = if let Some(metadata_bytes) = metadata_bytes { let metadata = crate::db::deser_repo_meta(metadata_bytes.as_ref())?; - db.pending + db.indexer.pending .get(keys::pending_key(metadata.index_id)) .into_diagnostic()? .is_some() @@ -509,7 +509,7 @@ impl FirehoseWorker { let db = &ctx.state.db; let prefix = keys::resync_buffer_prefix(did); - for guard in db.resync_buffer.prefix(&prefix) { + for guard in db.indexer.resync_buffer.prefix(&prefix) { let (key, value) = guard.into_inner().into_diagnostic()?; let commit: Commit = rmp_serde::from_slice(&value).into_diagnostic()?; @@ -525,14 +525,14 @@ impl FirehoseWorker { Ok(r) => r, Err(e) => { if !Self::check_if_retriable_failure(&e) { - ctx.batch.remove(&db.resync_buffer, key); + ctx.batch.remove(&db.indexer.resync_buffer, key); } return Err(e); } }; match res { RepoProcessResult::Ok(rs) => { - ctx.batch.remove(&db.resync_buffer, key); + ctx.batch.remove(&db.indexer.resync_buffer, key); repo_state = rs; } RepoProcessResult::NeedsBackfill(_) => { @@ -540,7 +540,7 @@ impl FirehoseWorker { return Ok(RepoProcessResult::NeedsBackfill(None)); } RepoProcessResult::Deleted => { - ctx.batch.remove(&db.resync_buffer, key); + ctx.batch.remove(&db.indexer.resync_buffer, key); return Ok(RepoProcessResult::Deleted); } } @@ -575,18 +575,18 @@ impl FirehoseWorker { let old_pkey = keys::pending_key(metadata.index_id); let was_pending = had_metadata && db - .pending + .indexer.pending .get(old_pkey.as_slice()) .into_diagnostic()? .is_some(); // remove old pending entry and insert new one with fresh index_id if had_metadata { // only remove if we had one so we dont delete a random entry - batch.remove(&db.pending, old_pkey); + batch.remove(&db.indexer.pending, old_pkey); } metadata.index_id = rand::random::(); - batch.insert(&db.pending, keys::pending_key(metadata.index_id), &repo_key); + batch.insert(&db.indexer.pending, keys::pending_key(metadata.index_id), &repo_key); batch.insert(&db.repo_metadata, &meta_key, ser_repo_meta(&metadata)?); if !was_pending { lifecycle_counts.transition(&mut batch, did, GaugeState::Pending)?; diff --git a/src/ingest/relay/sink/relay.rs b/src/ingest/relay/sink/relay.rs index b4f91b0..4a1a263 100644 --- a/src/ingest/relay/sink/relay.rs +++ b/src/ingest/relay/sink/relay.rs @@ -75,9 +75,9 @@ impl EventSink { let started = crate::ingest::firehose_stats::StatsInstant::now(); let result = (|| { let db = &state.db; - let seq = db.next_relay_seq.fetch_add(1, Ordering::SeqCst); + let seq = db.relay.next_seq.fetch_add(1, Ordering::SeqCst); let frame = make_frame(seq as i64)?; - batch.insert(&db.relay_events, keys::relay_event_key(seq), frame.as_ref()); + batch.insert(&db.relay.events, keys::relay_event_key(seq), frame.as_ref()); self.broadcasts.push(RelayBroadcast::Ephemeral(seq, frame)); self.broadcasts.push(RelayBroadcast::Persisted(seq)); Ok(seq) @@ -118,7 +118,7 @@ impl EventSink { { // skip car parse and record materialization when no jetstream subscribers are // connected — the stream thread re-reads from relay_events as a fallback. - let has_subscribers = state.db.jetstream_tx.receiver_count() > 0; + let has_subscribers = state.db.jetstream.tx.receiver_count() > 0; for (op_index, collection) in &jetstream_ops { let ephemeral = has_subscribers .then(|| build_relay_commit_ephemeral(&commit, *op_index, collection, &parsed_blocks)) @@ -210,7 +210,7 @@ impl EventSink { { let mut batch = batch; let mut jetstream_broadcasts = Vec::new(); - let _lock = state.db.jetstream_lock.lock(); + let _lock = state.db.jetstream.lock.lock(); for (event, ephemeral) in self.jetstream_events.drain(..) { jetstream_broadcasts.push(crate::jetstream::stage_event( &mut batch, &state.db, event, ephemeral, @@ -231,11 +231,11 @@ impl EventSink { pub(crate) fn flush(&mut self, state: &AppState, staged: Staged) { for broadcast in self.broadcasts.drain(..) { - let _ = state.db.relay_broadcast_tx.send(broadcast); + let _ = state.db.relay.broadcast_tx.send(broadcast); } #[cfg(feature = "jetstream")] for broadcast in staged.jetstream_broadcasts { - let _ = state.db.jetstream_tx.send(broadcast); + let _ = state.db.jetstream.tx.send(broadcast); } #[cfg(not(feature = "jetstream"))] let _ = staged; diff --git a/src/jetstream.rs b/src/jetstream.rs index e716241..fb331e9 100644 --- a/src/jetstream.rs +++ b/src/jetstream.rs @@ -31,7 +31,7 @@ pub(crate) fn stage_event( event: StoredJetstreamEvent<'_>, ephemeral: Option, ) -> Result { - let id = db.next_jetstream_id.fetch_add(1, Ordering::SeqCst); + let id = db.jetstream.next_id.fetch_add(1, Ordering::SeqCst); let time_us = next_time_us(db); let ephemeral = ephemeral.and_then(|data| { let json_event = crate::control::stream::JetstreamEvent { @@ -54,7 +54,7 @@ pub(crate) fn stage_event( let event = event.into_static(); let bytes = rmp_serde::to_vec(&event).into_diagnostic()?; batch.insert( - &db.jetstream_events, + &db.jetstream.events, keys::jetstream_event_key(time_us as u64, id), bytes, ); @@ -68,11 +68,11 @@ pub(crate) fn stage_event( fn next_time_us(db: &Db) -> i64 { loop { - let last = db.last_jetstream_time_us.load(Ordering::SeqCst); + let last = db.jetstream.last_time_us.load(Ordering::SeqCst); let now = chrono::Utc::now().timestamp_micros(); let next = now.max(last.saturating_add(1)); if db - .last_jetstream_time_us + .jetstream.last_time_us .compare_exchange(last, next, Ordering::SeqCst, Ordering::SeqCst) .is_ok() { diff --git a/src/ops.rs b/src/ops.rs index d9f535e..bd88ff7 100644 --- a/src/ops.rs +++ b/src/ops.rs @@ -33,7 +33,7 @@ use crate::types::{JetstreamBroadcast, StoredJetstreamEvent}; pub fn persist_to_resync_buffer(db: &Db, did: &Did, commit: &Commit) -> Result<()> { let key = keys::resync_buffer_key(did, DbTid::from(&commit.rev)); let value = rmp_serde::to_vec_named(commit).into_diagnostic()?; - db.resync_buffer.insert(key, value).into_diagnostic()?; + db.indexer.resync_buffer.insert(key, value).into_diagnostic()?; debug!( did = %did, seq = commit.seq, @@ -46,7 +46,7 @@ pub fn persist_to_resync_buffer(db: &Db, did: &Did, commit: &Commit) -> Result<( // we dont replay these, consumers can just fetch identity themselves if they need it #[cfg(feature = "indexer_stream")] pub fn make_identity_event(db: &Db, evt: IdentityEvt<'static>) -> BroadcastEvent { - let event_id = db.next_event_id.fetch_add(1, Ordering::SeqCst); + let event_id = db.stream.next_event_id.fetch_add(1, Ordering::SeqCst); let marshallable = MarshallableEvt { id: event_id, kind: crate::types::EventType::Identity, @@ -59,7 +59,7 @@ pub fn make_identity_event(db: &Db, evt: IdentityEvt<'static>) -> BroadcastEvent #[cfg(feature = "indexer_stream")] pub fn make_account_event(db: &Db, evt: AccountEvt<'static>) -> BroadcastEvent { - let event_id = db.next_event_id.fetch_add(1, Ordering::SeqCst); + let event_id = db.stream.next_event_id.fetch_add(1, Ordering::SeqCst); let marshallable = MarshallableEvt { id: event_id, kind: crate::types::EventType::Account, @@ -84,28 +84,28 @@ pub fn delete_repo( let metadata_bytes = db.repo_metadata.get(&metadata_key).into_diagnostic()?; if let Some(metadata_bytes) = metadata_bytes { let metadata = db::deser_repo_meta(&metadata_bytes)?; - batch.remove(&db.pending, keys::pending_key(metadata.index_id)); + batch.remove(&db.indexer.pending, keys::pending_key(metadata.index_id)); } // we don't delete from repos, relay uses it as a tombstone // todo: we should still delete it after some time - batch.remove(&db.resync, &repo_key); + batch.remove(&db.indexer.resync, &repo_key); batch.remove(&db.repo_metadata, &metadata_key); // 2. delete from resync buffer let resync_prefix = keys::resync_buffer_prefix(did); - for guard in db.resync_buffer.prefix(&resync_prefix) { + for guard in db.indexer.resync_buffer.prefix(&resync_prefix) { let k = guard.key().into_diagnostic()?; - batch.remove(&db.resync_buffer, k); + batch.remove(&db.indexer.resync_buffer, k); } // 3. delete from records // todo: figure out how we want to handle the blocks associated with these records // without delving too much into gc madness we had before let records_prefix = keys::record_prefix_did(did); - for guard in db.records.prefix(&records_prefix) { + for guard in db.indexer.records.prefix(&records_prefix) { let (k, _cid_bytes) = guard.into_inner().into_diagnostic()?; - batch.remove(&db.records, k); + batch.remove(&db.indexer.records, k); } // 4. reset collection counts @@ -149,7 +149,7 @@ pub fn transition_repo<'s>( match &new_status { RepoStatus::Synced => { lifecycle_transitions.push((did.clone().into_static(), GaugeState::Synced)); - batch.remove(&db.pending, pending_key.as_slice()); + batch.remove(&db.indexer.pending, pending_key.as_slice()); // we dont have to remove from resync here because it has to transition resync -> pending first } RepoStatus::Error(msg) => { @@ -158,7 +158,7 @@ pub fn transition_repo<'s>( did.clone().into_static(), GaugeState::Resync(Some(ResyncErrorKind::Generic)), )); - batch.remove(&db.pending, pending_key.as_slice()); + batch.remove(&db.indexer.pending, pending_key.as_slice()); // TODO: we need to make errors have kind instead of "message" in repo status // and then pass it to resync error kind let resync_state = crate::types::ResyncState::Error { @@ -167,7 +167,7 @@ pub fn transition_repo<'s>( next_retry: chrono::Utc::now().timestamp(), }; batch.insert( - &db.resync, + &db.indexer.resync, &repo_key, rmp_serde::to_vec(&resync_state).into_diagnostic()?, ); @@ -180,7 +180,7 @@ pub fn transition_repo<'s>( status: new_status.clone(), }; batch.insert( - &db.resync, + &db.indexer.resync, &repo_key, rmp_serde::to_vec(&resync_state).into_diagnostic()?, ); @@ -188,8 +188,8 @@ pub fn transition_repo<'s>( RepoStatus::Deleted => { lifecycle_transitions.push((did.clone().into_static(), GaugeState::Synced)); // terminal state: remove from queues, no resync entry needed - batch.remove(&db.pending, pending_key.as_slice()); - batch.remove(&db.resync, &repo_key); + batch.remove(&db.indexer.pending, pending_key.as_slice()); + batch.remove(&db.indexer.resync, &repo_key); } RepoStatus::Desynchronized | RepoStatus::Throttled => { lifecycle_transitions.push(( @@ -197,14 +197,14 @@ pub fn transition_repo<'s>( GaugeState::Resync(Some(ResyncErrorKind::Generic)), )); // like an error: remove from pending and schedule a resync attempt - batch.remove(&db.pending, pending_key.as_slice()); + batch.remove(&db.indexer.pending, pending_key.as_slice()); let resync_state = crate::types::ResyncState::Error { kind: crate::types::ResyncErrorKind::Generic, retry_count: 0, next_retry: chrono::Utc::now().timestamp(), }; batch.insert( - &db.resync, + &db.indexer.resync, &repo_key, rmp_serde::to_vec(&resync_state).into_diagnostic()?, ); @@ -257,13 +257,13 @@ pub fn apply_commit<'s>( let rev = DbTid::from(&commit.rev); #[cfg(feature = "indexer_stream")] - let should_broadcast_live = db.event_tx.receiver_count() > 0; + let should_broadcast_live = db.stream.event_tx.receiver_count() > 0; #[cfg(feature = "indexer_stream")] let mut live_events = Vec::new(); #[cfg(feature = "indexer_stream")] let mut last_event_id = None; #[cfg(feature = "jetstream")] - let should_stage_jetstream = db.jetstream_tx.receiver_count() > 0; + let should_stage_jetstream = db.jetstream.tx.receiver_count() > 0; #[cfg(feature = "jetstream")] let mut jetstream_events = Vec::new(); @@ -311,9 +311,9 @@ pub fn apply_commit<'s>( blocks_count += 1; if !ephemeral { if !only_index_links { - batch.insert(&db.blocks, block_key.clone(), bytes.as_ref()); + batch.insert(&db.indexer.blocks, block_key.clone(), bytes.as_ref()); } - batch.insert(&db.records, db_key.clone(), cid_raw); + batch.insert(&db.indexer.records, db_key.clone(), cid_raw); // accumulate counts if action == DbAction::Create { records_delta += 1; @@ -345,7 +345,7 @@ pub fn apply_commit<'s>( } DbAction::Delete => { if !ephemeral { - batch.remove(&db.records, db_key); + batch.remove(&db.indexer.records, db_key); // accumulate counts records_delta -= 1; @@ -375,7 +375,7 @@ pub fn apply_commit<'s>( }) .unwrap_or(StoredData::Nothing); - let event_id = db.next_event_id.fetch_add(1, Ordering::SeqCst); + let event_id = db.stream.next_event_id.fetch_add(1, Ordering::SeqCst); last_event_id = Some(event_id); let did_trimmed = TrimmedDid::from(did); let collection = CowStr::Borrowed(collection); @@ -390,7 +390,7 @@ pub fn apply_commit<'s>( data: data.clone(), }; let bytes = rmp_serde::to_vec(&evt).into_diagnostic()?; - batch.insert(&db.events, keys::event_key(event_id), bytes); + batch.insert(&db.stream.events, keys::event_key(event_id), bytes); #[cfg(feature = "jetstream")] { -- 2.51.2