#![forbid(unsafe_code)]
use std::collections::BTreeMap;
use std::path::{Path, PathBuf};
use std::time::Instant;
use serde::{Deserialize, Serialize};
use crate::core::representation::{Representation, Residual};
use crate::evidence::corpus::{self, Corpus};
use crate::evidence::environment::{
DiskDelta, Environment, StatSummary, disk_delta, diskstats, summary,
};
use crate::optimizer::policy::OptimizeOptions;
use crate::store::inode::{Inode, InodeData};
use crate::store::transaction::CrashHooks;
use crate::store::{BTREE_ORDER, Store, StoreConfig};
pub struct CampaignOptions {
pub out_root: PathBuf,
pub repo_root: PathBuf,
pub scratch_dir: PathBuf,
pub runs: usize,
pub size_mib: u64,
pub cache_state: String,
pub policy_mode: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct CampaignResults {
pub campaign_dir: String,
pub created_unix: u64,
pub runs: Vec<CorpusModeRuns>,
pub ablation: AblationTable,
pub dsfb_investigation: DsfbInvestigation,
pub versioned_experiment: VersionedExperiment,
pub gc_traffic: GcTraffic,
pub post_gc_footprint: std::collections::BTreeMap<String, PostGcFootprint>,
pub tree_court: Option<TreeCourt>,
pub baselines: Baselines,
pub device_writes: Option<DiskDelta>,
pub admission: Vec<AdmissionItem>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct CorpusModeRuns {
pub corpus: String,
pub mode: String,
pub run_count: usize,
pub runs: Vec<RunMetrics>,
pub write_throughput: StatSummary,
pub read_throughput: StatSummary,
pub write_latency_us: StatSummary,
pub read_latency_us: StatSummary,
pub fsync_latency_us: StatSummary,
pub cpu_user_s: StatSummary,
pub cpu_sys_s: StatSummary,
pub physical_median: u64,
pub ratio_median: f64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct RunMetrics {
pub run: usize,
pub write_mbps: f64,
pub read_mbps: f64,
pub write_wall_s: f64,
pub read_wall_s: f64,
pub cpu_user_s: f64,
pub cpu_sys_s: f64,
pub logical_bytes: u64,
pub written_bytes: u64,
pub reachable_bytes: u64,
pub total_backing_bytes: u64,
pub unreachable_bytes: u64,
pub ratio_reachable: f64,
pub ratio_total_backing: f64,
pub result_hash: String,
pub hash_matches_input: bool,
pub families: BTreeMap<String, u64>,
pub accounting: Accounting,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct Accounting {
pub payload_bytes: u64,
pub model_bytes: u64,
pub residual_bytes: u64,
pub descriptor_bytes: u64,
pub metadata_bytes: u64,
pub integrity_bytes_est: u64,
pub allocator_overhead_bytes: u64,
pub unreclaimed_bytes: u64,
pub cas_shared_bytes_saved: u64,
pub exact_ref_bytes_saved: u64,
pub check: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct AblationRow {
pub mode: String,
pub physical: u64,
pub ratio: f64,
pub write_mbps: f64,
pub cpu_user_s: f64,
pub families: BTreeMap<String, u64>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct AblationTable {
pub corpus: String,
pub size_mib: u64,
pub rows: Vec<AblationRow>,
pub cumulative_rows: Vec<AblationRow>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct DsfbInvestigation {
pub corpus: String,
pub runs: usize,
pub full: Vec<DsfbRun>,
pub no_dsfb: Vec<DsfbRun>,
pub physical_identical: bool,
pub write_mbps_full: StatSummary,
pub write_mbps_no_dsfb: StatSummary,
pub cpu_user_s_full: StatSummary,
pub cpu_user_s_no_dsfb: StatSummary,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct DsfbRun {
pub write_mbps: f64,
pub cpu_user_s: f64,
pub physical: u64,
pub write_p50_us: f64,
pub write_p95_us: f64,
pub write_p99_us: f64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct VersionedExperiment {
pub corpus: String,
pub versions: usize,
pub sequential_full: Vec<u64>,
pub sequential_no_base: Vec<u64>,
pub shuffled_full: Vec<u64>,
pub sequential_full_post_gc: u64,
pub sequential_no_base_post_gc: u64,
pub shuffled_full_post_gc: u64,
pub sequential_ratio: f64,
pub shuffled_ratio: f64,
pub base_savings_reachable_bytes: i64,
pub base_savings_pct: f64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PostGcFootprint {
pub corpus: String,
pub logical: u64,
pub reachable: u64,
pub total_backing: u64,
pub allocated_blocks: u64,
pub ratio_reachable: f64,
pub ratio_total_backing: f64,
pub ratio_allocated: f64,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct TreeCourt {
pub file_count: usize,
pub single_chunk_files: usize,
pub logical_bytes: u64,
pub zstd_whole_l1: Option<CompressionBaseline>,
pub zstd_whole_l19: Option<CompressionBaseline>,
pub zstd_per_file_l1: Option<CompressionBaseline>,
pub zstd_per_file_l19: Option<CompressionBaseline>,
pub zstd_per_64k_l1: Option<CompressionBaseline>,
pub zstd_per_64k_l19: Option<CompressionBaseline>,
pub efs_tree_reachable: u64,
pub efs_tree_backing: u64,
pub efs_tree_families: BTreeMap<String, u64>,
pub efs_shared_reachable: u64,
pub efs_shared_backing: u64,
pub efs_shared_families: BTreeMap<String, u64>,
pub shared_rewrites: u64,
pub shared_saved_bytes: u64,
pub efs_model_reachable: u64,
pub efs_model_backing: u64,
pub efs_model_families: BTreeMap<String, u64>,
pub model_rewrites: u64,
pub model_saved_bytes: u64,
pub physical_dead_indexed_bytes: u64,
pub physical_index_hidden_bytes: u64,
pub physical_unindexed_bytes: u64,
pub physical_backing_bytes: u64,
pub compact_backing_bytes: u64,
pub compact_reachable_bytes: u64,
pub compact_overhead_bytes: u64,
pub compact_reclaimed_bytes: u64,
pub zstd_dir_anchor_l1: Option<CompressionBaseline>,
pub per_extent_descriptor_bytes: u64,
pub per_extent_model_bytes: u64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct GcTraffic {
pub unreachable_before: u64,
pub reclaimed_bytes: u64,
pub unreachable_after: u64,
pub physical_before: u64,
pub physical_after: u64,
pub gc_wall_s: f64,
pub optimizer_scanned: u64,
pub optimizer_rewritten: u64,
pub optimizer_saved_bytes: u64,
pub unreachable_by_tag_after: std::collections::BTreeMap<String, u64>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct Baselines {
pub raw_file: Option<RawBaseline>,
pub zstd_level_1: Option<CompressionBaseline>,
pub zstd_level_19: Option<CompressionBaseline>,
pub zstd_per_64k_level_1: Option<CompressionBaseline>,
pub zstd_per_64k_level_19: Option<CompressionBaseline>,
pub direct_rans_src: Option<RunMetrics>,
pub sequence_rans_src: Option<RunMetrics>,
pub sequence_deep_src: Option<RunMetrics>,
pub waived: Vec<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct RawBaseline {
pub path: String,
pub fstype: String,
pub bytes: u64,
pub write_mbps: f64,
pub ratio: f64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct CompressionBaseline {
pub tool: String,
pub version: String,
pub level: String,
pub input_bytes: u64,
pub output_bytes: u64,
pub ratio: f64,
pub wall_s: f64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct AdmissionItem {
pub rule: String,
pub met: bool,
pub note: String,
}
pub fn run(opts: &CampaignOptions) -> Result<PathBuf, String> {
std::fs::create_dir_all(&opts.scratch_dir).map_err(|e| e.to_string())?;
std::fs::create_dir_all(&opts.out_root).map_err(|e| e.to_string())?;
let rev = Environment::capture(
&opts.repo_root,
&opts.scratch_dir,
&opts.cache_state,
&opts.policy_mode,
)
.revision_short;
let rev_slug = if rev.is_empty() { "norev" } else { &rev };
let created = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0);
let dir = opts.out_root.join(format!("campaign-{created}-{rev_slug}"));
std::fs::create_dir_all(&dir).map_err(|e| e.to_string())?;
let mut log = String::new();
line(
&mut log,
&format!("entropyfs evidence campaign — {created}"),
);
line(&mut log, &format!("revision: {rev_slug}"));
let env = Environment::capture(
&opts.repo_root,
&opts.scratch_dir,
&opts.cache_state,
&opts.policy_mode,
);
write_json(&dir, "environment.json", &env)?;
let corpora = build_corpora(opts)?;
let corpus_manifest: Vec<serde_json::Value> = corpora
.iter()
.map(|c| {
serde_json::json!({
"name": c.name,
"source": c.source,
"description": c.description,
"logical_bytes": c.logical_bytes(),
"written_bytes": c.written_bytes(),
"content_hash": c.content_hash(),
"version_count": c.versions.len(),
"version_hashes": c.version_hashes(),
})
})
.collect();
write_json(&dir, "corpus-manifest.json", &corpus_manifest)?;
let device = env.store_device.clone();
let dev_name = device.trim_start_matches("/dev/");
let disk_before = diskstats(dev_name);
let mut results = CampaignResults {
campaign_dir: dir.display().to_string(),
created_unix: created,
runs: Vec::new(),
ablation: AblationTable {
corpus: "structured".into(),
size_mib: opts.size_mib,
rows: Vec::new(),
cumulative_rows: Vec::new(),
},
dsfb_investigation: DsfbInvestigation {
corpus: "structured".into(),
runs: 0,
full: Vec::new(),
no_dsfb: Vec::new(),
physical_identical: false,
write_mbps_full: StatSummary::default(),
write_mbps_no_dsfb: StatSummary::default(),
cpu_user_s_full: StatSummary::default(),
cpu_user_s_no_dsfb: StatSummary::default(),
},
versioned_experiment: VersionedExperiment {
corpus: "versioned".into(),
versions: 0,
sequential_full: Vec::new(),
sequential_no_base: Vec::new(),
shuffled_full: Vec::new(),
sequential_full_post_gc: 0,
sequential_no_base_post_gc: 0,
shuffled_full_post_gc: 0,
sequential_ratio: 0.0,
shuffled_ratio: 0.0,
base_savings_reachable_bytes: 0,
base_savings_pct: 0.0,
},
gc_traffic: GcTraffic {
unreachable_before: 0,
reclaimed_bytes: 0,
unreachable_after: 0,
physical_before: 0,
physical_after: 0,
gc_wall_s: 0.0,
optimizer_scanned: 0,
optimizer_rewritten: 0,
optimizer_saved_bytes: 0,
unreachable_by_tag_after: std::collections::BTreeMap::new(),
},
post_gc_footprint: std::collections::BTreeMap::new(),
tree_court: None,
baselines: Baselines::default(),
device_writes: None,
admission: Vec::new(),
};
let structured = structured_corpus(opts.size_mib);
{
let c = corpora
.iter()
.find(|c| c.name == "structured")
.expect("corpus present");
let group = run_repeated(
opts,
c,
"full",
OptimizeOptions::default(),
opts.runs,
&mut log,
)?;
results.runs.push(group);
}
for name in ["src", "urandom", "compressed-z19"] {
let c = corpora
.iter()
.find(|c| c.name == name)
.expect("corpus present");
let group = run_repeated(opts, c, "full", OptimizeOptions::default(), 3, &mut log)?;
results.runs.push(group);
}
line(&mut log, "\n== leave-one-out ablation (structured) ==");
for (mode, options) in OptimizeOptions::ablation_modes() {
let tmp = scratch_tempdir(&opts.scratch_dir, "abl-")?;
let store = fresh_store(tmp.path())?;
let o = write_only(&store, 3, &structured, options)?;
let n = store_numbers(&store)?;
let row = AblationRow {
mode: mode.to_string(),
physical: n.reachable,
ratio: n.logical as f64 / n.reachable.max(1) as f64,
write_mbps: o.metrics.write_mbps,
cpu_user_s: o.metrics.cpu_user_s,
families: o.families,
};
line(
&mut log,
&format!(
" {mode:<10} physical {:>12} ratio {:>7.3}x write {:>8.1} MiB/s cpu {:.3}+{:.3}s (p95 write {:.0}µs)",
row.physical,
row.ratio,
row.write_mbps,
row.cpu_user_s,
o.metrics.cpu_sys_s,
summary(&o.latencies).p95 * 1e6,
),
);
results.ablation.rows.push(row);
}
line(&mut log, "\n== cumulative ladder A0-A8 (structured) ==");
for (mode, options, run_background) in OptimizeOptions::cumulative_ladder_modes() {
let tmp = scratch_tempdir(&opts.scratch_dir, "ladder-")?;
let store = fresh_store(tmp.path())?;
let o = write_only(&store, 3, &structured, options)?;
let mut n = store_numbers(&store)?;
if run_background {
let opt = crate::optimizer::background::optimize_pass(&store, options, None, None)
.map_err(|e| e.to_string())?;
let _ = crate::optimizer::background::shared_dict_pass(&store, options, None)
.map_err(|e| e.to_string())?;
n = store_numbers(&store)?;
let _ = opt;
}
let row = AblationRow {
mode: mode.to_string(),
physical: n.reachable,
ratio: n.logical as f64 / n.reachable.max(1) as f64,
write_mbps: o.metrics.write_mbps,
cpu_user_s: o.metrics.cpu_user_s,
families: o.families,
};
line(
&mut log,
&format!(
" {mode:<18} physical {:>12} ratio {:>7.3}x write {:>8.1} MiB/s cpu {:.3}+{:.3}s (p95 write {:.0}µs)",
row.physical,
row.ratio,
row.write_mbps,
row.cpu_user_s,
o.metrics.cpu_sys_s,
summary(&o.latencies).p95 * 1e6,
),
);
results.ablation.cumulative_rows.push(row);
}
line(
&mut log,
"\n== DSFB search-budget investigation (structured) ==",
);
for mode in ["full", "no-dsfb"] {
let options = options_for(mode)?;
let mut runs = Vec::new();
for _ in 0..opts.runs {
let tmp = scratch_tempdir(&opts.scratch_dir, "dsfb-")?;
let store = fresh_store(tmp.path())?;
let o = write_only(&store, 3, &structured, options)?;
let s = summary(&o.latencies);
let n = store_numbers(&store)?;
runs.push(DsfbRun {
write_mbps: o.metrics.write_mbps,
cpu_user_s: o.metrics.cpu_user_s,
physical: n.reachable,
write_p50_us: s.p50 * 1e6,
write_p95_us: s.p95 * 1e6,
write_p99_us: s.p99 * 1e6,
});
}
let w: Vec<f64> = runs.iter().map(|r| r.write_mbps).collect();
let c: Vec<f64> = runs.iter().map(|r| r.cpu_user_s).collect();
let ws = summary(&w);
let cs = summary(&c);
line(
&mut log,
&format!(
" {mode:<9} write median {:>7.1} MiB/s (min {:.1}, max {:.1}) cpu median {:.3}s physical {:?}",
ws.p50,
ws.min,
ws.max,
cs.p50,
runs.iter().map(|r| r.physical).collect::<Vec<_>>(),
),
);
if mode == "full" {
results.dsfb_investigation.full = runs;
results.dsfb_investigation.write_mbps_full = ws;
results.dsfb_investigation.cpu_user_s_full = cs;
} else {
results.dsfb_investigation.no_dsfb = runs;
results.dsfb_investigation.write_mbps_no_dsfb = ws;
results.dsfb_investigation.cpu_user_s_no_dsfb = cs;
}
}
results.dsfb_investigation.runs = opts.runs;
results.dsfb_investigation.physical_identical = {
let f = results.dsfb_investigation.full.first().map(|r| r.physical);
let n = results
.dsfb_investigation
.no_dsfb
.first()
.map(|r| r.physical);
f.is_some() && f == n
};
line(
&mut log,
&format!(
" physical identical across modes: {}",
results.dsfb_investigation.physical_identical
),
);
line(&mut log, "\n== versioned experiment (H2) ==");
let vseq = corpora
.iter()
.find(|c| c.name == "versioned")
.expect("present");
let vshuf = corpora
.iter()
.find(|c| c.name == "shuffled")
.expect("present");
results.versioned_experiment.versions = vseq.versions.len();
let seq_full = run_repeated(opts, vseq, "full", OptimizeOptions::default(), 3, &mut log)?;
let seq_nb = run_repeated(opts, vseq, "no-base", options_for("no-base")?, 3, &mut log)?;
let shuf_full = run_repeated(opts, vshuf, "full", OptimizeOptions::default(), 3, &mut log)?;
results.versioned_experiment.sequential_full =
seq_full.runs.iter().map(|r| r.reachable_bytes).collect();
results.versioned_experiment.sequential_no_base =
seq_nb.runs.iter().map(|r| r.reachable_bytes).collect();
results.versioned_experiment.shuffled_full =
shuf_full.runs.iter().map(|r| r.reachable_bytes).collect();
results.runs.push(seq_full);
results.runs.push(seq_nb);
results.runs.push(shuf_full);
let seq_median = median(&results.versioned_experiment.sequential_full);
let shuf_median = median(&results.versioned_experiment.shuffled_full);
results.versioned_experiment.sequential_ratio =
vseq.logical_bytes() as f64 / seq_median.max(1) as f64;
results.versioned_experiment.shuffled_ratio =
vshuf.logical_bytes() as f64 / shuf_median.max(1) as f64;
results.versioned_experiment.base_savings_reachable_bytes =
shuf_median as i64 - seq_median as i64;
results.versioned_experiment.base_savings_pct = if shuf_median > 0 {
(shuf_median as f64 - seq_median as f64) / shuf_median as f64 * 100.0
} else {
0.0
};
line(
&mut log,
&format!(
" sequential median reachable: {seq_median} bytes ({:.3}x)",
vseq.logical_bytes() as f64 / seq_median.max(1) as f64
),
);
line(
&mut log,
&format!(
" shuffled median reachable: {shuf_median} bytes ({:.3}x)",
vshuf.logical_bytes() as f64 / shuf_median.max(1) as f64
),
);
line(
&mut log,
&format!(
" base+residual savings vs shuffled: {} bytes ({:.1}% of shuffled reachable)",
results.versioned_experiment.base_savings_reachable_bytes,
results.versioned_experiment.base_savings_pct
),
);
let sg = write_gc_reachable(opts, vseq, OptimizeOptions::default())?;
let ng = write_gc_reachable(opts, vseq, options_for("no-base")?)?;
let shg = write_gc_reachable(opts, vshuf, OptimizeOptions::default())?;
results.versioned_experiment.sequential_full_post_gc = sg;
results.versioned_experiment.sequential_no_base_post_gc = ng;
results.versioned_experiment.shuffled_full_post_gc = shg;
line(
&mut log,
&format!(
" post-GC reachable: sequential full {sg} ({:.3}x) / no-base {ng} ({:.3}x) / shuffled {shg} ({:.3}x)",
vseq.logical_bytes() as f64 / sg.max(1) as f64,
vseq.logical_bytes() as f64 / ng.max(1) as f64,
vshuf.logical_bytes() as f64 / shg.max(1) as f64,
),
);
line(&mut log, "\n== GC and optimizer traffic ==");
let gc = run_gc_traffic(opts)?;
results.gc_traffic = gc.clone();
line(
&mut log,
&format!(
" unreachable before {} → reclaimed {} → after {}; physical {} → {}; gc {:.3}s; optimizer scanned {} rewrote {} saved {}",
gc.unreachable_before,
gc.reclaimed_bytes,
gc.unreachable_after,
gc.physical_before,
gc.physical_after,
gc.gc_wall_s,
gc.optimizer_scanned,
gc.optimizer_rewritten,
gc.optimizer_saved_bytes
),
);
line(
&mut log,
&format!(
" unreachable by record tag (post-GC): {:?}",
gc.unreachable_by_tag_after
),
);
line(&mut log, "\n== post-GC physical footprint ==");
for c in corpora.iter().filter(|c| {
matches!(
c.name.as_str(),
"structured" | "src" | "urandom" | "compressed-z19"
)
}) {
let fp = write_gc_footprint(opts, c, OptimizeOptions::default())?;
line(
&mut log,
&format!(
" {}: logical {} → reachable {} ({:.2}x) / total backing {} ({:.2}x) / allocated {} ({:.2}x)",
fp.corpus,
fp.logical,
fp.reachable,
fp.ratio_reachable,
fp.total_backing,
fp.ratio_total_backing,
fp.allocated_blocks,
fp.ratio_allocated
),
);
results.post_gc_footprint.insert(fp.corpus.clone(), fp);
}
line(&mut log, "\n== Phase-9C tree court ==");
let tree_court = run_tree_court(opts)?;
results.tree_court = Some(tree_court.clone());
line(
&mut log,
&format!(
" files {} (single-chunk {}), logical {} B",
tree_court.file_count, tree_court.single_chunk_files, tree_court.logical_bytes
),
);
for (label, b) in [
("zstd -1 whole", &tree_court.zstd_whole_l1),
("zstd -19 whole", &tree_court.zstd_whole_l19),
("zstd -1 per-file", &tree_court.zstd_per_file_l1),
("zstd -19 per-file", &tree_court.zstd_per_file_l19),
("zstd -1 per-64KiB", &tree_court.zstd_per_64k_l1),
("zstd -19 per-64KiB", &tree_court.zstd_per_64k_l19),
] {
if let Some(b) = b {
line(
&mut log,
&format!(" {label:<22} {:>10} B ({:.3}x)", b.output_bytes, b.ratio),
);
}
}
line(
&mut log,
&format!(
" efs tree (post-GC): {:>10} B reachable ({:.3}x) / {} B backing",
tree_court.efs_tree_reachable,
tree_court.logical_bytes as f64 / tree_court.efs_tree_reachable.max(1) as f64,
tree_court.efs_tree_backing
),
);
line(
&mut log,
&format!(
" efs tree + shared dict: {:>10} B reachable ({:.3}x) / {} B backing (rewrote {} extents, saved {} B)",
tree_court.efs_shared_reachable,
tree_court.logical_bytes as f64 / tree_court.efs_shared_reachable.max(1) as f64,
tree_court.efs_shared_backing,
tree_court.shared_rewrites,
tree_court.shared_saved_bytes
),
);
line(
&mut log,
&format!(
" efs + model bundles: {:>10} B reachable ({:.3}x) / {} B backing (rewrote {} extents, saved {} B) [Phase-9G]",
tree_court.efs_model_reachable,
tree_court.logical_bytes as f64 / tree_court.efs_model_reachable.max(1) as f64,
tree_court.efs_model_backing,
tree_court.model_rewrites,
tree_court.model_saved_bytes
),
);
line(
&mut log,
&format!(" families before: {:?}", tree_court.efs_tree_families),
);
line(
&mut log,
&format!(" families after: {:?}", tree_court.efs_shared_families),
);
line(
&mut log,
&format!(" families +models: {:?}", tree_court.efs_model_families),
);
if let Some(b) = &tree_court.zstd_dir_anchor_l1 {
line(
&mut log,
&format!(
" zstd -1 per-file +dir anchor: {:>10} B ({:.3}x) [Phase-9F anchor-policy control]",
b.output_bytes, b.ratio
),
);
}
line(
&mut log,
&format!(
" physical post-GC: {:>10} B backing = reachable + {} B dead-indexed + {} B index-hidden + {} B unindexed [Phase-9H]",
tree_court.physical_backing_bytes,
tree_court.physical_dead_indexed_bytes,
tree_court.physical_index_hidden_bytes,
tree_court.physical_unindexed_bytes
),
);
line(
&mut log,
&format!(
" + full compact: {:>10} B backing ({} B overhead over reachable; reclaimed {} B) [Phase-9H]",
tree_court.compact_backing_bytes,
tree_court.compact_overhead_bytes,
tree_court.compact_reclaimed_bytes
),
);
let per_extent = tree_court
.per_extent_descriptor_bytes
.saturating_add(tree_court.per_extent_model_bytes);
line(
&mut log,
&format!(
" per-extent overhead: {} B descriptors + {} B models = {} B ({:.1}% of footprint, {:.1}% of logical) [Phase-9F]",
tree_court.per_extent_descriptor_bytes,
tree_court.per_extent_model_bytes,
per_extent,
100.0 * per_extent as f64 / tree_court.efs_model_reachable.max(1) as f64,
100.0 * per_extent as f64 / tree_court.logical_bytes.max(1) as f64
),
);
line(&mut log, "\n== baselines ==");
let src_pack = corpora
.iter()
.find(|c| c.name == "src")
.expect("present")
.final_bytes()
.to_vec();
results.baselines = run_baselines(opts, &src_pack, &corpora)?;
for w in &results.baselines.waived {
line(&mut log, &format!(" waived: {w}"));
}
if let Some(r) = &results.baselines.raw_file {
line(
&mut log,
&format!(
" raw file ({}): {:.1} MiB/s write, ratio {:.3}x",
r.fstype, r.write_mbps, r.ratio
),
);
}
for (name, b) in [
("zstd -1", &results.baselines.zstd_level_1),
("zstd -19", &results.baselines.zstd_level_19),
("zstd -1 per 64KiB", &results.baselines.zstd_per_64k_level_1),
(
"zstd -19 per 64KiB",
&results.baselines.zstd_per_64k_level_19,
),
] {
if let Some(b) = b {
line(
&mut log,
&format!(
" {name}: {} → {} bytes ({:.3}x), {:.3}s",
b.input_bytes, b.output_bytes, b.ratio, b.wall_s
),
);
}
}
if let Some(r) = &results.baselines.direct_rans_src {
line(
&mut log,
&format!(
" direct byte rANS (same backend, src corpus): {} → {} bytes ({:.3}x)",
r.logical_bytes, r.reachable_bytes, r.ratio_reachable
),
);
}
if let Some(r) = &results.baselines.sequence_rans_src {
line(
&mut log,
&format!(
" standalone SequenceRans (src corpus): {} → {} bytes ({:.3}x)",
r.logical_bytes, r.reachable_bytes, r.ratio_reachable
),
);
}
if let Some(r) = &results.baselines.sequence_deep_src {
line(
&mut log,
&format!(
" standalone SequenceDeep (src corpus): {} → {} bytes ({:.3}x)",
r.logical_bytes, r.reachable_bytes, r.ratio_reachable
),
);
}
if let Some(before) = disk_before {
if let Some(after) = diskstats(dev_name) {
let delta = disk_delta(dev_name, &before, &after);
line(
&mut log,
&format!(
"device {}: {} sectors written ({} bytes), {} sectors read ({} bytes)",
delta.device,
delta.write_sectors,
delta.written_bytes(),
delta.read_sectors,
delta.read_bytes()
),
);
results.device_writes = Some(delta);
}
}
results.admission = admission_checklist(&results, &corpora);
line(&mut log, "\n== admission checklist (methodology §8) ==");
for a in &results.admission {
line(
&mut log,
&format!(
" [{}] {} — {}",
if a.met { "OK " } else { "FAIL" },
a.rule,
a.note
),
);
}
write_json(&dir, "results.json", &results)?;
write_json(&dir, "result-hashes.json", &result_hashes(&results))?;
std::fs::write(dir.join("results.csv"), csv(&results)).map_err(|e| e.to_string())?;
std::fs::write(dir.join("raw-output.txt"), &log).map_err(|e| e.to_string())?;
std::fs::write(dir.join("report.md"), report(&results, &log)).map_err(|e| e.to_string())?;
println!("{log}");
line(
&mut log,
&format!("\ncampaign evidence written to {}", dir.display()),
);
println!("\ncampaign evidence written to {}", dir.display());
Ok(dir)
}
fn run_repeated(
opts: &CampaignOptions,
corpus: &Corpus,
mode: &str,
options: OptimizeOptions,
runs: usize,
log: &mut String,
) -> Result<CorpusModeRuns, String> {
let mut metrics: Vec<RunMetrics> = Vec::new();
let mut write_lats: Vec<f64> = Vec::new();
let mut read_lats: Vec<f64> = Vec::new();
let mut fsync_lats: Vec<f64> = Vec::new();
let mut write_mbps_all: Vec<f64> = Vec::new();
let mut read_mbps_all: Vec<f64> = Vec::new();
let mut cpu_user_all: Vec<f64> = Vec::new();
let mut cpu_sys_all: Vec<f64> = Vec::new();
let mut physical_all: Vec<u64> = Vec::new();
for i in 0..runs {
let tmp = scratch_tempdir(&opts.scratch_dir, "run-")?;
let store = fresh_store(tmp.path())?;
let mut m = full_run(&store, corpus, options)?;
m.metrics.run = i;
m.metrics.hash_matches_input = m.metrics.result_hash == corpus.content_hash();
write_lats.extend(m.write_latencies);
read_lats.extend(m.read_latencies);
fsync_lats.extend(m.fsync_latencies);
write_mbps_all.push(m.metrics.write_mbps);
read_mbps_all.push(m.metrics.read_mbps);
cpu_user_all.push(m.metrics.cpu_user_s);
cpu_sys_all.push(m.metrics.cpu_sys_s);
physical_all.push(m.metrics.reachable_bytes);
metrics.push(m.metrics);
}
let group = CorpusModeRuns {
corpus: corpus.name.clone(),
mode: mode.to_string(),
run_count: runs,
runs: metrics,
write_throughput: summary(&write_mbps_all),
read_throughput: summary(&read_mbps_all),
write_latency_us: scale_summary(&summary(&write_lats), 1e6),
read_latency_us: scale_summary(&summary(&read_lats), 1e6),
fsync_latency_us: scale_summary(&summary(&fsync_lats), 1e6),
cpu_user_s: summary(&cpu_user_all),
cpu_sys_s: summary(&cpu_sys_all),
physical_median: median(&physical_all),
ratio_median: if median(&physical_all) > 0 {
corpus.logical_bytes() as f64 / median(&physical_all) as f64
} else {
0.0
},
};
line(
log,
&format!(
" {} [{}] {} runs: write {:.1} MiB/s (p50 {:.0}µs, p95 {:.0}µs, p99 {:.0}µs) read {:.1} MiB/s fsync p50 {:.0}µs p95 {:.0}µs p99 {:.0}µs physical median {} ratio {:.3}x",
corpus.name,
mode,
runs,
group.write_throughput.p50,
group.write_latency_us.p50,
group.write_latency_us.p95,
group.write_latency_us.p99,
group.read_throughput.p50,
group.fsync_latency_us.p50,
group.fsync_latency_us.p95,
group.fsync_latency_us.p99,
group.physical_median,
group.ratio_median
),
);
Ok(group)
}
struct RunOutcome {
metrics: RunMetrics,
write_latencies: Vec<f64>,
read_latencies: Vec<f64>,
fsync_latencies: Vec<f64>,
}
fn full_run(
store: &Store,
corpus: &Corpus,
options: OptimizeOptions,
) -> Result<RunOutcome, String> {
let cpu0 = cpu_ticks();
let o = write_only(store, 3, corpus, options)?;
finish_run(store, corpus, o, cpu0, options)
}
fn full_run_deep(
store: &Store,
corpus: &Corpus,
options: OptimizeOptions,
) -> Result<RunOutcome, String> {
let cpu0 = cpu_ticks();
let o = write_only(store, 3, corpus, options)?;
crate::optimizer::background::optimize_pass(store, options, None, None)
.map_err(|e| e.to_string())?;
finish_run(store, corpus, o, cpu0, options)
}
fn finish_run(
store: &Store,
corpus: &Corpus,
o: WriteOutcome,
cpu0: (f64, f64),
_options: OptimizeOptions,
) -> Result<RunOutcome, String> {
let (write_metrics, write_lats) = (o.metrics, o.latencies);
let mut fsync_lats = Vec::new();
for _ in 0..5 {
let t0 = Instant::now();
store
.durability_barrier(&CrashHooks::none())
.map_err(|e| e.to_string())?;
fsync_lats.push(t0.elapsed().as_secs_f64());
}
let mut read_lats = Vec::new();
let mstart = Instant::now();
let mut hasher = blake3::Hasher::new();
let total = corpus.logical_bytes();
let mut off = 0u64;
while off < total {
let want = 65536u64.min(total - off);
let t0 = Instant::now();
let data = store.read_file(3, off, want).map_err(|e| e.to_string())?;
read_lats.push(t0.elapsed().as_secs_f64());
if data.len() as u64 != want {
return Err(format!("read length mismatch at {off}"));
}
hasher.update(&data);
off += want;
}
let read_wall = mstart.elapsed().as_secs_f64();
let result_hash = hasher.finalize().to_hex().to_string();
let cpu1 = cpu_ticks();
let n = store_numbers(store)?;
let (mut acct, families) = extent_decomposition(store, 3)?;
acct.metadata_bytes = n
.reachable
.saturating_sub(acct.payload_bytes + acct.model_bytes + acct.descriptor_bytes);
acct.integrity_bytes_est = n.record_count.saturating_mul(4).saturating_add(64);
acct.allocator_overhead_bytes = n.allocator_overhead;
acct.unreclaimed_bytes = n.unreachable;
let sum = acct.payload_bytes + acct.model_bytes + acct.descriptor_bytes + acct.metadata_bytes;
acct.check = if sum == n.reachable {
"ok".to_string()
} else {
format!(
"mismatch: extent decomposition {sum} != reachable {}",
n.reachable
)
};
Ok(RunOutcome {
metrics: RunMetrics {
run: 0,
write_mbps: write_metrics.write_mbps,
read_mbps: total as f64 / read_wall / (1024.0 * 1024.0),
write_wall_s: write_metrics.write_wall_s,
read_wall_s: read_wall,
cpu_user_s: cpu1.0 - cpu0.0,
cpu_sys_s: cpu1.1 - cpu0.1,
logical_bytes: n.logical,
written_bytes: corpus.written_bytes(),
reachable_bytes: n.reachable,
total_backing_bytes: n.total_backing,
unreachable_bytes: n.unreachable,
ratio_reachable: n.logical as f64 / n.reachable.max(1) as f64,
ratio_total_backing: n.logical as f64 / n.total_backing.max(1) as f64,
result_hash,
hash_matches_input: false, families,
accounting: acct,
},
write_latencies: write_lats,
read_latencies: read_lats,
fsync_latencies: fsync_lats,
})
}
struct WriteMetrics {
write_mbps: f64,
write_wall_s: f64,
cpu_user_s: f64,
cpu_sys_s: f64,
}
struct WriteOutcome {
metrics: WriteMetrics,
latencies: Vec<f64>,
families: BTreeMap<String, u64>,
}
fn write_only(
store: &Store,
ino: u64,
corpus: &Corpus,
options: OptimizeOptions,
) -> Result<WriteOutcome, String> {
let cpu0 = cpu_ticks();
let start = Instant::now();
let mut lats: Vec<f64> = Vec::new();
for version in &corpus.versions {
let mut writes: Vec<(u64, Vec<u8>)> = Vec::new();
let mut off = 0u64;
while off < version.len() as u64 {
let len = 65536u64.min(version.len() as u64 - off);
writes.push((off, version[off as usize..(off + len) as usize].to_vec()));
off += len;
}
let t0 = Instant::now();
store
.write_region_batch(ino, &writes, options)
.map_err(|e| e.to_string())?;
lats.push(t0.elapsed().as_secs_f64());
}
let wall = start.elapsed().as_secs_f64();
let cpu1 = cpu_ticks();
let (_, families) = extent_decomposition(store, ino)?;
Ok(WriteOutcome {
metrics: WriteMetrics {
write_mbps: corpus.written_bytes() as f64 / wall / (1024.0 * 1024.0),
write_wall_s: wall,
cpu_user_s: cpu1.0 - cpu0.0,
cpu_sys_s: cpu1.1 - cpu0.1,
},
latencies: lats,
families,
})
}
fn write_gc_reachable(
opts: &CampaignOptions,
corpus: &Corpus,
options: OptimizeOptions,
) -> Result<u64, String> {
let tmp = scratch_tempdir(&opts.scratch_dir, "h2gc-")?;
let store = fresh_store(tmp.path())?;
write_only(&store, 3, corpus, options)?;
crate::store::gc::collect(&store, &crate::store::transaction::CrashHooks::none())
.map_err(|e| e.to_string())?;
let n = store_numbers(&store)?;
Ok(n.reachable)
}
fn write_gc_footprint(
opts: &CampaignOptions,
corpus: &Corpus,
options: OptimizeOptions,
) -> Result<PostGcFootprint, String> {
let tmp = scratch_tempdir(&opts.scratch_dir, "fpgc-")?;
let store = fresh_store(tmp.path())?;
write_only(&store, 3, corpus, options)?;
crate::store::gc::collect(&store, &crate::store::transaction::CrashHooks::none())
.map_err(|e| e.to_string())?;
let n = store_numbers(&store)?;
let backing = dir_bytes(store.dir());
let allocated = allocated_blocks(store.dir());
let logical = corpus.logical_bytes();
Ok(PostGcFootprint {
corpus: corpus.name.clone(),
logical,
reachable: n.reachable,
total_backing: backing,
allocated_blocks: allocated,
ratio_reachable: logical as f64 / n.reachable.max(1) as f64,
ratio_total_backing: logical as f64 / backing.max(1) as f64,
ratio_allocated: logical as f64 / allocated.max(1) as f64,
})
}
fn allocated_blocks(dir: &Path) -> u64 {
use std::os::unix::fs::MetadataExt;
let mut total = 0u64;
if let Ok(md) = std::fs::metadata(dir.join("superblock")) {
total = total.saturating_add(md.blocks().saturating_mul(512));
}
if let Ok(segments) = crate::store::segment::list_segments(dir) {
for seq in segments {
if let Ok(md) = std::fs::metadata(crate::store::segment::segment_path(dir, seq)) {
total = total.saturating_add(md.blocks().saturating_mul(512));
}
}
}
total
}
struct StoreNumbers {
logical: u64,
reachable: u64,
total_backing: u64,
unreachable: u64,
allocator_overhead: u64,
record_count: u64,
}
fn store_numbers(store: &Store) -> Result<StoreNumbers, String> {
let total_backing = dir_bytes(store.dir());
let unreachable = crate::store::gc::unreachable_bytes(store).map_err(|e| e.to_string())?;
let records_total: u64 = store
.object_index()
.iter()
.into_iter()
.map(|(_, loc)| loc.total_size())
.sum();
let allocator_overhead = total_backing.saturating_sub(records_total);
let reachable = records_total.saturating_sub(unreachable);
let logical = store.logical_bytes().map_err(|e| e.to_string())?;
Ok(StoreNumbers {
logical,
reachable,
total_backing,
unreachable,
allocator_overhead,
record_count: store.object_index().len() as u64,
})
}
fn extent_decomposition(
store: &Store,
ino: u64,
) -> Result<(Accounting, BTreeMap<String, u64>), String> {
let limits = *store.limits();
let inode = store
.get_inode(ino)
.map_err(|e| e.to_string())?
.ok_or("inode missing")?;
let root = match inode.data {
InodeData::File { extent_root } => extent_root,
_ => return Err("not a file".into()),
};
let mut acct = Accounting::default();
let mut families: BTreeMap<String, u64> = BTreeMap::new();
let mut descriptor_bytes = 0u64;
let mut payload_objs: std::collections::HashSet<crate::core::extent::ChunkId> =
std::collections::HashSet::new();
let mut model_objs: std::collections::HashSet<crate::core::extent::ChunkId> =
std::collections::HashSet::new();
let mut residual_bytes = 0u64;
let mut object_refs: std::collections::HashMap<crate::core::extent::ChunkId, u64> =
std::collections::HashMap::new();
let mut exact_ref_lens: Vec<u64> = Vec::new();
for (_, bytes) in
crate::store::extent_tree::scan_all(root, BTREE_ORDER, limits.max_fanout, store)
.map_err(|e| e.to_string())?
{
let d = crate::format::descriptor::decode(
&bytes,
limits.max_descriptor_bytes,
limits.max_inline_bytes,
limits.max_palette,
limits.max_period,
limits.max_chunk_size,
)
.map_err(|e| e.to_string())?;
descriptor_bytes += d.encoded_size();
*families.entry(d.family().to_string()).or_insert(0) += 1;
let mut refs: Vec<crate::core::extent::ChunkId> = Vec::new();
let residual_refs = |r: &Residual| -> Vec<crate::core::extent::ChunkId> {
match r {
Residual::RansCoded { enc_obj, model, .. }
| Residual::BaseSequence { enc_obj, model, .. } => vec![*enc_obj, *model],
_ => Vec::new(),
}
};
match &d {
Representation::Raw { obj, .. } => {
payload_objs.insert(*obj);
refs.push(*obj);
}
Representation::Rans { model, enc_obj, .. }
| Representation::SequenceRans { model, enc_obj, .. }
| Representation::SparseBlock64 { model, enc_obj, .. }
| Representation::SequenceDeep { model, enc_obj, .. } => {
model_objs.insert(*model);
payload_objs.insert(*enc_obj);
refs.push(*model);
refs.push(*enc_obj);
}
Representation::SequenceDict {
dictionary,
model,
enc_obj,
..
} => {
model_objs.insert(*model);
payload_objs.insert(*enc_obj);
refs.push(*model);
refs.push(*enc_obj);
payload_objs.insert(*dictionary);
refs.push(*dictionary);
}
Representation::SequenceSharedDict {
dictionary,
shared,
model,
enc_obj,
..
} => {
model_objs.insert(*model);
payload_objs.insert(*enc_obj);
refs.push(*model);
refs.push(*enc_obj);
if !dictionary.is_zero() {
payload_objs.insert(*dictionary);
refs.push(*dictionary);
}
payload_objs.insert(*shared);
refs.push(*shared);
}
Representation::ExactRef { target, len, .. } => {
payload_objs.insert(*target);
refs.push(*target);
exact_ref_lens.push(*len);
}
Representation::BaseResidual { base, residual, .. } => {
payload_objs.insert(*base);
refs.push(*base);
residual_bytes += residual.encoded_size();
refs.extend(residual_refs(residual));
}
Representation::EntropyRef { residual, .. } => {
refs.extend(residual_refs(residual));
}
_ => {}
}
for id in refs {
*object_refs.entry(id).or_insert(0) += 1;
}
}
let payload_bytes: u64 = payload_objs.iter().map(|id| object_size(store, id)).sum();
let model_bytes: u64 = model_objs.iter().map(|id| object_size(store, id)).sum();
let mut cas_shared = 0u64;
for (id, count) in &object_refs {
if *count >= 2 {
cas_shared = cas_shared.saturating_add((count - 1) * object_size(store, id));
}
}
let mut exact_ref_saved = 0u64;
for len in &exact_ref_lens {
exact_ref_saved = exact_ref_saved.saturating_add(*len);
}
exact_ref_saved = exact_ref_saved.saturating_sub(exact_ref_lens
.len()
.saturating_mul(41) as u64);
acct.descriptor_bytes = descriptor_bytes;
acct.payload_bytes = payload_bytes;
acct.model_bytes = model_bytes;
acct.residual_bytes = residual_bytes;
acct.cas_shared_bytes_saved = cas_shared;
acct.exact_ref_bytes_saved = exact_ref_saved;
Ok((acct, families))
}
fn object_size(store: &Store, id: &crate::core::extent::ChunkId) -> u64 {
store
.object_index()
.get(id)
.map(|loc| loc.total_size())
.unwrap_or(0)
}
fn dir_bytes(path: &Path) -> u64 {
let mut total = 0u64;
let mut stack = vec![path.to_path_buf()];
while let Some(dir) = stack.pop() {
if let Ok(rd) = std::fs::read_dir(&dir) {
for e in rd.flatten() {
let p = e.path();
if p.is_dir() {
stack.push(p);
} else if let Ok(md) = e.metadata() {
total += md.len();
}
}
}
}
total
}
fn run_gc_traffic(opts: &CampaignOptions) -> Result<GcTraffic, String> {
let config = StoreConfig {
segment_size: 8 * 1024 * 1024,
..StoreConfig::default()
};
let tmp = scratch_tempdir(&opts.scratch_dir, "gc-")?;
let store = Store::create(tmp.path(), &config, [0x66; 16]).map_err(|e| e.to_string())?;
let inode = Inode::new_file(1000, 1000, 0o644);
{
let mut tx = store.begin_tx().map_err(|e| e.to_string())?;
Store::put_inode_in_tx(&mut tx, 3, &inode).map_err(|e| e.to_string())?;
tx.commit(&CrashHooks::none()).map_err(|e| e.to_string())?;
}
let vseq = corpus::versioned(4, 8);
write_only(&store, 3, &vseq, OptimizeOptions::default())?;
let ur = corpus::urandom(opts.size_mib / 2, 0x1234_5678);
write_only(&store, 3, &ur, OptimizeOptions::default())?;
let before = crate::store::gc::unreachable_bytes(&store).map_err(|e| e.to_string())?;
let physical_before = store.physical_used();
let t0 = Instant::now();
let reclaimed =
crate::store::gc::collect(&store, &CrashHooks::none()).map_err(|e| e.to_string())?;
let gc_wall = t0.elapsed().as_secs_f64();
let after = crate::store::gc::unreachable_bytes(&store).map_err(|e| e.to_string())?;
let physical_after = store.physical_used();
let by_tag_after =
crate::store::gc::unreachable_bytes_by_record_tag(&store).map_err(|e| e.to_string())?;
let opt =
crate::optimizer::background::optimize_pass(&store, OptimizeOptions::default(), None, None)
.map_err(|e| e.to_string())?;
Ok(GcTraffic {
unreachable_before: before,
reclaimed_bytes: reclaimed,
unreachable_after: after,
physical_before,
physical_after,
gc_wall_s: gc_wall,
optimizer_scanned: opt.scanned,
optimizer_rewritten: opt.rewritten,
optimizer_saved_bytes: opt.saved_bytes,
unreachable_by_tag_after: by_tag_after,
})
}
fn run_baselines(
opts: &CampaignOptions,
src_pack: &[u8],
corpora: &[Corpus],
) -> Result<Baselines, String> {
let mut b = Baselines::default();
let raw_path = opts.scratch_dir.join("baseline-raw.bin");
let t0 = Instant::now();
std::fs::write(&raw_path, src_pack).map_err(|e| e.to_string())?;
let wall = t0.elapsed().as_secs_f64();
let (_, fstype) = crate::evidence::environment::mount_of(&opts.scratch_dir);
b.raw_file = Some(RawBaseline {
path: raw_path.display().to_string(),
fstype,
bytes: src_pack.len() as u64,
write_mbps: src_pack.len() as f64 / wall / (1024.0 * 1024.0),
ratio: 1.0,
});
b.zstd_level_1 = zstd_baseline(src_pack, 1);
b.zstd_level_19 = zstd_baseline(src_pack, 19);
b.zstd_per_64k_level_1 = zstd_per_64k_baseline(src_pack, 1);
b.zstd_per_64k_level_19 = zstd_per_64k_baseline(src_pack, 19);
if let Some(src) = corpora.iter().find(|c| c.name == "src") {
let tmp = scratch_tempdir(&opts.scratch_dir, "rans-")?;
let store = fresh_store(tmp.path())?;
let outcome = full_run(&store, src, OptimizeOptions::raw_rans())?;
b.direct_rans_src = Some(outcome.metrics);
let tmp2 = scratch_tempdir(&opts.scratch_dir, "seq-")?;
let store2 = fresh_store(tmp2.path())?;
let outcome2 = full_run(&store2, src, OptimizeOptions::raw_sequence())?;
b.sequence_rans_src = Some(outcome2.metrics);
let tmp3 = scratch_tempdir(&opts.scratch_dir, "deep-")?;
let store3 = fresh_store(tmp3.path())?;
let outcome3 = full_run_deep(&store3, src, OptimizeOptions::raw_sequence_deep())?;
b.sequence_deep_src = Some(outcome3.metrics);
}
b.waived.push("btrfs with compression: waived — writable compressed-FS baseline requires root for loop-mounting a test image".into());
b.waived.push("EROFS/SquashFS: waived — read-only compressed-image baseline deferred (requires root and/or mkfs.erofs)".into());
Ok(b)
}
fn zstd_baseline(input: &[u8], level: i32) -> Option<CompressionBaseline> {
let tmp = tempfile::NamedTempFile::new().ok()?;
std::fs::write(tmp.path(), input).ok()?;
let t0 = Instant::now();
let out = std::process::Command::new("zstd")
.args(["-q", &format!("-{level}"), "-c"])
.arg(tmp.path())
.output()
.ok()?;
let wall = t0.elapsed().as_secs_f64();
if !out.status.success() {
return None;
}
let version = std::process::Command::new("zstd")
.arg("--version")
.output()
.ok()
.map(|o| String::from_utf8_lossy(&o.stdout).trim().to_string())
.unwrap_or_default();
Some(CompressionBaseline {
tool: "zstd".into(),
version,
level: level.to_string(),
input_bytes: input.len() as u64,
output_bytes: out.stdout.len() as u64,
ratio: input.len() as f64 / out.stdout.len().max(1) as f64,
wall_s: wall,
})
}
fn zstd_per_64k_baseline(input: &[u8], level: i32) -> Option<CompressionBaseline> {
let chunk = 64 * 1024;
let mut total_out = 0u64;
let t0 = Instant::now();
let mut off = 0usize;
while off < input.len() {
let end = (off + chunk).min(input.len());
let tmp = tempfile::NamedTempFile::new().ok()?;
std::fs::write(tmp.path(), &input[off..end]).ok()?;
let out = std::process::Command::new("zstd")
.args(["-q", &format!("-{level}"), "-c"])
.arg(tmp.path())
.output()
.ok()?;
if !out.status.success() {
return None;
}
total_out = total_out.saturating_add(out.stdout.len() as u64);
off = end;
}
let wall = t0.elapsed().as_secs_f64();
Some(CompressionBaseline {
tool: "zstd".into(),
version: "per-64KiB".into(),
level: level.to_string(),
input_bytes: input.len() as u64,
output_bytes: total_out,
ratio: input.len() as f64 / total_out.max(1) as f64,
wall_s: wall,
})
}
fn zstd_per_file_baseline(files: &[(String, Vec<u8>)], level: i32) -> Option<CompressionBaseline> {
let logical: u64 = files.iter().map(|(_, b)| b.len() as u64).sum();
let mut total_out = 0u64;
let t0 = Instant::now();
for (_, bytes) in files {
let tmp = tempfile::NamedTempFile::new().ok()?;
std::fs::write(tmp.path(), bytes).ok()?;
let out = std::process::Command::new("zstd")
.args(["-q", &format!("-{level}"), "-c"])
.arg(tmp.path())
.output()
.ok()?;
if !out.status.success() {
return None;
}
total_out = total_out.saturating_add(out.stdout.len() as u64);
}
let wall = t0.elapsed().as_secs_f64();
Some(CompressionBaseline {
tool: "zstd".into(),
version: "per-file".into(),
level: level.to_string(),
input_bytes: logical,
output_bytes: total_out,
ratio: logical as f64 / total_out.max(1) as f64,
wall_s: wall,
})
}
fn zstd_dir_anchor_baseline(
files: &[(String, Vec<u8>)],
level: i32,
) -> Option<CompressionBaseline> {
use std::collections::HashMap;
let logical: u64 = files.iter().map(|(_, b)| b.len() as u64).sum();
let mut dirs: HashMap<String, Vec<(String, Vec<u8>)>> = HashMap::new();
for (name, bytes) in files {
let d = name.rsplit('/').nth(1).unwrap_or(".").to_string();
dirs.entry(d)
.or_default()
.push((name.clone(), bytes.clone()));
}
let tmp = tempfile::Builder::new().prefix("zdd-").tempdir().ok()?;
let mut total_out = 0u64;
let t0 = Instant::now();
for (dir_no, members) in dirs.values().enumerate() {
let anchor = members
.iter()
.max_by_key(|(_, b)| b.len())
.map(|(_, b)| b.clone());
let anchor_path = tmp.path().join(format!("anchor-{dir_no}.dict"));
if let Some(a) = &anchor {
std::fs::write(&anchor_path, a).ok()?;
}
for (name, bytes) in members {
let in_path = tmp
.path()
.join(format!("in-{}.bin", name.replace('/', "_")));
std::fs::write(&in_path, bytes).ok()?;
let plain = std::process::Command::new("zstd")
.args(["-q", &format!("-{level}"), "-c"])
.arg(&in_path)
.output()
.ok()?
.stdout
.len() as u64;
let is_anchor = anchor.as_ref().map(|a| a == bytes).unwrap_or(false);
let out = if is_anchor || anchor.is_none() {
plain } else {
let with = std::process::Command::new("zstd")
.args(["-q", &format!("-{level}"), "-c", "-D"])
.arg(&anchor_path)
.arg(&in_path)
.output()
.ok()?
.stdout
.len() as u64;
if with == 0 && !bytes.is_empty() {
plain } else {
with
}
};
total_out = total_out.saturating_add(out);
let _ = std::fs::remove_file(&in_path);
}
}
let wall = t0.elapsed().as_secs_f64();
Some(CompressionBaseline {
tool: "zstd".into(),
version: "per-file+dir-anchor".into(),
level: level.to_string(),
input_bytes: logical,
output_bytes: total_out,
ratio: logical as f64 / total_out.max(1) as f64,
wall_s: wall,
})
}
fn per_extent_overhead(store: &Store) -> Result<(u64, u64), String> {
use crate::core::representation::Representation as Rep;
let limits = *store.limits();
let mut descriptor_bytes = 0u64;
let mut model_ids: std::collections::HashSet<crate::core::extent::ChunkId> =
std::collections::HashSet::new();
for ino in store.all_inodes().map_err(|e| e.to_string())? {
let Some(inode) = store.get_inode(ino).map_err(|e| e.to_string())? else {
continue;
};
let root = match inode.data {
InodeData::File { extent_root } => extent_root,
_ => continue,
};
if root.is_zero() {
continue;
}
for (_, bytes) in
crate::store::extent_tree::scan_all(root, BTREE_ORDER, limits.max_fanout, store)
.map_err(|e| e.to_string())?
{
let Ok(d) = crate::format::descriptor::decode(
&bytes,
limits.max_descriptor_bytes,
limits.max_inline_bytes,
limits.max_palette,
limits.max_period,
limits.max_chunk_size,
) else {
continue;
};
descriptor_bytes = descriptor_bytes.saturating_add(d.encoded_size());
match &d {
Rep::Rans { model, .. }
| Rep::SequenceRans { model, .. }
| Rep::SequenceDeep { model, .. }
| Rep::SparseBlock64 { model, .. }
| Rep::SequenceDict { model, .. }
| Rep::SequenceSharedDict { model, .. } => {
model_ids.insert(*model);
}
_ => {}
}
}
}
let mut model_bytes = 0u64;
for id in &model_ids {
if let Some(loc) = store.object_index().get(id) {
model_bytes = model_bytes.saturating_add(loc.stored_len);
}
}
Ok((descriptor_bytes, model_bytes))
}
fn write_tree(
store: &Store,
files: &[(String, Vec<u8>)],
options: OptimizeOptions,
) -> Result<(), String> {
let mut dir_cache: std::collections::HashMap<String, u64> = std::collections::HashMap::new();
dir_cache.insert(String::new(), store.current_root().root_dir_ino);
for (rel, bytes) in files {
let (dir_part, name) = match rel.rsplit_once('/') {
Some((d, n)) => (d.to_string(), n.to_string()),
None => (String::new(), rel.clone()),
};
if !dir_cache.contains_key(&dir_part) {
let mut cur = String::new();
let mut cur_ino = store.current_root().root_dir_ino;
for comp in dir_part.split('/') {
if comp.is_empty() {
continue;
}
let next_path = if cur.is_empty() {
comp.to_string()
} else {
format!("{cur}/{comp}")
};
let ino = match dir_cache.get(&next_path) {
Some(&cached) => cached,
None => {
let existing = store
.dir_lookup(cur_ino, comp.as_bytes())
.map_err(|e| e.to_string())?;
let ino = match existing {
Some(entry) => entry.ino,
None => store
.create_entry(
cur_ino,
comp.as_bytes(),
crate::store::NewEntry::dir(0o755, 1000, 1000),
&CrashHooks::none(),
)
.map_err(|e| e.to_string())?,
};
dir_cache.insert(next_path.clone(), ino);
ino
}
};
cur = next_path;
cur_ino = ino;
}
dir_cache.insert(dir_part.clone(), cur_ino);
}
let dir_ino = *dir_cache.get(&dir_part).expect("dir cached");
let ino = store
.create_entry(
dir_ino,
name.as_bytes(),
crate::store::NewEntry::file(0o644, 1000, 1000),
&CrashHooks::none(),
)
.map_err(|e| e.to_string())?;
let mut writes: Vec<(u64, Vec<u8>)> = Vec::new();
let mut off = 0u64;
while off < bytes.len() as u64 {
let len = 65536u64.min(bytes.len() as u64 - off);
writes.push((off, bytes[off as usize..(off + len) as usize].to_vec()));
off += len;
}
store
.write_region_batch(ino, &writes, options)
.map_err(|e| e.to_string())?;
}
Ok(())
}
fn tree_families(store: &Store) -> Result<BTreeMap<String, u64>, String> {
let limits = *store.limits();
let mut families: BTreeMap<String, u64> = BTreeMap::new();
for ino in store.all_inodes().map_err(|e| e.to_string())? {
let Some(inode) = store.get_inode(ino).map_err(|e| e.to_string())? else {
continue;
};
let root = match inode.data {
InodeData::File { extent_root } => extent_root,
_ => continue,
};
if root.is_zero() {
continue;
}
for (_, bytes) in
crate::store::extent_tree::scan_all(root, BTREE_ORDER, limits.max_fanout, store)
.map_err(|e| e.to_string())?
{
if let Ok(d) = crate::format::descriptor::decode(
&bytes,
limits.max_descriptor_bytes,
limits.max_inline_bytes,
limits.max_palette,
limits.max_period,
limits.max_chunk_size,
) {
*families.entry(d.family().to_string()).or_insert(0) += 1;
}
}
}
Ok(families)
}
fn run_tree_court(opts: &CampaignOptions) -> Result<TreeCourt, String> {
let files = corpus::source_tree_files(&opts.repo_root)?;
let pack = corpus::source_tree_pack(&opts.repo_root)?;
let logical: u64 = files.iter().map(|(_, b)| b.len() as u64).sum();
let single_chunk = files.iter().filter(|(_, b)| b.len() <= 65536).count();
let zw1 = zstd_baseline(&pack, 1);
let zw19 = zstd_baseline(&pack, 19);
let zf1 = zstd_per_file_baseline(&files, 1);
let zf19 = zstd_per_file_baseline(&files, 19);
let zc1 = zstd_per_64k_baseline(&pack, 1);
let zc19 = zstd_per_64k_baseline(&pack, 19);
let za1 = zstd_dir_anchor_baseline(&files, 1);
let tmp = scratch_tempdir(&opts.scratch_dir, "tree-")?;
let store = fresh_store(tmp.path())?;
write_tree(&store, &files, OptimizeOptions::default())?;
crate::store::gc::collect(&store, &CrashHooks::none()).map_err(|e| e.to_string())?;
let n1 = store_numbers(&store)?;
let fam1 = tree_families(&store)?;
let shared =
crate::optimizer::background::shared_dict_pass(&store, OptimizeOptions::default(), None)
.map_err(|e| e.to_string())?;
crate::store::gc::collect(&store, &CrashHooks::none()).map_err(|e| e.to_string())?;
let n2 = store_numbers(&store)?;
let fam2 = tree_families(&store)?;
let models =
crate::optimizer::background::model_bundle_pass(&store, OptimizeOptions::default(), None)
.map_err(|e| e.to_string())?;
crate::store::gc::collect(&store, &CrashHooks::none()).map_err(|e| e.to_string())?;
let n3 = store_numbers(&store)?;
let fam3 = tree_families(&store)?;
let physical = crate::store::physical::physical_report(&store).map_err(|e| e.to_string())?;
let compact_reclaimed =
crate::store::gc::compact_full(&store, &CrashHooks::none()).map_err(|e| e.to_string())?;
let n4 = store_numbers(&store)?;
let physical_after =
crate::store::physical::physical_report(&store).map_err(|e| e.to_string())?;
let (desc_bytes, model_bytes) = per_extent_overhead(&store)?;
Ok(TreeCourt {
file_count: files.len(),
single_chunk_files: single_chunk,
logical_bytes: logical,
zstd_whole_l1: zw1,
zstd_whole_l19: zw19,
zstd_per_file_l1: zf1,
zstd_per_file_l19: zf19,
zstd_per_64k_l1: zc1,
zstd_per_64k_l19: zc19,
efs_tree_reachable: n1.reachable,
efs_tree_backing: n1.total_backing,
efs_tree_families: fam1,
efs_shared_reachable: n2.reachable,
efs_shared_backing: n2.total_backing,
efs_shared_families: fam2,
shared_rewrites: shared.rewritten,
shared_saved_bytes: shared.saved_bytes,
efs_model_reachable: n3.reachable,
efs_model_backing: n3.total_backing,
efs_model_families: fam3,
model_rewrites: models.rewritten,
model_saved_bytes: models.saved_bytes,
physical_dead_indexed_bytes: physical.dead_indexed_bytes,
physical_index_hidden_bytes: physical.index_hidden_bytes,
physical_unindexed_bytes: physical.unindexed_bytes,
physical_backing_bytes: physical.file_bytes,
compact_backing_bytes: physical_after.file_bytes,
compact_reachable_bytes: n4.reachable,
compact_overhead_bytes: physical_after.file_bytes.saturating_sub(n4.reachable),
compact_reclaimed_bytes: compact_reclaimed,
zstd_dir_anchor_l1: za1,
per_extent_descriptor_bytes: desc_bytes,
per_extent_model_bytes: model_bytes,
})
}
fn build_corpora(opts: &CampaignOptions) -> Result<Vec<Corpus>, String> {
let pack = corpus::source_tree_pack(&opts.repo_root)?;
let mut v = vec![Corpus::single(
pack.clone(),
"src",
&format!(
"EntropyFS source tree pack (revision {})",
opts.repo_root.display()
),
"docs + src + evidence + manifests, length-prefixed; deterministic per revision",
)];
v.push(structured_corpus(opts.size_mib));
v.push(corpus::versioned(4, 8));
v.push(corpus::shuffled_versioned(4, 8));
v.push(corpus::urandom(opts.size_mib / 2, 0xdead_beef_cafe_f00d));
match corpus::compressed_zstd(&pack, 19) {
Ok(c) => v.push(c),
Err(e) => eprintln!("entropyfs: campaign: compressed control skipped: {e}"),
}
Ok(v)
}
fn structured_corpus(size_mib: u64) -> Corpus {
corpus::structured(size_mib)
}
fn options_for(mode: &str) -> Result<OptimizeOptions, String> {
OptimizeOptions::ablation_modes()
.into_iter()
.find(|(name, _)| *name == mode)
.map(|(_, o)| o)
.ok_or_else(|| format!("unknown mode {mode}"))
}
fn scratch_tempdir(scratch: &Path, prefix: &str) -> Result<tempfile::TempDir, String> {
tempfile::Builder::new()
.prefix(prefix)
.tempdir_in(scratch)
.map_err(|e| e.to_string())
}
fn fresh_store(dir: &Path) -> Result<Store, String> {
let config = StoreConfig::default();
Store::create(dir, &config, [0x66; 16]).map_err(|e| e.to_string())?;
let store = Store::open(dir, &config).map_err(|e| e.to_string())?;
let inode = Inode::new_file(1000, 1000, 0o644);
{
let mut tx = store.begin_tx().map_err(|e| e.to_string())?;
Store::put_inode_in_tx(&mut tx, 3, &inode).map_err(|e| e.to_string())?;
tx.commit(&CrashHooks::none()).map_err(|e| e.to_string())?;
}
Ok(store)
}
fn write_json<T: Serialize>(dir: &Path, name: &str, value: &T) -> Result<(), String> {
let json = serde_json::to_string_pretty(value).map_err(|e| e.to_string())?;
crate::store::write_atomic(&dir.join(name), json.as_bytes()).map_err(|e| e.to_string())
}
fn line(log: &mut String, s: &str) {
println!("{s}");
log.push_str(s);
log.push('\n');
}
fn scale_summary(s: &StatSummary, scale: f64) -> StatSummary {
StatSummary {
count: s.count,
mean: s.mean * scale,
min: s.min * scale,
p50: s.p50 * scale,
p95: s.p95 * scale,
p99: s.p99 * scale,
max: s.max * scale,
}
}
fn median(v: &[u64]) -> u64 {
if v.is_empty() {
return 0;
}
let mut s = v.to_vec();
s.sort_unstable();
s[s.len() / 2]
}
fn cpu_ticks() -> (f64, f64) {
let body = std::fs::read_to_string("/proc/self/stat").unwrap_or_default();
let rest = body.rsplit_once(')').map(|(_, r)| r).unwrap_or_default();
let tok: Vec<&str> = rest.split_whitespace().collect();
if tok.len() < 15 {
return (0.0, 0.0);
}
let utime: f64 = tok[11].parse().unwrap_or(0.0);
let stime: f64 = tok[12].parse().unwrap_or(0.0);
const HZ: f64 = 100.0;
(utime / HZ, stime / HZ)
}
fn result_hashes(results: &CampaignResults) -> serde_json::Value {
let mut m = serde_json::Map::new();
for group in &results.runs {
let mut runs = serde_json::Map::new();
for r in &group.runs {
runs.insert(
r.run.to_string(),
serde_json::json!({ "result_hash": r.result_hash, "matches_input": r.hash_matches_input }),
);
}
m.insert(
format!("{}[{}]", group.corpus, group.mode),
serde_json::Value::Object(runs),
);
}
serde_json::Value::Object(m)
}
fn csv(results: &CampaignResults) -> String {
let mut out = String::from(
"corpus,mode,run,logical_bytes,written_bytes,reachable_bytes,total_backing_bytes,unreachable_bytes,ratio_reachable,write_mbps,read_mbps,write_wall_s,read_wall_s,cpu_user_s,cpu_sys_s,result_hash\n",
);
for g in &results.runs {
for r in &g.runs {
out.push_str(&format!(
"{},{},{},{},{},{},{},{},{:.4},{:.3},{:.3},{:.4},{:.4},{:.4},{:.4},{}\n",
g.corpus,
g.mode,
r.run,
r.logical_bytes,
r.written_bytes,
r.reachable_bytes,
r.total_backing_bytes,
r.unreachable_bytes,
r.ratio_reachable,
r.write_mbps,
r.read_mbps,
r.write_wall_s,
r.read_wall_s,
r.cpu_user_s,
r.cpu_sys_s,
r.result_hash
));
}
}
out
}
fn admission_checklist(results: &CampaignResults, corpora: &[Corpus]) -> Vec<AdmissionItem> {
let mut items = Vec::new();
items.push(AdmissionItem {
rule: "benchmark context complete (revision, Cargo.lock, kernel, CPU, governor, device, cache state, command)".to_string(),
met: true,
note: "environment.json (revision, cargo_lock_hash, kernel_*, cpu_*, governor, store_device, cache_state, command) + corpus-manifest.json archived in this directory".into(),
});
let accounting_ok = results.runs.iter().all(|g| {
g.runs.iter().all(|r| {
r.accounting.check == "ok"
&& r.accounting.payload_bytes
+ r.accounting.model_bytes
+ r.accounting.descriptor_bytes
+ r.accounting.metadata_bytes
> 0
})
});
items.push(AdmissionItem {
rule: "every required byte is counted (payload/models/residuals/descriptors/metadata/integrity/allocator/unreclaimed)".to_string(),
met: accounting_ok,
note: if accounting_ok { "per-run Accounting tables pass the reachable-bytes cross-check".into() } else { "check per-run accounting.check fields".into() },
});
let has_baselines = results.baselines.raw_file.is_some()
&& (results.baselines.zstd_level_1.is_some() || results.baselines.zstd_level_19.is_some())
&& results.baselines.direct_rans_src.is_some();
items.push(AdmissionItem {
rule: "all listed baselines run or explicitly waived".to_string(),
met: has_baselines,
note: format!(
"raw file {}; zstd {}; direct rANS {}; waivers: {}",
if results.baselines.raw_file.is_some() {
"present"
} else {
"MISSING"
},
if results.baselines.zstd_level_1.is_some() {
"present"
} else {
"MISSING"
},
if results.baselines.direct_rans_src.is_some() {
"present"
} else {
"MISSING"
},
results.baselines.waived.len()
),
});
let ablation_ok =
results.ablation.rows.len() >= 8 && results.ablation.cumulative_rows.len() >= 9;
items.push(AdmissionItem {
rule: "ablations identify which mechanism caused the gain (cumulative ladder A0-A8 + leave-one-out)".to_string(),
met: ablation_ok,
note: format!(
"leave-one-out table {} rows; cumulative ladder {} rows (A0-A8)",
results.ablation.rows.len(),
results.ablation.cumulative_rows.len()
),
});
let has_urandom = results
.runs
.iter()
.any(|g| g.corpus == "urandom" && g.ratio_median < 1.5);
let has_compressed = results.runs.iter().any(|g| g.corpus == "compressed-z19");
let has_shuffled = !results.versioned_experiment.shuffled_full.is_empty();
items.push(AdmissionItem {
rule: "negative controls included (random → RAW, compressed → no gain, shuffled history → temporal gains disappear)".to_string(),
met: has_urandom && has_compressed && has_shuffled,
note: format!("urandom ratio {:.3}x (expected ≤1.5x); compressed present {}; shuffled present {}",
results.runs.iter().find(|g| g.corpus == "urandom").map(|g| g.ratio_median).unwrap_or(0.0),
has_compressed, has_shuffled),
});
let mut all_match = true;
for c in corpora {
for g in &results.runs {
if g.corpus != c.name {
continue;
}
for r in &g.runs {
if !r.hash_matches_input {
all_match = false;
}
}
}
}
items.push(AdmissionItem {
rule: "materialized output hashes match the input corpus hashes".to_string(),
met: all_match,
note: if all_match {
"result-hashes.json: all runs match corpus content hashes".into()
} else {
"result-hashes.json: at least one mismatch".into()
},
});
items.push(AdmissionItem {
rule: "raw result artifacts are archived".to_string(),
met: true,
note: "raw-output.txt, results.json, results.csv, environment.json, corpus-manifest.json, result-hashes.json, report.md".into(),
});
items
}
fn report(results: &CampaignResults, log: &str) -> String {
let mut s = String::new();
s.push_str("# EntropyFS evidence campaign\n\n");
s.push_str(&format!(
"- campaign dir: `{}`\n- created: unix {}\n",
results.campaign_dir, results.created_unix
));
s.push_str("\n## Admission checklist (methodology §8)\n\n");
for a in &results.admission {
s.push_str(&format!(
"- [{}] {} — {}\n",
if a.met { "x" } else { " " },
a.rule,
a.note
));
}
s.push_str("\n## Summary\n\n");
for g in &results.runs {
s.push_str(&format!(
"- `{}[{}]`: {} runs — write {:.1} MiB/s (p50 {:.0}µs, p95 {:.0}µs, p99 {:.0}µs), read {:.1} MiB/s, fsync p50 {:.0}µs, physical median {} bytes, ratio {:.3}x\n",
g.corpus,
g.mode,
g.run_count,
g.write_throughput.p50,
g.write_latency_us.p50,
g.write_latency_us.p95,
g.write_latency_us.p99,
g.read_throughput.p50,
g.fsync_latency_us.p50,
g.physical_median,
g.ratio_median
));
}
if let Some(d) = &results.device_writes {
s.push_str(&format!(
"\nDevice writes during campaign window ({}): {} bytes written, {} bytes read.\n",
d.device,
d.written_bytes(),
d.read_bytes()
));
}
s.push_str("\n## Raw output\n\n```text\n");
s.push_str(log);
s.push_str("```\n");
s
}