From 9c1e0f5ffcfa0604ccd148c0ffdef9f05e189a6a Mon Sep 17 00:00:00 2001 From: dawn <90008@klbr.net> Date: Tue, 6 Oct 2026 00:52:16 +0300 Subject: [PATCH] [db] count from the row a txn already read instead of reading it again put_repo read and decoded the row it replaces on every write, so apply_commit paid a second point read and msgpack decode per live commit for a row the shard had read under the same lock moments before, only to find its active and pds unchanged. Txn::read_repo now takes the repo lock, reads and decodes the row and keeps the pds it counts for, and a later put_repo of it in the same txn moves the count from that without a read. the shard's commit load, the relay's indexer branches, the crawler's reconcile and the backfill's inactive path read through it. put_repo also compares the pds as stored first, so the usual write that keeps it parses no url. --- src/backfill/worker/process.rs | 10 ++-- src/crawler/worker.rs | 3 +- src/db/counts.rs | 15 +++--- src/db/txn.rs | 93 +++++++++++++++++++++++++++++----- src/ingest/indexer/shard.rs | 23 +++------ src/ingest/relay/context.rs | 17 ++----- 6 files changed, 108 insertions(+), 53 deletions(-) diff --git a/src/backfill/worker/process.rs b/src/backfill/worker/process.rs index 932da0b..fe16848 100644 --- a/src/backfill/worker/process.rs +++ b/src/backfill/worker/process.rs @@ -258,12 +258,14 @@ pub(super) async fn process_did( db.indexer .stage_pending_remove(&mut txn.batch, &did, &pending_key); if applied { - let did_key = keys::repo_key(&did); - if let Some(bytes) = db.repos.get(&did_key).into_diagnostic()? { - let mut state = crate::db::deser_repo_state(&bytes)?.into_static(); + if let Some(mut state) = txn.read_repo(&did)? { state.active = false; state.status = status; - txn.batch.insert(&db.indexer.resync, &did_key, resync_bytes); + txn.batch.insert( + &db.indexer.resync, + keys::repo_key(&did), + resync_bytes, + ); txn.put_repo(&did, &state)?; } #[cfg(feature = "indexer_stream")] diff --git a/src/crawler/worker.rs b/src/crawler/worker.rs index 86c45cf..b331fd1 100644 --- a/src/crawler/worker.rs +++ b/src/crawler/worker.rs @@ -336,7 +336,7 @@ fn reconcile_listing( return Ok(ReconcileOutcome::Unchanged); } let did_key = keys::repo_key(&listing.did); - let Some(state_bytes) = db.repos.get(&did_key).into_diagnostic()? else { + let Some(mut state) = txn.read_repo(&listing.did)? else { if listing.active || !listing.create_if_missing { return Ok(ReconcileOutcome::Unchanged); } @@ -359,7 +359,6 @@ fn reconcile_listing( txn.counts.add_repos(1); return Ok(ReconcileOutcome::StateOnly); }; - let mut state = crate::db::deser_repo_state(&state_bytes)?.into_static(); if listing.active { let listing_is_newer = state diff --git a/src/db/counts.rs b/src/db/counts.rs index 0ce57b6..7e3e236 100644 --- a/src/db/counts.rs +++ b/src/db/counts.rs @@ -119,12 +119,15 @@ impl CountDeltas { } } -/// the pds host a repo row counts one account for, which is what the v6 rebuild counts too -pub(crate) fn pds_account_host(row: &RepoState<'_>) -> Option { - row.pds - .as_deref() - .filter(|_| row.active) - .and_then(|pds| Url::parse(pds).ok()) +/// the pds a repo row counts one account for, `None` while it's inactive. the count is kept +/// per [`pds_host`] of it, which is what the v6 rebuild counts too +pub(crate) fn counted_pds<'r>(row: &'r RepoState<'_>) -> Option<&'r str> { + row.pds.as_deref().filter(|_| row.active) +} + +pub(crate) fn pds_host(pds: &str) -> Option { + Url::parse(pds) + .ok() .and_then(|url| url.host_str().map(SmolStr::new)) } diff --git a/src/db/txn.rs b/src/db/txn.rs index e86471f..44a169c 100644 --- a/src/db/txn.rs +++ b/src/db/txn.rs @@ -1,6 +1,8 @@ use std::collections::HashMap; use std::time::Duration; +#[cfg(feature = "indexer")] +use jacquard_common::IntoStatic; use jacquard_common::types::did::Did; #[cfg(feature = "indexer")] use jacquard_common::types::string::Tid; @@ -9,7 +11,7 @@ use smol_str::SmolStr; #[cfg(feature = "indexer")] use crate::db::LifecycleCountBatch; -use crate::db::counts::{PdsAccountMove, pds_account_host}; +use crate::db::counts::{PdsAccountMove, counted_pds, pds_host}; use crate::db::outbox::{Assigned, Outbox}; #[cfg(feature = "indexer")] use crate::db::types::{DbAction, DbRkey, DbTid}; @@ -229,21 +231,39 @@ impl<'db> Txn<'db> { }) } + /// reads `did`'s repo row under its lock and keeps the pds it counts for, so a + /// [`Self::put_repo`] of it later in this transaction needn't read it again + #[cfg(feature = "indexer")] + pub(crate) fn read_repo(&mut self, did: &Did) -> Result>> { + self.hold_repo_write_lock(did); + let key = keys::repo_key(did); + let row = self + .db + .repos + .get(&key) + .into_diagnostic()? + .map(|bytes| crate::db::deser_repo_state(&bytes).map(IntoStatic::into_static)) + .transpose()?; + let counted = row.as_ref().and_then(counted_pds).map(SmolStr::new); + self.pds_accounts.insert(key, counted); + Ok(row) + } + /// puts `row` as `did`'s repo row. every repo row write goes through here or /// [`Self::remove_repo`], because the per-pds account counts move by the row a write /// replaces, and only the repo lock this holds keeps that row from changing until the commit pub(crate) fn put_repo(&mut self, did: &Did, row: &RepoState<'_>) -> Result { let key = keys::repo_key(did); - let moved = self.move_pds_account(did, &key, pds_account_host(row))?; - self.batch - .insert(&self.db.repos, key, crate::db::ser_repo_state(row)?); + let value = crate::db::ser_repo_state(row)?; + let moved = self.move_pds_account(did, key.clone(), counted_pds(row))?; + self.batch.insert(&self.db.repos, key, value); Ok(moved) } #[cfg(feature = "indexer")] pub(crate) fn remove_repo(&mut self, did: &Did) -> Result { let key = keys::repo_key(did); - let moved = self.move_pds_account(did, &key, None)?; + let moved = self.move_pds_account(did, key.clone(), None)?; self.batch.remove(&self.db.repos, key); Ok(moved) } @@ -252,24 +272,34 @@ impl<'db> Txn<'db> { fn move_pds_account( &mut self, did: &Did, - key: &[u8], - to: Option, + key: Vec, + to: Option<&str>, ) -> Result { // without the indexer there are no repo locks to take #[cfg(feature = "indexer")] self.hold_repo_write_lock(did); - let from = match self.pds_accounts.get(key) { - Some(staged) => staged.clone(), + let from = match self.pds_accounts.remove(&key) { + Some(staged) => staged, None => self .db .repos - .get(key) + .get(&key) .into_diagnostic()? - .map(|bytes| crate::db::deser_repo_state(&bytes).map(|row| pds_account_host(&row))) + .map(|bytes| { + crate::db::deser_repo_state(&bytes) + .map(|row| counted_pds(&row).map(SmolStr::new)) + }) .transpose()? .flatten(), }; - self.pds_accounts.insert(key.to_vec(), to.clone()); + // the same pds most writes keep, which needs no url parse + if from.as_deref() == to { + self.pds_accounts.insert(key, from); + return Ok(PdsAccountMove::default()); + } + self.pds_accounts.insert(key, to.map(SmolStr::new)); + let from = from.as_deref().and_then(pds_host); + let to = to.and_then(pds_host); if from == to { return Ok(PdsAccountMove::default()); } @@ -810,6 +840,45 @@ mod tests { assert_eq!(accounts("three.example"), 0); } + #[test] + fn a_put_after_a_read_counts_from_the_row_the_read_saw() { + let (_tmp, state) = test_state(); + let did = did(); + let accounts = |host| state.db.get_count_sync(&keys::pds_account_count_key(host)); + let row = |pds: &'static str| RepoState { + pds: Some(jacquard_common::CowStr::Borrowed(pds)), + ..RepoState::synced() + }; + let mut txn = Txn::new(&state.db); + txn.put_repo(&did, &row("https://one.example/")).unwrap(); + txn.commit().unwrap(); + + let mut txn = Txn::new(&state.db); + let read = txn.read_repo(&did).unwrap().expect("repo row"); + assert_eq!(read.pds.as_deref(), Some("https://one.example/")); + // nothing else can write the row under the lock the read took, so a put after it has + // to count from what the read saw. this stored row, written past the lock, shows the + // put didn't read the row again + state + .db + .repos + .insert( + keys::repo_key(&did), + crate::db::ser_repo_state(&RepoState::backfilling()).unwrap(), + ) + .unwrap(); + txn.put_repo(&did, &row("https://two.example/")).unwrap(); + txn.commit().unwrap(); + assert_eq!((accounts("one.example"), accounts("two.example")), (0, 1)); + + let other = Did::new_static("did:plc:aaaaaaaaaaaaaaaaaaaaaaaa").unwrap(); + let mut txn = Txn::new(&state.db); + assert!(txn.read_repo(&other).unwrap().is_none()); + txn.put_repo(&other, &row("https://two.example/")).unwrap(); + txn.commit().unwrap(); + assert_eq!(accounts("two.example"), 2); + } + #[test] fn create_writes_head_without_history() { let (_tmp, state) = test_state(); diff --git a/src/ingest/indexer/shard.rs b/src/ingest/indexer/shard.rs index 58059cb..1285433 100644 --- a/src/ingest/indexer/shard.rs +++ b/src/ingest/indexer/shard.rs @@ -134,25 +134,14 @@ impl FirehoseWorker { } } - let repo_bytes = { - let repo_key = keys::repo_key(did); - match state.db.repos.get(&repo_key).into_diagnostic() { - Ok(Some(b)) => b, - Ok(None) => { - checkpoint(); - continue; - } - Err(e) => { - error!(err = %e, "failed to get repo state"); - checkpoint(); - continue; - } + let repo_state = match ctx.txn.read_repo(did) { + Ok(Some(s)) => s, + Ok(None) => { + checkpoint(); + continue; } - }; - let repo_state = match crate::db::deser_repo_state(&repo_bytes) { - Ok(s) => s, Err(e) => { - error!(err = %e, "failed to deser repo state"); + error!(err = %e, "failed to get repo state"); checkpoint(); continue; } diff --git a/src/ingest/relay/context.rs b/src/ingest/relay/context.rs index f33b857..3dc77de 100644 --- a/src/ingest/relay/context.rs +++ b/src/ingest/relay/context.rs @@ -72,35 +72,28 @@ impl WorkerContext<'_> { let Some(staged) = self.staged.take() else { return Ok(()); }; - #[cfg(feature = "indexer")] - let key = keys::repo_key(did); let row = match loaded { #[cfg(feature = "indexer")] Some(before) => { - txn.hold_repo_write_lock(did); - let Some(bytes) = self.state.db.repos.get(&key).into_diagnostic()? else { + let Some(current) = txn.read_repo(did)? else { // erased while this message was in flight return Ok(()); }; - staged.rebased(&before, db::deser_repo_state(&bytes)?.into_static()) + staged.rebased(&before, current) } // the row `load_repo_state` made up for a repo it hadn't seen #[cfg(feature = "indexer")] None => { - txn.hold_repo_write_lock(did); - match self.state.db.repos.get(&key).into_diagnostic()? { + match txn.read_repo(did)? { None => staged, // a track, the crawler or a backfill made the row while this message // resolved the doc, and counted the repo already. what this message // learned goes onto that row as if it had started from a blank one. whoever // made it saw to its backfill, and a second NewRepo would requeue it - Some(bytes) => { + Some(current) => { txn.counts.add_repos(-1); self.sink.forget_new_repo(did); - staged.rebased( - &RepoState::backfilling(), - db::deser_repo_state(&bytes)?.into_static(), - ) + staged.rebased(&RepoState::backfilling(), current) } } } -- 2.51.2