diff --git a/src/backfill/manager.rs b/src/backfill/manager.rs index 1fc82f2..4ff28ba 100644 --- a/src/backfill/manager.rs +++ b/src/backfill/manager.rs @@ -244,6 +244,7 @@ mod tests { .unwrap() ); assert_eq!(state.db.indexer.pending_iter().count(), 0); + state.db.indexer.assert_pending_dids_match(); } #[test] @@ -266,5 +267,6 @@ mod tests { .unwrap() ); assert_eq!(state.db.indexer.pending_iter().count(), 1); + state.db.indexer.assert_pending_dids_match(); } } diff --git a/src/backfill/worker/task.rs b/src/backfill/worker/task.rs index 64e5880..fe75324 100644 --- a/src/backfill/worker/task.rs +++ b/src/backfill/worker/task.rs @@ -617,6 +617,15 @@ mod tests { _tmp: tempfile::TempDir, } + // whatever a test does to the queue, its two keyspaces have to end in step + impl Drop for Fixture { + fn drop(&mut self) { + if !std::thread::panicking() { + self.state.db.indexer.assert_pending_dids_match(); + } + } + } + impl Fixture { async fn new() -> miette::Result { let pds = PdsStub::spawn().await; diff --git a/src/control/repos/indexer.rs b/src/control/repos/indexer.rs index a1474f4..efb47fe 100644 --- a/src/control/repos/indexer.rs +++ b/src/control/repos/indexer.rs @@ -828,6 +828,7 @@ mod tests { .is_pending(&keys::pending_key(7)) .into_diagnostic()? ); + db.indexer.assert_pending_dids_match(); assert!(gone(db.crawler.raw(), keys::crawler_retry_key(&did))); assert!(gone( db.counts.raw(), @@ -857,6 +858,36 @@ mod tests { Ok(()) } + #[tokio::test] + async fn queue_moves_keep_the_did_order_copy_in_step() -> miette::Result<()> { + let tmp = tempfile::tempdir().into_diagnostic()?; + let config = crate::config::Config { + database_path: tmp.path().to_path_buf(), + ..Default::default() + }; + let state = Arc::new(AppState::new(&config)?); + let repos = ReposControl(state.clone()); + let did = Did::new_static("did:plc:testoperatorrepo").into_diagnostic()?; + let other = Did::new_static("did:web:example.com").into_diagnostic()?; + let in_step = || state.db.indexer.assert_pending_dids_match(); + + repos.track([did.clone(), other.clone()]).await?; + in_step(); + // already queued, so nothing moves + repos.track([did.clone()]).await?; + in_step(); + repos.untrack([did.clone(), other.clone()]).await?; + in_step(); + assert_eq!(state.db.indexer.pending_iter().count(), 0); + repos.track([did.clone()]).await?; + in_step(); + // an untracked repo comes back under a fresh pending key + assert_eq!(repos.resync([other.clone()]).await?, vec![other.clone()]); + in_step(); + assert_eq!(state.db.indexer.pending_iter().count(), 2); + Ok(()) + } + #[tokio::test] #[cfg(not(feature = "backlinks"))] async fn test_generate_car_ephemeral_rejected() -> miette::Result<()> { diff --git a/src/crawler/list_repos/producer.rs b/src/crawler/list_repos/producer.rs index 7418df2..43c34fa 100644 --- a/src/crawler/list_repos/producer.rs +++ b/src/crawler/list_repos/producer.rs @@ -482,6 +482,7 @@ mod tests { assert!(!repo.active); assert_eq!(repo.status, crate::types::RepoStatus::Takendown); assert_eq!(state.db.indexer.pending_iter().count(), 0); + state.db.indexer.assert_pending_dids_match(); producer.abort(); worker.abort(); diff --git a/src/crawler/worker.rs b/src/crawler/worker.rs index c223675..7321247 100644 --- a/src/crawler/worker.rs +++ b/src/crawler/worker.rs @@ -662,6 +662,7 @@ mod tests { assert!(!state.active); assert_eq!(state.status, RepoStatus::Takendown); assert_eq!(db.indexer.pending_iter().count(), 0); + db.indexer.assert_pending_dids_match(); assert!( db.indexer .resync @@ -686,6 +687,7 @@ mod tests { assert_eq!(apply(&db, &listing)?, ReconcileOutcome::Unchanged); assert_eq!(db.indexer.pending_iter().count(), 0); + db.indexer.assert_pending_dids_match(); Ok(()) } @@ -706,6 +708,7 @@ mod tests { assert_eq!(apply(&db, &listing)?, ReconcileOutcome::StateOnly); assert_eq!(db.indexer.pending_iter().count(), 0); + db.indexer.assert_pending_dids_match(); Ok(()) } @@ -723,6 +726,7 @@ mod tests { assert_eq!(apply(&db, &listing)?, ReconcileOutcome::StateOnly); assert_eq!(db.indexer.pending_iter().count(), 1); + db.indexer.assert_pending_dids_match(); let state = read_state(&db)?; assert!(state.active); assert_eq!(state.status, RepoStatus::Desynchronized); @@ -861,6 +865,7 @@ mod tests { assert!(events.try_recv().is_err()); assert_eq!(state.db.indexer.pending_iter().count(), 1); + state.db.indexer.assert_pending_dids_match(); Ok(()) } } diff --git a/src/db/keys/indexer.rs b/src/db/keys/indexer.rs index e01a498..7727594 100644 --- a/src/db/keys/indexer.rs +++ b/src/db/keys/indexer.rs @@ -11,6 +11,16 @@ pub fn pending_key(id: u64) -> [u8; 8] { id.to_be_bytes() } +// key format: {DID} 00 {ID} (DID trimmed). the NUL keeps did order, because a +// did only prefixes another when both are text, and text never holds a NUL +pub fn pending_did_key(repo_key: &[u8], pending_key: &[u8]) -> Vec { + let mut key = Vec::with_capacity(repo_key.len() + 1 + pending_key.len()); + key.extend_from_slice(repo_key); + key.push(REC_SEP); + key.extend_from_slice(pending_key); + key +} + #[cfg(feature = "indexer_stream")] pub fn event_watermark_key(timestamp_secs: u64) -> Vec { let mut key = Vec::with_capacity(EVENT_WATERMARK_PREFIX.len() + 8); @@ -474,6 +484,19 @@ mod tests { let rev = DbTid::from(&Tid::new("3kzbif5moe22m").unwrap()); assert_eq!(pending_key(42), 42_u64.to_be_bytes()); + let repo_key = crate::db::keys::repo_key(&did); + assert_eq!( + pending_did_key(&repo_key, &pending_key(42)), + [repo_key.as_slice(), &[REC_SEP], &pending_key(42)].concat() + ); + // a did that prefixes another sorts first here too, the same as in `resync` + let short = crate::db::keys::repo_key(&Did::new_static("did:web:a.com").unwrap()); + let long = crate::db::keys::repo_key(&Did::new_static("did:web:a.com.au").unwrap()); + assert!(short < long); + assert!( + pending_did_key(&short, &pending_key(u64::MAX)) + < pending_did_key(&long, &pending_key(0)) + ); #[cfg(feature = "indexer_stream")] assert_eq!( event_watermark_key(42), diff --git a/src/db/keyspaces.rs b/src/db/keyspaces.rs index 0316878..58996b5 100644 --- a/src/db/keyspaces.rs +++ b/src/db/keyspaces.rs @@ -62,9 +62,12 @@ pub struct IndexerDb { pub(super) history: Ks, /// operator-redacted versions: `{record key} 00 {cid}` -> empty pub(super) redactions: Ks, - /// backfill queue of `{ID}` -> `{DID}`, written only by `stage_pending_insert` - /// and `stage_pending_remove` + /// backfill queue of `{ID}` -> `{DID}` pub(super) pending: Ks, + /// the same queue in did order, `{DID} 00 {ID}` -> empty. both are written + /// only by `stage_pending_insert` and `stage_pending_remove`, which keep them + /// in step + pub(super) pending_dids: Ks, /// per-repo resync/retry state pub resync: Ks, /// live events buffered during backfill @@ -81,6 +84,7 @@ impl IndexerDb { history: Ks::open(cx)?, redactions: Ks::open(cx)?, pending: Ks::open(cx)?, + pending_dids: Ks::open(cx)?, resync: Ks::open(cx)?, resync_buffer: Ks::open(cx)?, lifecycle_count_lock: Arc::new(std::sync::Mutex::new(())), @@ -144,19 +148,48 @@ impl IndexerDb { did: &Did, pending_key: &[u8], ) { - batch.insert(&self.pending, pending_key, super::keys::repo_key(did)); + let repo_key = super::keys::repo_key(did); + let did_key = super::keys::pending_did_key(&repo_key, pending_key); + batch.insert(&self.pending_dids, did_key, []); + batch.insert(&self.pending, pending_key, repo_key); } - /// drops `did`'s queue entry under `pending_key` + /// drops `did`'s queue entry under `pending_key`. the did has to be the + /// entry's own, or its did-order copy stays behind pub(crate) fn stage_pending_remove( &self, batch: &mut fjall::OwnedWriteBatch, - _did: &Did, + did: &Did, pending_key: &[u8], ) { + let repo_key = super::keys::repo_key(did); + batch.remove( + &self.pending_dids, + super::keys::pending_did_key(&repo_key, pending_key), + ); batch.remove(&self.pending, pending_key); } + /// panics unless `pending_dids` holds exactly the did-order copy of every + /// `pending` entry + #[cfg(test)] + pub(crate) fn assert_pending_dids_match(&self) { + let expected = self + .pending + .iter() + .map(|guard| { + let (pending_key, repo_key) = guard.into_inner().unwrap(); + super::keys::pending_did_key(&repo_key, &pending_key) + }) + .collect::>(); + let actual = self + .pending_dids + .iter() + .map(|guard| guard.key().unwrap().to_vec()) + .collect::>(); + assert_eq!(actual, expected, "pending_dids drifted from pending"); + } + pub(crate) fn is_pending(&self, pending_key: &[u8]) -> fjall::Result { self.pending.contains_key(pending_key) } diff --git a/src/db/lifecycle_counts.rs b/src/db/lifecycle_counts.rs index 9012669..b94a0b3 100644 --- a/src/db/lifecycle_counts.rs +++ b/src/db/lifecycle_counts.rs @@ -316,6 +316,7 @@ mod tests { assert!(!complete_pending(&db, &did, stale_pending)?); assert_eq!(db.get_count_sync("pending"), 1); assert!(db.indexer.is_pending(¤t_pending).into_diagnostic()?); + db.indexer.assert_pending_dids_match(); Ok(()) } @@ -356,6 +357,7 @@ mod tests { .is_pending(&keys::pending_key(1)) .into_diagnostic()? ); + db.indexer.assert_pending_dids_match(); } Ok(()) diff --git a/src/db/migration/mod.rs b/src/db/migration/mod.rs index 01f8a73..4291d92 100644 --- a/src/db/migration/mod.rs +++ b/src/db/migration/mod.rs @@ -65,6 +65,7 @@ fn legacy_records(db: &Db) -> Result> { mod v1; mod v10; mod v11; +mod v12; mod v2; mod v3; mod v4; @@ -242,6 +243,13 @@ const MIGRATIONS: &[(&str, Migration)] = &[ finalize: None, }, ), + ( + "pending_dids", + Migration::Chunked { + passes: v12::PASSES, + finalize: None, + }, + ), ]; pub(crate) const LATEST_VERSION: u64 = MIGRATIONS.len() as u64; @@ -275,7 +283,9 @@ fn migration_applied_key(name: &str) -> Vec { /// test helper: rewind the db to `version` as if later migrations never ran. /// opening a fresh test db already runs (and marks) every migration, so -/// simulating an old database means clearing the applied markers too. +/// simulating an old database means clearing the applied markers too. 9 is +/// as far forward as this can go, because a v10 database without its marker +/// is rejected on purpose. #[cfg(test)] pub(crate) fn rewind_version_for_test(db: &Db, version: u64) -> Result<()> { let mut batch = db.inner.batch(); diff --git a/src/db/migration/v11.rs b/src/db/migration/v11.rs index 27d5afb..3363da9 100644 --- a/src/db/migration/v11.rs +++ b/src/db/migration/v11.rs @@ -89,7 +89,7 @@ fn rewrite_status( mod tests { use super::*; use crate::config::Config; - use crate::db::migration::{LATEST_VERSION, rewind_version_for_test}; + use crate::db::migration::rewind_version_for_test; use tempfile::tempdir; #[test] @@ -115,7 +115,7 @@ mod tests { rmp_serde::to_vec(&HostStatus::Active).into_diagnostic()?, ); batch.commit().into_diagnostic()?; - rewind_version_for_test(&db, LATEST_VERSION - 2)?; + rewind_version_for_test(&db, 9)?; db.persist()?; } diff --git a/src/db/migration/v12.rs b/src/db/migration/v12.rs new file mode 100644 index 0000000..9a68dca --- /dev/null +++ b/src/db/migration/v12.rs @@ -0,0 +1,82 @@ +//! v12 copies the backfill queue into `pending_dids`, so a queue filled before +//! the did-order copy existed can be listed by did too. the pass writes a +//! different keyspace than it scans and inserting the same key twice changes +//! nothing, so a crash that re-visits a chunk is harmless. + +#[cfg(feature = "indexer")] +use fjall::{Keyspace, OwnedWriteBatch}; +#[cfg(feature = "indexer")] +use miette::Result; + +#[cfg(feature = "indexer")] +use crate::db::{Db, keys}; + +#[cfg(feature = "indexer")] +use super::ChunkBudget; +use super::Pass; + +pub(super) const PASSES: &[Pass] = &[ + #[cfg(feature = "indexer")] + Pass { + name: "pending_dids", + scan: "pending", + visit: copy_pending_entry, + budget: ChunkBudget::DEFAULT, + optional_on_absent: false, + }, +]; + +#[cfg(feature = "indexer")] +fn copy_pending_entry( + db: &Db, + batch: &mut OwnedWriteBatch, + _scanned: &Keyspace, + key: &[u8], + value: &[u8], +) -> Result { + let did_key = keys::pending_did_key(value, key); + let staged = did_key.len(); + batch.insert(&db.indexer.pending_dids, did_key, []); + Ok(staged) +} + +#[cfg(all(test, feature = "indexer"))] +mod tests { + use super::*; + use crate::config::Config; + use crate::db::migration::rewind_version_for_test; + use jacquard_common::types::did::Did; + use miette::IntoDiagnostic; + use tempfile::tempdir; + + #[test] + fn copies_a_queue_that_has_no_did_order_copy_yet() -> Result<()> { + let tmp = tempdir().into_diagnostic()?; + let cfg = Config { + database_path: tmp.path().to_path_buf(), + ..Default::default() + }; + let plc = Did::new_static("did:plc:yk4q3id7id6p5z3bypvshc64").into_diagnostic()?; + let web = Did::new_static("did:web:example.com").into_diagnostic()?; + { + let db = Db::open(&cfg)?; + let mut batch = db.inner.batch(); + // an older build leaves queue entries with no did-order copy + for (id, did) in [(3, &plc), (9, &web), (1, &web)] { + batch.insert( + &db.indexer.pending, + keys::pending_key(id), + keys::repo_key(did), + ); + } + batch.commit().into_diagnostic()?; + rewind_version_for_test(&db, 9)?; + db.persist()?; + } + + let db = Db::open(&cfg)?; + db.indexer.assert_pending_dids_match(); + assert_eq!(db.indexer.pending_dids.iter().count(), 3); + Ok(()) + } +} diff --git a/src/db/registry.rs b/src/db/registry.rs index e493ac1..7928352 100644 --- a/src/db/registry.rs +++ b/src/db/registry.rs @@ -104,6 +104,8 @@ registry! { #[cfg(feature = "indexer")] indexer.pending => schema::Pending, #[cfg(feature = "indexer")] + indexer.pending_dids => schema::PendingDids, + #[cfg(feature = "indexer")] indexer.resync => schema::Resync, #[cfg(feature = "indexer")] indexer.resync_buffer => schema::ResyncBuffer, diff --git a/src/db/schema.rs b/src/db/schema.rs index dad69ee..5868328 100644 --- a/src/db/schema.rs +++ b/src/db/schema.rs @@ -8,6 +8,8 @@ use std::marker::PhantomData; +#[cfg(feature = "indexer")] +use fjall::config::{FilterPolicy, PinningPolicy}; use fjall::{ CompressionType, Keyspace, KeyspaceCreateOptions, config::{BlockSizePolicy, CompressionPolicy, RestartIntervalPolicy}, @@ -547,6 +549,36 @@ impl Schema for Pending { } } +/// the backfill queue again in did order, `{DID} 00 {ID}` -> empty, so it can be +/// paged by did next to `resync` +#[cfg(feature = "indexer")] +pub struct PendingDids; + +#[cfg(feature = "indexer")] +impl Schema for PendingDids { + const NAME: &'static str = "pending_dids"; + + fn options(_cx: &OpenCx) -> KeyspaceCreateOptions { + KeyspaceCreateOptions::default() + // only ever scanned or written blind, so no level needs a filter and + // nothing has to stay pinned outside the block cache + .filter_policy(FilterPolicy::disabled()) + .filter_block_pinning_policy(PinningPolicy::disabled()) + .index_block_pinning_policy(PinningPolicy::disabled()) + // the same churn as `pending`, but entries are tiny so flushing more + // often is cheap + .max_memtable_size(mb(4)) + .data_block_size_policy(BlockSizePolicy::all(kb(8))) + // dids barely compress and entries leave once their backfill is done + .data_block_compression_policy(CompressionPolicy::disabled()) + .data_block_restart_interval_policy(RestartIntervalPolicy::all(16)) + } + + fn render_value(_value: &[u8]) -> Value { + Value::Null + } +} + /// per-repo resync/retry state, `{DID}` -> `ResyncState` (msgpack) #[cfg(feature = "indexer")] pub struct Resync; @@ -797,6 +829,7 @@ mod tests { "history", "redactions", "pending", + "pending_dids", "resync", "resync_buffer", ]); diff --git a/src/ingest/indexer/shard.rs b/src/ingest/indexer/shard.rs index 1cee149..04a05d4 100644 --- a/src/ingest/indexer/shard.rs +++ b/src/ingest/indexer/shard.rs @@ -630,6 +630,7 @@ mod tests { .is_pending(&pending_key) .into_diagnostic()? ); + state.db.indexer.assert_pending_dids_match(); Ok(()) } diff --git a/tests/representation_inventory.tsv b/tests/representation_inventory.tsv index 74a8ac6..fe790e0 100644 --- a/tests/representation_inventory.tsv +++ b/tests/representation_inventory.tsv @@ -290,6 +290,7 @@ db-key-codec:src/db/keys/indexer.rs::history_hit_for rust:src/db/keys/indexer.rs db-key-codec:src/db/keys/indexer.rs::history_key rust:src/db/keys/indexer.rs::indexer_key_codec_inventory_contracts db-key-codec:src/db/keys/indexer.rs::history_range_after rust:src/db/keys/indexer.rs::indexer_key_codec_inventory_contracts db-key-codec:src/db/keys/indexer.rs::parse_rkey_text rust:src/db/keys/indexer.rs::indexer_key_codec_inventory_contracts +db-key-codec:src/db/keys/indexer.rs::pending_did_key rust:src/db/keys/indexer.rs::indexer_key_codec_inventory_contracts db-key-codec:src/db/keys/indexer.rs::pending_key rust:src/db/keys/indexer.rs::indexer_key_codec_inventory_contracts db-key-codec:src/db/keys/indexer.rs::record_cid_key rust:src/db/keys/indexer.rs::indexer_key_codec_inventory_contracts db-key-codec:src/db/keys/indexer.rs::record_key rust:src/db/keys/indexer.rs::indexer_key_codec_inventory_contracts @@ -379,6 +380,7 @@ db-keyspace:src/db/schema.rs::Heads rust:src/db/schema.rs::keyspace_names_match_ db-keyspace:src/db/schema.rs::History rust:src/db/schema.rs::keyspace_names_match_the_compiled_storage_schema db-keyspace:src/db/schema.rs::JetstreamEvents rust:src/db/schema.rs::keyspace_names_match_the_compiled_storage_schema db-keyspace:src/db/schema.rs::Pending rust:src/db/schema.rs::keyspace_names_match_the_compiled_storage_schema +db-keyspace:src/db/schema.rs::PendingDids rust:src/db/schema.rs::keyspace_names_match_the_compiled_storage_schema db-keyspace:src/db/schema.rs::Redactions rust:src/db/schema.rs::keyspace_names_match_the_compiled_storage_schema db-keyspace:src/db/schema.rs::RelayEvents rust:src/db/schema.rs::keyspace_names_match_the_compiled_storage_schema db-keyspace:src/db/schema.rs::RepoMetadata rust:src/db/schema.rs::keyspace_names_match_the_compiled_storage_schema @@ -420,6 +422,7 @@ migration-transform:src/db/migration/v10/events.rs::decode_event rust:src/db/mig migration-transform:src/db/migration/v10/events.rs::materialize_pointer_event rust:src/db/migration/v10/events.rs::cid_mismatch_aborts_without_clearing_blocks migration-transform:src/db/migration/v10/events.rs::retire_blocks rust:src/db/migration/v10/events.rs::cid_mismatch_aborts_without_clearing_blocks migration-transform:src/db/migration/v10/events.rs::stage_archive_body rust:src/db/migration/v10/events.rs::cid_mismatch_aborts_without_clearing_blocks +migration-transform:src/db/migration/v12.rs::copy_pending_entry rust:src/db/migration/v12.rs::copies_a_queue_that_has_no_did_order_copy_yet migration-transform:src/db/migration/v2.rs::repo_state_root_commit rust:src/db/migration/mod.rs::v0_database_migrates_all_early_representation_boundaries migration-transform:src/db/migration/v3.rs::firehose_source_is_pds rust:src/db/migration/mod.rs::v0_database_migrates_all_early_representation_boundaries migration-transform:src/db/migration/v3.rs::is_known_relay rust:src/db/migration/mod.rs::v0_database_migrates_all_early_representation_boundaries