use std::fs::File;
use std::io::BufRead;
use std::path::{Path, PathBuf};
use engramdb_core::count_index::CountIndex;
use engramdb_core::layout::Layout;
use engramdb_keygen::PleSpec;
const FNV_OFFSET: u64 = 0xcbf29ce484222325;
const FNV_PRIME: u64 = 0x100000001b3;
fn fnv(bytes: &[u8]) -> u64 {
let mut h = FNV_OFFSET;
for &b in bytes {
h ^= b as u64;
h = h.wrapping_mul(FNV_PRIME);
}
h
}
fn ple_layout() -> Layout {
Layout::new(128, 2_500_012, 160, 1)
}
fn main() {
let mut args = std::env::args().skip(1);
let Some(cmd) = args.next() else {
println!("usage: engramdb <build|index|gather|verify|bench-real|warm> [args...]");
return;
};
let rest = args;
let out = match cmd.as_str() {
"build" => cmd_build(rest),
"index" => cmd_index(rest),
"gather" => cmd_gather(rest),
"verify" => cmd_verify(rest),
"bench-real" => cmd_bench_real(rest),
"warm" => cmd_warm(rest),
_ => Err(format!("unknown command: {cmd}")),
};
if let Err(e) = out {
eprintln!("error: {e}");
std::process::exit(1);
}
}
fn cmd_build(mut args: impl Iterator<Item = String>) -> Result<(), String> {
let src = PathBuf::from(args.next().ok_or("需要 <shard_dir>")?);
let dst = PathBuf::from(args.next().ok_or("需要 <out_dir>")?);
let layout = ple_layout();
std::fs::create_dir_all(&dst).map_err(|e| e.to_string())?;
for shard in 0..layout.shards {
let sp = src.join(format!("shard_{:03}.bin", shard));
if !sp.exists() {
return Err(format!("缺分片 {sp:?}({}/{})", shard, layout.shards));
}
let mut f = File::open(&sp).map_err(io_err)?;
let mut out = File::create(dst.join(format!("badge_{:03}.bin", shard))).map_err(io_err)?;
let badge_bytes = layout.badge_bytes() as usize;
let mut badge = vec![0u8; badge_bytes];
loop {
let mut got = 0usize;
while got < badge_bytes {
let n = std::io::Read::read(&mut f, &mut badge[got..]).map_err(io_err)?;
if n == 0 {
break;
}
got += n;
}
if got == 0 {
break;
}
if got < badge_bytes {
badge[got..].fill(0);
}
std::io::Write::write_all(&mut out, &badge).map_err(io_err)?;
}
drop(out);
}
let manifest = serde_json::json!({
"layout": { "shards": layout.shards, "rows_per_shard": layout.rows_per_shard,
"width": layout.width, "elem_bytes": 1, "badge_rows": layout.badge_rows,
"badge_bytes": layout.badge_bytes(), "total_rows": layout.total_rows() },
"source": src.to_string_lossy(),
});
std::fs::write(
dst.join("manifest.json"),
serde_json::to_vec_pretty(&manifest).unwrap(),
)
.map_err(io_err)?;
println!(
"built {dst:?} ({} shards × {} rows × {})",
layout.shards, layout.rows_per_shard, layout.width
);
Ok(())
}
fn cmd_index(mut args: impl Iterator<Item = String>) -> Result<(), String> {
let src = PathBuf::from(args.next().ok_or("需要 <rowids.bin>")?);
let dst = args.next().unwrap_or_else(|| "index".to_string());
let dst = PathBuf::from(dst);
std::fs::create_dir_all(&dst).map_err(io_err)?;
let f = File::open(&src).map_err(io_err)?;
let idx = CountIndex::build_from_bin_stream(std::io::BufReader::new(f)).map_err(io_err)?;
idx.write_bin(&dst.join("counts.bin")).map_err(io_err)?;
idx.write_dump(&dst.join("counts.dump.txt"))
.map_err(io_err)?;
println!("indexed {} unique rows -> {dst:?}", idx.iter().count());
Ok(())
}
fn cmd_warm(mut rest: impl Iterator<Item = String>) -> Result<(), String> {
let dir = PathBuf::from(rest.next().ok_or("需要 <rows_dir|badge_dir>")?);
let mut budget_gb = 1.0f64;
let mut rest2 = rest;
while let Some(arg) = rest2.next() {
if arg == "--budget" {
budget_gb = rest2
.next()
.ok_or("budget 值")?
.parse()
.map_err(|e: std::num::ParseFloatError| e.to_string())?;
}
}
let layout = ple_layout();
let mut warmed: u64 = 0;
let budget = (budget_gb * 1e9) as u64;
for shard in 0..layout.shards {
let p = dir.join(format!("shard_{:03}.bin", shard));
if !p.exists() {
let p2 = dir.join(format!("badge_{:03}.bin", shard));
if !p2.exists() {
continue;
}
let mut f = File::open(&p2).map_err(io_err)?;
let mut probe = 0u8;
let _gone = 0u64;
while warmed < budget {
let mut buf = [0u8; 1 << 20];
let n = std::io::Read::read(&mut f, &mut buf).map_err(io_err)?;
if n == 0 {
break;
}
warmed += n as u64;
let gone = warmed;
let _ = gone;
probe = buf[0];
}
let _ = probe;
} else {
let mut f = File::open(&p).map_err(io_err)?;
let mut probe = 0u8;
while warmed < budget {
let mut buf = [0u8; 1 << 20];
let n = std::io::Read::read(&mut f, &mut buf).map_err(io_err)?;
if n == 0 {
break;
}
warmed += n as u64;
probe = buf[0];
}
let _ = probe;
}
if warmed >= budget {
break;
}
}
println!("warmed {:.2} GB of OS page cache", warmed as f64 / 1e9);
Ok(())
}
fn cmd_gather(mut args: impl Iterator<Item = String>) -> Result<(), String> {
let dir = PathBuf::from(args.next().ok_or("需要 <badge_dir>")?);
let layout = ple_layout();
let batch = engramdb_io::batch::BadgeGather::open(&dir, &layout).map_err(io_err)?;
let stdin = std::io::stdin();
let mut keys = Vec::new();
for line in stdin.lock().lines() {
let line = line.map_err(|e| e.to_string())?;
let line = line.trim();
if line.is_empty() {
continue;
}
keys.push(line.parse::<u64>().map_err(|e| e.to_string())?);
}
let w = layout.width as usize;
let mut out = vec![0u8; keys.len() * w];
batch.gather_planned(&keys, &mut out).map_err(io_err)?;
println!("{}", fnv(&out));
Ok(())
}
fn cmd_verify(mut args: impl Iterator<Item = String>) -> Result<(), String> {
let dir = PathBuf::from(args.next().ok_or("需要 <badge_dir>")?);
let input = PathBuf::from(args.next().ok_or("需要 <rowids.txt>")?);
let layout = ple_layout();
let batch = engramdb_io::batch::BadgeGather::open(&dir, &layout).map_err(io_err)?;
let mut keys = Vec::new();
let text = std::fs::read_to_string(&input).map_err(io_err)?;
for line in text.lines() {
let line = line.trim();
if !line.is_empty() {
keys.push(line.parse::<u64>().map_err(|e| e.to_string())?);
}
}
let w = layout.width as usize;
if keys.is_empty() {
println!("fnv=0");
return Ok(());
}
let mut out = vec![0u8; keys.len() * w];
batch.gather_planned(&keys, &mut out).map_err(io_err)?;
println!("fnv={} keys={}", fnv(&out), keys.len());
Ok(())
}
fn cmd_bench_real(mut args: impl Iterator<Item = String>) -> Result<(), String> {
let dir = PathBuf::from(args.next().ok_or("需要 <badge_dir>")?);
let spec = PleSpec::real();
let n_tokens = 4096usize;
let mut state: u64 = 0x1234_5678_9ABC_DEF0;
let mut tokens: Vec<u32> = Vec::with_capacity(n_tokens);
for _ in 0..n_tokens {
state = state
.wrapping_mul(6_364_136_223_846_793_005)
.wrapping_add(1_442_695_040_888_963_407);
tokens.push((state % 248_320) as u32);
}
let mut rowids: Vec<u32> = Vec::with_capacity(n_tokens * 16);
for t in 0..n_tokens {
let c = tokens[t];
let triple = [tokens[t.saturating_sub(2)], tokens[t.saturating_sub(1)], c];
let ids = spec.rowids_for_seq(&triple);
rowids.extend_from_slice(&ids[0]);
}
let keys: Vec<u64> = rowids.iter().map(|&x| x as u64).collect();
let threads: usize = 8;
let layout = ple_layout();
let batch = engramdb_io::batch::BadgeGather::open(&dir, &layout).map_err(io_err)?;
let w = layout.width as usize;
let mut out = vec![0u8; keys.len() * w];
let t0 = std::time::Instant::now();
for _ in 0..8 {
batch.gather_pp(&keys, &mut out, threads).map_err(io_err)?;
black_box(&out);
}
let dt = t0.elapsed().as_secs_f64() / 8.0;
let rows_per_s = keys.len() as f64 / dt;
println!(
"rows/s={:.0} keys/batch={} badge-rows={} payload=KB/tok={:.1}",
rows_per_s,
keys.len(),
keys.len(),
keys.len() * 160 / 1024 / (keys.len() / 16)
);
Ok(())
}
fn black_box<T>(x: &T) {
unsafe {
std::ptr::read_volatile(&(x as *const T as usize));
}
}
fn io_err(e: std::io::Error) -> String {
e.to_string()
}
pub fn exported_fnv(bytes: &[u8]) -> u64 {
fnv(bytes)
}
#[allow(dead_code)]
fn _keep_serde(_p: &Path) {
let _ = serde_json::Value::Null;
}