From b45547fcc2bd0956ba0da9c69337da20fea80a0a Mon Sep 17 00:00:00 2001 From: dawn <90008@klbr.net> Date: Tue, 6 Oct 2026 00:24:49 +0300 Subject: [PATCH] [db] count a repo for its pds from the row each write replaces the per-pds account counts are meant to equal the repo rows that are active with that pds, which is what the v6 rebuild counts. only the relay moved them, and from the row it read, so every other writer of `active` or `pds` let them drift: a backfill sets the pds from the doc and can force a repo active, the backfill finish flips active from the status, the crawler's listings deactivate and reactivate rows, a failed fetch marks a repo inactive, the minidoc refresh rewrites the pds, and erasing or discarding a repo removes its row. eg a deactivation took one off, a backfill forced the repo active without adding it back, and the next active event saw it active already, so the host stayed one short for good. the relay also counted the row it made up or read rather than the rebased row it wrote, so a row made inactive meanwhile still counted once its probe failed. every repo row write now goes through Txn::put_repo or Txn::remove_repo. they take the repo lock, read the row the write replaces (or what the same txn already put), and move the count from that row's host to the new one, counting a row the same way the v6 rebuild does. the relay applies the host limit status for the hosts a write moved. AGENTS.md names the rule. --- AGENTS.md | 3 +- src/backfill/worker/process.rs | 22 +++-- src/backfill/worker/task.rs | 9 +-- src/control/repos/indexer.rs | 8 +- src/control/repos/mod.rs | 15 +--- src/crawler/worker.rs | 41 +++++++--- src/db/counts.rs | 25 ++++++ src/db/indexer.rs | 37 ++------- src/db/txn.rs | 141 ++++++++++++++++++++++++++------- src/ingest/indexer/shard.rs | 3 +- src/ingest/relay/context.rs | 22 +++-- src/ingest/relay/handlers.rs | 55 ------------- src/ingest/relay/worker.rs | 97 ++++++++++++++++++++++- src/ops.rs | 7 +- 14 files changed, 301 insertions(+), 184 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index 28aef2a..817aac4 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -104,6 +104,7 @@ Modes (`indexer`, `relay`) are compile-time choices and mutually exclusive (enfo - **State**: Use `rmp-serde` (MessagePack) for all internal state (`RepoState`, `ErrorState`, `StoredEvent`). - **Blocks**: Record bodies are stored inline in `heads`; superseded bodies move to `history`. v10 materializes unresolved legacy event bodies into `event_bodies`, moves bodies into the compressed `heads` keyspace, and deletes the old `blocks` and `records` keyspaces. - **Cursors**: Store cursors as big-endian bytes (`u64`/`i64`). +- **Repo rows**: Write `repos` rows only through `Txn::put_repo`/`Txn::remove_repo`. They move the per-pds account counts (`k|p|{host}`, the active rows with that pds) by the row the write replaces, read under the repo lock they hold, so a raw `batch.insert(&db.repos, ..)` leaves a host's count off for good. - **Compression**: Configurable via `HYDRANT_DATA_COMPRESSION` (`lz4`, `zstd`, `none`). Per-keyspace zstd dictionaries can be trained via `POST /db/train` and are stored as `dict_{keyspace}.bin` in the database directory. - **Keyspaces**: Use the `keys.rs` module to maintain consistent composite key formats. - **Stream positions**: `/stream` event ids, jetstream keys, and relay seqs are only assigned inside `db::sequencer::Sequencer::commit`. Stage stream events in the transaction's `outbox` (`txn.outbox.stream`, `.jetstream`, `.relay`) and let `Txn::commit` give them positions and broadcast them; never write event rows or broadcast stream events directly. This keeps position order equal to commit order, which the shared stream engine (`control::stream::engine`) relies on. @@ -129,7 +130,7 @@ Hydrant uses multiple `fjall` keyspaces: - `pending_dids`: The same queue in DID order, `{DID} 00 {ID}` -> empty, so `GET /repos?queued=true` can page it by DID next to `resync`. Only `IndexerDb::stage_pending_insert`/`stage_pending_remove` write either keyspace, always both in one batch. - `resync`: Maps `{DID}` -> `ResyncState` (MessagePack) for retry logic/tombstones. - `resync_buffer`: Maps `{DID}|{Rev}` -> `Commit` (MessagePack). Used to buffer live events during backfill. -- `counts`: Maps `k|{NAME}` or `r|{DID}|{COL}` -> `Count` (u64 BE Bytes). +- `counts`: Maps `k|{NAME}` or `r|{DID}|{COL}` -> `Count` (u64 BE Bytes). `k|p|{host}` counts the active repo rows whose pds is on that host. - `filter`: Stores filter config. Handled by the `db::filter` and `db::pds_meta` modules. Includes: - Mode key `m` -> `FilterMode` (MessagePack). - Storage mode marker `storage_mode` -> `StorageModeMarker` (MessagePack): the immutable `only_index_links`/`ephemeral` layout the database was adopted under, written on first open and enforced by `db::storage_mode`; a config flip against an existing database is rejected at startup (markerless legacy databases pass one inference check before adoption). diff --git a/src/backfill/worker/process.rs b/src/backfill/worker/process.rs index d1f6172..932da0b 100644 --- a/src/backfill/worker/process.rs +++ b/src/backfill/worker/process.rs @@ -258,17 +258,14 @@ pub(super) async fn process_did( db.indexer .stage_pending_remove(&mut txn.batch, &did, &pending_key); if applied { - crate::db::Db::update_repo_state( - &mut txn.batch, - &db.repos, - &did, - move |state, (key, batch)| { - state.active = false; - state.status = status; - batch.insert(&db.indexer.resync, key, resync_bytes); - Ok((true, ())) - }, - )?; + 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(); + state.active = false; + state.status = status; + txn.batch.insert(&db.indexer.resync, &did_key, resync_bytes); + txn.put_repo(&did, &state)?; + } #[cfg(feature = "indexer_stream")] txn.outbox.stream.push_account(account); } @@ -550,7 +547,6 @@ async fn remove_discarded_repo( did: &Did, pending_key: &Slice, ) -> Result<(), BackfillError> { - let did_key = keys::repo_key(did); let metadata_key = keys::repo_metadata_key(did); let did = did.clone(); let pending_key = pending_key.clone(); @@ -566,7 +562,7 @@ async fn remove_discarded_repo( db.indexer .stage_pending_remove(&mut txn.batch, &did, &pending_key); if applied { - txn.batch.remove(&db.repos, &did_key); + txn.remove_repo(&did)?; txn.batch.remove(&db.repo_metadata, &metadata_key); txn.counts.add_repos(-1); } diff --git a/src/backfill/worker/task.rs b/src/backfill/worker/task.rs index 71c603d..62a9ebf 100644 --- a/src/backfill/worker/task.rs +++ b/src/backfill/worker/task.rs @@ -332,11 +332,7 @@ pub(super) async fn did_task( if repo.active { repo.status = RepoStatus::Error(error_string.into()); } - txn.batch.insert( - &db.repos, - &did_key, - rmp_serde::to_vec(&repo).into_diagnostic()?, - ); + txn.put_repo(&did, &repo)?; } } txn.commit()?; @@ -786,8 +782,7 @@ mod tests { let mut txn = DbTxn::new(&self.state.db); crate::ops::delete_repo(&mut txn, &self.state.db, &self.did, repo)?; crate::ops::transition_repo( - &mut txn.batch, - &self.state.db, + &mut txn, &mut Vec::new(), &self.did, repo.clone(), diff --git a/src/control/repos/indexer.rs b/src/control/repos/indexer.rs index 7205d76..435c491 100644 --- a/src/control/repos/indexer.rs +++ b/src/control/repos/indexer.rs @@ -250,7 +250,6 @@ impl ReposControl { { continue; } - let did_key = keys::repo_key(&did); let metadata_key = keys::repo_metadata_key(&did); let metadata_bytes = db.repo_metadata.get(&metadata_key).into_diagnostic()?; @@ -263,13 +262,8 @@ impl ReposControl { queued.push(did); } } else { - let repo_state = RepoState::backfilling(); let metadata = RepoMetadata::backfilling(rand::random()); - txn.batch.insert( - &db.repos, - &did_key, - crate::db::ser_repo_state(&repo_state)?, - ); + txn.put_repo(&did, &RepoState::backfilling())?; txn.batch.insert( &db.repo_metadata, &metadata_key, diff --git a/src/control/repos/mod.rs b/src/control/repos/mod.rs index acaead1..21d4fd8 100644 --- a/src/control/repos/mod.rs +++ b/src/control/repos/mod.rs @@ -906,12 +906,9 @@ impl RepoHandle { .state .db .run(move |db| { - #[cfg(feature = "indexer")] let mut txn = crate::db::Txn::new(db); #[cfg(feature = "indexer")] txn.hold_repo_write_lock(&did_for_db); - #[cfg(not(feature = "indexer"))] - let mut batch = db.inner.batch(); let Some(MiniDocSnapshot { mut state, @@ -921,7 +918,6 @@ impl RepoHandle { else { return Ok(None); }; - let did_key = keys::repo_key(&did_for_db); let raced = state.identity() != initial_identity; let current_trigger = @@ -937,14 +933,8 @@ impl RepoHandle { false }; - #[cfg(feature = "indexer")] - if identity_changed { - txn.batch - .insert(&db.repos, &did_key, crate::db::ser_repo_state(&state)?); - } - #[cfg(not(feature = "indexer"))] if identity_changed { - batch.insert(&db.repos, &did_key, crate::db::ser_repo_state(&state)?); + txn.put_repo(&did_for_db, &state)?; } // a DID document never revives lifecycle state. gone recovery @@ -961,10 +951,7 @@ impl RepoHandle { let queued_backfill = false; if identity_changed || queued_backfill { - #[cfg(feature = "indexer")] txn.commit()?; - #[cfg(not(feature = "indexer"))] - batch.commit().into_diagnostic()?; } Ok(Some(MiniDocReconcile { diff --git a/src/crawler/worker.rs b/src/crawler/worker.rs index 7321247..86c45cf 100644 --- a/src/crawler/worker.rs +++ b/src/crawler/worker.rs @@ -220,10 +220,8 @@ impl CrawlerWorker { if db.repos.contains_key(&did_key).into_diagnostic()? { continue; } - let state = RepoState::backfilling(); let metadata = RepoMetadata::backfilling(rng.next_u64()); - txn.batch - .insert(&db.repos, &did_key, ser_repo_state(&state)?); + txn.put_repo(&guard, &RepoState::backfilling())?; txn.batch.insert( &db.repo_metadata, &metadata_key, @@ -357,8 +355,7 @@ fn reconcile_listing( ); txn.transition_lifecycle(&listing.did, GaugeState::Resync(None))?; } - txn.batch - .insert(&db.repos, &did_key, ser_repo_state(&state)?); + txn.put_repo(&listing.did, &state)?; txn.counts.add_repos(1); return Ok(ReconcileOutcome::StateOnly); }; @@ -388,8 +385,7 @@ fn reconcile_listing( } state.active = true; state.touch(); - txn.batch - .insert(&db.repos, did_key, ser_repo_state(&state)?); + txn.put_repo(&listing.did, &state)?; return Ok(if was_active { ReconcileOutcome::StateOnly } else { @@ -440,8 +436,7 @@ fn reconcile_listing( state.active = false; state.status = status; state.touch(); - txn.batch - .insert(&db.repos, did_key, ser_repo_state(&state)?); + txn.put_repo(&listing.did, &state)?; Ok(if account_changed { ReconcileOutcome::AccountChanged } else { @@ -621,6 +616,34 @@ mod tests { Ok(crate::db::deser_repo_state(&bytes)?.into_static()) } + #[test] + fn listings_move_the_repo_between_counted_and_not() -> Result<()> { + let (_tmp, db) = open_db()?; + let mut state = state_with_root(rev(2)); + state.pds = Some(jacquard_common::CowStr::Borrowed("https://pds.example/")); + let mut txn = crate::db::Txn::new(&db); + txn.put_repo(&did(), &state)?; + txn.commit()?; + let accounts = || db.get_count_sync(&keys::pds_account_count_key("pds.example")); + assert_eq!(accounts(), 1); + + let mut listing = RepoListing { + did: did(), + rev: rev(2), + active: false, + status: Some(ListedRepoStatus::Deactivated), + create_if_missing: false, + }; + assert_eq!(apply(&db, &listing)?, ReconcileOutcome::AccountChanged); + assert_eq!(accounts(), 0); + + listing.active = true; + listing.status = None; + assert_eq!(apply(&db, &listing)?, ReconcileOutcome::AccountChanged); + assert_eq!(accounts(), 1); + Ok(()) + } + #[test] fn inactive_listing_removes_pending_backfill() -> Result<()> { let (_tmp, db) = open_db()?; diff --git a/src/db/counts.rs b/src/db/counts.rs index 2bd27df..0ce57b6 100644 --- a/src/db/counts.rs +++ b/src/db/counts.rs @@ -1,4 +1,5 @@ use crate::db::{Db, keys}; +use crate::types::RepoState; use fjall::OwnedWriteBatch; use miette::{Context, IntoDiagnostic, Result}; use smol_str::SmolStr; @@ -7,6 +8,7 @@ use std::sync::Arc; use std::sync::Mutex; use std::sync::atomic::Ordering; use tracing::error; +use url::Url; #[derive(Debug, Clone, Default)] pub struct CountDeltas { @@ -117,6 +119,29 @@ 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()) + .and_then(|url| url.host_str().map(SmolStr::new)) +} + +/// the per-pds account counts one repo row write moved, `from` losing an account and `to` +/// gaining one +#[derive(Debug, Default)] +pub(crate) struct PdsAccountMove { + pub(crate) from: Option, + pub(crate) to: Option, +} + +impl PdsAccountMove { + pub(crate) fn hosts(self) -> impl Iterator { + self.from.into_iter().chain(self.to) + } +} + pub(crate) struct CountDeltaReservation { in_flight: Arc>>, start_id: u64, diff --git a/src/db/indexer.rs b/src/db/indexer.rs index 6ecc47a..2390651 100644 --- a/src/db/indexer.rs +++ b/src/db/indexer.rs @@ -1,5 +1,4 @@ use fjall::{Keyspace, OwnedWriteBatch}; -use jacquard_common::IntoStatic; #[cfg(feature = "indexer_stream")] use jacquard_common::types::cid::IpldCid; use jacquard_common::types::string::Did; @@ -7,33 +6,11 @@ use miette::{IntoDiagnostic, Result, WrapErr}; use url::Url; use crate::db::types::DbTid; -use crate::db::{CountDeltas, Db, Txn, deser_repo_state, keys, ser_repo_state}; -use crate::types::{DeleteBodiesReport, DeleteBodyTarget, RepoState}; +use crate::db::{CountDeltas, Db, Txn, keys}; +use crate::types::{DeleteBodiesReport, DeleteBodyTarget}; use std::collections::BTreeMap; impl Db { - pub(crate) fn update_repo_state( - batch: &mut OwnedWriteBatch, - repos: &Keyspace, - did: &Did, - f: F, - ) -> Result, T)>> - where - F: FnOnce(&mut RepoState, (&[u8], &mut fjall::OwnedWriteBatch)) -> Result<(bool, T)>, - { - let key = keys::repo_key(did); - if let Some(bytes) = repos.get(&key).into_diagnostic()? { - let mut state: RepoState = deser_repo_state(bytes.as_ref())?.into_static(); - let (changed, result) = f(&mut state, (key.as_slice(), batch))?; - if changed { - batch.insert(repos, key, ser_repo_state(&state)?); - } - Ok(Some((state, result))) - } else { - Ok(None) - } - } - /// resolve one stored event pointer from the body-bearing layouts that can /// legitimately own it. every candidate is content-verified before use. #[cfg(feature = "indexer_stream")] @@ -308,10 +285,9 @@ pub(crate) fn redact_record_bodies( /// replaying them resolves no body. the caller untracks the repo first, which /// keeps ingestion from refilling it while the chunks go. pub(crate) fn erase_repo(db: &Db, did: &Did) -> Result { - let lock_index = crate::db::record_lock_index_for_did(did); - let _guard = db.record_write_locks[lock_index as usize] - .lock() - .unwrap_or_else(|e| e.into_inner()); + // the repo lock stays held across the chunked erases below and the final commit + let mut txn = Txn::new(db); + txn.hold_repo_write_lock(did); let prefix = keys::record_prefix_did(did); let mut report = DeleteBodiesReport::default(); @@ -371,7 +347,6 @@ pub(crate) fn erase_repo(db: &Db, did: &Did) -> Result { report.buffered_commits_deleted = count; report.buffered_commit_bytes_deleted = bytes; - let mut txn = Txn::new(db); // the gauge is read from the pending and resync entries, so it goes first txn.forget_lifecycle(did)?; let metadata_key = keys::repo_metadata_key(did); @@ -383,7 +358,7 @@ pub(crate) fn erase_repo(db: &Db, did: &Did) -> Result { } let repo_key = keys::repo_key(did); if db.repos.contains_key(&repo_key).into_diagnostic()? { - txn.batch.remove(&db.repos, &repo_key); + txn.remove_repo(did)?; txn.counts.add_repos(-1); } txn.batch.remove(&db.indexer.resync, &repo_key); diff --git a/src/db/txn.rs b/src/db/txn.rs index 282859c..e86471f 100644 --- a/src/db/txn.rs +++ b/src/db/txn.rs @@ -1,29 +1,26 @@ -#[cfg(feature = "indexer")] use std::collections::HashMap; use std::time::Duration; -#[cfg(feature = "indexer")] -#[cfg(feature = "indexer")] -#[cfg(feature = "indexer")] use jacquard_common::types::did::Did; #[cfg(feature = "indexer")] use jacquard_common::types::string::Tid; -#[cfg(feature = "indexer")] -use miette::IntoDiagnostic; -use miette::Result; +use miette::{IntoDiagnostic, Result}; +use smol_str::SmolStr; +#[cfg(feature = "indexer")] +use crate::db::LifecycleCountBatch; +use crate::db::counts::{PdsAccountMove, pds_account_host}; use crate::db::outbox::{Assigned, Outbox}; #[cfg(feature = "indexer")] use crate::db::types::{DbAction, DbRkey, DbTid}; -use crate::db::{CountDeltas, Db}; -#[cfg(feature = "indexer")] -use crate::db::{LifecycleCountBatch, keys}; +use crate::db::{CountDeltas, Db, keys}; #[cfg(feature = "indexer")] use crate::ops::record_events::{EmitOp, RecordEmitter, RecordEventOrigin}; #[cfg(feature = "indexer")] use crate::state::AppState; #[cfg(feature = "indexer")] -use crate::types::{GaugeState, RepoState}; +use crate::types::GaugeState; +use crate::types::RepoState; #[cfg(feature = "firehose-diagnostics")] type TxnInstant = std::time::Instant; @@ -61,6 +58,8 @@ pub(crate) struct Txn<'db> { pub(crate) counts: CountDeltas, /// stream events, given positions and broadcast when the batch commits pub(crate) outbox: Outbox, + /// the pds host each repo row this transaction put counts an account for + pds_accounts: HashMap, Option>, #[cfg(feature = "indexer")] lifecycle: Option>, /// record-write guards retained from the first record/purge touch until @@ -75,18 +74,7 @@ pub(crate) struct Txn<'db> { impl<'db> Txn<'db> { #[allow(dead_code)] pub(crate) fn new(db: &'db Db) -> Self { - Self { - batch: db.inner.batch(), - db, - counts: CountDeltas::default(), - outbox: Outbox::default(), - #[cfg(feature = "indexer")] - lifecycle: None, - #[cfg(feature = "indexer")] - record_lock_guards: Vec::new(), - #[cfg(feature = "indexer")] - held_record_locks: std::collections::BTreeSet::new(), - } + Self::from_parts(db, db.inner.batch(), CountDeltas::default()) } pub(crate) fn from_parts( @@ -99,6 +87,7 @@ impl<'db> Txn<'db> { db, counts, outbox: Outbox::default(), + pds_accounts: HashMap::new(), #[cfg(feature = "indexer")] lifecycle: None, #[cfg(feature = "indexer")] @@ -240,6 +229,59 @@ impl<'db> Txn<'db> { }) } + /// 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)?); + 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)?; + self.batch.remove(&self.db.repos, key); + Ok(moved) + } + + #[cfg_attr(not(feature = "indexer"), allow(unused_variables))] + fn move_pds_account( + &mut self, + did: &Did, + key: &[u8], + to: Option, + ) -> 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(), + None => self + .db + .repos + .get(key) + .into_diagnostic()? + .map(|bytes| crate::db::deser_repo_state(&bytes).map(|row| pds_account_host(&row))) + .transpose()? + .flatten(), + }; + self.pds_accounts.insert(key.to_vec(), to.clone()); + if from == to { + return Ok(PdsAccountMove::default()); + } + if let Some(host) = &from { + self.counts.add_pds_account(host, -1); + } + if let Some(host) = &to { + self.counts.add_pds_account(host, 1); + } + Ok(PdsAccountMove { from, to }) + } + #[cfg(feature = "indexer")] pub(crate) fn transition_lifecycle(&mut self, did: &Did, gauge: GaugeState) -> Result { let lifecycle = self @@ -621,12 +663,7 @@ impl RecordTxn<'_, '_, '_> { if self.is_excluded() { return Ok(()); } - self.txn.batch.insert( - &self.txn.db.repos, - keys::repo_key(self.did), - crate::db::ser_repo_state(state)?, - ); - Ok(()) + self.txn.put_repo(self.did, state).map(drop) } pub(crate) fn finish(self) -> Result<()> { @@ -727,6 +764,52 @@ mod tests { .collect() } + #[test] + fn repo_rows_count_for_the_pds_they_are_active_on() { + let (_tmp, state) = test_state(); + let did = did(); + let accounts = |host| state.db.get_count_sync(&keys::pds_account_count_key(host)); + let row = |active, pds: &'static str| RepoState { + active, + pds: Some(jacquard_common::CowStr::Borrowed(pds)), + ..RepoState::backfilling() + }; + let put = |row: &RepoState| { + let mut txn = Txn::new(&state.db); + txn.put_repo(&did, row).unwrap(); + txn.commit().unwrap(); + }; + + put(&row(true, "https://one.example/")); + put(&row(true, "https://one.example/")); + assert_eq!(accounts("one.example"), 1); + put(&row(false, "https://one.example/")); + assert_eq!(accounts("one.example"), 0); + put(&row(true, "https://two.example/")); + assert_eq!((accounts("one.example"), accounts("two.example")), (0, 1)); + + // the second put moves the count from what the first one staged, not the stored row + let mut txn = Txn::new(&state.db); + txn.put_repo(&did, &row(true, "https://one.example/")) + .unwrap(); + txn.put_repo(&did, &row(true, "https://three.example/")) + .unwrap(); + txn.commit().unwrap(); + assert_eq!( + ( + accounts("one.example"), + accounts("two.example"), + accounts("three.example") + ), + (0, 0, 1) + ); + + let mut txn = Txn::new(&state.db); + txn.remove_repo(&did).unwrap(); + txn.commit().unwrap(); + assert_eq!(accounts("three.example"), 0); + } + #[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 c2548cf..d03ffac 100644 --- a/src/ingest/indexer/shard.rs +++ b/src/ingest/indexer/shard.rs @@ -471,8 +471,7 @@ impl FirehoseWorker { repo_state.status.clone() }; ops::transition_repo( - &mut ctx.txn.batch, - &ctx.state.db, + &mut ctx.txn, ctx.lifecycle_transitions, did, repo_state, diff --git a/src/ingest/relay/context.rs b/src/ingest/relay/context.rs index 5c84c54..7a7e182 100644 --- a/src/ingest/relay/context.rs +++ b/src/ingest/relay/context.rs @@ -21,8 +21,8 @@ use crate::types::{RepoState, RepoStatus}; use crate::util; use super::{ - AuthorityOutcome, RelayWorker, WRONG_HOST_AUTHORITY_CACHE_PRUNE_AT, - WRONG_HOST_AUTHORITY_RECHECK_INTERVAL, WorkerMessage, map_repo_status_probe, + AuthorityOutcome, WRONG_HOST_AUTHORITY_CACHE_PRUNE_AT, WRONG_HOST_AUTHORITY_RECHECK_INTERVAL, + WorkerMessage, map_repo_status_probe, }; use crate::ingest::stream::AccountStatus; #[cfg(feature = "indexer")] @@ -72,6 +72,7 @@ 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")] @@ -105,8 +106,13 @@ impl WorkerContext<'_> { #[cfg(not(feature = "indexer"))] _ => staged, }; - txn.batch - .insert(&self.state.db.repos, key, db::ser_repo_state(&row)?); + for host in txn.put_repo(did, &row)?.hosts() { + let count = txn + .counts + .projected_pds_account_count(&self.state.db, &host); + self.state + .apply_host_limit_status(&mut txn.batch, &host, count); + } Ok(()) } @@ -537,14 +543,6 @@ impl WorkerContext<'_> { let mut repo_state = RepoState::synced(); repo_state.update_from_doc(doc); - RelayWorker::update_pds_account_count( - self, - false, - None, - repo_state.active, - RelayWorker::pds_host(repo_state.pds.as_deref()).as_deref(), - ); - self.stage_repo_state(&repo_state); self.sink.new_repo(did.clone()); diff --git a/src/ingest/relay/handlers.rs b/src/ingest/relay/handlers.rs index 87a8955..ca25461 100644 --- a/src/ingest/relay/handlers.rs +++ b/src/ingest/relay/handlers.rs @@ -1,5 +1,4 @@ use miette::Result; -use smol_str::SmolStr; use tokio::runtime::Handle; use tracing::{debug, warn}; use url::Url; @@ -98,8 +97,6 @@ impl RelayWorker { if event_is_new { repo_state.advance_identity_time(event_ms); } - let was_active = repo_state.active; - let was_pds_host = Self::pds_host(repo_state.pds.as_deref()); // refresh did doc if a pds sent this event // or if there is no handle specified. the doc is the authority, so what @@ -125,14 +122,6 @@ impl RelayWorker { identity.handle = None; } - Self::update_pds_account_count( - ctx, - was_active, - was_pds_host.as_deref(), - repo_state.active, - Self::pds_host(repo_state.pds.as_deref()).as_deref(), - ); - ctx.sink .identity(ctx.state, &mut ctx.batch, firehose, identity, changed, true)?; @@ -156,9 +145,7 @@ impl RelayWorker { repo_state.advance_account_time(event_ms); - // always capture was_active for count tracking, not just in indexer mode let was_active = repo_state.active; - let was_pds_host = Self::pds_host(repo_state.pds.as_deref()); let snapshot = ctx.sink.account_snapshot(repo_state); repo_state.active = account.active; @@ -185,14 +172,6 @@ impl RelayWorker { }; } - Self::update_pds_account_count( - ctx, - was_active, - was_pds_host.as_deref(), - repo_state.active, - Self::pds_host(repo_state.pds.as_deref()).as_deref(), - ); - ctx.sink.account( ctx.state, &mut ctx.batch, @@ -208,38 +187,4 @@ impl RelayWorker { Ok(()) } - - pub(super) fn pds_host(pds: Option<&str>) -> Option { - pds.and_then(|pds| Url::parse(pds).ok()) - .and_then(|url| url.host_str().map(SmolStr::new)) - } - - pub(super) fn update_pds_account_count( - ctx: &mut WorkerContext, - old_active: bool, - old_host: Option<&str>, - new_active: bool, - new_host: Option<&str>, - ) { - if old_active && old_host == new_host && new_active { - return; - } - - let mut update_host = |host: &str, delta| { - ctx.count_deltas.add_pds_account(host, delta); - let count = ctx - .count_deltas - .projected_pds_account_count(&ctx.state.db, host); - ctx.state - .apply_host_limit_status(&mut ctx.batch, host, count); - }; - - if old_active && let Some(host) = old_host { - update_host(host, -1); - } - - if new_active && let Some(host) = new_host { - update_host(host, 1); - } - } } diff --git a/src/ingest/relay/worker.rs b/src/ingest/relay/worker.rs index 8a67a86..37d64b3 100644 --- a/src/ingest/relay/worker.rs +++ b/src/ingest/relay/worker.rs @@ -344,7 +344,7 @@ mod tests { use crate::db::{self, keys}; use crate::ingest::indexer::{IndexerMessage, IndexerTx}; use crate::ingest::mailbox::{ShardedReceiver, ShardedSender}; - use crate::ingest::stream::{Datetime, Identity}; + use crate::ingest::stream::{Account, AccountStatus, Datetime, Identity}; use crate::types::RepoState; const DID: &str = "did:plc:aaaaaaaaaaaaaaaaaaaaaaaa"; @@ -509,6 +509,38 @@ mod tests { Ok(()) } + /// sends an account event `secs` seconds from now, so later calls are newer, and comes + /// back once the relay has committed it + async fn account( + &mut self, + active: bool, + status: Option>, + secs: i64, + ) -> miette::Result<()> { + let time = chrono::Utc::now() + chrono::Duration::seconds(secs); + self.buffer_tx + .send(IngestMessage::Firehose { + url: Url::parse("wss://relay.example").into_diagnostic()?, + is_pds: false, + msg: SubscribeReposMessage::Account(Box::new(Account { + active, + did: self.did.clone(), + seq: secs, + status, + time: Datetime(time.into()), + })), + }) + .await + .into_diagnostic()?; + self.release().await + } + + fn pds_accounts(&self) -> u64 { + self.state + .db + .get_count_sync(&keys::pds_account_count_key("pds.example")) + } + fn put_repo(&self, row: &RepoState<'_>) -> miette::Result<()> { let mut batch = self.state.db.inner.batch(); batch.insert( @@ -612,6 +644,69 @@ mod tests { Ok(()) } + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn a_backfill_that_finds_a_repo_active_again_counts_it_for_its_pds() -> miette::Result<()> + { + let mut relay = HeldRelay::start_without_a_row().await?; + let mut row = RepoState::synced(); + row.pds = Some(DOC_PDS.into()); + let mut txn = crate::db::Txn::new(&relay.state.db); + txn.put_repo(&relay.did, &row)?; + txn.commit()?; + assert_eq!(relay.pds_accounts(), 1); + + relay + .account(false, Some(AccountStatus::Deactivated), 1) + .await?; + assert_eq!(relay.pds_accounts(), 0); + + // a backfill of the deactivated repo gets its car anyway, which says it's active, and + // writes the row the way the backfill's persist does + let before = relay.repo()?.expect("repo row"); + let mut fetched = before.clone(); + fetched.root = Some(backfilled_root()?); + let mut txn = crate::db::Txn::new(&relay.state.db); + txn.hold_repo_write_lock(&relay.did); + let after = + crate::backfill::backfilled_state(&relay.state.db, &relay.did, &before, fetched)? + .expect("repo row"); + assert!(after.active); + let rev = after.root.as_ref().expect("root").rev; + let mut records = txn.backfill_records(&relay.state, &rev, &relay.did)?; + records.update_repo_state(&after)?; + records.finish()?; + txn.commit()?; + assert_eq!(relay.pds_accounts(), 1); + + relay.account(true, None, 2).await?; + assert_eq!(relay.pds_accounts(), 1); + Ok(()) + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn a_repo_made_inactive_while_the_relay_took_it_for_new_is_not_counted() + -> miette::Result<()> { + let mut relay = HeldRelay::start_without_a_row().await?; + relay.identity_waiting_on_the_doc().await?; + + // the crawler lists the repo as deactivated and makes its row while the relay resolves + // the doc. the relay's probe then fails, so the row it made up says active + let mut listed = RepoState::backfilling(); + listed.active = false; + listed.status = crate::types::RepoStatus::Deactivated; + let mut txn = crate::db::Txn::new(&relay.state.db); + txn.put_repo(&relay.did, &listed)?; + txn.counts.add_repos(1); + txn.commit()?; + + relay.release().await?; + let after = relay.repo()?.expect("repo row"); + assert!(!after.active); + assert_eq!(after.pds.as_deref(), Some(DOC_PDS)); + assert_eq!(relay.pds_accounts(), 0); + Ok(()) + } + /// lets the relay's doc through while this test holds the repo lock, and checks the relay /// commits only once the lock is free. for a new row the relay also probes the pds, its last /// network call before the lock, so `probes` makes the window below start after that diff --git a/src/ops.rs b/src/ops.rs index 73ec84d..76e3f65 100644 --- a/src/ops.rs +++ b/src/ops.rs @@ -129,8 +129,7 @@ pub(crate) fn delete_repo( } pub fn transition_repo<'s>( - batch: &mut OwnedWriteBatch, - db: &Db, + txn: &mut Txn<'_>, lifecycle_transitions: &mut Vec<(Did, GaugeState)>, did: &Did, mut repo_state: RepoState<'s>, @@ -138,6 +137,8 @@ pub fn transition_repo<'s>( ) -> Result> { debug!(did = %did, status = ?new_status, "updating repo status"); + let db = txn.db; + let batch = &mut txn.batch; let repo_key = keys::repo_key(did); let metadata_key = keys::repo_metadata_key(did); @@ -215,7 +216,7 @@ pub fn transition_repo<'s>( repo_state.active = matches!(new_status, RepoStatus::Synced | RepoStatus::Error(_)); repo_state.status = new_status; repo_state.touch(); - batch.insert(&db.repos, repo_key, db::ser_repo_state(&repo_state)?); + txn.put_repo(did, &repo_state)?; Ok(repo_state) } -- 2.51.2