From 8f52f23a3b8c1fcf57d14c54bafc6b97f925b3f9 Mon Sep 17 00:00:00 2001 From: dawn <90008@klbr.net> Date: Sat, 11 Jul 2026 21:18:27 +0300 Subject: [PATCH] [db] unify record mutation transactions --- src/api/debug.rs | 16 ++- src/backfill/sparse.rs | 152 ++------------------ src/backfill/worker/process.rs | 161 ++++----------------- src/backlinks/mod.rs | 2 +- src/backlinks/store.rs | 6 +- src/control/indexer.rs | 3 +- src/control/repos/indexer.rs | 41 ++++-- src/control/stats.rs | 2 +- src/control/stream/indexer.rs | 26 ++-- src/control/stream/jetstream.rs | 11 +- src/db/indexer.rs | 12 ++ src/db/keyspaces.rs | 50 ++++++- src/db/mod.rs | 14 +- src/db/txn.rs | 243 ++++++++++++++++++++++++++++++++ src/ingest/indexer/shard.rs | 82 ++++------- src/ops.rs | 104 +++----------- src/ops/record_events.rs | 60 +++++--- 17 files changed, 495 insertions(+), 490 deletions(-) create mode 100644 src/db/txn.rs diff --git a/src/api/debug.rs b/src/api/debug.rs index d4bbe3f..e06d6a0 100644 --- a/src/api/debug.rs +++ b/src/api/debug.rs @@ -60,7 +60,9 @@ pub async fn handle_debug_count( .map_err(|_| StatusCode::BAD_REQUEST)?; let db = &state.db; - let ks = db.indexer.records.clone(); + let ks = db + .keyspace_by_name("records") + .expect("records keyspace exists in indexer mode"); // {TrimmedDid}|{collection}| let prefix = keys::record_prefix_collection(&did, &req.collection); @@ -303,9 +305,14 @@ pub async fn handle_debug_seed_events( for _ in 0..req.count { let seq = state .db - .stream.next_event_id + .stream + .next_event_id .fetch_add(1, std::sync::atomic::Ordering::SeqCst); - batch.insert(&state.db.stream.events, crate::db::keys::event_key(seq), b"dummy"); + state.db.stream.stage_event( + &mut batch, + crate::db::keys::event_key(seq), + b"dummy", + ); } } } else if req.partition == "relay_events" { @@ -314,7 +321,8 @@ pub async fn handle_debug_seed_events( for _ in 0..req.count { let seq = state .db - .relay.next_seq + .relay + .next_seq .fetch_add(1, std::sync::atomic::Ordering::SeqCst); batch.insert( &state.db.relay.events, diff --git a/src/backfill/sparse.rs b/src/backfill/sparse.rs index b357540..014e436 100644 --- a/src/backfill/sparse.rs +++ b/src/backfill/sparse.rs @@ -2,21 +2,16 @@ use crate::backfill::client::ThrottledHttpClient; use crate::backfill::error::BackfillError; use crate::config::{BackfillStrategy, RateTier}; use crate::db::types::{DbAction, DbRkey}; -#[cfg(feature = "indexer_stream")] -use crate::db::types::TrimmedDid; -use crate::db::{self, CountDeltas, keys, ser_repo_state}; +use crate::db::{self, Txn as DbTxn, keys}; use crate::filter::{FilterConfig, FilterMode}; use crate::ops; use crate::sparse_mst::{SparseScanner, mst_node_layer, sparse_probe_collection, sparse_ranges}; use crate::state::AppState; use crate::types::{Commit, RepoState}; -use fjall::Slice; use futures::{StreamExt, stream}; use jacquard_api::com_atproto::sync::get_blocks::GetBlocksError; use jacquard_api::com_atproto::sync::get_record::{GetRecord, GetRecordError}; -#[cfg(feature = "jetstream")] -use jacquard_common::IntoStatic; use jacquard_common::types::cid::{Cid as AtCid, IpldCid}; use jacquard_common::types::did::Did; use jacquard_common::types::string::{Nsid, RecordKey}; @@ -27,16 +22,10 @@ use reqwest::StatusCode; use smol_str::{SmolStr, ToSmolStr}; use std::collections::{BTreeMap, HashMap}; use std::sync::Arc; -#[cfg(feature = "indexer_stream")] -use std::sync::atomic::Ordering; use tracing::{debug, trace, warn}; -#[cfg(feature = "indexer_stream")] -use crate::types::{StoredData, StoredEvent}; use crate::util::throttle::ThrottleHandle; use crate::util::url_to_fluent_uri; -#[cfg(feature = "indexer_stream")] -use jacquard_common::CowStr; #[derive(Debug)] pub(crate) struct SparseBackfillSuccess { @@ -425,18 +414,14 @@ async fn persist_sparse_backfill( tokio::task::spawn_blocking(move || { let filter = app_state.filter.load(); let ephemeral = app_state.ephemeral; - let only_index_links = app_state.only_index_links; let mut count = 0; - let mut delta = 0; - let mut added_blocks = 0; let mut collection_counts: HashMap = HashMap::new(); - let mut batch = app_state.db.inner.batch(); let prefix = keys::record_prefix_did(&did); let mut existing_cids: HashMap<(SmolStr, DbRkey), SmolStr> = HashMap::new(); if !ephemeral { - for guard in app_state.db.indexer.records.prefix(&prefix) { + for guard in app_state.db.indexer.record_prefix(&prefix) { let (key, cid_bytes) = guard.into_inner().into_diagnostic()?; let mut remaining = key[prefix.len()..].splitn(2, |b| keys::SEP.eq(b)); let collection_raw = remaining @@ -463,6 +448,8 @@ async fn persist_sparse_backfill( } let mut signal_seen = filter.mode == FilterMode::Full || filter.signals.is_empty(); + let mut txn = DbTxn::new(&app_state.db); + let mut record_txn = txn.backfill_records(&app_state, &root_commit.rev, &did); for (key, cid) in leaves { let (collection, rkey) = ops::parse_path(&key)?; @@ -471,7 +458,7 @@ async fn persist_sparse_backfill( continue; } - let Some(val) = blocks.get(&cid).cloned() else { + let Some(val) = blocks.get(&cid) else { return Err(miette::miette!("missing sparse record block {cid}").into()); }; @@ -498,114 +485,13 @@ async fn persist_sparse_backfill( }; trace!(collection = %collection, rkey = %rkey, cid = %cid, ?action, "action sparse record"); - let db_key = keys::record_key(&did, collection, &rkey); - let cid_raw = cid.to_bytes(); - let block_key = Slice::from(keys::block_key(collection, &cid_raw)); - if !ephemeral { - if !only_index_links { - batch.insert(&app_state.db.indexer.blocks, block_key.clone(), val.as_ref()); - } - 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( - &mut batch, - &app_state.db.backlinks, - did.as_str(), - collection, - &rkey.to_smolstr(), - &value, - )?; - } - } - - added_blocks += 1; - if action == DbAction::Create { - delta += 1; - } - - #[cfg(feature = "indexer_stream")] - { - let event_id = app_state.db.stream.next_event_id.fetch_add(1, Ordering::SeqCst); - let evt = StoredEvent { - live: false, - did: TrimmedDid::from(&did), - rev: root_commit.rev, - collection: CowStr::Borrowed(collection), - rkey, - action, - data: if ephemeral { - StoredData::Block(val) - } else if only_index_links { - StoredData::Nothing - } else { - StoredData::Ptr(cid_obj.to_ipld().expect("valid cid")) - }, - }; - let bytes = rmp_serde::to_vec(&evt).into_diagnostic()?; - batch.insert(&app_state.db.stream.events, keys::event_key(event_id), bytes); - - #[cfg(feature = "jetstream")] - { - let jetstream = crate::types::StoredJetstreamEvent::Commit { - did: TrimmedDid::from(&did).into_static(), - collection: CowStr::Borrowed(collection).into_static(), - event_id, - live: false, - }; - crate::jetstream::stage_event(&mut batch, &app_state.db, jetstream, None)?; - } - } - + record_txn.put_record(collection, &rkey, cid, val, action)?; count += 1; } for ((collection, rkey), cid) in existing_cids { trace!(collection = %collection, rkey = %rkey, cid = %cid, "remove sparse-stale record"); - - batch.remove( - &app_state.db.indexer.records, - keys::record_key(&did, &collection, &rkey), - ); - #[cfg(feature = "backlinks")] - crate::backlinks::store::delete_record( - &mut batch, - &app_state.db.backlinks, - did.as_str(), - &collection, - &rkey.to_smolstr(), - )?; - - #[cfg(feature = "indexer_stream")] - { - let event_id = app_state.db.stream.next_event_id.fetch_add(1, Ordering::SeqCst); - let evt = StoredEvent { - live: false, - did: TrimmedDid::from(&did), - rev: root_commit.rev, - collection: CowStr::Borrowed(&collection), - rkey, - action: DbAction::Delete, - data: StoredData::Nothing, - }; - let bytes = rmp_serde::to_vec(&evt).into_diagnostic()?; - batch.insert(&app_state.db.stream.events, keys::event_key(event_id), bytes); - - #[cfg(feature = "jetstream")] - { - let jetstream = crate::types::StoredJetstreamEvent::Commit { - did: TrimmedDid::from(&did).into_static(), - collection: CowStr::Borrowed(&collection).into_static(), - event_id, - live: false, - }; - crate::jetstream::stage_event(&mut batch, &app_state.db, jetstream, None)?; - } - } - - delta -= 1; + record_txn.delete_record(&collection, &rkey)?; count += 1; } @@ -616,12 +502,8 @@ async fn persist_sparse_backfill( state.root = Some(root_commit); state.touch(); - - batch.insert( - &app_state.db.repos, - keys::repo_key(&did), - ser_repo_state(&state)?, - ); + record_txn.update_repo_state(&state)?; + let _events = record_txn.finish()?; let metadata_key = keys::repo_metadata_key(&did); let metadata_bytes = app_state @@ -632,7 +514,7 @@ async fn persist_sparse_backfill( .ok_or_else(|| miette::miette!("repo metadata not found for {}", did))?; let mut metadata = crate::db::deser_repo_meta(&metadata_bytes)?; metadata.tracked = true; - batch.insert( + txn.batch.insert( &app_state.db.repo_metadata, &metadata_key, crate::db::ser_repo_meta(&metadata)?, @@ -640,7 +522,7 @@ async fn persist_sparse_backfill( if !ephemeral { db::replace_record_counts_matching( - &mut batch, + &mut txn.batch, &app_state.db, &did, |collection| filter.matches_collection(collection), @@ -648,17 +530,7 @@ async fn persist_sparse_backfill( )?; } - let mut count_deltas = CountDeltas::default(); - if delta != 0 { - count_deltas.add_records(delta); - } - if added_blocks > 0 { - count_deltas.add_blocks(added_blocks); - } - let reservation = app_state.db.stage_count_deltas(&mut batch, &count_deltas); - batch.commit().into_diagnostic()?; - app_state.db.apply_count_deltas(&count_deltas); - drop(reservation); + txn.commit()?; Ok::<_, miette::Report>(Some((count, state))) }) diff --git a/src/backfill/worker/process.rs b/src/backfill/worker/process.rs index 225f7f8..165fdb6 100644 --- a/src/backfill/worker/process.rs +++ b/src/backfill/worker/process.rs @@ -20,9 +20,7 @@ use crate::backfill::error::BackfillError; use crate::backfill::sparse::{SparseBackfillResult, process_did_sparse}; use crate::config::BackfillStrategy; use crate::db::types::{DbAction, DbRkey}; -#[cfg(feature = "indexer_stream")] -use crate::db::types::TrimmedDid; -use crate::db::{self, CountDeltas, Db, keys, ser_repo_state}; +use crate::db::{self, CountDeltas, Db, Txn as DbTxn, keys}; use crate::filter::FilterMode; use crate::ops; use crate::sparse_mst::sparse_probe_collection; @@ -31,9 +29,7 @@ use crate::types::{Commit, GaugeState, RepoState, RepoStatus, ResyncState}; use crate::util::url_to_fluent_uri; #[cfg(feature = "indexer_stream")] -use crate::types::{AccountEvt, BroadcastEvent, StoredData, StoredEvent}; -#[cfg(feature = "indexer_stream")] -use jacquard_common::CowStr; +use crate::types::{AccountEvt, BroadcastEvent}; #[cfg(feature = "indexer_stream")] use std::sync::atomic::Ordering; @@ -91,7 +87,11 @@ pub(crate) async fn process_did( .then_some(None) .unwrap_or_else(|| Some(status.into())), }; - let _ = app_state.db.stream.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 { @@ -387,18 +387,15 @@ pub(crate) async fn process_did( let result = { let app_state = app_state.clone(); let did = did.clone(); - #[cfg(feature = "indexer_stream")] let rev = root_commit.rev; tokio::task::spawn_blocking(move || { let filter = app_state.filter.load(); let ephemeral = app_state.ephemeral; - let only_index_links = app_state.only_index_links; let mut count = 0; - let mut delta = 0; - let mut added_blocks = 0; let mut collection_counts: HashMap = HashMap::new(); - let mut batch = app_state.db.inner.batch(); + let mut txn = DbTxn::new(&app_state.db); + let mut record_txn = txn.backfill_records(&app_state, &rev, &did); // clone the Arc so we hold an independent reference to the block store, // allowing mst (and its entire loaded node tree) to be freed immediately // rather than surviving until the end of spawn_blocking. @@ -409,7 +406,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.indexer.records.prefix(&prefix) { + for guard in app_state.db.indexer.record_prefix(&prefix) { let (key, cid_bytes) = guard.into_inner().into_diagnostic()?; // key is did|collection|rkey // skip did| @@ -473,66 +470,13 @@ pub(crate) async fn process_did( }; trace!(collection = %collection, rkey = %rkey, cid = %cid, ?action, "action record"); - // key is did|collection|rkey - let db_key = keys::record_key(&did, collection, &rkey); - - let cid_raw = cid.to_bytes(); - let block_key = Slice::from(keys::block_key(collection, &cid_raw)); - if !ephemeral { - if !only_index_links { - batch.insert(&app_state.db.indexer.blocks, block_key.clone(), val.as_ref()); - } - 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( - &mut batch, - &app_state.db.backlinks, - did.as_str(), - collection, - &rkey.to_smolstr(), - &value, - )?; - } - } - - added_blocks += 1; - if action == DbAction::Create { - delta += 1; - } - - #[cfg(feature = "indexer_stream")] - { - let event_id = app_state.db.stream.next_event_id.fetch_add(1, Ordering::SeqCst); - let evt = StoredEvent { - live: false, - did: TrimmedDid::from(&did), - rev, - collection: CowStr::Borrowed(collection), - rkey, - action, - data: if ephemeral { - StoredData::Block(val) - } else if only_index_links { - StoredData::Nothing - } else { - StoredData::Ptr(cid_obj.to_ipld().expect("valid cid")) - }, - }; - let bytes = rmp_serde::to_vec(&evt).into_diagnostic()?; - batch.insert(&app_state.db.stream.events, keys::event_key(event_id), bytes); - - #[cfg(feature = "jetstream")] - { - let jetstream = crate::types::StoredJetstreamEvent::Commit { - did: TrimmedDid::from(&did).into_static(), - collection: CowStr::Borrowed(collection).into_static(), - event_id, - live: false, - }; - crate::jetstream::stage_event(&mut batch, &app_state.db, jetstream, None)?; - } - } + record_txn.put_record( + collection, + &rkey, + cid_obj.to_ipld().expect("valid cid"), + &val, + action, + )?; count += 1; } @@ -546,49 +490,8 @@ pub(crate) async fn process_did( for ((collection, rkey), cid) in existing_cids { trace!(collection = %collection, rkey = %rkey, cid = %cid, "remove existing record"); - // we dont have to put if ephemeral around here since - // existing_cids will be empty anyway - batch.remove( - &app_state.db.indexer.records, - keys::record_key(&did, &collection, &rkey), - ); - #[cfg(feature = "backlinks")] - crate::backlinks::store::delete_record( - &mut batch, - &app_state.db.backlinks, - did.as_str(), - &collection, - &rkey.to_smolstr(), - )?; - - #[cfg(feature = "indexer_stream")] - { - let event_id = app_state.db.stream.next_event_id.fetch_add(1, Ordering::SeqCst); - let evt = StoredEvent { - live: false, - did: TrimmedDid::from(&did), - rev, - collection: CowStr::Borrowed(&collection), - rkey, - action: DbAction::Delete, - data: StoredData::Nothing, - }; - let bytes = rmp_serde::to_vec(&evt).into_diagnostic()?; - batch.insert(&app_state.db.stream.events, keys::event_key(event_id), bytes); - - #[cfg(feature = "jetstream")] - { - let jetstream = crate::types::StoredJetstreamEvent::Commit { - did: TrimmedDid::from(&did).into_static(), - collection: CowStr::Borrowed(&collection).into_static(), - event_id, - live: false, - }; - crate::jetstream::stage_event(&mut batch, &app_state.db, jetstream, None)?; - } - } - - delta -= 1; + // existing_cids is empty in ephemeral mode. + record_txn.delete_record(&collection, &rkey)?; count += 1; } @@ -596,16 +499,10 @@ pub(crate) async fn process_did( trace!(signals = ?filter.signals, "no signal-matching records found, discarding repo"); return Ok::<_, miette::Report>(None); } - - // 6. update data, status is updated in worker shard state.root = Some(root_commit); state.touch(); - - batch.insert( - &app_state.db.repos, - keys::repo_key(&did), - ser_repo_state(&state)?, - ); + record_txn.update_repo_state(&state)?; + record_txn.finish()?; let metadata_key = keys::repo_metadata_key(&did); let metadata_bytes = app_state @@ -616,7 +513,7 @@ pub(crate) async fn process_did( .ok_or_else(|| miette::miette!("repo metadata not found for {}", did))?; let mut metadata = crate::db::deser_repo_meta(&metadata_bytes)?; metadata.tracked = true; - batch.insert( + txn.batch.insert( &app_state.db.repo_metadata, &metadata_key, crate::db::ser_repo_meta(&metadata)?, @@ -625,24 +522,14 @@ pub(crate) async fn process_did( // add the counts if !ephemeral { db::replace_record_counts( - &mut batch, + &mut txn.batch, &app_state.db, &did, collection_counts.iter().map(|(col, cnt)| (col.as_str(), *cnt)), )?; } - let mut count_deltas = CountDeltas::default(); - if delta != 0 { - count_deltas.add_records(delta); - } - if added_blocks > 0 { - count_deltas.add_blocks(added_blocks); - } - let reservation = app_state.db.stage_count_deltas(&mut batch, &count_deltas); - batch.commit().into_diagnostic()?; - app_state.db.apply_count_deltas(&count_deltas); - drop(reservation); + txn.commit()?; Ok::<_, miette::Report>(Some(count)) }) diff --git a/src/backlinks/mod.rs b/src/backlinks/mod.rs index 6f7cca0..7b188ab 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.indexer.records, did, col, rkey) else { + let Some(cid) = store::lookup_cid(&db.indexer, did, col, rkey) else { // record deleted after backlink was written continue; }; diff --git a/src/backlinks/store.rs b/src/backlinks/store.rs index 130dee6..aa1e225 100644 --- a/src/backlinks/store.rs +++ b/src/backlinks/store.rs @@ -217,8 +217,8 @@ pub fn delete_repo(batch: &mut OwnedWriteBatch, backlinks_ks: &Keyspace, did: &D /// look up the CID string for a record from the records keyspace. /// returns `None` if the record is not found (e.g. deleted after backlink was written). -pub fn lookup_cid_from_ks( - records: &Keyspace, +pub fn lookup_cid( + indexer: &crate::db::IndexerDb, did: &str, collection: &str, rkey: &str, @@ -226,7 +226,7 @@ pub fn lookup_cid_from_ks( let did = Did::new(did).ok()?; let db_rkey = DbRkey::new(rkey); let record_key = keys::record_key(&did, collection, &db_rkey); - let cid_bytes = records.get(record_key).ok()??; + let cid_bytes = indexer.record(record_key).ok()??; Cid::new(&cid_bytes).ok().map(|c| c.as_str().to_string()) } diff --git a/src/control/indexer.rs b/src/control/indexer.rs index d0e2110..aa7c212 100644 --- a/src/control/indexer.rs +++ b/src/control/indexer.rs @@ -128,7 +128,8 @@ impl Hydrant { }; let bytes = rmp_serde::to_vec(&evt).expect("msgpack serialization cannot fail"); total_bytes += bytes.len(); - batch.insert(&db.stream.events, keys::event_key(event_id), bytes); + db.stream + .stage_event(&mut batch, 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 23a7656..52555e4 100644 --- a/src/control/repos/indexer.rs +++ b/src/control/repos/indexer.rs @@ -21,7 +21,8 @@ impl ReposControl { let state = self.0.clone(); self.0 .db - .indexer.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 +72,8 @@ impl ReposControl { let state = self.0.clone(); self.0 .db - .indexer.resync + .indexer + .resync .range((start_bound, std::ops::Bound::Unbounded)) .map(move |g| { let did_key = g.key().into_diagnostic()?; @@ -124,7 +126,8 @@ impl ReposControl { // skip if already in pending queue let is_pending = db - .indexer.pending + .indexer + .pending .get(keys::pending_key(metadata.index_id)) .into_diagnostic()? .is_some(); @@ -134,7 +137,11 @@ impl ReposControl { let old_pending = keys::pending_key(metadata.index_id); batch.remove(&db.indexer.pending, old_pending); metadata.index_id = rand::Rng::next_u64(&mut rand::rng()); - batch.insert(&db.indexer.pending, keys::pending_key(metadata.index_id), &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, @@ -234,7 +241,11 @@ impl ReposControl { &metadata_key, crate::db::ser_repo_meta(&metadata)?, ); - batch.insert(&db.indexer.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); @@ -332,14 +343,14 @@ impl<'i> RepoHandle<'i> { tokio::task::spawn_blocking(move || { use miette::WrapErr; - let cid_bytes = state.db.indexer.records.get(db_key).into_diagnostic()?; + let cid_bytes = state.db.indexer.record(db_key).into_diagnostic()?; let Some(cid_bytes) = cid_bytes else { return Ok(None); }; // lookup block using col|cid key let block_key = keys::block_key(&collection, &cid_bytes); - let Some(block_bytes) = state.db.indexer.blocks.get(block_key).into_diagnostic()? else { + let Some(block_bytes) = state.db.indexer.block(block_key).into_diagnostic()? else { miette::bail!("block {cid_bytes:?} not found, this is a bug!!"); }; @@ -398,8 +409,8 @@ impl<'i> RepoHandle<'i> { Box::new( state .db - .indexer.records - .range(prefix.as_slice()..end_key.as_slice()) + .indexer + .record_range(prefix.as_slice()..end_key.as_slice()) .rev(), ) } else { @@ -413,7 +424,7 @@ impl<'i> RepoHandle<'i> { prefix.clone() }; - Box::new(state.db.indexer.records.range(start_key.as_slice()..)) + Box::new(state.db.indexer.record_range(start_key.as_slice()..)) }; for item in iter { @@ -432,8 +443,8 @@ 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 - .indexer.blocks - .get(keys::block_key(collection.as_str(), &cid_bytes)) + .indexer + .block(keys::block_key(collection.as_str(), &cid_bytes)) { let value: Data = serde_ipld_dagcbor::from_slice(&block_bytes).unwrap_or(Data::Null); @@ -509,7 +520,7 @@ impl<'i> RepoHandle<'i> { let mut mst = mst; let prefix = keys::record_prefix_did(&did); - for guard in app_state.db.indexer.records.prefix(&prefix) { + for guard in app_state.db.indexer.record_prefix(&prefix) { let (key, cid_bytes) = guard.into_inner().into_diagnostic()?; let rest = &key[prefix.len()..]; @@ -534,8 +545,8 @@ impl<'i> RepoHandle<'i> { let block_key = keys::block_key(collection, cid_bytes.as_ref()); let block_bytes = app_state .db - .indexer.blocks - .get(&block_key) + .indexer + .block(&block_key) .into_diagnostic()? .ok_or_else(|| miette::miette!("block missing for record {mst_key}"))?; diff --git a/src/control/stats.rs b/src/control/stats.rs index 5b0fcf7..56978e3 100644 --- a/src/control/stats.rs +++ b/src/control/stats.rs @@ -51,7 +51,7 @@ impl Hydrant { .collect(); #[cfg(feature = "indexer_stream")] - counts.insert("events", state.db.stream.events.approximate_len() as u64); + counts.insert("events", state.db.stream.approximate_event_count() as u64); #[cfg(feature = "relay")] counts.insert( diff --git a/src/control/stream/indexer.rs b/src/control/stream/indexer.rs index 281c03d..87b6dce 100644 --- a/src/control/stream/indexer.rs +++ b/src/control/stream/indexer.rs @@ -24,13 +24,21 @@ pub(crate) fn event_stream_thread( ) { let db = &state.db; 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.stream.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.stream.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(); @@ -41,7 +49,7 @@ pub(crate) fn event_stream_thread( catch_up_target, opts, move |current_id, target, chunk_size| { - read_event_replay_chunk(&replay_state, &ks, current_id, target, chunk_size) + read_event_replay_chunk(&replay_state, current_id, target, chunk_size) }, move |event| broadcast_to_event(&state, event), ); @@ -49,7 +57,6 @@ pub(crate) fn event_stream_thread( fn read_event_replay_chunk( state: &AppState, - ks: &fjall::Keyspace, current_id: Option, target: u64, chunk_size: usize, @@ -68,7 +75,10 @@ fn read_event_replay_chunk( let mut exhausted = false; let max_scanned = chunk_size.saturating_mul(4).max(chunk_size); let mut scanned = 0usize; - let mut iter = ks.range(keys::event_key(start)..=keys::event_key(target)); + let mut iter = state + .db + .stream + .event_range(keys::event_key(start)..=keys::event_key(target)); while events.len() < chunk_size && scanned < max_scanned { let Some(item) = iter.next() else { @@ -157,8 +167,8 @@ pub(crate) fn stored_to_event( } else { let block = state .db - .indexer.blocks - .get(keys::block_key(collection.as_str(), &cid.to_bytes())); + .indexer + .block(keys::block_key(collection.as_str(), &cid.to_bytes())); match block { Ok(Some(bytes)) => { match serde_ipld_dagcbor::from_slice::(bytes.as_ref()) { diff --git a/src/control/stream/jetstream.rs b/src/control/stream/jetstream.rs index e30540a..dace9c4 100644 --- a/src/control/stream/jetstream.rs +++ b/src/control/stream/jetstream.rs @@ -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.stream.events.get(keys::event_key(*event_id)).ok()??; + let bytes = state.db.stream.event(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,8 @@ 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 +412,8 @@ 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 +444,8 @@ 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/db/indexer.rs b/src/db/indexer.rs index 12a60e6..9c3c469 100644 --- a/src/db/indexer.rs +++ b/src/db/indexer.rs @@ -31,6 +31,18 @@ impl Db { } } +pub(crate) fn delete_repo_records( + batch: &mut OwnedWriteBatch, + db: &Db, + did: &Did<'_>, +) -> Result<()> { + for guard in db.indexer.records.prefix(keys::record_prefix_did(did)) { + let (key, _) = guard.into_inner().into_diagnostic()?; + batch.remove(&db.indexer.records, key); + } + Ok(()) +} + pub fn set_record_count( batch: &mut OwnedWriteBatch, db: &Db, diff --git a/src/db/keyspaces.rs b/src/db/keyspaces.rs index c6a6ede..b9ce42f 100644 --- a/src/db/keyspaces.rs +++ b/src/db/keyspaces.rs @@ -51,9 +51,9 @@ impl OpenCx<'_> { #[cfg(feature = "indexer")] pub struct IndexerDb { /// maps `{DID}|{COL}|{RKey}` -> record CID - pub records: Ks, + pub(super) records: Ks, /// content-addressable storage of raw DAG-CBOR blocks - pub blocks: Ks, + pub(super) blocks: Ks, /// backfill queue of `{ID}` -> empty pub pending: Ks, /// per-repo resync/retry state @@ -76,12 +76,31 @@ impl IndexerDb { lifecycle_count_lock: Arc::new(std::sync::Mutex::new(())), }) } + pub(crate) fn record>(&self, key: K) -> fjall::Result> { + self.records.get(key) + } + + pub(crate) fn record_prefix>(&self, prefix: K) -> fjall::Iter { + self.records.prefix(prefix) + } + + pub(crate) fn record_range(&self, range: R) -> fjall::Iter + where + K: AsRef<[u8]>, + R: std::ops::RangeBounds, + { + self.records.range(range) + } + + pub(crate) fn block>(&self, key: K) -> fjall::Result> { + self.blocks.get(key) + } } #[cfg(feature = "indexer_stream")] pub struct StreamDb { /// maps `{ID}` (u64 BE) -> `StoredEvent`, the source for the json stream api - pub events: Ks, + pub(super) events: Ks, pub(crate) event_tx: tokio::sync::broadcast::Sender, pub next_event_id: Arc, } @@ -98,6 +117,31 @@ impl StreamDb { }) } + #[cfg(feature = "jetstream")] + pub(crate) fn event>(&self, key: K) -> fjall::Result> { + self.events.get(key) + } + + pub(crate) fn event_range(&self, range: R) -> fjall::Iter + where + K: AsRef<[u8]>, + R: std::ops::RangeBounds, + { + self.events.range(range) + } + + pub(crate) fn stage_event(&self, batch: &mut fjall::OwnedWriteBatch, key: K, value: V) + where + K: Into, + V: Into, + { + batch.insert(&self.events, key, value); + } + + pub(crate) fn approximate_event_count(&self) -> usize { + self.events.approximate_len() + } + /// resume event ids after the last stored event. pub(super) fn init(&self) -> Result<()> { let mut last_id = 0; diff --git a/src/db/mod.rs b/src/db/mod.rs index 8bad2d5..cff00aa 100644 --- a/src/db/mod.rs +++ b/src/db/mod.rs @@ -26,10 +26,12 @@ pub mod pds_meta; pub mod types; pub mod keyspaces; +mod open; pub mod registry; pub mod schema; -mod open; mod train; +#[cfg(feature = "indexer")] +mod txn; #[cfg(feature = "indexer")] pub use keyspaces::IndexerDb; @@ -44,11 +46,11 @@ pub use schema::Ks; #[cfg(feature = "indexer")] pub(crate) use lifecycle_counts::LifecycleCountBatch; +#[cfg(feature = "indexer")] +pub(crate) use txn::Txn; use tracing::error; - - pub struct Db { pub inner: Arc, pub path: std::path::PathBuf, @@ -67,7 +69,7 @@ pub struct Db { #[cfg(feature = "relay")] pub(crate) relay: RelayDb, #[cfg(feature = "backlinks")] - pub backlinks: Ks, + pub(crate) backlinks: Ks, pub counts_map: HashMap, next_count_delta_id: Arc, count_delta_checkpoint_watermark: Arc, @@ -123,9 +125,7 @@ impl Db { pub async fn get(ks: Keyspace, key: impl Into) -> Result> { let key = key.into(); tokio::task::spawn_blocking(move || { - ks.get(key) - .inspect_err(check_poisoned) - .into_diagnostic() + ks.get(key).inspect_err(check_poisoned).into_diagnostic() }) .await .into_diagnostic()? diff --git a/src/db/txn.rs b/src/db/txn.rs new file mode 100644 index 0000000..3a2c45f --- /dev/null +++ b/src/db/txn.rs @@ -0,0 +1,243 @@ +use std::collections::HashMap; + +use bytes::Bytes; +use jacquard_common::IntoStatic; +use jacquard_common::types::cid::IpldCid; +use jacquard_common::types::did::Did; +use miette::{IntoDiagnostic, Result}; + +use crate::db::types::{DbAction, DbRkey, DbTid}; +use crate::db::{CountDeltas, Db, keys}; +use crate::ops::record_events::{EmitOp, RecordEmitter, RecordEventOrigin, RecordEvents}; +use crate::state::AppState; +#[cfg(feature = "indexer")] +use crate::types::GaugeState; +use crate::types::RepoState; + +/// one atomic database write, including its in-memory count projections. +/// +/// count deltas are persisted in the same fjall batch and applied in memory only +/// after that batch commits. dropping a transaction discards both. +pub(crate) struct Txn<'db> { + pub(crate) batch: fjall::OwnedWriteBatch, + pub(crate) db: &'db Db, + pub(crate) counts: CountDeltas, + #[cfg(feature = "indexer")] + lifecycle_transitions: Vec<(Did<'static>, GaugeState)>, +} + +impl<'db> Txn<'db> { + pub(crate) fn new(db: &'db Db) -> Self { + Self { + batch: db.inner.batch(), + db, + counts: CountDeltas::default(), + #[cfg(feature = "indexer")] + lifecycle_transitions: Vec::new(), + } + } + + pub(crate) fn records<'txn, 'did, 'repo>( + &'txn mut self, + state: &AppState, + commit_rev: &DbTid, + did: &'did Did<'repo>, + ) -> RecordTxn<'txn, 'db, 'did, 'repo> { + self.record_scope(state, commit_rev, did, RecordEventOrigin::Live, false) + } + + pub(crate) fn backfill_records<'txn, 'did, 'repo>( + &'txn mut self, + state: &AppState, + commit_rev: &DbTid, + did: &'did Did<'repo>, + ) -> RecordTxn<'txn, 'db, 'did, 'repo> { + self.record_scope(state, commit_rev, did, RecordEventOrigin::Backfill, true) + } + + fn record_scope<'txn, 'did, 'repo>( + &'txn mut self, + state: &AppState, + commit_rev: &DbTid, + did: &'did Did<'repo>, + origin: RecordEventOrigin, + count_ephemeral_records: bool, + ) -> RecordTxn<'txn, 'db, 'did, 'repo> { + RecordTxn { + txn: self, + emitter: RecordEmitter::new(state, commit_rev, origin), + did, + ephemeral: state.ephemeral, + only_index_links: state.only_index_links, + count_ephemeral_records, + records_delta: 0, + blocks_count: 0, + collection_deltas: HashMap::new(), + } + } + + #[cfg(feature = "indexer")] + pub(crate) fn transition_lifecycle(&mut self, did: &Did<'_>, gauge: GaugeState) { + self.lifecycle_transitions + .push((did.clone().into_static(), gauge)); + } + + pub(crate) fn commit(mut self) -> Result<()> { + #[cfg(feature = "indexer")] + let lifecycle_reservation = { + let mut lifecycle = self.db.lifecycle_counts(); + for (did, gauge) in self.lifecycle_transitions { + lifecycle.transition(&mut self.batch, &did, gauge)?; + } + lifecycle.stage(&mut self.batch) + }; + let count_reservation = self.db.stage_count_deltas(&mut self.batch, &self.counts); + + self.batch.commit().into_diagnostic()?; + + self.db.apply_count_deltas(&self.counts); + drop(count_reservation); + #[cfg(feature = "indexer")] + self.db.apply_lifecycle_counts(lifecycle_reservation); + Ok(()) + } +} + +/// one commit's record mutations within a larger atomic transaction. +pub(crate) struct RecordTxn<'txn, 'db, 'did, 'repo> { + txn: &'txn mut Txn<'db>, + emitter: RecordEmitter, + did: &'did Did<'repo>, + ephemeral: bool, + only_index_links: bool, + count_ephemeral_records: bool, + records_delta: i64, + blocks_count: i64, + collection_deltas: HashMap, +} + +impl RecordTxn<'_, '_, '_, '_> { + pub(crate) fn put_record( + &mut self, + collection: &str, + rkey: &DbRkey, + cid: IpldCid, + block: &Bytes, + action: DbAction, + ) -> Result<()> { + self.blocks_count += 1; + if action == DbAction::Create && (!self.ephemeral || self.count_ephemeral_records) { + self.records_delta += 1; + } + + #[cfg(feature = "indexer")] + if !self.ephemeral { + let cid_bytes = cid.to_bytes(); + if !self.only_index_links { + self.txn.batch.insert( + &self.txn.db.indexer.blocks, + keys::block_key(collection, &cid_bytes), + block.as_ref(), + ); + } + self.txn.batch.insert( + &self.txn.db.indexer.records, + keys::record_key(self.did, collection, rkey), + cid_bytes, + ); + if action == DbAction::Create { + *self + .collection_deltas + .entry(collection.to_owned()) + .or_default() += 1; + } + crate::ops::backlink_ops::index_record( + &mut self.txn.batch, + self.txn.db, + self.did, + collection, + &rkey.to_smolstr(), + block, + )?; + } + + self.emitter.emit( + &mut self.txn.batch, + self.txn.db, + EmitOp { + did: self.did, + collection, + rkey, + action, + cid: Some(cid), + block: Some(block), + }, + ) + } + + pub(crate) fn delete_record(&mut self, collection: &str, rkey: &DbRkey) -> Result<()> { + if !self.ephemeral || self.count_ephemeral_records { + self.records_delta -= 1; + } + + #[cfg(feature = "indexer")] + if !self.ephemeral { + self.txn.batch.remove( + &self.txn.db.indexer.records, + keys::record_key(self.did, collection, rkey), + ); + *self + .collection_deltas + .entry(collection.to_owned()) + .or_default() -= 1; + crate::ops::backlink_ops::delete_record( + &mut self.txn.batch, + self.txn.db, + self.did, + collection, + &rkey.to_smolstr(), + )?; + } + + self.emitter.emit( + &mut self.txn.batch, + self.txn.db, + EmitOp { + did: self.did, + collection, + rkey, + action: DbAction::Delete, + cid: None, + block: None, + }, + ) + } + + pub(crate) fn update_repo_state(&mut self, state: &RepoState<'_>) -> Result<()> { + self.txn.batch.insert( + &self.txn.db.repos, + keys::repo_key(self.did), + crate::db::ser_repo_state(state)?, + ); + Ok(()) + } + + pub(crate) fn finish(self) -> Result { + #[cfg(feature = "indexer")] + if !self.ephemeral { + for (collection, delta) in &self.collection_deltas { + crate::db::update_record_count( + &mut self.txn.batch, + self.txn.db, + self.did, + collection, + *delta, + )?; + } + } + self.txn.counts.add_records(self.records_delta); + self.txn.counts.add_blocks(self.blocks_count); + + Ok(self.emitter.finish()) + } +} diff --git a/src/ingest/indexer/shard.rs b/src/ingest/indexer/shard.rs index 926f1d3..f30311e 100644 --- a/src/ingest/indexer/shard.rs +++ b/src/ingest/indexer/shard.rs @@ -1,4 +1,3 @@ -use fjall::OwnedWriteBatch; use jacquard_common::IntoStatic; use jacquard_common::types::did::Did; use miette::{IntoDiagnostic, Result}; @@ -7,7 +6,7 @@ use std::sync::atomic::Ordering::SeqCst; use tokio::runtime::Handle as TokioHandle; use tracing::{debug, error, warn}; -use crate::db::{self, CountDeltas, keys, ser_repo_meta}; +use crate::db::{self, Txn, keys, ser_repo_meta}; use crate::ingest::stream::types::AccountStatus; use crate::ingest::stream::{Account, Commit, Identity}; use crate::ingest::validation; @@ -35,11 +34,8 @@ use super::worker::{FirehoseWorker, IngestError, RepoProcessResult}; struct WorkerContext<'a> { state: &'a AppState, - batch: OwnedWriteBatch, + txn: Txn<'a>, lifecycle_transitions: &'a mut Vec<(Did<'static>, GaugeState)>, - added_blocks: &'a mut i64, - records_delta: &'a mut i64, - count_deltas: &'a mut CountDeltas, #[cfg(feature = "indexer_stream")] broadcast_events: &'a mut Vec, #[cfg(feature = "jetstream")] @@ -58,24 +54,17 @@ impl FirehoseWorker { let mut lifecycle_transitions = Vec::new(); while let Some(msg) = rx.blocking_recv() { - let batch = state.db.inner.batch(); + let txn = Txn::new(&state.db); #[cfg(feature = "indexer_stream")] broadcast_events.clear(); #[cfg(feature = "jetstream")] jetstream_events.clear(); lifecycle_transitions.clear(); - let mut added_blocks = 0; - let mut records_delta = 0; - let mut count_deltas = CountDeltas::default(); - let mut ctx = WorkerContext { state: &state, - batch, + txn, lifecycle_transitions: &mut lifecycle_transitions, - added_blocks: &mut added_blocks, - records_delta: &mut records_delta, - count_deltas: &mut count_deltas, #[cfg(feature = "indexer_stream")] broadcast_events: &mut broadcast_events, #[cfg(feature = "jetstream")] @@ -96,7 +85,7 @@ impl FirehoseWorker { match Self::drain_resync_buffer(&mut ctx, &did, repo_state) { Ok(RepoProcessResult::Ok(s)) => { let res = ops::transition_repo( - &mut ctx.batch, + &mut ctx.txn.batch, &state.db, ctx.lifecycle_transitions, &did, @@ -205,7 +194,7 @@ impl FirehoseWorker { ) { Ok(RepoProcessResult::Ok(_)) => {} Ok(RepoProcessResult::Deleted) => { - ctx.count_deltas.add_repos(-1); + ctx.txn.counts.add_repos(-1); } Ok(RepoProcessResult::NeedsBackfill(Some(commit))) => { try_persist(commit); @@ -255,36 +244,13 @@ impl FirehoseWorker { } } - let mut batch = ctx.batch; - if added_blocks > 0 { - count_deltas.add_blocks(added_blocks); - } - if records_delta != 0 { - count_deltas.add_records(records_delta); - } - let mut lifecycle_counts = state.db.lifecycle_counts(); - let mut lifecycle_failed = false; - for (did, gauge) in lifecycle_transitions.drain(..) { - if let Err(e) = lifecycle_counts.transition(&mut batch, &did, gauge) { - error!(did = %did, err = %e, "failed to stage lifecycle transition"); - lifecycle_failed = true; - break; - } - } - if lifecycle_failed { - continue; + for (did, gauge) in ctx.lifecycle_transitions.drain(..) { + ctx.txn.transition_lifecycle(&did, gauge); } - let reservation = state.db.stage_count_deltas(&mut batch, &count_deltas); - let lifecycle_reservation = lifecycle_counts.stage(&mut batch); - if let Err(e) = batch.commit() { - error!(shard = id, err = %e, "failed to commit batch"); - drop(lifecycle_reservation); - drop(reservation); + if let Err(e) = ctx.txn.commit() { + error!(shard = id, err = %e, "failed to commit transaction"); continue; } - state.db.apply_lifecycle_counts(lifecycle_reservation); - state.db.apply_count_deltas(&count_deltas); - drop(reservation); #[cfg(feature = "indexer_stream")] for evt in broadcast_events.drain(..) { let _ = state.db.stream.event_tx.send(evt); @@ -337,7 +303,8 @@ 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.indexer.pending + db.indexer + .pending .get(keys::pending_key(metadata.index_id)) .into_diagnostic()? .is_some() @@ -372,15 +339,13 @@ impl FirehoseWorker { }; let res = ops::apply_commit( - &mut ctx.batch, + &mut ctx.txn, ctx.state, repo_state, validated, &ctx.state.filter.load(), )?; let repo_state = res.repo_state; - *ctx.added_blocks += res.blocks_count; - *ctx.records_delta += res.records_delta; #[cfg(feature = "indexer_stream")] { use std::sync::Arc; @@ -414,7 +379,7 @@ impl FirehoseWorker { let did = &identity.did; #[cfg(feature = "jetstream")] ctx.jetstream_events.push(crate::jetstream::stage_event( - &mut ctx.batch, + &mut ctx.txn.batch, db, StoredJetstreamEvent::Identity { did: TrimmedDid::from(did).into_static(), @@ -455,7 +420,7 @@ impl FirehoseWorker { }; #[cfg(feature = "jetstream")] ctx.jetstream_events.push(crate::jetstream::stage_event( - &mut ctx.batch, + &mut ctx.txn.batch, db, StoredJetstreamEvent::Account { did: TrimmedDid::from(did).into_static(), @@ -473,7 +438,7 @@ impl FirehoseWorker { debug!("account deleted, wiping data"); ctx.lifecycle_transitions .push((did.clone().into_static(), GaugeState::Synced)); - crate::ops::delete_repo(&mut ctx.batch, db, did, &repo_state)?; + crate::ops::delete_repo(&mut ctx.txn.batch, db, did, &repo_state)?; return Ok(RepoProcessResult::Deleted); } _ => { @@ -525,14 +490,14 @@ impl FirehoseWorker { Ok(r) => r, Err(e) => { if !Self::check_if_retriable_failure(&e) { - ctx.batch.remove(&db.indexer.resync_buffer, key); + ctx.txn.batch.remove(&db.indexer.resync_buffer, key); } return Err(e); } }; match res { RepoProcessResult::Ok(rs) => { - ctx.batch.remove(&db.indexer.resync_buffer, key); + ctx.txn.batch.remove(&db.indexer.resync_buffer, key); repo_state = rs; } RepoProcessResult::NeedsBackfill(_) => { @@ -540,7 +505,7 @@ impl FirehoseWorker { return Ok(RepoProcessResult::NeedsBackfill(None)); } RepoProcessResult::Deleted => { - ctx.batch.remove(&db.indexer.resync_buffer, key); + ctx.txn.batch.remove(&db.indexer.resync_buffer, key); return Ok(RepoProcessResult::Deleted); } } @@ -575,7 +540,8 @@ impl FirehoseWorker { let old_pkey = keys::pending_key(metadata.index_id); let was_pending = had_metadata && db - .indexer.pending + .indexer + .pending .get(old_pkey.as_slice()) .into_diagnostic()? .is_some(); @@ -586,7 +552,11 @@ impl FirehoseWorker { } metadata.index_id = rand::random::(); - batch.insert(&db.indexer.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/ops.rs b/src/ops.rs index 6ad9d06..8744213 100644 --- a/src/ops.rs +++ b/src/ops.rs @@ -1,10 +1,8 @@ use fjall::OwnedWriteBatch; -use fjall::Slice; use jacquard_common::IntoStatic; use jacquard_common::types::did::Did; use miette::{Context, IntoDiagnostic, Result}; -use std::collections::HashMap; use tracing::debug; use crate::db::types::{DbAction, DbRkey, DbTid, TrimmedDid}; @@ -20,16 +18,18 @@ use { crate::types::{AccountEvt, BroadcastEvent, IdentityEvt, MarshallableEvt}, std::sync::atomic::Ordering, }; +pub(crate) mod backlink_ops; +use record_events::RecordEvents; -mod backlink_ops; -pub(crate) mod record_events; - -use record_events::{EmitOp, RecordEmitter, RecordEvents}; +pub mod record_events; 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.indexer.resync_buffer.insert(key, value).into_diagnostic()?; + db.indexer + .resync_buffer + .insert(key, value) + .into_diagnostic()?; debug!( did = %did, seq = commit.seq, @@ -98,11 +98,7 @@ pub fn delete_repo( // 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.indexer.records.prefix(&records_prefix) { - let (k, _cid_bytes) = guard.into_inner().into_diagnostic()?; - batch.remove(&db.indexer.records, k); - } + db::delete_repo_records(batch, db, did)?; // 4. reset collection counts let mut count_prefix = Vec::new(); @@ -216,22 +212,17 @@ pub fn transition_repo<'s>( pub struct ApplyCommitResults<'s> { pub repo_state: RepoState<'s>, - pub records_delta: i64, - pub blocks_count: i64, #[cfg_attr(not(feature = "indexer_stream"), allow(dead_code))] pub events: RecordEvents, } pub fn apply_commit<'s>( - batch: &mut OwnedWriteBatch, + txn: &mut crate::db::Txn<'_>, state: &AppState, mut repo_state: RepoState<'s>, validated: ValidatedCommit<'_>, filter: &FilterConfig, ) -> Result> { - let db = &state.db; - let ephemeral = state.ephemeral; - let only_index_links = state.only_index_links; let commit = validated.commit; let parsed = validated.parsed_blocks; let did = &commit.repo; @@ -240,11 +231,7 @@ pub fn apply_commit<'s>( repo_state.root = Some(validated.commit_obj.into()); repo_state.touch(); - // 2. iterate ops and update records index - let mut records_delta = 0; - let mut blocks_count = 0; - let mut collection_deltas: HashMap<&str, i64> = HashMap::new(); - let mut emitter = RecordEmitter::new(state, &commit); + let mut record_txn = txn.records(state, &DbTid::from(&commit.rev), did); for op in &commit.ops { let (collection, rkey) = parse_path(&op.path)?; @@ -254,13 +241,8 @@ pub fn apply_commit<'s>( } let rkey = DbRkey::new(rkey); - let db_key = keys::record_key(did, collection, &rkey); - let action = DbAction::try_from(op.action.as_str())?; - let mut event_cid = None; - let mut event_block = None; - match action { DbAction::Create | DbAction::Update => { let Some(cid) = &op.cid else { @@ -270,77 +252,23 @@ pub fn apply_commit<'s>( .to_ipld() .into_diagnostic() .wrap_err("expected valid cid from relay")?; - event_cid = Some(cid_ipld); let Some(bytes) = parsed.blocks.get(&cid_ipld) else { return Err(miette::miette!( "block {cid} not found in CAR for record {did}/{collection}/{rkey}" )); }; - event_block = Some(bytes); - let cid_raw = cid_ipld.to_bytes(); - let block_key = Slice::from(keys::block_key(collection, &cid_raw)); - - blocks_count += 1; - if !ephemeral { - if !only_index_links { - batch.insert(&db.indexer.blocks, block_key.clone(), bytes.as_ref()); - } - batch.insert(&db.indexer.records, db_key.clone(), cid_raw); - // accumulate counts - if action == DbAction::Create { - records_delta += 1; - *collection_deltas.entry(collection).or_default() += 1; - } - backlink_ops::index_record( - batch, - db, - did, - collection, - &rkey.to_smolstr(), - bytes.as_ref(), - )?; - } + record_txn.put_record(collection, &rkey, cid_ipld, bytes, action)?; } DbAction::Delete => { - if !ephemeral { - batch.remove(&db.indexer.records, db_key); - - // accumulate counts - records_delta -= 1; - *collection_deltas.entry(collection).or_default() -= 1; - - backlink_ops::delete_record(batch, db, did, collection, &rkey.to_smolstr())?; - } + record_txn.delete_record(collection, &rkey)?; } - }; - - emitter.emit( - batch, - db, - EmitOp { - did, - collection, - rkey: &rkey, - action, - cid: event_cid, - block: event_block, - }, - )?; - } - - // update counts - if !ephemeral { - for (col, delta) in collection_deltas { - db::update_record_count(batch, db, did, col, delta)?; } } - Ok(ApplyCommitResults { - repo_state, - records_delta, - blocks_count, - events: emitter.finish(), - }) + record_txn.update_repo_state(&repo_state)?; + let events = record_txn.finish()?; + + Ok(ApplyCommitResults { repo_state, events }) } pub fn parse_path(path: &str) -> Result<(&str, &str)> { diff --git a/src/ops/record_events.rs b/src/ops/record_events.rs index 20de4e0..e8596a8 100644 --- a/src/ops/record_events.rs +++ b/src/ops/record_events.rs @@ -17,7 +17,6 @@ mod enabled { use crate::db::types::{DbAction, DbRkey, DbTid, TrimmedDid}; use crate::db::{Db, keys}; - use crate::ingest::stream::Commit; use crate::state::AppState; use crate::types::{LiveRecordEvent, StoredData, StoredEvent}; @@ -33,10 +32,17 @@ mod enabled { pub(crate) block: Option<&'a Bytes>, } + #[derive(Clone, Copy)] + pub(crate) enum RecordEventOrigin { + Live, + Backfill, + } + pub(crate) struct RecordEmitter { rev: DbTid, ephemeral: bool, only_index_links: bool, + live: bool, should_broadcast_live: bool, live_events: Vec, last_event_id: Option, @@ -47,16 +53,18 @@ mod enabled { } impl RecordEmitter { - pub(crate) fn new(state: &AppState, commit: &Commit) -> Self { + pub(crate) fn new(state: &AppState, commit_rev: &DbTid, origin: RecordEventOrigin) -> Self { + let live = matches!(origin, RecordEventOrigin::Live); Self { - rev: DbTid::from(&commit.rev), + rev: *commit_rev, ephemeral: state.ephemeral, only_index_links: state.only_index_links, - should_broadcast_live: state.db.stream.event_tx.receiver_count() > 0, + live, + should_broadcast_live: live && state.db.stream.event_tx.receiver_count() > 0, live_events: Vec::new(), last_event_id: None, #[cfg(feature = "jetstream")] - should_stage_jetstream: state.db.jetstream.tx.receiver_count() > 0, + should_stage_jetstream: live && state.db.jetstream.tx.receiver_count() > 0, #[cfg(feature = "jetstream")] jetstream_events: Vec::new(), } @@ -69,15 +77,12 @@ mod enabled { op: EmitOp<'_>, ) -> Result<()> { // in ephemeral mode, the event payload is the only place we persist the record. - let block_inline_for_event = (self.ephemeral) - .then(|| op.block.cloned()) - .flatten(); + let block_inline_for_event = (self.ephemeral).then(|| op.block.cloned()).flatten(); // inline record bytes for live tailing so we don't have to load from blocks. - let inline_block = (!self.ephemeral - && self.should_broadcast_live - && !self.only_index_links) - .then(|| op.block.cloned()) - .flatten(); + let inline_block = + (!self.ephemeral && self.should_broadcast_live && !self.only_index_links) + .then(|| op.block.cloned()) + .flatten(); let data = block_inline_for_event .map(StoredData::Block) @@ -94,7 +99,7 @@ mod enabled { let collection = CowStr::Borrowed(op.collection); let evt = StoredEvent { - live: true, + live: self.live, did: did_trimmed.clone(), rev: self.rev, collection: collection.clone(), @@ -103,7 +108,8 @@ mod enabled { data: data.clone(), }; let bytes = rmp_serde::to_vec(&evt).into_diagnostic()?; - batch.insert(&db.stream.events, keys::event_key(event_id), bytes); + db.stream + .stage_event(batch, keys::event_key(event_id), bytes); #[cfg(feature = "jetstream")] { @@ -111,7 +117,7 @@ mod enabled { did: did_trimmed.clone().into_static(), collection: collection.clone().into_static(), event_id, - live: true, + live: self.live, }; let ephemeral = self .should_stage_jetstream @@ -124,13 +130,14 @@ mod enabled { op.rkey.to_smolstr().as_str(), &data, inline_block.as_ref(), - true, + self.live, ) }) .flatten(); - self.jetstream_events.push(crate::jetstream::stage_event( - batch, db, jetstream, ephemeral, - )?); + let broadcast = crate::jetstream::stage_event(batch, db, jetstream, ephemeral)?; + if self.live { + self.jetstream_events.push(broadcast); + } } if self.should_broadcast_live { @@ -184,9 +191,14 @@ mod noop { use crate::db::Db; use crate::db::types::{DbAction, DbRkey}; - use crate::ingest::stream::Commit; use crate::state::AppState; + #[derive(Clone, Copy)] + pub(crate) enum RecordEventOrigin { + Live, + Backfill, + } + #[allow(dead_code)] pub(crate) struct EmitOp<'a> { pub(crate) did: &'a Did<'a>, @@ -201,7 +213,11 @@ mod noop { impl RecordEmitter { #[inline(always)] - pub(crate) fn new(_state: &AppState, _commit: &Commit) -> Self { + pub(crate) fn new( + _state: &AppState, + _commit_rev: &crate::db::types::DbTid, + _origin: RecordEventOrigin, + ) -> Self { Self } -- 2.51.2