use std::fs::File;
use std::io::{BufWriter, Write};
use std::path::{Path, PathBuf};
use flate2::write::GzEncoder;
use flate2::Compression;
use serde::{Deserialize, Serialize};
use sha2::{Digest, Sha256};
use super::metrics::RequestSample;
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct SamplesFile {
pub path: PathBuf,
pub sha256: String,
pub bytes: u64,
pub rows: usize,
}
impl SamplesFile {
#[must_use]
pub fn exceeds_budget(&self, budget_bytes: u64) -> bool {
self.bytes > budget_bytes
}
}
pub fn write_samples_gz(path: &Path, samples: &[RequestSample]) -> std::io::Result<SamplesFile> {
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent)?;
}
{
let file = File::create(path)?;
let mut gz = GzEncoder::new(BufWriter::new(file), Compression::default());
for s in samples {
let line = serde_json::to_string(s)
.map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidData, e))?;
gz.write_all(line.as_bytes())?;
gz.write_all(b"\n")?;
}
gz.finish()?.flush()?;
}
let bytes = std::fs::read(path)?;
let mut hasher = Sha256::new();
hasher.update(&bytes);
Ok(SamplesFile {
path: path
.file_name()
.map_or_else(|| path.to_path_buf(), PathBuf::from),
sha256: format!("{:x}", hasher.finalize()),
bytes: bytes.len() as u64,
rows: samples.len(),
})
}
pub fn read_samples_gz(path: &Path) -> std::io::Result<Vec<RequestSample>> {
use std::io::BufRead;
let file = File::open(path)?;
let gz = flate2::read::GzDecoder::new(file);
let reader = std::io::BufReader::new(gz);
let mut out = Vec::new();
for line in reader.lines() {
let line = line?;
if line.trim().is_empty() {
continue;
}
out.push(
serde_json::from_str(&line)
.map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidData, e))?,
);
}
Ok(out)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::perf_gate::bootstrap::bootstrap_agg_tok_s_ci;
use crate::perf_gate::protocol::Outcome;
fn deck(n: usize) -> Vec<RequestSample> {
(0..n)
.map(|i| RequestSample {
index: i,
worker: i % 4,
start_s: i as f64 * 0.25,
end_s: i as f64 * 0.25 + 1.0 + f64::from((i % 3) as u32) * 0.1,
token_times_s: vec![i as f64 * 0.25 + 0.05, i as f64 * 0.25 + 0.9],
generated_tokens: 128,
prompt_tokens: 512,
outcome: Outcome::Completed,
in_flight_at_start: 4,
drained: false,
})
.collect()
}
fn tmpdir(name: &str) -> PathBuf {
let d = std::env::temp_dir().join(format!("perf024-{name}-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&d);
d
}
#[test]
fn samples_round_trip_through_gzip() {
let dir = tmpdir("roundtrip");
let path = dir.join("samples.jsonl.gz");
let want = deck(25);
let meta = write_samples_gz(&path, &want).expect("write");
assert_eq!(meta.rows, 25);
assert!(meta.bytes > 0);
assert_eq!(meta.sha256.len(), 64);
assert_eq!(meta.path, PathBuf::from("samples.jsonl.gz"));
let got = read_samples_gz(&path).expect("read");
assert_eq!(got, want, "retained samples must survive the round trip");
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn the_ci_is_reproducible_from_the_retained_file_alone() {
let dir = tmpdir("rederive");
let path = dir.join("samples.jsonl.gz");
let original = deck(40);
write_samples_gz(&path, &original).expect("write");
let from_disk = read_samples_gz(&path).expect("read");
let a = bootstrap_agg_tok_s_ci(&original, 0.95).expect("n >= 2");
let b = bootstrap_agg_tok_s_ci(&from_disk, 0.95).expect("n >= 2");
assert_eq!(a, b, "a receipt is only evidence if its CI re-derives");
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn the_file_is_actually_gzip() {
let dir = tmpdir("magic");
let path = dir.join("samples.jsonl.gz");
write_samples_gz(&path, &deck(3)).expect("write");
let bytes = std::fs::read(&path).expect("read raw");
assert_eq!(&bytes[..2], &[0x1f, 0x8b], "gzip magic bytes absent");
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn digest_matches_the_bytes_on_disk() {
let dir = tmpdir("digest");
let path = dir.join("samples.jsonl.gz");
let meta = write_samples_gz(&path, &deck(5)).expect("write");
let bytes = std::fs::read(&path).expect("read raw");
let mut h = Sha256::new();
h.update(&bytes);
assert_eq!(meta.sha256, format!("{:x}", h.finalize()));
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn one_json_object_per_line() {
let dir = tmpdir("lines");
let path = dir.join("samples.jsonl.gz");
write_samples_gz(&path, &deck(7)).expect("write");
let file = File::open(&path).expect("open");
let mut text = String::new();
std::io::Read::read_to_string(&mut flate2::read::GzDecoder::new(file), &mut text)
.expect("decode");
assert_eq!(text.lines().count(), 7);
for line in text.lines() {
let _: RequestSample = serde_json::from_str(line).expect("each line stands alone");
}
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn a_samples_record_with_an_unknown_key_is_refused() {
let honest = r#"{"path":"samples.c1.r1.jsonl.gz","sha256":"ab","bytes":10,"rows":3}"#;
let parsed: SamplesFile = serde_json::from_str(honest).expect("the known shape parses");
assert_eq!(parsed.rows, 3);
let extra = r#"{"path":"s.gz","sha256":"ab","bytes":10,"rows":3,"rows_dropped":7}"#;
let err = serde_json::from_str::<SamplesFile>(extra)
.expect_err("an unknown key must be refused, not dropped");
assert!(err.to_string().contains("rows_dropped"), "{err}");
}
#[test]
fn budget_check_takes_the_budget_as_an_argument() {
let dir = tmpdir("budget");
let path = dir.join("samples.jsonl.gz");
let meta = write_samples_gz(&path, &deck(10)).expect("write");
assert!(meta.exceeds_budget(0));
assert!(!meta.exceeds_budget(u64::MAX));
let _ = std::fs::remove_dir_all(&dir);
}
}