diff --git a/ufos/src/storage_fjall.rs b/ufos/src/storage_fjall.rs index 5a9ffc4..e482788 100644 --- a/ufos/src/storage_fjall.rs +++ b/ufos/src/storage_fjall.rs @@ -224,6 +224,23 @@ impl StorageWhatever for sketch_secret }; + for (partition, name) in [ + (&global, "global"), + (&feeds, "feeds"), + (&records, "records"), + (&rollups, "rollups"), + (&queues, "queues"), + ] { + let size0 = partition.disk_space(); + log::info!("beggining major compaction for {name} (original size: {size0})"); + let t0 = Instant::now(); + partition.major_compact().expect("compact better work 😬"); + let dt = t0.elapsed(); + let sizef = partition.disk_space(); + let dsize = (sizef as i64) - (size0 as i64); + log::info!("completed compaction for {name} in {dt:?} (new size: {sizef}, {dsize})"); + } + let reader = FjallReader { keyspace: keyspace.clone(), global: global.clone(), -- 2.51.2 From 4b4b627a8d3c8366e77ccfeedb0d13b0d0c08179 Mon Sep 17 00:00:00 2001 From: phil Date: Fri, 15 Aug 2025 13:22:50 -0400 Subject: [PATCH 2/5] =?UTF-8?q?remove=20records=20from=20the=20records=20c?= =?UTF-8?q?ollection=20=F0=9F=98=AD?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit noooooooooooooooooooooooo --- ufos/src/storage_fjall.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/ufos/src/storage_fjall.rs b/ufos/src/storage_fjall.rs index e482788..f71fc92 100644 --- a/ufos/src/storage_fjall.rs +++ b/ufos/src/storage_fjall.rs @@ -1622,7 +1622,7 @@ impl StoreWriter for FjallWriter { candidate_new_feed_lower_cursor = Some(feed_key.cursor()); } - self.feeds.remove(&location_key_bytes)?; + self.records.remove(&location_key_bytes)?; self.feeds.remove(key_bytes)?; records_deleted += 1; } -- 2.51.2 From 877c09750ec0805f1ae2f2f3fc8956a2857efb76 Mon Sep 17 00:00:00 2001 From: phil Date: Fri, 15 Aug 2025 16:37:44 -0400 Subject: [PATCH 3/5] big gc bigtime previously: forgot to delete all records this probably takes a lot of memory hopefully won't be needed again --- ufos/src/main.rs | 24 ++++++++- ufos/src/storage_fjall.rs | 102 +++++++++++++++++++++++++++++++------- 2 files changed, 107 insertions(+), 19 deletions(-) diff --git a/ufos/src/main.rs b/ufos/src/main.rs index 949db21..60601ee 100644 --- a/ufos/src/main.rs +++ b/ufos/src/main.rs @@ -9,7 +9,7 @@ use ufos::consumer; use ufos::file_consumer; use ufos::server; use ufos::storage::{StorageWhatever, StoreBackground, StoreReader, StoreWriter}; -use ufos::storage_fjall::FjallStorage; +use ufos::storage_fjall::{FjallConfig, FjallStorage}; use ufos::store_types::SketchSecretPrefix; use ufos::{nice_duration, ConsumerInfo}; @@ -55,6 +55,12 @@ struct Args { /// DEBUG: interpret jetstream as a file fixture #[arg(long, action)] jetstream_fixture: bool, + /// HOPEFULLY only needed once + /// + /// brute-force garbage-collect all dangling records because we weren't deleting + /// them before at all (oops) + #[arg(long, action)] + fjall_records_gc: bool, } #[tokio::main] @@ -67,8 +73,22 @@ async fn main() -> anyhow::Result<()> { args.data.clone(), jetstream, args.jetstream_force, - Default::default(), + FjallConfig { + major_compact: !args.fjall_records_gc, + }, )?; + + if args.fjall_records_gc { + log::info!("beginning brute-force records gc"); + let t0 = std::time::Instant::now(); + let (n, m) = write_store.records_brute_gc_danger()?; + let dt = t0.elapsed(); + log::info!( + "completed brute-force records gc in {dt:?}, removed {n} and retained {m} records." + ); + return Ok(()); + } + go(args, read_store, write_store, cursor, sketch_secret).await?; Ok(()) } diff --git a/ufos/src/storage_fjall.rs b/ufos/src/storage_fjall.rs index f71fc92..ccd027a 100644 --- a/ufos/src/storage_fjall.rs +++ b/ufos/src/storage_fjall.rs @@ -148,6 +148,10 @@ pub struct FjallConfig { /// this is only meant for tests #[cfg(test)] pub temp: bool, + /// do major compaction on startup + /// + /// default is false. probably a good thing unless it's too slow. + pub major_compact: bool, } impl StorageWhatever for FjallStorage { @@ -155,7 +159,7 @@ impl StorageWhatever for path: impl AsRef, endpoint: String, force_endpoint: bool, - _config: FjallConfig, + config: FjallConfig, ) -> StorageResult<(FjallReader, FjallWriter, Option, SketchSecretPrefix)> { let keyspace = { let config = Config::new(path); @@ -224,21 +228,27 @@ impl StorageWhatever for sketch_secret }; - for (partition, name) in [ - (&global, "global"), - (&feeds, "feeds"), - (&records, "records"), - (&rollups, "rollups"), - (&queues, "queues"), - ] { - let size0 = partition.disk_space(); - log::info!("beggining major compaction for {name} (original size: {size0})"); - let t0 = Instant::now(); - partition.major_compact().expect("compact better work 😬"); - let dt = t0.elapsed(); - let sizef = partition.disk_space(); - let dsize = (sizef as i64) - (size0 as i64); - log::info!("completed compaction for {name} in {dt:?} (new size: {sizef}, {dsize})"); + if config.major_compact { + for (partition, name) in [ + (&global, "global"), + (&feeds, "feeds"), + (&records, "records"), + (&rollups, "rollups"), + (&queues, "queues"), + ] { + let size0 = partition.disk_space(); + log::info!("beggining major compaction for {name} (original size: {size0})"); + let t0 = Instant::now(); + partition.major_compact().expect("compact better work 😬"); + let dt = t0.elapsed(); + let sizef = partition.disk_space(); + let dsize = (sizef as i64) - (size0 as i64); + log::info!( + "completed compaction for {name} in {dt:?} (new size: {sizef}, {dsize})" + ); + } + } else { + log::info!("skipping major compaction on startup"); } let reader = FjallReader { @@ -1366,6 +1376,61 @@ impl FjallWriter { batch.commit()?; Ok((cursors_advanced, dirty_nsids)) } + pub fn records_brute_gc_danger(&self) -> StorageResult<(usize, usize)> { + let (mut removed, mut retained) = (0, 0); + let mut to_retain = HashSet::>::new(); + + // Partition: 'feed' + // + // - Per-collection list of record references ordered by jetstream cursor + // - key: nullstr || u64 (collection nsid null-terminated, jetstream cursor) + // - val: nullstr || nullstr || nullstr (did, rkey, rev. rev is mostly a sanity-check for now.) + // + // + // Partition: 'records' + // + // - Actual records by their atproto location + // - key: nullstr || nullstr || nullstr (did, collection, rkey) + // - val: u64 || bool || nullstr || rawval (js_cursor, is_update, rev, actual record) + // + // + + log::warn!("loading *all* record keys from feed into memory (yikes)"); + let t0 = Instant::now(); + for (i, kv) in self.feeds.iter().enumerate() { + if i > 0 && (i % 100000 == 0) { + log::info!("{i}..."); + } + let (key_bytes, val_bytes) = kv?; + let key = db_complete::(&key_bytes)?; + let val = db_complete::(&val_bytes)?; + let record_key: RecordLocationKey = (&key, &val).into(); + to_retain.insert(record_key.to_db_bytes()?); + } + log::warn!( + "loaded. wow. took {:?}, found {} keys", + t0.elapsed(), + to_retain.len() + ); + + log::warn!("warmup OVER, iterating some billions of record keys now"); + let t0 = Instant::now(); + for (i, k) in self.records.keys().enumerate() { + let key_bytes = k?; + if to_retain.contains(&*key_bytes) { + retained += 1; + } else { + self.records.remove(key_bytes)?; + removed += 1; + } + if i > 0 && (i % 10_000_000) == 0 { + log::info!("{i}: {retained} retained, {removed} removed."); + } + } + log::warn!("whew! that took {:?}", t0.elapsed()); + + Ok((removed, retained)) + } } impl StoreWriter for FjallWriter { @@ -1817,7 +1882,10 @@ mod tests { tempfile::tempdir().unwrap(), "offline test (no real jetstream endpoint)".to_string(), false, - FjallConfig { temp: true }, + FjallConfig { + temp: true, + ..Default::default() + }, ) .unwrap(); (read, write) -- 2.51.2 From 4b462c828ecca4a8d42ca1ab175ce56de1f6e3c8 Mon Sep 17 00:00:00 2001 From: phil Date: Sat, 16 Aug 2025 13:31:38 -0400 Subject: [PATCH 4/5] log a reasonable amount --- ufos/src/storage_fjall.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/ufos/src/storage_fjall.rs b/ufos/src/storage_fjall.rs index ccd027a..3d68aae 100644 --- a/ufos/src/storage_fjall.rs +++ b/ufos/src/storage_fjall.rs @@ -1398,7 +1398,7 @@ impl FjallWriter { log::warn!("loading *all* record keys from feed into memory (yikes)"); let t0 = Instant::now(); for (i, kv) in self.feeds.iter().enumerate() { - if i > 0 && (i % 100000 == 0) { + if i > 0 && (i % 10_000_000 == 0) { log::info!("{i}..."); } let (key_bytes, val_bytes) = kv?; @@ -1423,7 +1423,7 @@ impl FjallWriter { self.records.remove(key_bytes)?; removed += 1; } - if i > 0 && (i % 10_000_000) == 0 { + if i > 0 && (i % 100_000_000) == 0 { log::info!("{i}: {retained} retained, {removed} removed."); } } -- 2.51.2 From 204bb40c58c9062dd2a88f6fe4bc365382870072 Mon Sep 17 00:00:00 2001 From: phil Date: Mon, 11 May 2026 17:44:17 -0400 Subject: [PATCH 5/5] just do major compact and fix the default-run config --- ufos/Cargo.toml | 5 ++ ufos/src/bin/major-compact.rs | 34 +++++++++++++ ufos/src/main.rs | 24 +--------- ufos/src/storage_fjall.rs | 89 +---------------------------------- 4 files changed, 43 insertions(+), 109 deletions(-) create mode 100644 ufos/src/bin/major-compact.rs diff --git a/ufos/Cargo.toml b/ufos/Cargo.toml index bc7bec0..2c000ce 100644 --- a/ufos/Cargo.toml +++ b/ufos/Cargo.toml @@ -2,6 +2,7 @@ name = "ufos" version = "0.1.0" edition = "2021" +default-run = "ufos" [dependencies] anyhow = "1.0.97" @@ -38,5 +39,9 @@ tikv-jemallocator = "0.6.0" name = "analyze" path = "src/bin/analyze.rs" +[[bin]] +name = "major-compact" +path = "src/bin/major-compact.rs" + [dev-dependencies] tempfile = "3.19.1" diff --git a/ufos/src/bin/major-compact.rs b/ufos/src/bin/major-compact.rs new file mode 100644 index 0000000..a9c17fe --- /dev/null +++ b/ufos/src/bin/major-compact.rs @@ -0,0 +1,34 @@ +use clap::Parser; +use fjall::{Config, PartitionCreateOptions}; +use std::path::PathBuf; +use std::time::Instant; + +#[derive(Parser)] +#[command(about = "Run a major compaction over every ufos partition")] +struct Cli { + /// path to the fjall data directory + /// + /// WARNING: MUST NOT RUN WHILE ANOTHER UFOS PROCESS IS USING IT + data: PathBuf, +} + +fn main() -> anyhow::Result<()> { + let cli = Cli::parse(); + + eprintln!("opening db at {:?}...", cli.data); + let keyspace = Config::new(&cli.data).open()?; + + for name in ["global", "feeds", "records", "rollups", "queues"] { + let partition = keyspace.open_partition(name, PartitionCreateOptions::default())?; + let size0 = partition.disk_space(); + eprintln!("beginning major compaction for {name} (original size: {size0})"); + let t0 = Instant::now(); + partition.major_compact()?; + let dt = t0.elapsed(); + let sizef = partition.disk_space(); + let dsize = (sizef as i64) - (size0 as i64); + eprintln!("completed compaction for {name} in {dt:?} (new size: {sizef}, {dsize})"); + } + + Ok(()) +} diff --git a/ufos/src/main.rs b/ufos/src/main.rs index 1b320bb..20f7cfe 100644 --- a/ufos/src/main.rs +++ b/ufos/src/main.rs @@ -9,7 +9,7 @@ use ufos::consumer; use ufos::file_consumer; use ufos::server; use ufos::storage::{StorageWhatever, StoreBackground, StoreReader, StoreWriter}; -use ufos::storage_fjall::{FjallConfig, FjallStorage}; +use ufos::storage_fjall::FjallStorage; use ufos::store_types::SketchSecretPrefix; use ufos::{nice_duration, ConsumerInfo}; @@ -59,12 +59,6 @@ struct Args { /// DEBUG: interpret jetstream as a file fixture #[arg(long, action, env = "UFOS_JETSTREAM_FIXTURE")] jetstream_fixture: bool, - /// HOPEFULLY only needed once - /// - /// brute-force garbage-collect all dangling records because we weren't deleting - /// them before at all (oops) - #[arg(long, action)] - fjall_records_gc: bool, /// enable metrics collection and serving #[arg(long, action, env = "UFOS_COLLECT_METRICS")] collect_metrics: bool, @@ -84,22 +78,8 @@ async fn main() -> anyhow::Result<()> { args.data.clone(), jetstream, args.jetstream_force, - FjallConfig { - major_compact: !args.fjall_records_gc, - }, + Default::default(), )?; - - if args.fjall_records_gc { - log::info!("beginning brute-force records gc"); - let t0 = std::time::Instant::now(); - let (n, m) = write_store.records_brute_gc_danger()?; - let dt = t0.elapsed(); - log::info!( - "completed brute-force records gc in {dt:?}, removed {n} and retained {m} records." - ); - return Ok(()); - } - go(args, read_store, write_store, cursor, sketch_secret).await?; Ok(()) } diff --git a/ufos/src/storage_fjall.rs b/ufos/src/storage_fjall.rs index 09694e9..dc96c45 100644 --- a/ufos/src/storage_fjall.rs +++ b/ufos/src/storage_fjall.rs @@ -148,10 +148,6 @@ pub struct FjallConfig { /// this is only meant for tests #[cfg(test)] pub temp: bool, - /// do major compaction on startup - /// - /// default is false. probably a good thing unless it's too slow. - pub major_compact: bool, } impl StorageWhatever for FjallStorage { @@ -159,7 +155,7 @@ impl StorageWhatever for path: impl AsRef, endpoint: String, force_endpoint: bool, - config: FjallConfig, + _config: FjallConfig, ) -> StorageResult<(FjallReader, FjallWriter, Option, SketchSecretPrefix)> { let keyspace = { let config = Config::new(path); @@ -228,29 +224,6 @@ impl StorageWhatever for sketch_secret }; - if config.major_compact { - for (partition, name) in [ - (&global, "global"), - (&feeds, "feeds"), - (&records, "records"), - (&rollups, "rollups"), - (&queues, "queues"), - ] { - let size0 = partition.disk_space(); - log::info!("beggining major compaction for {name} (original size: {size0})"); - let t0 = Instant::now(); - partition.major_compact().expect("compact better work 😬"); - let dt = t0.elapsed(); - let sizef = partition.disk_space(); - let dsize = (sizef as i64) - (size0 as i64); - log::info!( - "completed compaction for {name} in {dt:?} (new size: {sizef}, {dsize})" - ); - } - } else { - log::info!("skipping major compaction on startup"); - } - let reader = FjallReader { keyspace: keyspace.clone(), global: global.clone(), @@ -1381,61 +1354,6 @@ impl FjallWriter { batch.commit()?; Ok((cursors_advanced, dirty_nsids)) } - pub fn records_brute_gc_danger(&self) -> StorageResult<(usize, usize)> { - let (mut removed, mut retained) = (0, 0); - let mut to_retain = HashSet::>::new(); - - // Partition: 'feed' - // - // - Per-collection list of record references ordered by jetstream cursor - // - key: nullstr || u64 (collection nsid null-terminated, jetstream cursor) - // - val: nullstr || nullstr || nullstr (did, rkey, rev. rev is mostly a sanity-check for now.) - // - // - // Partition: 'records' - // - // - Actual records by their atproto location - // - key: nullstr || nullstr || nullstr (did, collection, rkey) - // - val: u64 || bool || nullstr || rawval (js_cursor, is_update, rev, actual record) - // - // - - log::warn!("loading *all* record keys from feed into memory (yikes)"); - let t0 = Instant::now(); - for (i, kv) in self.feeds.iter().enumerate() { - if i > 0 && (i % 10_000_000 == 0) { - log::info!("{i}..."); - } - let (key_bytes, val_bytes) = kv?; - let key = db_complete::(&key_bytes)?; - let val = db_complete::(&val_bytes)?; - let record_key: RecordLocationKey = (&key, &val).into(); - to_retain.insert(record_key.to_db_bytes()?); - } - log::warn!( - "loaded. wow. took {:?}, found {} keys", - t0.elapsed(), - to_retain.len() - ); - - log::warn!("warmup OVER, iterating some billions of record keys now"); - let t0 = Instant::now(); - for (i, k) in self.records.keys().enumerate() { - let key_bytes = k?; - if to_retain.contains(&*key_bytes) { - retained += 1; - } else { - self.records.remove(key_bytes)?; - removed += 1; - } - if i > 0 && (i % 100_000_000) == 0 { - log::info!("{i}: {retained} retained, {removed} removed."); - } - } - log::warn!("whew! that took {:?}", t0.elapsed()); - - Ok((removed, retained)) - } } impl StoreWriter for FjallWriter { @@ -1887,10 +1805,7 @@ mod tests { tempfile::tempdir().unwrap(), "offline test (no real jetstream endpoint)".to_string(), false, - FjallConfig { - temp: true, - ..Default::default() - }, + FjallConfig { temp: true }, ) .unwrap(); (read, write)