use std::fs::{self, File};
use std::io::BufWriter;
use std::path::PathBuf;
use std::time::Instant;
use anyhow::{Context, Result};
use fst::{IntoStreamer, Streamer};
use memmap2::MmapOptions;
use trie_rs::{Trie, TrieBuilder};
const MEMORY_TARGET_BYTES: u64 = 1_000_000_000;
const LATENCY_TARGET_MEAN_NS: f64 = 1_000.0;
const OFFSET_STEP: u64 = 4096;
const KEYS_PER_TENANT: u64 = 1000;
#[cfg(feature = "dev-mode")]
const TRIE_SHARD_KEYS: u64 = 500_000;
#[cfg(not(feature = "dev-mode"))]
const TRIE_SHARD_KEYS: u64 = 2_000_000;
#[derive(Debug, Clone)]
pub struct ModuleEFstConfig {
pub keys: u64,
pub lookups: u64,
pub list_queries: u64,
pub index_path: PathBuf,
}
#[derive(Debug, Clone)]
pub struct ModuleETrieConfig {
pub keys: u64,
pub lookups: u64,
pub list_queries: u64,
}
#[derive(Debug, Clone)]
pub struct ModuleECompareConfig {
pub keys: u64,
pub lookups: u64,
pub list_queries: u64,
pub index_path: PathBuf,
}
#[derive(Debug, Clone)]
pub struct ModuleEFstStats {
pub keys: u64,
pub lookups: u64,
pub list_queries: u64,
pub build_elapsed_ms: f64,
pub query_elapsed_ms: f64,
pub index_size_bytes: u64,
pub vm_rss_bytes: u64,
pub lookup_mean_ns: f64,
pub lookup_p50_ns: u64,
pub lookup_p99_ns: u64,
pub list_seek_mean_ns: f64,
pub list_seek_p50_ns: u64,
pub list_seek_p99_ns: u64,
pub memory_target_met: bool,
pub latency_target_met: bool,
pub index_path: PathBuf,
}
impl ModuleEFstStats {
pub fn to_json(&self) -> String {
format!(
"{{\"module\":\"E_FST\",\"keys\":{},\"lookups\":{},\"list_queries\":{},\"build_elapsed_ms\":{:.3},\"query_elapsed_ms\":{:.3},\"index_size_bytes\":{},\"vm_rss_bytes\":{},\"lookup_mean_ns\":{:.2},\"lookup_p50_ns\":{},\"lookup_p99_ns\":{},\"list_seek_mean_ns\":{:.2},\"list_seek_p50_ns\":{},\"list_seek_p99_ns\":{},\"memory_target_met\":{},\"latency_target_met\":{},\"index_path\":\"{}\"}}",
self.keys,
self.lookups,
self.list_queries,
self.build_elapsed_ms,
self.query_elapsed_ms,
self.index_size_bytes,
self.vm_rss_bytes,
self.lookup_mean_ns,
self.lookup_p50_ns,
self.lookup_p99_ns,
self.list_seek_mean_ns,
self.list_seek_p50_ns,
self.list_seek_p99_ns,
self.memory_target_met,
self.latency_target_met,
self.index_path.display()
)
}
}
#[derive(Debug, Clone)]
pub struct ModuleETrieStats {
pub keys: u64,
pub lookups: u64,
pub list_queries: u64,
pub build_elapsed_ms: f64,
pub query_elapsed_ms: f64,
pub vm_rss_bytes: u64,
pub lookup_mean_ns: f64,
pub lookup_p50_ns: u64,
pub lookup_p99_ns: u64,
pub list_seek_mean_ns: f64,
pub list_seek_p50_ns: u64,
pub list_seek_p99_ns: u64,
pub memory_target_met: bool,
pub latency_target_met: bool,
}
impl ModuleETrieStats {
pub fn to_json(&self) -> String {
format!(
"{{\"module\":\"E_TRIE\",\"keys\":{},\"lookups\":{},\"list_queries\":{},\"build_elapsed_ms\":{:.3},\"query_elapsed_ms\":{:.3},\"vm_rss_bytes\":{},\"lookup_mean_ns\":{:.2},\"lookup_p50_ns\":{},\"lookup_p99_ns\":{},\"list_seek_mean_ns\":{:.2},\"list_seek_p50_ns\":{},\"list_seek_p99_ns\":{},\"memory_target_met\":{},\"latency_target_met\":{}}}",
self.keys,
self.lookups,
self.list_queries,
self.build_elapsed_ms,
self.query_elapsed_ms,
self.vm_rss_bytes,
self.lookup_mean_ns,
self.lookup_p50_ns,
self.lookup_p99_ns,
self.list_seek_mean_ns,
self.list_seek_p50_ns,
self.list_seek_p99_ns,
self.memory_target_met,
self.latency_target_met
)
}
}
#[derive(Debug, Clone)]
pub struct ModuleECompareStats {
pub fst: ModuleEFstStats,
pub trie: ModuleETrieStats,
pub recommendation: String,
}
impl ModuleECompareStats {
pub fn to_json(&self) -> String {
format!(
"{{\"module\":\"E_COMPARE\",\"recommendation\":\"{}\",\"fst\":{},\"trie\":{}}}",
self.recommendation,
self.fst.to_json(),
self.trie.to_json()
)
}
}
pub fn run_fst(config: ModuleEFstConfig) -> Result<ModuleEFstStats> {
validate_workload(config.keys, config.lookups, config.list_queries)?;
if let Some(parent) = config.index_path.parent() {
if !parent.as_os_str().is_empty() {
fs::create_dir_all(parent)
.with_context(|| format!("failed creating index parent {}", parent.display()))?;
}
}
let build_start = Instant::now();
{
let file = File::create(&config.index_path)
.with_context(|| format!("failed creating index {}", config.index_path.display()))?;
let writer = BufWriter::with_capacity(8 * 1024 * 1024, file);
let mut builder =
fst::MapBuilder::new(writer).context("failed creating fst map builder")?;
for i in 0..config.keys {
let key = key_for(i);
let value = i
.checked_mul(OFFSET_STEP)
.context("offset overflow while building fst")?;
builder
.insert(key, value)
.with_context(|| format!("failed inserting fst key at index {}", i))?;
}
builder.finish().context("failed finishing fst map")?;
}
let build_elapsed_ms = build_start.elapsed().as_secs_f64() * 1_000.0;
let index_size_bytes = fs::metadata(&config.index_path)
.with_context(|| format!("failed to stat fst index {}", config.index_path.display()))?
.len();
let query_start = Instant::now();
let file = File::open(&config.index_path)
.with_context(|| format!("failed opening fst index {}", config.index_path.display()))?;
let mapped = unsafe { MmapOptions::new().map(&file) }
.with_context(|| format!("failed mmap fst index {}", config.index_path.display()))?;
let map = fst::Map::new(&mapped[..]).context("failed loading mmap-backed fst map")?;
let mut lookup_latencies = Vec::with_capacity(config.lookups as usize);
let mut lookup_rng = fastrand::Rng::with_seed(0xC001_D00D_0001);
for _ in 0..config.lookups {
let idx = lookup_rng.u64(..config.keys);
let key = key_for(idx);
let start = Instant::now();
let value = map
.get(key.as_bytes())
.ok_or_else(|| anyhow::anyhow!("fst lookup missing key {}", key))?;
let elapsed_ns = start.elapsed().as_nanos() as u64;
let expected = idx
.checked_mul(OFFSET_STEP)
.context("offset overflow while validating fst lookup")?;
if value != expected {
anyhow::bail!(
"fst lookup value mismatch: expected {}, got {} for key {}",
expected,
value,
key
);
}
lookup_latencies.push(elapsed_ns);
}
let mut list_seek_latencies = Vec::with_capacity(config.list_queries as usize);
let mut list_rng = fastrand::Rng::with_seed(0xC001_D00D_0002);
for _ in 0..config.list_queries {
let idx = list_rng.u64(..config.keys);
let prefix = prefix_for_index(idx);
let start = Instant::now();
let mut range = map.range().ge(prefix.as_bytes());
if let Some(upper) = prefix_upper_bound(prefix.as_bytes()) {
range = range.lt(upper);
}
let mut stream = range.into_stream();
let elapsed_ns = start.elapsed().as_nanos() as u64;
let Some((key, _value)) = stream.next() else {
anyhow::bail!("fst prefix scan returned no result for prefix {}", prefix);
};
if !key.starts_with(prefix.as_bytes()) {
anyhow::bail!(
"fst prefix scan returned key outside prefix: prefix={}, key={}",
prefix,
String::from_utf8_lossy(key)
);
}
list_seek_latencies.push(elapsed_ns);
}
let query_elapsed_ms = query_start.elapsed().as_secs_f64() * 1_000.0;
let vm_rss_bytes = read_vm_rss_bytes();
let lookup = summarize_latencies(&mut lookup_latencies);
let list_seek = summarize_latencies(&mut list_seek_latencies);
let memory_target_met = vm_rss_bytes < MEMORY_TARGET_BYTES;
let latency_target_met =
lookup.mean_ns < LATENCY_TARGET_MEAN_NS && list_seek.mean_ns < LATENCY_TARGET_MEAN_NS;
Ok(ModuleEFstStats {
keys: config.keys,
lookups: config.lookups,
list_queries: config.list_queries,
build_elapsed_ms,
query_elapsed_ms,
index_size_bytes,
vm_rss_bytes,
lookup_mean_ns: lookup.mean_ns,
lookup_p50_ns: lookup.p50_ns,
lookup_p99_ns: lookup.p99_ns,
list_seek_mean_ns: list_seek.mean_ns,
list_seek_p50_ns: list_seek.p50_ns,
list_seek_p99_ns: list_seek.p99_ns,
memory_target_met,
latency_target_met,
index_path: config.index_path,
})
}
pub fn run_trie(config: ModuleETrieConfig) -> Result<ModuleETrieStats> {
validate_workload(config.keys, config.lookups, config.list_queries)?;
let shard_count = config.keys.div_ceil(TRIE_SHARD_KEYS) as usize;
let mut build_elapsed_ms = 0.0_f64;
let mut query_elapsed_ms = 0.0_f64;
let mut vm_rss_peak_bytes = 0_u64;
let mut lookup_latencies = Vec::with_capacity(config.lookups as usize);
let mut list_seek_latencies = Vec::with_capacity(config.list_queries as usize);
let mut lookup_rng = fastrand::Rng::with_seed(0xC001_D00D_0003);
let mut list_rng = fastrand::Rng::with_seed(0xC001_D00D_0004);
for shard_idx in 0..shard_count {
let shard_start = shard_idx as u64 * TRIE_SHARD_KEYS;
let shard_end = (shard_start + TRIE_SHARD_KEYS).min(config.keys);
let shard_keys = shard_end.saturating_sub(shard_start);
if shard_keys == 0 {
continue;
}
let build_start = Instant::now();
let mut builder = TrieBuilder::<u8>::new();
for i in shard_start..shard_end {
let key = key_for(i);
builder.push(key.as_bytes());
}
let trie: Trie<u8> = builder.build();
build_elapsed_ms += build_start.elapsed().as_secs_f64() * 1_000.0;
vm_rss_peak_bytes = vm_rss_peak_bytes.max(read_vm_rss_bytes());
let shard_lookup_queries = split_queries(config.lookups, shard_idx, shard_count);
let shard_list_queries = split_queries(config.list_queries, shard_idx, shard_count);
let query_start = Instant::now();
for _ in 0..shard_lookup_queries {
let idx = shard_start + lookup_rng.u64(..shard_keys);
let key = key_for(idx);
let start = Instant::now();
let found = trie.exact_match(key.as_bytes());
let elapsed_ns = start.elapsed().as_nanos() as u64;
if !found {
anyhow::bail!("trie lookup missing key {}", key);
}
lookup_latencies.push(elapsed_ns);
}
for _ in 0..shard_list_queries {
let idx = shard_start + list_rng.u64(..shard_keys);
let prefix = prefix_for_index(idx);
let start = Instant::now();
let first = trie
.predictive_search::<Vec<u8>, _>(prefix.as_bytes())
.next();
let elapsed_ns = start.elapsed().as_nanos() as u64;
let Some(key) = first else {
anyhow::bail!("trie prefix scan returned no result for prefix {}", prefix);
};
if !key.starts_with(prefix.as_bytes()) {
anyhow::bail!(
"trie prefix scan returned key outside prefix: prefix={}, key={}",
prefix,
String::from_utf8_lossy(&key)
);
}
list_seek_latencies.push(elapsed_ns);
}
query_elapsed_ms += query_start.elapsed().as_secs_f64() * 1_000.0;
}
if lookup_latencies.len() as u64 != config.lookups {
anyhow::bail!(
"trie lookup sample mismatch: expected {}, got {}",
config.lookups,
lookup_latencies.len()
);
}
if list_seek_latencies.len() as u64 != config.list_queries {
anyhow::bail!(
"trie list sample mismatch: expected {}, got {}",
config.list_queries,
list_seek_latencies.len()
);
}
let lookup = summarize_latencies(&mut lookup_latencies);
let list_seek = summarize_latencies(&mut list_seek_latencies);
let memory_target_met = vm_rss_peak_bytes < MEMORY_TARGET_BYTES;
let latency_target_met =
lookup.mean_ns < LATENCY_TARGET_MEAN_NS && list_seek.mean_ns < LATENCY_TARGET_MEAN_NS;
Ok(ModuleETrieStats {
keys: config.keys,
lookups: config.lookups,
list_queries: config.list_queries,
build_elapsed_ms,
query_elapsed_ms,
vm_rss_bytes: vm_rss_peak_bytes,
lookup_mean_ns: lookup.mean_ns,
lookup_p50_ns: lookup.p50_ns,
lookup_p99_ns: lookup.p99_ns,
list_seek_mean_ns: list_seek.mean_ns,
list_seek_p50_ns: list_seek.p50_ns,
list_seek_p99_ns: list_seek.p99_ns,
memory_target_met,
latency_target_met,
})
}
fn split_queries(total_queries: u64, shard_idx: usize, shard_count: usize) -> u64 {
if shard_count == 0 {
return 0;
}
let base = total_queries / shard_count as u64;
let remainder = total_queries % shard_count as u64;
if (shard_idx as u64) < remainder {
base + 1
} else {
base
}
}
pub fn run_compare(config: ModuleECompareConfig) -> Result<ModuleECompareStats> {
let fst = run_fst(ModuleEFstConfig {
keys: config.keys,
lookups: config.lookups,
list_queries: config.list_queries,
index_path: config.index_path,
})?;
let trie = run_trie(ModuleETrieConfig {
keys: config.keys,
lookups: config.lookups,
list_queries: config.list_queries,
})?;
let recommendation = pick_recommendation(&fst, &trie);
Ok(ModuleECompareStats {
fst,
trie,
recommendation,
})
}
fn validate_workload(keys: u64, lookups: u64, list_queries: u64) -> Result<()> {
if keys == 0 {
anyhow::bail!("keys must be > 0");
}
if lookups == 0 {
anyhow::bail!("lookups must be > 0");
}
if list_queries == 0 {
anyhow::bail!("list_queries must be > 0");
}
Ok(())
}
fn key_for(index: u64) -> String {
let tenant = index / KEYS_PER_TENANT;
format!("tenant-{tenant:08}/production/images/2026/02/25/image-{index:012}.jpg")
}
fn prefix_for_index(index: u64) -> String {
let tenant = index / KEYS_PER_TENANT;
format!("tenant-{tenant:08}/production/images/2026/02/25/image-")
}
fn prefix_upper_bound(prefix: &[u8]) -> Option<Vec<u8>> {
let mut upper = prefix.to_vec();
for idx in (0..upper.len()).rev() {
if upper[idx] != 0xFF {
upper[idx] += 1;
upper.truncate(idx + 1);
return Some(upper);
}
}
None
}
fn read_vm_rss_bytes() -> u64 {
let Ok(status) = fs::read_to_string("/proc/self/status") else {
return 0;
};
for line in status.lines() {
if !line.starts_with("VmRSS:") {
continue;
}
let mut parts = line.split_whitespace();
let _ = parts.next();
let Some(value_kib) = parts.next() else {
return 0;
};
if let Ok(kib) = value_kib.parse::<u64>() {
return kib.saturating_mul(1024);
}
return 0;
}
0
}
#[derive(Debug, Clone, Copy)]
struct LatencySummary {
mean_ns: f64,
p50_ns: u64,
p99_ns: u64,
}
fn summarize_latencies(samples: &mut [u64]) -> LatencySummary {
if samples.is_empty() {
return LatencySummary {
mean_ns: 0.0,
p50_ns: 0,
p99_ns: 0,
};
}
let sum = samples.iter().copied().map(|v| v as f64).sum::<f64>();
let mean_ns = sum / samples.len() as f64;
samples.sort_unstable();
let p50_idx = percentile_index(samples.len(), 50);
let p99_idx = percentile_index(samples.len(), 99);
LatencySummary {
mean_ns,
p50_ns: samples[p50_idx],
p99_ns: samples[p99_idx],
}
}
fn percentile_index(len: usize, percentile: usize) -> usize {
if len == 0 {
return 0;
}
let rank = ((len as f64) * (percentile as f64 / 100.0)).ceil() as usize;
rank.saturating_sub(1).min(len - 1)
}
fn pick_recommendation(fst: &ModuleEFstStats, trie: &ModuleETrieStats) -> String {
match (fst.memory_target_met, trie.memory_target_met) {
(true, false) => return "fst".to_string(),
(false, true) => return "trie-rs".to_string(),
_ => {}
}
if fst.lookup_p50_ns < trie.lookup_p50_ns {
return "fst".to_string();
}
if trie.lookup_p50_ns < fst.lookup_p50_ns {
return "trie-rs".to_string();
}
if fst.build_elapsed_ms <= trie.build_elapsed_ms {
"fst".to_string()
} else {
"trie-rs".to_string()
}
}
#[cfg(test)]
mod tests {
use std::time::{SystemTime, UNIX_EPOCH};
use super::*;
#[test]
fn keys_are_sorted_and_unique() {
let mut prev = key_for(0);
for idx in 1..20_000 {
let next = key_for(idx);
assert!(next > prev, "keys must be strictly increasing");
prev = next;
}
}
#[test]
fn prefix_upper_bound_increments_suffix() {
let upper = prefix_upper_bound(b"abc").expect("upper bound should exist");
assert_eq!(upper, b"abd");
}
#[test]
fn fst_small_round_trip() {
let path = unique_index_path("module-e-fst-small");
let stats = run_fst(ModuleEFstConfig {
keys: 10_000,
lookups: 2_000,
list_queries: 2_000,
index_path: path.clone(),
})
.expect("small fst run should succeed");
assert_eq!(stats.keys, 10_000);
assert!(stats.index_size_bytes > 0);
assert!(stats.lookup_p99_ns >= stats.lookup_p50_ns);
std::fs::remove_file(path).expect("temporary fst file should be removable");
}
#[test]
fn trie_small_round_trip() {
let stats = run_trie(ModuleETrieConfig {
keys: 10_000,
lookups: 2_000,
list_queries: 2_000,
})
.expect("small trie run should succeed");
assert_eq!(stats.keys, 10_000);
assert!(stats.lookup_p99_ns >= stats.lookup_p50_ns);
}
fn unique_index_path(prefix: &str) -> PathBuf {
let nanos = SystemTime::now()
.duration_since(UNIX_EPOCH)
.expect("clock should be after unix epoch")
.as_nanos();
std::env::temp_dir().join(format!("{}-{}-{}.fst", prefix, std::process::id(), nanos))
}
}