From 79ae05d1e56157aa801f4c4c5e0d121e1ae865ac Mon Sep 17 00:00:00 2001 From: dawn <90008@klbr.net> Date: Thu, 6 Aug 2026 14:44:42 +0300 Subject: [PATCH] [db] write v10 bodies into the compressed heads keyspace v9 created records with compression disabled at every level and fjall persists keyspace options on disk, so bodies written there stay uncompressed no matter what the schema says at open (measured on the prod db copy: 71.3MB stored vs 16MB zstd:3 standalone on the same bodies). the schema registers the store under a new name, heads, with zstd at every level, and v10 stages rewritten bodies straight into it, retiring both legacy stores in its finalize. chunked passes now resolve scan keyspaces through an on-disk fallback for stores the registry has dropped, pass the scanned keyspace to visit so in-place rewrites keep writing the store they read, and treat a never-existing source as a no-op when the pass is optional_on_absent. verified on a copy of the prod database: 136,311 entries migrated in ~500ms, heads compacts to 13.1MB (5.4x smaller), record resolution unchanged. --- AGENTS.md | 6 +- docs/specs/block-gc.md | 4 +- src/db/keyspaces.rs | 2 +- src/db/migration/mod.rs | 80 ++++++++++++++++-- src/db/migration/v10.rs | 149 +++++++++++++++++++++++++++------ src/db/migration/v10/events.rs | 3 + src/db/registry.rs | 6 +- src/db/schema.rs | 31 ++++--- tests/debug_endpoints.nu | 6 +- 9 files changed, 229 insertions(+), 58 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index 7172914..b0fa2be 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -96,7 +96,7 @@ Modes (`indexer`, `relay`) are compile-time choices and mutually exclusive (enfo ### Storage and serialization - **State**: Use `rmp-serde` (MessagePack) for all internal state (`RepoState`, `ErrorState`, `StoredEvent`). -- **Blocks**: Record bodies are stored inline in `records`; superseded bodies move to `history`. v10 materializes unresolved legacy event bodies into `event_bodies`, then deletes the old `blocks` keyspace. +- **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`). - **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. @@ -112,10 +112,10 @@ Modes (`indexer`, `relay`) are compile-time choices and mutually exclusive (enfo Hydrant uses multiple `fjall` keyspaces: - `repos`: Maps `{DID}` -> `RepoState` (MessagePack). - `repo_metadata`: Maps `rm|{DID}` -> `RepoMetadata` (MessagePack). -- `records`: Maps `{DID} 00 {COL} 2f {RKey}` (MST order, trimmed DID) -> record body (raw DAG-CBOR). Operator-redacted heads retain only their raw CID bytes; links-only mode stores CIDs for every head. Values were CIDs before v10. +- `heads`: Maps `{DID} 00 {COL} 2f {RKey}` (MST order, trimmed DID) -> record body (raw DAG-CBOR), zstd-compressed at every level. Operator-redacted heads retain only their raw CID bytes; links-only mode stores CIDs for every head. Replaces v9's `records` keyspace in the v10 migration (v9 persisted "never compress" into the old keyspace's on-disk config); values were CIDs before v10. - `history`: Maps `{record key} 00 {death_rev}` (rev = 8-byte BE TID) -> superseded record body; raw CID bytes = operator-redacted body; empty value = delete tombstone ("deleted at rev", vs never existed). Retention is bounded by `HYDRANT_HISTORY_TTL` via a key-only compaction filter (`db::compaction::HistoryRetentionFilterFactory`); unset = permanent. Ephemeral mode stores no history. - `redactions`: Maps `{record key} 00 {CID bytes}` -> empty for operator-erased versions. Markers are permanent and suppress the same record+CID on replay even after an intervening version or history compaction. -- `event_bodies` (indexer stream): Finite compatibility archive, `{record key} 00 {CID bytes}` -> raw DAG-CBOR. v10 uses it only for legacy events that cannot resolve through `records` or `history`; new writes never add entries. +- `event_bodies` (indexer stream): Finite compatibility archive, `{record key} 00 {CID bytes}` -> raw DAG-CBOR. v10 uses it only for legacy events that cannot resolve through `heads` or `history`; new writes never add entries. - `events`: Maps `{ID}` (u64 BE) -> `StoredEvent` (MessagePack). This is the source for the JSON stream API. Ephemeral mode stores create/update bodies inline as `StoredData::Block`, so event TTL removes identity and body bytes together without creating record heads or history. - `cursors`: Maps per-relay cursor keys -> `Value` (u64/i64 BE Bytes). Keys: `firehose_cursor|{relay}`, `crawler_cursor|{relay}`, `by_collection_cursor|{url}|{collection}`. - `pending`: Queue of `{ID}` (u64 BE) -> `Empty` (Backfill queue). diff --git a/docs/specs/block-gc.md b/docs/specs/block-gc.md index ad516c5..a768811 100644 --- a/docs/specs/block-gc.md +++ b/docs/specs/block-gc.md @@ -6,7 +6,9 @@ Tracking: `hydrant-6vo` (description predates this framing); related `hydrant-71 ## Decision update -The permanent layout below landed: current bodies live inline in `records`, +The permanent layout below landed: current bodies live inline in `heads` +(the v10 migration builds that keyspace with zstd at every level — v9 had +persisted "never compress" into the old `records` keyspace's on-disk config), superseded bodies move to `history`, and `HYDRANT_HISTORY_TTL` optionally bounds that history. The queryable-ephemeral extension did not land. Ephemeral remains the shipped separate architecture: it stores create/update bodies only inside diff --git a/src/db/keyspaces.rs b/src/db/keyspaces.rs index 0e825fc..1b95773 100644 --- a/src/db/keyspaces.rs +++ b/src/db/keyspaces.rs @@ -52,7 +52,7 @@ impl OpenCx<'_> { #[derive(Clone)] pub struct IndexerDb { /// maps `{DID} 00 {COL} 2f {RKey}` -> DAG-CBOR body or redacted CID marker - pub(super) records: Ks, + pub(super) records: Ks, /// superseded bodies/markers: `{record key} 00 {death_rev}` -> body/CID/empty pub(super) history: Ks, /// operator-redacted versions: `{record key} 00 {cid}` -> empty diff --git a/src/db/migration/mod.rs b/src/db/migration/mod.rs index 866f2de..2b07e7f 100644 --- a/src/db/migration/mod.rs +++ b/src/db/migration/mod.rs @@ -38,6 +38,30 @@ fn legacy_blocks_for_test(db: &Db) -> Result { .into_diagnostic() } +/// the pre-v12 record store. v12's schema drops `records` from the registry, +/// but its on-disk tables persist until the `recompress_record_heads` +/// migration copies them into `heads` and deletes the keyspace. +#[cfg(feature = "indexer")] +const LEGACY_RECORDS_KEYSPACE: &str = "records"; + +#[cfg(all(test, feature = "indexer"))] +fn legacy_records_for_test(db: &Db) -> Result { + db.inner + .keyspace(LEGACY_RECORDS_KEYSPACE, Default::default) + .into_diagnostic() +} + +#[cfg(feature = "indexer")] +fn legacy_records(db: &Db) -> Result> { + if !db.inner.keyspace_exists(LEGACY_RECORDS_KEYSPACE) { + return Ok(None); + } + db.inner + .keyspace(LEGACY_RECORDS_KEYSPACE, Default::default) + .into_diagnostic() + .map(Some) +} + mod v1; mod v10; mod v2; @@ -55,7 +79,13 @@ 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; +/// visit one scanned entry, staging its rewrite into the chunk batch, and +/// return the staged byte count for the chunk budget. +/// +/// the third argument is the keyspace being scanned: in-place passes stage +/// into it, so a pass keeps writing the same store it reads even when that +/// store is a legacy keyspace the registry no longer knows. +type VisitFn = fn(&Db, &mut OwnedWriteBatch, &Keyspace, &[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 @@ -127,8 +157,14 @@ struct Pass { 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. + /// the registry resolves names from this build first; a name it no longer + /// knows falls back to opening an existing on-disk keyspace, so a pass can + /// still finish migrating data out of a store a later schema dropped. + /// + /// a pass whose keyspace exists nowhere is skipped, so a migration can + /// name mode-specific keyspaces unconditionally — unless + /// `optional_on_absent` is set, in which case absence means the database + /// never had the source layout and the pass completes as a no-op. /// /// a pass scanning `counts` also sees its own resume cursor, which sorts /// between the `k|` and `r|` prefixes; `visit` must ignore keys under @@ -136,6 +172,7 @@ struct Pass { scan: &'static str, visit: VisitFn, budget: ChunkBudget, + optional_on_absent: bool, } /// ordered list of schema migrations. @@ -194,7 +231,7 @@ const MIGRATIONS: &[(&str, Migration)] = &[ "inline_record_bodies", Migration::Chunked { passes: v10::PASSES, - finalize: Some(v10::retire_blocks), + finalize: Some(v10::retire_legacy_stores), }, ), ]; @@ -320,7 +357,7 @@ fn run_chunk(db: &Db, cursor_key: &[u8], ks: &Keyspace, pass: &Pass) -> Result Result Result { - let Some(ks) = db.keyspace_by_name(pass.scan) else { + let ks = match db.keyspace_by_name(pass.scan) { + Some(ks) => Some(ks), + // a keyspace dropped from the registry by a later schema may still + // exist on disk; open it so the pass can finish migrating its data + None if db.inner.keyspace_exists(pass.scan) => Some( + db.inner + .keyspace(pass.scan, Default::default) + .into_diagnostic()?, + ), + None => None, + }; + let Some(ks) = ks else { + if pass.optional_on_absent { + // the source layout never existed on this database, so there is + // genuinely nothing to migrate + tracing::debug!( + "db: migration pass {} found no {} keyspace, nothing to migrate", + pass.name, + pass.scan + ); + return Ok(true); + } // keyspace is not compiled into this build, so there is nothing to migrate tracing::debug!( "db: migration pass {} skipped, keyspace {} not in this build", @@ -492,6 +550,7 @@ mod tests { fn rewrite_prefix( db: &Db, batch: &mut OwnedWriteBatch, + scanned: &Keyspace, key: &[u8], value: &[u8], ) -> Result { @@ -500,9 +559,10 @@ mod tests { }; let mut new_key = MIGRATED_PREFIX.to_vec(); new_key.extend_from_slice(suffix); + let _ = db; let staged_bytes = new_key.len() + value.len(); - batch.insert(&db.cursors, new_key, value); - batch.remove(&db.cursors, key); + batch.insert(scanned, new_key, value); + batch.remove(scanned, key); Ok(staged_bytes) } @@ -515,6 +575,7 @@ mod tests { entries, bytes: usize::MAX, }, + optional_on_absent: false, } } @@ -637,6 +698,7 @@ mod tests { fn fail_partway( db: &Db, batch: &mut OwnedWriteBatch, + scanned: &Keyspace, key: &[u8], value: &[u8], ) -> Result { @@ -647,7 +709,7 @@ mod tests { if n == 4 { miette::bail!("simulated failure at entry 4"); } - rewrite_prefix(db, batch, key, value) + rewrite_prefix(db, batch, scanned, key, value) } #[test] fn a_chunk_that_fails_partway_commits_nothing() -> Result<()> { diff --git a/src/db/migration/v10.rs b/src/db/migration/v10.rs index 1861cf1..b648d1a 100644 --- a/src/db/migration/v10.rs +++ b/src/db/migration/v10.rs @@ -1,16 +1,22 @@ -//! v10: inline record bodies and retire the legacy blocks CAS. +//! v10: inline record bodies into the compressed `heads` keyspace, and retire +//! the legacy `records` store and `blocks` CAS. //! -//! legacy `{DID}|{COL}|{tag}{rkey}` -> cid entries become -//! `{DID} 00 {COL} 2f {rkey text}` -> body, resolved from the blocks CAS. -//! new keys sort before legacy keys within a repo (NUL < `|` after the -//! trimmed DID), so the resuming chunked scan only ever meets legacy keys. +//! legacy `{DID}|{COL}|{tag}{rkey}` -> cid entries in `records` become +//! `{DID} 00 {COL} 2f {rkey text}` -> body in `heads`, resolved from the +//! blocks CAS. the schema drops `records` from the registry for this: v9 +//! persisted "never compress" into its on-disk config (fair when values were +//! cids), and fjall reads keyspace options from disk at open, so bodies could +//! only ever be compressed in a fresh keyspace. the pass scans `records` +//! through the on-disk fallback and stages into a different keyspace than it +//! scans, so chunks never re-observe their own output and re-inserting the +//! same key/value is a fixpoint. //! //! the same bounded migration archives otherwise-unresolved legacy permanent -//! pointer bodies before deleting `blocks`. ephemeral inline events stay -//! inline so their TTL still removes body bytes atomically. no ingestion can -//! run during `Db::open`, so every pointer visible to this migration is -//! necessarily legacy; an event watermark would only model versions that -//! never shipped. +//! pointer bodies before deleting `blocks` and `records`. ephemeral inline +//! events stay inline so their TTL still removes body bytes atomically. no +//! ingestion can run during `Db::open`, so every pointer visible to this +//! migration is necessarily legacy; an event watermark would only model +//! versions that never shipped. mod events; @@ -43,6 +49,7 @@ pub(super) const PASSES: &[Pass] = &[ scan: "filter", visit: rewrite_filter_key, budget: ChunkBudget::DEFAULT, + optional_on_absent: false, }, #[cfg(feature = "indexer")] Pass { @@ -53,6 +60,8 @@ pub(super) const PASSES: &[Pass] = &[ entries: 100_000, bytes: DEFAULT_RECORD_CHUNK_BYTES, }, + // v12 databases register `heads` instead and never had `records` + optional_on_absent: true, }, #[cfg(feature = "indexer_stream")] events::POINTER_EVENT_BODIES, @@ -73,8 +82,9 @@ fn is_current_key(key: &[u8]) -> bool { #[cfg(feature = "indexer")] fn rewrite_filter_key( - db: &Db, + _db: &Db, batch: &mut OwnedWriteBatch, + scanned: &fjall::Keyspace, key: &[u8], value: &[u8], ) -> Result { @@ -88,8 +98,8 @@ fn rewrite_filter_key( let mut new_key = exclude_prefix.to_vec(); trimmed.write_to_vec(&mut new_key); let staged = new_key.len(); - batch.insert(&db.filter, new_key, []); - batch.remove(&db.filter, key); + batch.insert(scanned, new_key, []); + batch.remove(scanned, key); return Ok(staged); } @@ -110,12 +120,31 @@ fn rewrite_filter_key( new_key.extend_from_slice(current_prefix); new_key.extend_from_slice(host); let staged = new_key.len() + value.len(); - batch.insert(&db.filter, new_key, value); - batch.remove(&db.filter, key); + batch.insert(scanned, new_key, value); + batch.remove(scanned, key); Ok(staged) } -pub(super) fn retire_blocks(db: &crate::db::Db) -> miette::Result { +#[cfg(feature = "indexer")] +/// delete the legacy `records` keyspace only after every copy chunk and its +/// resume cursor are durable. chunk commits use the data WAL while keyspace +/// deletion updates fjall's metadata independently, so force the copy durable +/// before making the old store unreachable. fjall removes the on-disk +/// directory on the next open, when recovery finds it unreferenced. +pub(super) fn retire_legacy_stores(db: &crate::db::Db) -> miette::Result { + if !events::retire_blocks(db)? { + return Ok(false); + } + if let Some(records) = super::legacy_records(db)? { + db.persist()?; + db.inner.delete_keyspace(records).into_diagnostic()?; + db.persist()?; + } + Ok(true) +} + +#[cfg(not(feature = "indexer"))] +pub(super) fn retire_legacy_stores(db: &crate::db::Db) -> miette::Result { events::retire_blocks(db) } @@ -158,9 +187,18 @@ fn parse_legacy_key(key: &[u8]) -> Result<(usize, &str, DbRkey)> { } #[cfg(feature = "indexer")] -fn rewrite_record(db: &Db, batch: &mut OwnedWriteBatch, key: &[u8], value: &[u8]) -> Result { +fn rewrite_record( + db: &Db, + batch: &mut OwnedWriteBatch, + _scanned: &fjall::Keyspace, + key: &[u8], + value: &[u8], +) -> Result { + // already-current key (defense against out-of-order surprises): the body + // is final, copy it verbatim into heads if is_current_key(key) { - return Ok(0); + batch.insert(&db.indexer.records, key, value); + return Ok(key.len() + value.len()); } let (did_len, collection, rkey) = parse_legacy_key(key)?; @@ -205,7 +243,6 @@ fn rewrite_record(db: &Db, batch: &mut OwnedWriteBatch, key: &[u8], value: &[u8] let staged_bytes = new_value.len(); batch.insert(&db.indexer.records, new_key, new_value); - batch.remove(&db.indexer.records, key); Ok(staged_bytes) } @@ -279,13 +316,14 @@ mod tests { let db = Db::open(&cfg)?; let mut batch = db.inner.batch(); let blocks = super::super::legacy_blocks_for_test(&db)?; + let records = super::super::legacy_records_for_test(&db)?; let post_legacy = legacy_record_key(&did(), "app.bsky.feed.post", &tid_rkey); let profile_legacy = legacy_record_key(&did(), "app.bsky.actor.profile", &str_rkey); let orphan_legacy = legacy_record_key(&did(), "app.bsky.feed.like", &tid_rkey); - batch.insert(&db.indexer.records, &post_legacy, &cid1); - batch.insert(&db.indexer.records, &profile_legacy, &cid2); - batch.insert(&db.indexer.records, &orphan_legacy, &cid3); + batch.insert(&records, &post_legacy, &cid1); + batch.insert(&records, &profile_legacy, &cid2); + batch.insert(&records, &orphan_legacy, &cid3); // bodies live in the blocks cas batch.insert(&blocks, keys::block_key("app.bsky.feed.post", &cid1), body1); @@ -297,7 +335,7 @@ mod tests { // an already-current record: the pass must leave it alone let current_key = keys::record_key(&did(), "app.bsky.feed.repost", &str_rkey); - batch.insert(&db.indexer.records, ¤t_key, b"repost body"); + batch.insert(&records, ¤t_key, b"repost body"); // a body-less legacy event remains valid through the event passes #[cfg(feature = "indexer_stream")] @@ -365,6 +403,8 @@ mod tests { // v10 retires the frozen cas after event compatibility is established assert!(!db.inner.keyspace_exists("blocks")); + // and the legacy records store, whose contents moved to heads + assert!(!db.inner.keyspace_exists("records")); Ok(()) } @@ -423,8 +463,9 @@ mod tests { { let db = Db::open(&cfg)?; let blocks = super::super::legacy_blocks_for_test(&db)?; + let records = super::super::legacy_records_for_test(&db)?; let mut batch = db.inner.batch(); - batch.insert(&db.indexer.records, &legacy_key, &cid); + batch.insert(&records, &legacy_key, &cid); batch.insert( &blocks, keys::block_key("app.bsky.actor.profile", &cid), @@ -456,6 +497,60 @@ mod tests { Ok(()) } + #[test] + fn v10_heads_are_actually_compressed() -> Result<()> { + // v9's records keyspace persisted "never compress" into its on-disk + // config; the migrated bodies must land in a store that compresses + let tmp = tempdir().into_diagnostic()?; + let cfg = Config { + database_path: tmp.path().to_path_buf(), + ..Default::default() + }; + + // highly compressible bodies, larger than one data block + let body = vec![7u8; 4096]; + let count = 2000u32; + let raw: u64; + { + let db = Db::open(&cfg)?; + let records = super::super::legacy_records_for_test(&db)?; + let mut batch = db.inner.batch(); + for n in 0..count { + let key = keys::record_key( + &Did::new(&format!("did:plc:{n:016x}")).unwrap(), + "app.bsky.feed.post", + &DbRkey::new(&format!("r{n}")), + ); + batch.insert(&records, key, &body); + } + batch.commit().into_diagnostic()?; + raw = (0..count) + .map(|n| { + (keys::record_key( + &Did::new(&format!("did:plc:{n:016x}")).unwrap(), + "app.bsky.feed.post", + &DbRkey::new(&format!("r{n}")), + ) + .len() + + body.len()) as u64 + }) + .sum(); + crate::db::migration::rewind_version_for_test(&db, 9)?; + db.persist()?; + } + + let db = Db::open(&cfg)?; + db.indexer.records.major_compact().into_diagnostic()?; + db.persist()?; + + let stored = db.indexer.records.disk_space(); + assert!( + stored < raw / 2, + "heads should compress repeated bodies: stored {stored} vs raw {raw}" + ); + Ok(()) + } + #[test] fn v10_byte_budget_chunking_limits_output_bytes() -> Result<()> { let tmp = tempdir().into_diagnostic()?; @@ -467,6 +562,7 @@ mod tests { let db = Db::open(&cfg)?; let mut batch = db.inner.batch(); let blocks = super::super::legacy_blocks_for_test(&db)?; + let records = super::super::legacy_records_for_test(&db)?; let large_body = vec![7u8; 10_000]; // 10 KB body // tid alphabet is base32-sortable (no 0/1/8/9 digits) @@ -480,7 +576,7 @@ mod tests { .into_diagnostic()? .to_bytes(); let legacy_key = legacy_record_key(&did(), "app.bsky.feed.post", &tid_rkey); - batch.insert(&db.indexer.records, &legacy_key, &cid); + batch.insert(&records, &legacy_key, &cid); batch.insert( &blocks, keys::block_key("app.bsky.feed.post", &cid), @@ -497,9 +593,10 @@ mod tests { entries: usize::MAX, bytes: 25_000, // 25 KB budget }, + optional_on_absent: true, }; let cursor_key = super::super::pass_cursor_key(9, pass.name); - let ks = db.keyspace_by_name("records").expect("records keyspace"); + let ks = super::super::legacy_records_for_test(&db)?; let outcome = super::super::run_chunk(&db, &cursor_key, &ks, &pass)?; assert_eq!(outcome.seen, 3); diff --git a/src/db/migration/v10/events.rs b/src/db/migration/v10/events.rs index ecb1c9d..46a558a 100644 --- a/src/db/migration/v10/events.rs +++ b/src/db/migration/v10/events.rs @@ -27,6 +27,7 @@ pub(super) const POINTER_EVENT_BODIES: Pass = Pass { scan: "events", visit: materialize_pointer_event, budget: ChunkBudget::DEFAULT, + optional_on_absent: false, }; #[cfg(feature = "indexer_stream")] @@ -73,6 +74,7 @@ fn decode_event<'a>(value: &'a [u8]) -> Result> { fn materialize_pointer_event( db: &Db, batch: &mut OwnedWriteBatch, + _scanned: &fjall::Keyspace, _key: &[u8], value: &[u8], ) -> Result { @@ -590,6 +592,7 @@ mod tests { entries: 1, bytes: usize::MAX, }, + optional_on_absent: false, }; let cursor = super::super::super::pass_cursor_key(9, pass.name); let events = db.keyspace_by_name("events").expect("events keyspace"); diff --git a/src/db/registry.rs b/src/db/registry.rs index 595ac17..e493ac1 100644 --- a/src/db/registry.rs +++ b/src/db/registry.rs @@ -96,7 +96,7 @@ registry! { #[cfg(feature = "backlinks")] backlinks => schema::Backlinks, #[cfg(feature = "indexer")] - indexer.records => schema::Records, + indexer.records => schema::Heads, #[cfg(feature = "indexer")] indexer.history => schema::History, #[cfg(feature = "indexer")] @@ -122,11 +122,11 @@ mod tests { #[test] fn record_debug_key_roundtrips() { let parsed = super::debug_parse_key( - "records", + "heads", "web:guestbook.gaze.systems|app.bsky.actor.profile|self", ); assert!(parsed.is_some(), "parse failed"); - let rendered = super::debug_render_key("records", &parsed.unwrap()); + let rendered = super::debug_render_key("heads", &parsed.unwrap()); assert_eq!( rendered, "web:guestbook.gaze.systems|app.bsky.actor.profile|self" diff --git a/src/db/schema.rs b/src/db/schema.rs index 7c2903b..9a810da 100644 --- a/src/db/schema.rs +++ b/src/db/schema.rs @@ -323,11 +323,11 @@ impl Schema for Crawler { /// maps `{DID} 00 {COL} 2f {RKey}` -> DAG-CBOR body or redacted CID marker. #[cfg(feature = "indexer")] -pub struct Records; +pub struct Heads; #[cfg(feature = "indexer")] -impl Schema for Records { - const NAME: &'static str = "records"; +impl Schema for Heads { + const NAME: &'static str = "heads"; fn options(cx: &OpenCx) -> KeyspaceCreateOptions { KeyspaceCreateOptions::default() @@ -341,13 +341,20 @@ impl Schema for Records { // grouped by repo: author/time locality beats the old col|cid // content addressing (-56% measured on a live corpus, zstd:3@64k) .data_block_size_policy(BlockSizePolicy::new([kb(16), kb(64), kb(128)])) - // leave the first level uncompressed; head bodies of fresh records - // are hot for stream inflation + // compress from the first level: bulk loads (migrations) write in + // sorted order, so their tables never overlap and compaction + // trivial-moves them without ever re-compressing — a None first + // level sticks to the whole tree (measured: 71MB stored vs 16MB + // zstd:3 standalone on the same bodies). heads stay hot in the + // memtable regardless. v10 migrates pre-v10 databases into this + // keyspace because v9's `records` persisted "never compress" into + // fjall's on-disk config, which no open-time schema change can + // undo. .data_block_compression_policy(CompressionPolicy::new([ - CompressionType::None, - (cx.compression)("records", 3), - (cx.compression)("records", 3), - (cx.compression)("records", 5), + (cx.compression)("heads", 3), + (cx.compression)("heads", 3), + (cx.compression)("heads", 3), + (cx.compression)("heads", 5), ])) .data_block_restart_interval_policy(RestartIntervalPolicy::new([8, 16, 32])) } @@ -429,7 +436,7 @@ impl Schema for History { } let rev_bytes: [u8; 8] = key[key.len() - 8..].try_into().expect("length checked"); let rev = crate::db::types::DbTid::new_from_bytes(rev_bytes); - format!("{}@{rev}", Records::render_debug_key(&key[..key.len() - 9])) + format!("{}@{rev}", Heads::render_debug_key(&key[..key.len() - 9])) } } @@ -467,7 +474,7 @@ impl Schema for Redactions { let Some((record_key, cid)) = super::keys::indexer::split_record_cid_key(key) else { return String::from_utf8_lossy(key).into_owned(); }; - format!("{}#{cid}", Records::render_debug_key(record_key)) + format!("{}#{cid}", Heads::render_debug_key(record_key)) } } @@ -507,7 +514,7 @@ impl Schema for EventBodies { let Some((record_key, cid)) = super::keys::indexer::split_record_cid_key(key) else { return String::from_utf8_lossy(key).into_owned(); }; - format!("{}#{cid}", Records::render_debug_key(record_key)) + format!("{}#{cid}", Heads::render_debug_key(record_key)) } } diff --git a/tests/debug_endpoints.nu b/tests/debug_endpoints.nu index 26a6934..40819cf 100644 --- a/tests/debug_endpoints.nu +++ b/tests/debug_endpoints.nu @@ -24,8 +24,8 @@ def main [] { print "backfill complete, testing debug endpoints" # 1. Test /debug/iter to find a key - print "testing /debug/iter on records partition" - let records = http get $"($debug_url)/debug/iter?partition=records&limit=1" + print "testing /debug/iter on heads partition" + let records = http get $"($debug_url)/debug/iter?partition=heads&limit=1" if ($records.items | is-empty) { print "FAILED: /debug/iter returned empty items" @@ -47,7 +47,7 @@ def main [] { # 2. Test /debug/get with that key (sent as string) print "testing /debug/get" let encoded_key = ($key_str | url encode) - let get_res = http get $"($debug_url)/debug/get?partition=records&key=($encoded_key)" + let get_res = http get $"($debug_url)/debug/get?partition=heads&key=($encoded_key)" if $get_res.value != $value_cid { print $"FAILED: /debug/get returned different value. expected: ($value_cid), got: ($get_res.value)" -- 2.51.2