use std::collections::BTreeMap;
use std::hint::black_box;
use std::io::{self, Write};
use std::path::{Path, PathBuf};
use subms::{SubMsFeatureManifest, SubMsP99Source, SubMsPerfHarness, classify_feature, summarize};
use subms_lsm_tree::LsmTree;
const SIZES: [usize; 3] = [8_192, 65_536, 524_288];
const CANON_N: usize = SIZES[SIZES.len() - 1];
const OPS: usize = 20_000;
const BULK_REPS: usize = 256;
const BULK_WARM_NANOS: u64 = 300_000_000;
const BULK_WARM_MAX_REPS: usize = 5_000;
const KEYED_WARM_NANOS: u64 = 300_000_000;
const KEYED_WARM_MAX_REPS: usize = 200_000;
const BLOCK_BYTES: usize = 4096;
const ENTRIES_PER_RUN: usize = 4_096;
const KEYS_PER_BLOCK: usize = 8;
const KEYS_PER_SSTABLE: usize = 4_096;
const FLUSH_BYTES: usize = 1_000_000;
const VALUE: &[u8] = b"value-payload-bytes-24ch";
fn key(i: usize) -> String {
format!("k{i:09}")
}
fn probe(i: usize, n: usize) -> usize {
(i.wrapping_mul(2_654_435_761)) % n
}
fn stat(h: &SubMsPerfHarness, median: bool) -> u64 {
summarize(h)
.stages
.iter()
.find(|s| s.name == "op")
.map_or(0, |s| if median { s.p50_ns } else { s.p99_ns })
}
fn keyed(mut op: impl FnMut(usize), median: bool) -> u64 {
let start = std::time::Instant::now();
for i in 0..KEYED_WARM_MAX_REPS {
op(i % OPS);
if start.elapsed().as_nanos() as u64 >= KEYED_WARM_NANOS {
break;
}
}
let mut h = SubMsPerfHarness::new("lsm-feature", "rust");
let st = h.stage("op", OPS);
for i in 0..OPS {
st.time(|| op(i));
}
stat(&h, median)
}
fn bulk(mut op: impl FnMut(), median: bool) -> u64 {
let start = std::time::Instant::now();
for _ in 0..BULK_WARM_MAX_REPS {
op();
if start.elapsed().as_nanos() as u64 >= BULK_WARM_NANOS {
break;
}
}
let mut h = SubMsPerfHarness::new("lsm-feature", "rust");
let st = h.stage("op", BULK_REPS);
for _ in 0..BULK_REPS {
st.time(&mut op);
}
stat(&h, median)
}
fn bulk_each<T>(mut setup: impl FnMut() -> T, mut op: impl FnMut(&mut T), median: bool) -> u64 {
let start = std::time::Instant::now();
for _ in 0..BULK_WARM_MAX_REPS {
let mut input = setup();
op(&mut input);
if start.elapsed().as_nanos() as u64 >= BULK_WARM_NANOS {
break;
}
}
let mut h = SubMsPerfHarness::new("lsm-feature", "rust");
let st = h.stage("op", BULK_REPS);
for _ in 0..BULK_REPS {
let mut input = setup();
st.time(|| op(&mut input));
}
stat(&h, median)
}
fn sweep(label: &str, mut at: impl FnMut(usize) -> u64) -> Vec<(usize, u64)> {
let rows: Vec<(usize, u64)> = SIZES.iter().map(|&n| (n, at(n))).collect();
eprintln!("sweep {label}: {rows:?}");
rows
}
fn main() -> io::Result<()> {
let path = PathBuf::from(env!("CARGO_MANIFEST_DIR"))
.join("..")
.join(".subms")
.join("features")
.join("rust.json");
let existing = std::fs::read_to_string(&path).unwrap_or_default();
let mut manifest = SubMsFeatureManifest::load_str("rust", &existing);
let (source, instance) = SubMsP99Source::from_env();
manifest.set_p99_source(source, instance.as_deref());
let tmp = TempDir::new("subms-lsm-features");
let base_p50 = {
let mut tree = LsmTree::open(tmp.path().join("base"), FLUSH_BYTES)?;
for i in 0..CANON_N {
tree.put(&key(i), VALUE)?;
}
tree.flush()?;
let probes: Vec<String> = (0..OPS).map(|i| key(probe(i, CANON_N))).collect();
keyed(
|i| {
black_box(tree.get(&probes[i]).expect("get"));
},
true,
)
};
eprintln!("base get p50: {base_p50}ns ({CANON_N} live keys)");
#[cfg(feature = "wal")]
feature_wal(&mut manifest, base_p50, tmp.path());
#[cfg(feature = "tiered-compaction")]
feature_tiered(&mut manifest, base_p50);
#[cfg(feature = "leveled-compaction")]
feature_leveled(&mut manifest, base_p50);
#[cfg(feature = "snapshot")]
feature_snapshot(&mut manifest, base_p50);
#[cfg(feature = "lz4")]
feature_lz4(&mut manifest, base_p50);
#[cfg(feature = "zstd")]
feature_zstd(&mut manifest, base_p50);
#[cfg(feature = "block-cache-integration")]
feature_block_cache(&mut manifest, base_p50);
drop(tmp);
std::fs::create_dir_all(path.parent().unwrap())?;
std::fs::write(&path, manifest.to_json())?;
io::stdout().write_all(manifest.to_json().as_bytes())?;
Ok(())
}
#[cfg(feature = "wal")]
fn wal_of(dir: &Path, n: usize) -> PathBuf {
use subms_lsm_tree::WriteAheadLog;
let path = dir.join(format!("replay-{n}.wal"));
let _ = std::fs::remove_file(&path);
let mut wal = WriteAheadLog::open(&path).expect("open wal");
for i in 0..n {
wal.log_put(&key(i), VALUE).expect("log_put");
}
wal.sync().expect("sync");
path
}
#[cfg(feature = "wal")]
fn feature_wal(manifest: &mut SubMsFeatureManifest, base_p50: u64, dir: &Path) {
use subms_lsm_tree::WriteAheadLog;
let sw = sweep("wal/replay", |n| {
let path = wal_of(dir, n);
let out = bulk(
|| {
black_box(WriteAheadLog::replay(&path).expect("replay").len());
},
true,
);
let _ = std::fs::remove_file(&path);
out
});
let (cat, reason) = classify_feature(&sw, Some(base_p50), None);
let replay_path = wal_of(dir, CANON_N);
let mut p99 = BTreeMap::new();
p99.insert(
"replay".to_string(),
bulk(
|| {
black_box(WriteAheadLog::replay(&replay_path).expect("replay").len());
},
false,
),
);
let _ = std::fs::remove_file(&replay_path);
let append_path = dir.join("append.wal");
let _ = std::fs::remove_file(&append_path);
{
let mut wal = WriteAheadLog::open(&append_path).expect("open wal");
let keys: Vec<String> = (0..OPS).map(key).collect();
p99.insert(
"log_put".to_string(),
keyed(|i| wal.log_put(&keys[i], VALUE).expect("log_put"), false),
);
}
let _ = std::fs::remove_file(&append_path);
manifest.set_feature("wal", cat, &p99, &reason);
}
#[cfg(feature = "tiered-compaction")]
fn feature_tiered(manifest: &mut SubMsFeatureManifest, base_p50: u64) {
use subms_lsm_tree::{TieredCompactionPlanner, TieredManifest, TieredRun};
fn runs(n: usize) -> Vec<TieredRun> {
(0..n.div_ceil(ENTRIES_PER_RUN))
.map(|r| {
let entries: Vec<(String, Option<Vec<u8>>)> = (0..ENTRIES_PER_RUN)
.map(|j| (key(r * ENTRIES_PER_RUN + j), Some(VALUE.to_vec())))
.collect();
TieredRun::new(r as u64, entries)
})
.collect()
}
let planner = TieredCompactionPlanner::new(2);
let sw = sweep("tiered-compaction/merge", |n| {
let template = runs(n);
bulk_each(
|| TieredManifest {
levels: vec![template.clone()],
},
|m| planner.merge(m, 0, 9_999),
true,
)
});
let (cat, reason) = classify_feature(&sw, Some(base_p50), None);
let template = runs(CANON_N);
let mut p99 = BTreeMap::new();
p99.insert(
"merge".to_string(),
bulk_each(
|| TieredManifest {
levels: vec![template.clone()],
},
|m| planner.merge(m, 0, 9_999),
false,
),
);
let planned = TieredManifest {
levels: vec![template],
};
p99.insert(
"plan".to_string(),
keyed(|_| _ = black_box(planner.pick_level(&planned)), false),
);
manifest.set_feature("tiered-compaction", cat, &p99, &reason);
}
#[cfg(feature = "leveled-compaction")]
fn feature_leveled(manifest: &mut SubMsFeatureManifest, base_p50: u64) {
use subms_lsm_tree::{LeveledCompactionPlanner, LeveledManifest, LeveledRun};
fn halves(n: usize) -> (Vec<LeveledRun>, Vec<LeveledRun>) {
let per_level = n / 2;
let build = |parity: usize| -> Vec<LeveledRun> {
(0..per_level.div_ceil(ENTRIES_PER_RUN))
.map(|r| {
let entries: Vec<(String, Option<Vec<u8>>)> = (0..ENTRIES_PER_RUN)
.map(|j| {
let idx = 2 * (r * ENTRIES_PER_RUN + j) + parity;
(key(idx), Some(VALUE.to_vec()))
})
.collect();
LeveledRun::new((r * 2 + parity) as u64, entries)
})
.collect()
};
(build(0), build(1))
}
let planner = LeveledCompactionPlanner::new(64_000, 10, 4);
let sw = sweep("leveled-compaction/compact", |n| {
let (l0, l1) = halves(n);
bulk_each(
|| LeveledManifest {
levels: vec![l0.clone(), l1.clone()],
},
|m| planner.compact(m, 0, 9_999),
true,
)
});
let (cat, reason) = classify_feature(&sw, Some(base_p50), None);
let (l0, l1) = halves(CANON_N);
let mut p99 = BTreeMap::new();
p99.insert(
"compact".to_string(),
bulk_each(
|| LeveledManifest {
levels: vec![l0.clone(), l1.clone()],
},
|m| planner.compact(m, 0, 9_999),
false,
),
);
let planned = LeveledManifest {
levels: vec![l0, l1],
};
p99.insert(
"plan".to_string(),
keyed(|_| _ = black_box(planner.pick_level(&planned)), false),
);
manifest.set_feature("leveled-compaction", cat, &p99, &reason);
}
#[cfg(feature = "snapshot")]
fn feature_snapshot(manifest: &mut SubMsFeatureManifest, base_p50: u64) {
use subms_lsm_tree::{SnapshotManager, SnapshotManifest};
fn manager(n: usize) -> SnapshotManager {
let ids: Vec<u64> = (0..(n / KEYS_PER_SSTABLE).max(1) as u64).collect();
SnapshotManager::with_initial(SnapshotManifest::new(ids))
}
let sw = sweep("snapshot/snapshot", |n| {
let mgr = manager(n);
keyed(|_| _ = black_box(mgr.snapshot()), true)
});
let (cat, reason) = classify_feature(&sw, Some(base_p50), None);
let mgr = manager(CANON_N);
let mut p99 = BTreeMap::new();
p99.insert(
"snapshot".to_string(),
keyed(|_| _ = black_box(mgr.snapshot()), false),
);
let held = mgr.snapshot();
let ids = held.sstable_ids();
let targets: Vec<u64> = (0..OPS)
.map(|i| probe(i, ids.len().max(1) * 2) as u64)
.collect();
p99.insert(
"get_on_snapshot".to_string(),
keyed(
|i| {
let t = targets[i];
black_box(ids.iter().rev().any(|&id| id == t));
},
false,
),
);
manifest.set_feature("snapshot", cat, &p99, &reason);
}
#[cfg(any(feature = "lz4", feature = "zstd", feature = "block-cache-integration"))]
fn representative_block() -> Vec<u8> {
let pattern = b"key-0000042\x00present\x00value-payload-bytes-for-block|";
let mut out = Vec::with_capacity(BLOCK_BYTES + pattern.len());
while out.len() < BLOCK_BYTES {
out.extend_from_slice(pattern);
}
out.truncate(BLOCK_BYTES);
out
}
#[cfg(feature = "lz4")]
fn feature_lz4(manifest: &mut SubMsFeatureManifest, base_p50: u64) {
use subms_lsm_tree::Lz4BlockCompressor;
let c = Lz4BlockCompressor::new();
let block = representative_block();
let sw = sweep("lz4/compress", |_| {
keyed(|_| _ = black_box(c.compress(&block)), true)
});
let (cat, reason) = classify_feature(&sw, Some(base_p50), None);
let encoded = c.compress(&block);
let mut p99 = BTreeMap::new();
p99.insert(
"compress_block".to_string(),
keyed(|_| _ = black_box(c.compress(&block)), false),
);
p99.insert(
"decompress_block".to_string(),
keyed(
|_| _ = black_box(c.decompress(&encoded).expect("lz4 decode")),
false,
),
);
manifest.set_feature("lz4", cat, &p99, &reason);
}
#[cfg(feature = "zstd")]
fn feature_zstd(manifest: &mut SubMsFeatureManifest, base_p50: u64) {
use subms_lsm_tree::ZstdBlockCompressor;
let c = ZstdBlockCompressor::new();
let block = representative_block();
let sw = sweep("zstd/compress", |_| {
keyed(
|_| _ = black_box(c.compress(&block).expect("zstd encode")),
true,
)
});
let (cat, reason) = classify_feature(&sw, Some(base_p50), None);
let encoded = c.compress(&block).expect("zstd encode");
let mut p99 = BTreeMap::new();
p99.insert(
"compress_block".to_string(),
keyed(
|_| _ = black_box(c.compress(&block).expect("zstd encode")),
false,
),
);
p99.insert(
"decompress_block".to_string(),
keyed(
|_| _ = black_box(c.decompress(&encoded).expect("zstd decode")),
false,
),
);
manifest.set_feature("zstd", cat, &p99, &reason);
}
#[cfg(feature = "block-cache-integration")]
fn feature_block_cache(manifest: &mut SubMsFeatureManifest, base_p50: u64) {
use std::sync::Arc;
use subms_lsm_tree::{Block, BlockCache, BlockKey, LruBlockCache};
fn filled(n: usize) -> (LruBlockCache, usize) {
let cap = (n / KEYS_PER_BLOCK).max(64);
let cache = LruBlockCache::new(cap);
let block: Block = Arc::from(representative_block().into_boxed_slice());
for i in 0..cap as u64 {
cache.put(BlockKey::new(i % 8, i * BLOCK_BYTES as u64), block.clone());
}
(cache, cap)
}
let sw = sweep("block-cache-integration/get_cached", |n| {
let (cache, cap) = filled(n);
let keys: Vec<BlockKey> = (0..OPS)
.map(|i| {
let k = probe(i, cap) as u64;
BlockKey::new(k % 8, k * BLOCK_BYTES as u64)
})
.collect();
keyed(|i| _ = black_box(cache.get(&keys[i])), true)
});
let (cat, reason) = classify_feature(&sw, Some(base_p50), None);
let (cache, cap) = filled(CANON_N);
let hits: Vec<BlockKey> = (0..OPS)
.map(|i| {
let k = probe(i, cap) as u64;
BlockKey::new(k % 8, k * BLOCK_BYTES as u64)
})
.collect();
let misses: Vec<BlockKey> = (0..OPS)
.map(|i| BlockKey::new(999, probe(i, cap) as u64))
.collect();
let mut p99 = BTreeMap::new();
p99.insert(
"get_cached".to_string(),
keyed(|i| _ = black_box(cache.get(&hits[i])), false),
);
p99.insert(
"get_miss".to_string(),
keyed(|i| _ = black_box(cache.get(&misses[i])), false),
);
manifest.set_feature("block-cache-integration", cat, &p99, &reason);
}
struct TempDir {
path: PathBuf,
}
impl TempDir {
fn new(label: &str) -> Self {
let path = std::env::temp_dir().join(format!("{}-{}", label, std::process::id()));
let _ = std::fs::remove_dir_all(&path);
std::fs::create_dir_all(&path).expect("create temp dir");
Self { path }
}
fn path(&self) -> &Path {
&self.path
}
}
impl Drop for TempDir {
fn drop(&mut self) {
let _ = std::fs::remove_dir_all(&self.path);
}
}