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) }