diff --git a/src/db/migration/mod.rs b/src/db/migration/mod.rs index e5477b8..125ec72 100644 --- a/src/db/migration/mod.rs +++ b/src/db/migration/mod.rs @@ -1,8 +1,10 @@ -use fjall::OwnedWriteBatch; +use std::ops::Bound; + +use fjall::{Keyspace, OwnedWriteBatch, Readable}; use miette::{Context, IntoDiagnostic, Result}; use crate::db::Db; -use crate::db::keys::VERSIONING_KEY; +use crate::db::keys::{SEP, VERSIONING_KEY}; mod v1; mod v2; @@ -14,7 +16,96 @@ mod v7; mod v8; mod v9; -type MigrationFn = fn(&Db, &mut OwnedWriteBatch) -> Result<()>; +/// a migration whose entire write set is buffered into one batch. +/// returns `Ok(true)` if applied, or `Ok(false)` if skipped due to missing keyspaces in this build. +type AtomicFn = fn(&Db, &mut OwnedWriteBatch) -> Result; + +/// transform one scanned entry. staged writes may target any keyspace. +/// returns the number of output bytes staged into the write batch for this entry. +type VisitFn = fn(&Db, &mut OwnedWriteBatch, &[u8], &[u8]) -> Result; + +/// destructive or otherwise non-batched work that runs only after every +/// chunked pass is durable. returning false means this build does not contain +/// the keyspaces needed by the migration. +type FinalizeFn = fn(&Db) -> Result; + +/// how a migration is applied. +enum Migration { + /// the whole write set is buffered into one batch, committed atomically with + /// the version bump. + /// + /// a batch is held entirely in memory, so this is only for migrations whose + /// write set is small and bounded. it is *not* bounded by how much the + /// migration reads: streaming a large keyspace to accumulate a handful of + /// counters is fine, staging one write per entry of a large keyspace is not. + Atomic(AtomicFn), + /// ordered resumable passes, each committed in chunks with bounded memory. + /// + /// the only way to rewrite a keyspace larger than memory. passes run in + /// order, and each runs to completion before the next starts. + // first user lands with `inline_record_bodies` + #[allow(dead_code)] + Chunked { + passes: &'static [Pass], + finalize: Option, + }, +} + +/// per-chunk limits. both bound peak batch memory; whichever is hit first ends +/// the chunk. they also bound how much work a crash re-does. +#[derive(Clone, Copy)] +struct ChunkBudget { + /// entries scanned per chunk. + entries: usize, + /// staged output bytes per chunk. + /// + /// bounds peak write batch memory when converting small input keys/references + /// into large staged bodies (e.g. v10 record body expansion). + bytes: usize, +} + +impl ChunkBudget { + /// what a pass should use unless it has a measured reason not to. + #[allow(dead_code)] + const DEFAULT: Self = Self { + entries: 100_000, + bytes: 64 * 1024 * 1024, + }; +} + +/// one resumable key-ordered scan over a single keyspace. +/// +/// the framework owns iteration, chunk commits, and the resume cursor; a pass +/// supplies only [`Pass::visit`]. +/// +/// # visit must be idempotent +/// +/// the resume cursor advances only on a committed chunk, so a crash re-visits +/// every entry of the in-flight chunk. +/// +/// more subtly, a pass that writes into the keyspace it scans can observe its +/// own output. each chunk reads a fresh snapshot, so output that sorts *after* +/// the resume cursor is scanned again in a later chunk. `visit` must therefore +/// be safe to apply twice, and re-visiting its own output must reach a fixpoint +/// — normally by making the migrated form recognizable and doing nothing when it +/// is seen. a transform whose output sorts before its input (as key-shortening +/// rewrites usually do) is never rescanned, but that is a property of the +/// encoding, not a guarantee of the framework. +struct Pass { + /// stable cursor identity. THIS SHOULD ALWAYS BE STABLE. DO NOT CHANGE. + name: &'static str, + /// on-disk name of the keyspace to scan. + /// + /// a pass whose keyspace is not compiled into this build is skipped, so a + /// migration can name mode-specific keyspaces unconditionally. + /// + /// a pass scanning `counts` also sees its own resume cursor, which sorts + /// between the `k|` and `r|` prefixes; `visit` must ignore keys under + /// [`MIGRATION_CURSOR_PREFIX`]. + scan: &'static str, + visit: VisitFn, + budget: ChunkBudget, +} /// ordered list of schema migrations. /// @@ -26,26 +117,113 @@ type MigrationFn = fn(&Db, &mut OwnedWriteBatch) -> Result<()>; /// - if a stored type changes shape, define a new versioned type in /// `src/types.rs`, export the newest one for live code, and have the new /// migration explicitly deserialize the previous version and write the new one +/// - a migration that stages one write per entry of a high-cardinality keyspace +/// must be [`Migration::Chunked`]; [`Migration::Atomic`] would buffer the +/// whole rewrite in memory /// /// example: if `RepoState` changes again after `types::v7::RepoState`, add a new /// `types::v8::RepoState`, keep `v7` frozen, and add a migration that rewrites /// `v7 -> v8`. /// ordered list of migrations. migration at index `i` upgrades the schema from version `i` to `i+1`. -const MIGRATIONS: &[(&str, MigrationFn)] = &[ - ("stable_firehose_cursors", v1::stable_firehose_cursors), - ("repo_state_root_commit", v2::repo_state_root_commit), - ("firehose_source_is_pds", v3::firehose_source_is_pds), - ("repo_state_active", v4::repo_state_active), - ("pds_meta_layout", v5::pds_meta_layout), - ("rebuild_pds_account_counts", v6::rebuild_pds_account_counts), - ("repo_state_event_clocks", v7::repo_state_event_clocks), - ("rebuild_lifecycle_counts", v8::rebuild_lifecycle_counts), - ("migrate_excludes_and_pds_keys", v9::migrate_v9), +const MIGRATIONS: &[(&str, Migration)] = &[ + ( + "stable_firehose_cursors", + Migration::Atomic(v1::stable_firehose_cursors), + ), + ( + "repo_state_root_commit", + Migration::Atomic(v2::repo_state_root_commit), + ), + ( + "firehose_source_is_pds", + Migration::Atomic(v3::firehose_source_is_pds), + ), + ( + "repo_state_active", + Migration::Atomic(v4::repo_state_active), + ), + ("pds_meta_layout", Migration::Atomic(v5::pds_meta_layout)), + ( + "rebuild_pds_account_counts", + Migration::Atomic(v6::rebuild_pds_account_counts), + ), + ( + "repo_state_event_clocks", + Migration::Atomic(v7::repo_state_event_clocks), + ), + ( + "rebuild_lifecycle_counts", + Migration::Atomic(v8::rebuild_lifecycle_counts), + ), + ( + "migrate_excludes_and_pds_keys", + Migration::Atomic(v9::migrate_v9), + ), ]; -#[cfg(test)] pub(crate) const LATEST_VERSION: u64 = MIGRATIONS.len() as u64; +/// resume cursors for chunked passes, in the `counts` keyspace: +/// `mig|{version}|{pass}` -> last committed key. +/// +/// the `counts` compaction filter only drops `r|` and old `d|` entries, so these +/// survive compaction. +const MIGRATION_CURSOR_PREFIX: &[u8] = b"mig|"; + +/// applied markers for migrations, in the `counts` keyspace: +/// `mig_applied|{name}` -> `b"1"`. +const MIGRATION_APPLIED_PREFIX: &[u8] = b"mig_applied|"; + +fn pass_cursor_key(version: usize, pass: &str) -> Vec { + let mut key = Vec::with_capacity(MIGRATION_CURSOR_PREFIX.len() + 8 + 1 + pass.len()); + key.extend_from_slice(MIGRATION_CURSOR_PREFIX); + key.extend_from_slice(&(version as u64).to_be_bytes()); + key.push(SEP); + key.extend_from_slice(pass.as_bytes()); + key +} + +fn migration_applied_key(name: &str) -> Vec { + let mut key = Vec::with_capacity(MIGRATION_APPLIED_PREFIX.len() + name.len()); + key.extend_from_slice(MIGRATION_APPLIED_PREFIX); + key.extend_from_slice(name.as_bytes()); + key +} + +/// 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. +#[cfg(test)] +pub(crate) fn rewind_version_for_test(db: &Db, version: u64) -> Result<()> { + let mut batch = db.inner.batch(); + for (name, _) in MIGRATIONS { + batch.remove(&db.counts, migration_applied_key(name)); + } + batch.insert(&db.counts, VERSIONING_KEY, version.to_be_bytes()); + batch.commit().into_diagnostic() +} + +fn is_migration_applied(db: &Db, index: usize, name: &str, stored_version: u64) -> Result { + let key = migration_applied_key(name); + if db.counts.contains_key(&key).into_diagnostic()? { + return Ok(true); + } + + // Legacy versions predate per-migration applied keys. Their global version + // proves ordinary migrations ran, while feature-gated migrations need a + // cheap migration-specific proof: reduced-feature builds advanced the old + // version even when they skipped the work. + if (index as u64) < stored_version { + return match name { + "rebuild_lifecycle_counts" => v8::legacy_migration_applied(db), + "migrate_excludes_and_pds_keys" => Ok(true), + _ => Ok(true), + }; + } + + Ok(false) +} + fn read_version(db: &Db) -> Result { db.counts .get(VERSIONING_KEY) @@ -61,20 +239,557 @@ fn read_version(db: &Db) -> Result { .map(|v| v.unwrap_or(0)) } +struct ChunkOutcome { + /// entries visited and committed. + seen: usize, + /// the scan reached the end of the keyspace. + exhausted: bool, +} + +/// visit one chunk of `pass`, committing the transformed entries and the +/// advanced resume cursor as one atomic batch. +fn run_chunk(db: &Db, cursor_key: &[u8], ks: &Keyspace, pass: &Pass) -> Result { + let resume = db.counts.get(cursor_key).into_diagnostic()?; + let start = resume + .as_ref() + .map_or(Bound::Unbounded, |k| Bound::Excluded(k.to_vec())); + + let mut batch = db.inner.batch(); + let mut seen = 0usize; + let mut bytes = 0usize; + let mut last = None; + let mut exhausted = true; + + { + // a fresh snapshot per chunk keeps the scan consistent while the batch is + // staged, without pinning superseded versions for the whole pass. + let snapshot = db.inner.snapshot(); + for guard in snapshot.range(ks, (start, Bound::>::Unbounded)) { + // checked before visiting, so breaking here proves an entry remains. + // a chunk always visits at least one entry, so neither a degenerate + // budget nor an entry larger than the byte budget can stall the pass + // or make an unfinished scan look exhausted. + if seen > 0 && (seen >= pass.budget.entries || bytes >= pass.budget.bytes) { + exhausted = false; + break; + } + let (key, value) = guard.into_inner().into_diagnostic()?; + let staged_bytes = (pass.visit)(db, &mut batch, &key, &value) + .wrap_err_with(|| format!("migration pass {} failed on an entry", pass.name))?; + bytes += staged_bytes; + last = Some(key); + seen += 1; + } + } + + let Some(last) = last else { + // nothing left to do; leave the cursor for `run` to clear with the + // version bump + return Ok(ChunkOutcome { + seen: 0, + exhausted: true, + }); + }; + + batch.insert(&db.counts, cursor_key, last); + batch.commit().into_diagnostic()?; + + Ok(ChunkOutcome { seen, exhausted }) +} + +/// run one pass to completion, resuming from its stored cursor if present. +fn run_pass(db: &Db, version: usize, pass: &Pass) -> Result { + let Some(ks) = db.keyspace_by_name(pass.scan) else { + // keyspace is not compiled into this build, so there is nothing to migrate + tracing::debug!( + "db: migration pass {} skipped, keyspace {} not in this build", + pass.name, + pass.scan + ); + return Ok(false); + }; + let cursor_key = pass_cursor_key(version, pass.name); + + if db.counts.contains_key(&cursor_key).into_diagnostic()? { + tracing::info!("db: resuming migration pass {} from its cursor", pass.name); + } + + let mut total: u64 = 0; + loop { + let outcome = run_chunk(db, &cursor_key, &ks, pass)?; + total += outcome.seen as u64; + if outcome.seen > 0 { + tracing::info!( + "db: migration pass {} committed {total} entries so far", + pass.name + ); + } + if outcome.exhausted { + break; + } + } + + tracing::info!( + "db: migration pass {} complete ({total} entries)", + pass.name + ); + Ok(true) +} + /// run all pending database migrations in order. /// -/// each migration and its version bump are committed atomically. safe to run on a fresh -/// database (all migrations are no-ops when no relevant data exists). called during [`Db::open`]. +/// an atomic migration commits its write set and applied marker together. a +/// chunked migration commits each chunk with its resume cursor, durably runs +/// any finalizer, then commits cursor cleanup with the applied marker. the +/// global version is derived from the contiguous markers afterward. a crash +/// before the marker therefore resumes from the durable pass cursors. +/// +/// safe to run on a fresh database (all migrations are no-ops when no relevant +/// data exists). called during [`Db::open`]. pub(super) fn run(db: &Db) -> Result<()> { - let version = read_version(db)? as usize; - for (i, (name, migration)) in MIGRATIONS.iter().enumerate().skip(version) { + let stored_version = read_version(db)?; + if stored_version > LATEST_VERSION { + miette::bail!( + "database schema version {stored_version} is newer than supported version {LATEST_VERSION}" + ); + } + for (i, (name, migration)) in MIGRATIONS.iter().enumerate() { + if is_migration_applied(db, i, name, stored_version)? { + let key = migration_applied_key(name); + if !db.counts.contains_key(&key).into_diagnostic()? { + db.counts.insert(key, b"1").into_diagnostic()?; + tracing::info!("db: adopted legacy applied marker for {name}"); + } + tracing::debug!("db: migration {name} already applied"); + continue; + } + tracing::info!("db: running migration {name} (v{i} -> v{})", i + 1); + match migration { + Migration::Atomic(f) => { + let mut batch = db.inner.batch(); + if f(db, &mut batch)? { + batch.insert(&db.counts, migration_applied_key(name), b"1"); + batch.commit().into_diagnostic()?; + tracing::info!("db: migration {name} complete"); + } else { + tracing::info!( + "db: migration {name} skipped (required keyspace not in this build)" + ); + } + } + Migration::Chunked { passes, finalize } => { + let mut all_passed = !(passes.is_empty() && finalize.is_none()); + for pass in *passes { + if !run_pass(db, i, pass)? { + all_passed = false; + break; + } + } + if all_passed + && let Some(finalize) = finalize + && !finalize(db)? + { + all_passed = false; + } + if all_passed { + let mut batch = db.inner.batch(); + for pass in *passes { + batch.remove(&db.counts, pass_cursor_key(i, pass.name)); + } + batch.insert(&db.counts, migration_applied_key(name), b"1"); + batch.commit().into_diagnostic()?; + tracing::info!("db: migration {name} complete"); + } else { + tracing::info!( + "db: migration {name} skipped (required keyspace not in this build)" + ); + } + } + } + } + + // Advance VERSIONING_KEY to highest contiguous applied migration version + let mut max_contiguous_version = 0u64; + for (i, (name, _)) in MIGRATIONS.iter().enumerate() { + if is_migration_applied(db, i, name, stored_version)? { + max_contiguous_version = (i + 1) as u64; + } else { + break; + } + } + if max_contiguous_version != stored_version { let mut batch = db.inner.batch(); - migration(db, &mut batch)?; - let new_version = (i + 1) as u64; - batch.insert(&db.counts, VERSIONING_KEY, new_version.to_be_bytes()); + batch.insert( + &db.counts, + VERSIONING_KEY, + max_contiguous_version.to_be_bytes(), + ); batch.commit().into_diagnostic()?; - tracing::info!("db: migration {name} complete"); } + Ok(()) } + +#[cfg(test)] +mod tests { + use super::*; + + const TEST_PREFIX: &[u8] = b"chunktest|"; + const MIGRATED_PREFIX: &[u8] = b"chunkdone|"; + + /// rewrites `chunktest|{n}` to `chunkdone|{n}`, keeping the value. + /// + /// idempotent: already-migrated keys carry the new prefix and are ignored. + /// note the output sorts *before* the input (`chunkd` < `chunkt`), which is + /// the usual shape for a key rewrite and means output is never rescanned. + fn rewrite_prefix( + db: &Db, + batch: &mut OwnedWriteBatch, + key: &[u8], + value: &[u8], + ) -> Result { + let Some(suffix) = key.strip_prefix(TEST_PREFIX) else { + return Ok(0); + }; + let mut new_key = MIGRATED_PREFIX.to_vec(); + new_key.extend_from_slice(suffix); + let staged_bytes = new_key.len() + value.len(); + batch.insert(&db.cursors, new_key, value); + batch.remove(&db.cursors, key); + Ok(staged_bytes) + } + + fn test_pass(entries: usize) -> Pass { + Pass { + name: "rewrite_prefix", + scan: "cursors", + visit: rewrite_prefix, + budget: ChunkBudget { + entries, + bytes: usize::MAX, + }, + } + } + + fn open_db() -> Result<(tempfile::TempDir, Db)> { + let tmp = tempfile::tempdir().into_diagnostic()?; + let cfg = crate::config::Config { + database_path: tmp.path().to_path_buf(), + ..Default::default() + }; + let db = Db::open(&cfg)?; + Ok((tmp, db)) + } + + fn seed(db: &Db, count: usize) -> Result<()> { + let mut batch = db.inner.batch(); + for i in 0..count { + let mut key = TEST_PREFIX.to_vec(); + key.extend_from_slice(&(i as u64).to_be_bytes()); + batch.insert(&db.cursors, key, (i as u64).to_be_bytes()); + } + batch.commit().into_diagnostic()?; + Ok(()) + } + + fn migrated(db: &Db) -> Result> { + db.cursors + .prefix(MIGRATED_PREFIX) + .map(|guard| { + let (k, v) = guard.into_inner().into_diagnostic()?; + let n = + u64::from_be_bytes(k[MIGRATED_PREFIX.len()..].try_into().into_diagnostic()?); + let v = u64::from_be_bytes(v.as_ref().try_into().into_diagnostic()?); + Ok((n, v)) + }) + .collect() + } + + fn expected(count: usize) -> Vec<(u64, u64)> { + (0..count as u64).map(|i| (i, i)).collect() + } + + #[test] + fn chunked_pass_rewrites_every_entry_across_many_chunks() -> Result<()> { + let (_tmp, db) = open_db()?; + seed(&db, 50)?; + + let pass = test_pass(7); + run_pass(&db, 99, &pass)?; + + assert_eq!(migrated(&db)?, expected(50)); + assert_eq!(db.cursors.prefix(TEST_PREFIX).count(), 0); + Ok(()) + } + + #[test] + fn interrupted_pass_resumes_and_matches_an_uninterrupted_run() -> Result<()> { + let (_tmp, db) = open_db()?; + seed(&db, 50)?; + + let pass = test_pass(7); + let cursor_key = pass_cursor_key(99, pass.name); + let ks = db.keyspace_by_name("cursors").expect("cursors keyspace"); + + // one chunk only, standing in for a crash mid-pass + let first = run_chunk(&db, &cursor_key, &ks, &pass)?; + assert_eq!(first.seen, 7); + assert!(!first.exhausted); + assert_eq!(migrated(&db)?, expected(7)); + assert!(db.counts.contains_key(&cursor_key).into_diagnostic()?); + + // reopening resumes from the cursor rather than restarting + run_pass(&db, 99, &pass)?; + + assert_eq!(migrated(&db)?, expected(50)); + assert_eq!(db.cursors.prefix(TEST_PREFIX).count(), 0); + Ok(()) + } + + #[test] + fn completed_pass_is_a_cheap_no_op_when_rerun() -> Result<()> { + let (_tmp, db) = open_db()?; + seed(&db, 20)?; + + let pass = test_pass(7); + run_pass(&db, 99, &pass)?; + let after_first = migrated(&db)?; + + // the cursor survives until the version bump, so a rerun scans nothing + let cursor_key = pass_cursor_key(99, pass.name); + let ks = db.keyspace_by_name("cursors").expect("cursors keyspace"); + let outcome = run_chunk(&db, &cursor_key, &ks, &pass)?; + assert_eq!(outcome.seen, 0); + assert!(outcome.exhausted); + + assert_eq!(migrated(&db)?, after_first); + Ok(()) + } + + #[test] + fn visiting_migrated_entries_again_is_a_fixpoint() -> Result<()> { + let (_tmp, db) = open_db()?; + seed(&db, 20)?; + + let pass = test_pass(7); + run_pass(&db, 99, &pass)?; + let after_first = migrated(&db)?; + + // drop the cursor to force a full rescan over already-migrated entries, + // the worst case the idempotency contract has to survive + db.counts + .remove(pass_cursor_key(99, pass.name)) + .into_diagnostic()?; + run_pass(&db, 99, &pass)?; + + assert_eq!(migrated(&db)?, after_first); + Ok(()) + } + + /// fails partway through a chunk, standing in for a crash mid-chunk. + fn fail_partway( + db: &Db, + batch: &mut OwnedWriteBatch, + key: &[u8], + value: &[u8], + ) -> Result { + let Some(suffix) = key.strip_prefix(TEST_PREFIX) else { + return Ok(0); + }; + let n = u64::from_be_bytes(suffix.try_into().into_diagnostic()?); + if n == 4 { + miette::bail!("simulated failure at entry 4"); + } + rewrite_prefix(db, batch, key, value) + } + #[test] + fn a_chunk_that_fails_partway_commits_nothing() -> Result<()> { + let (_tmp, db) = open_db()?; + seed(&db, 50)?; + + let pass = Pass { + visit: fail_partway, + ..test_pass(7) + }; + let cursor_key = pass_cursor_key(99, pass.name); + let ks = db.keyspace_by_name("cursors").expect("cursors keyspace"); + + assert!(run_chunk(&db, &cursor_key, &ks, &pass).is_err()); + + // entries 0..=3 were staged before the failure, but the batch is dropped + // uncommitted, so neither they nor the cursor survive + assert_eq!(migrated(&db)?, vec![]); + assert!(!db.counts.contains_key(&cursor_key).into_diagnostic()?); + assert_eq!(db.cursors.prefix(TEST_PREFIX).count(), 50); + + // retrying with a working visit starts over from the beginning + run_pass(&db, 99, &test_pass(7))?; + assert_eq!(migrated(&db)?, expected(50)); + Ok(()) + } + + #[test] + fn byte_budget_ends_a_chunk_before_the_entry_budget() -> Result<()> { + let (_tmp, db) = open_db()?; + seed(&db, 50)?; + + // each entry is 26 scanned bytes (10-byte prefix + 8-byte key + 8-byte + // value), and the budget is checked before visiting, so a 60-byte budget + // admits 3 entries (26, 52, 78) before it trips + let pass = Pass { + budget: ChunkBudget { + entries: usize::MAX, + bytes: 60, + }, + ..test_pass(0) + }; + let cursor_key = pass_cursor_key(99, pass.name); + let ks = db.keyspace_by_name("cursors").expect("cursors keyspace"); + + let outcome = run_chunk(&db, &cursor_key, &ks, &pass)?; + assert_eq!(outcome.seen, 3); + assert!(!outcome.exhausted); + + run_pass(&db, 99, &pass)?; + assert_eq!(migrated(&db)?, expected(50)); + Ok(()) + } + + #[test] + fn a_degenerate_budget_still_makes_progress() -> Result<()> { + let (_tmp, db) = open_db()?; + seed(&db, 5)?; + + // a zero budget would otherwise trip before visiting anything, leaving + // the chunk with no cursor to advance and reporting itself exhausted + let pass = Pass { + budget: ChunkBudget { + entries: 0, + bytes: 0, + }, + ..test_pass(0) + }; + let cursor_key = pass_cursor_key(99, pass.name); + let ks = db.keyspace_by_name("cursors").expect("cursors keyspace"); + + let outcome = run_chunk(&db, &cursor_key, &ks, &pass)?; + assert_eq!(outcome.seen, 1); + assert!(!outcome.exhausted); + + run_pass(&db, 99, &pass)?; + assert_eq!(migrated(&db)?, expected(5)); + Ok(()) + } + + #[test] + fn pass_on_a_keyspace_missing_from_this_build_is_skipped() -> Result<()> { + let (_tmp, db) = open_db()?; + + let pass = Pass { + scan: "not_a_keyspace", + ..test_pass(7) + }; + run_pass(&db, 99, &pass)?; + + let cursor_key = pass_cursor_key(99, pass.name); + assert!(!db.counts.contains_key(&cursor_key).into_diagnostic()?); + Ok(()) + } + + #[test] + fn cursor_keys_are_distinct_per_version_and_pass() { + assert_ne!(pass_cursor_key(1, "a"), pass_cursor_key(2, "a")); + assert_ne!(pass_cursor_key(1, "a"), pass_cursor_key(1, "b")); + assert!(pass_cursor_key(1, "a").starts_with(MIGRATION_CURSOR_PREFIX)); + } + + #[test] + fn staged_byte_budget_bounds_chunk_write_size() -> Result<()> { + let (_tmp, db) = open_db()?; + let mut batch = db.inner.batch(); + // Insert small input keys/values that expand into large staged outputs + for i in 0..10 { + let mut key = TEST_PREFIX.to_vec(); + key.extend_from_slice(&(i as u64).to_be_bytes()); + batch.insert(&db.cursors, key, vec![0u8; 1000]); // 1 KB payload + } + batch.commit().into_diagnostic()?; + + // Set byte budget to 2500 bytes. Each output staged is ~1018 bytes. + // First entry: 1018 bytes. Second entry: 2036 bytes. Third entry: 3054 bytes (>= 2500). + // 4th check trips because 3054 >= 2500. So outcome.seen == 3. + let pass = Pass { + budget: ChunkBudget { + entries: usize::MAX, + bytes: 2500, + }, + ..test_pass(0) + }; + let cursor_key = pass_cursor_key(99, pass.name); + let ks = db.keyspace_by_name("cursors").expect("cursors keyspace"); + + let outcome = run_chunk(&db, &cursor_key, &ks, &pass)?; + assert_eq!(outcome.seen, 3); + assert!(!outcome.exhausted); + Ok(()) + } + + #[test] + fn applied_migration_keys_track_per_migration_status() -> Result<()> { + let (_tmp, db) = open_db()?; + assert!( + db.counts + .contains_key(&migration_applied_key("stable_firehose_cursors")) + .into_diagnostic()? + ); + Ok(()) + } + + #[cfg(feature = "indexer")] + #[test] + fn legacy_marker_adoption_does_not_rerun_lifecycle_rebuild() -> Result<()> { + use jacquard_common::types::string::Did; + + let tmp = tempfile::tempdir().into_diagnostic()?; + let cfg = crate::config::Config { + database_path: tmp.path().to_path_buf(), + ..Default::default() + }; + let repo = Did::new("did:plc:ewvi7nxzyoun6zhxrhs64oiz").into_diagnostic()?; + let stale = Did::new("did:plc:yk4q3id7id6p5z3bypvshc64").into_diagnostic()?; + let mut stale_gauge = crate::db::keys::COUNT_GAUGE_PREFIX.to_vec(); + stale_gauge.extend_from_slice(&crate::db::keys::repo_key(&stale)); + + { + let db = Db::open(&cfg)?; + let mut batch = db.inner.batch(); + batch.insert( + &db.repos, + crate::db::keys::repo_key(&repo), + crate::db::ser_repo_state(&crate::types::RepoState::synced())?, + ); + let mut gauge = crate::db::keys::COUNT_GAUGE_PREFIX.to_vec(); + gauge.extend_from_slice(&crate::db::keys::repo_key(&repo)); + let synced = + crate::db::lifecycle_counts::encode_gauge(crate::types::GaugeState::Synced); + batch.insert(&db.counts, gauge, synced); + // v8 would clear this stale-but-valid row. retaining it proves the + // legacy global version was adopted rather than rebuilt. + batch.insert(&db.counts, &stale_gauge, synced); + batch.commit().into_diagnostic()?; + // v9 is the last schema that existed before this change. + rewind_version_for_test(&db, 9)?; + db.persist()?; + } + + let db = Db::open(&cfg)?; + assert!(db.counts.contains_key(stale_gauge).into_diagnostic()?); + assert!( + db.counts + .contains_key(migration_applied_key("rebuild_lifecycle_counts")) + .into_diagnostic()? + ); + Ok(()) + } +} diff --git a/src/db/migration/v1.rs b/src/db/migration/v1.rs index 7165e40..03b191a 100644 --- a/src/db/migration/v1.rs +++ b/src/db/migration/v1.rs @@ -7,7 +7,7 @@ use url::Url; use crate::db::{Db, keys}; /// migrates firehose cursors from `firehose_cursor|{url}` to `firehose_cursor|{host}`. -pub(super) fn stable_firehose_cursors(db: &Db, batch: &mut OwnedWriteBatch) -> Result<()> { +pub(super) fn stable_firehose_cursors(db: &Db, batch: &mut OwnedWriteBatch) -> Result { let prefix = keys::FIREHOSE_CURSOR_PREFIX; for item in db.cursors.prefix(prefix) { let (old_key, value) = item.into_inner().into_diagnostic()?; @@ -29,5 +29,5 @@ pub(super) fn stable_firehose_cursors(db: &Db, batch: &mut OwnedWriteBatch) -> R batch.remove(&db.cursors, old_key); } - Ok(()) + Ok(true) } diff --git a/src/db/migration/v2.rs b/src/db/migration/v2.rs index fa2033d..0317261 100644 --- a/src/db/migration/v2.rs +++ b/src/db/migration/v2.rs @@ -31,7 +31,7 @@ pub(crate) struct OldRepoState<'i> { pub handle: Option>, } -pub(super) fn repo_state_root_commit(db: &Db, batch: &mut OwnedWriteBatch) -> Result<()> { +pub(super) fn repo_state_root_commit(db: &Db, batch: &mut OwnedWriteBatch) -> Result { for item in db.repos.iter() { let (k, v) = item.into_inner().into_diagnostic()?; let old: OldRepoState = rmp_serde::from_slice(&v) @@ -66,5 +66,5 @@ pub(super) fn repo_state_root_commit(db: &Db, batch: &mut OwnedWriteBatch) -> Re ); } - Ok(()) + Ok(true) } diff --git a/src/db/migration/v3.rs b/src/db/migration/v3.rs index 091bb5e..384c27f 100644 --- a/src/db/migration/v3.rs +++ b/src/db/migration/v3.rs @@ -39,7 +39,7 @@ fn is_known_relay(url_bytes: &[u8]) -> bool { .is_some_and(|h| KNOWN_RELAY_HOSTS.contains(&h)) } -pub(super) fn firehose_source_is_pds(db: &Db, batch: &mut OwnedWriteBatch) -> Result<()> { +pub(super) fn firehose_source_is_pds(db: &Db, batch: &mut OwnedWriteBatch) -> Result { let relay_bytes = rmp_serde::to_vec(&FirehoseSourceMeta { is_pds: false }) .map_err(|e| miette::miette!("failed to serialize meta: {e}"))?; let pds_bytes = rmp_serde::to_vec(&FirehoseSourceMeta { is_pds: true }) @@ -59,5 +59,5 @@ pub(super) fn firehose_source_is_pds(db: &Db, batch: &mut OwnedWriteBatch) -> Re batch.insert(&db.crawler, key, value); } - Ok(()) + Ok(true) } diff --git a/src/db/migration/v4.rs b/src/db/migration/v4.rs index bb28b1d..c4f535c 100644 --- a/src/db/migration/v4.rs +++ b/src/db/migration/v4.rs @@ -21,7 +21,7 @@ pub(crate) struct OldRepoState<'i> { pub handle: Option>, } -pub(super) fn repo_state_active(db: &Db, batch: &mut OwnedWriteBatch) -> Result<()> { +pub(super) fn repo_state_active(db: &Db, batch: &mut OwnedWriteBatch) -> Result { for item in db.repos.iter() { let (k, v) = item.into_inner().into_diagnostic()?; let old: OldRepoState = rmp_serde::from_slice(&v) @@ -85,5 +85,5 @@ pub(super) fn repo_state_active(db: &Db, batch: &mut OwnedWriteBatch) -> Result< ); } - Ok(()) + Ok(true) } diff --git a/src/db/migration/v5.rs b/src/db/migration/v5.rs index a506a91..e97ebcc 100644 --- a/src/db/migration/v5.rs +++ b/src/db/migration/v5.rs @@ -7,7 +7,7 @@ pub mod v4 { pub const PDS_BANNED_PREFIX: &[u8] = b"pb|"; } -pub(crate) fn pds_meta_layout(db: &Db, batch: &mut OwnedWriteBatch) -> Result<()> { +pub(crate) fn pds_meta_layout(db: &Db, batch: &mut OwnedWriteBatch) -> Result { for guard in db.filter.prefix(v4::PDS_TIER_PREFIX) { let (k, v) = guard.into_inner().into_diagnostic()?; let host = std::str::from_utf8(&k[v4::PDS_TIER_PREFIX.len()..]) @@ -34,5 +34,5 @@ pub(crate) fn pds_meta_layout(db: &Db, batch: &mut OwnedWriteBatch) -> Result<() batch.remove(&db.filter, k); } - Ok(()) + Ok(true) } diff --git a/src/db/migration/v6.rs b/src/db/migration/v6.rs index 21a05f7..feeb803 100644 --- a/src/db/migration/v6.rs +++ b/src/db/migration/v6.rs @@ -8,7 +8,7 @@ use miette::{Context, IntoDiagnostic, Result}; use smol_str::SmolStr; use url::Url; -pub(crate) fn rebuild_pds_account_counts(db: &Db, batch: &mut OwnedWriteBatch) -> Result<()> { +pub(crate) fn rebuild_pds_account_counts(db: &Db, batch: &mut OwnedWriteBatch) -> Result { let mut counts: BTreeMap = BTreeMap::new(); for guard in db.repos.iter() { @@ -54,5 +54,5 @@ pub(crate) fn rebuild_pds_account_counts(db: &Db, batch: &mut OwnedWriteBatch) - set_ks_count(batch, db, &keys::pds_account_count_key(&host), count); } - Ok(()) + Ok(true) } diff --git a/src/db/migration/v7.rs b/src/db/migration/v7.rs index 792b14e..23aa8e9 100644 --- a/src/db/migration/v7.rs +++ b/src/db/migration/v7.rs @@ -4,7 +4,7 @@ use miette::{Context, IntoDiagnostic, Result}; use crate::db::Db; use crate::types::{RepoState, v4}; -pub(crate) fn repo_state_event_clocks(db: &Db, batch: &mut OwnedWriteBatch) -> Result<()> { +pub(crate) fn repo_state_event_clocks(db: &Db, batch: &mut OwnedWriteBatch) -> Result { for guard in db.repos.iter() { let (key, value) = guard.into_inner().into_diagnostic()?; let old: v4::RepoState = rmp_serde::from_slice(value.as_ref()) @@ -33,7 +33,7 @@ pub(crate) fn repo_state_event_clocks(db: &Db, batch: &mut OwnedWriteBatch) -> R ); } - Ok(()) + Ok(true) } #[cfg(test)] @@ -71,8 +71,8 @@ mod tests { let mut batch = db.inner.batch(); batch.insert(&db.repos, &repo_key, &legacy_bytes); - batch.insert(&db.counts, keys::VERSIONING_KEY, 6_u64.to_be_bytes()); batch.commit().into_diagnostic()?; + crate::db::migration::rewind_version_for_test(&db, 6)?; db.persist()?; legacy_bytes diff --git a/src/db/migration/v8.rs b/src/db/migration/v8.rs index ef4ed0f..3040158 100644 --- a/src/db/migration/v8.rs +++ b/src/db/migration/v8.rs @@ -11,7 +11,10 @@ use crate::db::{deser_repo_meta, keys, set_ks_count}; use crate::types::{GaugeState, ResyncErrorKind}; #[cfg(feature = "indexer")] -pub(crate) fn rebuild_lifecycle_counts(db: &Db, batch: &mut OwnedWriteBatch) -> Result<()> { +pub(crate) fn rebuild_lifecycle_counts(db: &Db, batch: &mut OwnedWriteBatch) -> Result { + if db.keyspace_by_name("repo_metadata").is_none() { + return Ok(false); + } for guard in db.counts.prefix(keys::COUNT_GAUGE_PREFIX) { let key = guard.key().into_diagnostic()?; batch.remove(&db.counts, key); @@ -65,12 +68,34 @@ pub(crate) fn rebuild_lifecycle_counts(db: &Db, batch: &mut OwnedWriteBatch) -> set_ks_count(batch, db, "error_transport", error_transport); set_ks_count(batch, db, "error_generic", error_generic); - Ok(()) + Ok(true) +} + +/// whether the old atomic v8 migration committed in a legacy global-version +/// database. one gauge row was written for every repo in the same batch as the +/// version bump, so any row proves the all-or-nothing migration ran; an empty +/// repo set had no work to do. +#[cfg(feature = "indexer")] +pub(super) fn legacy_migration_applied(db: &Db) -> Result { + if db.repos.is_empty().into_diagnostic()? { + return Ok(true); + } + db.counts + .prefix(crate::db::keys::COUNT_GAUGE_PREFIX) + .next() + .map(|guard| guard.into_inner().into_diagnostic().map(|_| true)) + .transpose() + .map(|found| found.unwrap_or(false)) +} + +#[cfg(not(feature = "indexer"))] +pub(super) fn legacy_migration_applied(_db: &Db) -> Result { + Ok(false) } #[cfg(not(feature = "indexer"))] -pub(crate) fn rebuild_lifecycle_counts(_db: &Db, _batch: &mut OwnedWriteBatch) -> Result<()> { - Ok(()) +pub(crate) fn rebuild_lifecycle_counts(_db: &Db, _batch: &mut OwnedWriteBatch) -> Result { + Ok(false) } #[cfg(feature = "indexer")] @@ -227,8 +252,8 @@ mod tests { keys::count_delta_key(2, "error_ratelimited"), 10_i64.to_be_bytes(), ); - batch.insert(&db.counts, keys::VERSIONING_KEY, 7_u64.to_be_bytes()); batch.commit().into_diagnostic()?; + crate::db::migration::rewind_version_for_test(&db, 7)?; db.persist()?; } diff --git a/src/db/migration/v9.rs b/src/db/migration/v9.rs index d5c1c8f..f02d639 100644 --- a/src/db/migration/v9.rs +++ b/src/db/migration/v9.rs @@ -7,7 +7,7 @@ use crate::db::Db; use {crate::db::types::TrimmedDid, jacquard_common::types::did::Did, miette::IntoDiagnostic}; #[cfg(feature = "indexer")] -pub(crate) fn migrate_v9(db: &Db, batch: &mut OwnedWriteBatch) -> Result<()> { +pub(crate) fn migrate_v9(db: &Db, batch: &mut OwnedWriteBatch) -> Result { // 1. Migrate excludes let exclude_prefix = [crate::db::filter::EXCLUDE_PREFIX, crate::db::keys::SEP]; for guard in db.filter.prefix(exclude_prefix) { @@ -54,12 +54,12 @@ pub(crate) fn migrate_v9(db: &Db, batch: &mut OwnedWriteBatch) -> Result<()> { } } - Ok(()) + Ok(true) } #[cfg(not(feature = "indexer"))] -pub(crate) fn migrate_v9(_db: &Db, _batch: &mut OwnedWriteBatch) -> Result<()> { - Ok(()) +pub(crate) fn migrate_v9(_db: &Db, _batch: &mut OwnedWriteBatch) -> Result { + Ok(false) } #[cfg(all(test, feature = "indexer"))] @@ -103,12 +103,8 @@ mod tests { let legacy_tier_key = "example.com|tier"; batch.insert(&db.filter, legacy_tier_key.as_bytes(), b"tier1"); - batch.insert( - &db.counts, - crate::db::keys::VERSIONING_KEY, - 8_u64.to_be_bytes(), - ); batch.commit().into_diagnostic()?; + crate::db::migration::rewind_version_for_test(&db, 8)?; db.persist()?; }