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)