Files
Aura/crates/aura-registry/src/trace_store.rs
T
Brummel 5616aa6c6a fix(registry,cli): chart resolves the depth-2 sweep --trace family layout
TraceStore::name_kind and read_family resolved taps exactly one
directory below the family handle, so the per-member fan-out layout
sweep/walkforward --trace writes since #224
(<name>/<cell>/<member>/index.json) classified as NotFound and
`aura chart <printed handle>` exited 1. Both sides now resolve depth-1
(campaign nominee layout) and depth-2, with depth-2 members keyed
<cell>/<member> (deterministic sort, C1); the on-disk layout is
unchanged. emit_chart's not-found message no longer suggests re-running
with the unknown handle as --trace (the data-creating command cannot
take an output handle as input); the message prefix stays pinned.
Milestone-fieldtest finding B1 (fixtures: 09da04f); verified against
the fieldtest's own on-disk family and the campaign depth-1 sibling.
2026-07-11 15:05:31 +02:00

448 lines
18 KiB
Rust

//! The per-run trace store: a directory of columnar tap files under
//! `<runs_dir>/traces/<name>/`, beside the run registry's `runs.jsonl`. One
//! subdirectory per run name; one `<tap>.json` (a `ColumnarTrace`) per tap; plus
//! `index.json` carrying the run's `RunManifest` and the tap order. Mirrors the
//! registry's directory-co-located sibling-store discipline; a missing run reads as
//! `NotFound` (treat-as-absent), never a panic.
use std::fmt;
use std::fs;
use std::path::{Path, PathBuf};
use aura_engine::{ColumnarTrace, RunManifest};
use serde::{Deserialize, Serialize};
/// A per-run trace directory rooted at `<runs_dir>/traces`. One subdir per run name.
pub struct TraceStore {
dir: PathBuf,
}
/// A run's read-back traces: its manifest plus its taps in index (write) order.
#[derive(Debug)]
pub struct RunTraces {
pub manifest: RunManifest,
pub taps: Vec<ColumnarTrace>,
}
/// The on-disk shape of a trace name: a single run, a family of members, or absent.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum NameKind {
/// `traces/<name>/index.json` is present — a single recorded run.
Run,
/// No top-level `index.json`, but ≥1 immediate subdir has one (depth-1), or
/// ≥1 immediate subdir's own immediate subdir has one (depth-2, the #224
/// sweep/walk-forward fan-out) — a family.
Family,
/// Neither — no recorded run or family of this name.
NotFound,
}
/// Does `dir` hold ≥1 family member, at depth-1 (`dir/<key>/index.json`) or
/// depth-2 (`dir/<cell>/<member>/index.json`, the #224 fan-out)? Shared by
/// `name_kind`'s classification.
fn has_family_member(dir: &Path) -> bool {
let Ok(entries) = fs::read_dir(dir) else { return false };
for entry in entries.flatten() {
let path = entry.path();
if path.join("index.json").is_file() {
return true;
}
if let Ok(sub_entries) = fs::read_dir(&path) {
for sub_entry in sub_entries.flatten() {
if sub_entry.path().join("index.json").is_file() {
return true;
}
}
}
}
false
}
/// One member of a family, read back: its key (the member-dir name) + its traces.
#[derive(Debug)]
pub struct FamilyMember {
pub key: String,
pub traces: RunTraces,
}
/// Which kind of write a name is about to receive — the write-guard's intent.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum WriteKind {
Run,
Family,
}
/// `index.json`: the run manifest + the tap names in write order (= read order).
#[derive(Serialize, Deserialize)]
struct Index {
manifest: RunManifest,
taps: Vec<String>,
}
impl TraceStore {
/// Bind to `<runs_dir>/traces`. No I/O — directories are created on write.
pub fn open(runs_dir: impl AsRef<Path>) -> TraceStore {
TraceStore { dir: runs_dir.as_ref().join("traces") }
}
/// Write one run's taps + index under `traces/<name>/`, creating directories as
/// needed. Each tap lands in `<tap>.json`; the manifest + tap order in
/// `index.json`.
pub fn write(
&self,
name: &str,
manifest: &RunManifest,
taps: &[ColumnarTrace],
) -> Result<(), TraceStoreError> {
let run_dir = self.dir.join(name);
fs::create_dir_all(&run_dir)?;
for tap in taps {
let path = run_dir.join(format!("{}.json", tap.tap));
let body = serde_json::to_string(tap).map_err(|source| TraceStoreError::Parse {
file: path.display().to_string(),
source,
})?;
fs::write(&path, body)?;
}
let index = Index {
manifest: manifest.clone(),
taps: taps.iter().map(|t| t.tap.clone()).collect(),
};
let index_path = run_dir.join("index.json");
let body = serde_json::to_string(&index).map_err(|source| TraceStoreError::Parse {
file: index_path.display().to_string(),
source,
})?;
fs::write(&index_path, body)?;
Ok(())
}
/// Read a run's index + taps back, in index order. A missing run directory (no
/// `index.json`) is `NotFound`, not a panic; a malformed file is `Parse{file,..}`.
pub fn read(&self, name: &str) -> Result<RunTraces, TraceStoreError> {
let run_dir = self.dir.join(name);
let index_path = run_dir.join("index.json");
let index_text = match fs::read_to_string(&index_path) {
Ok(t) => t,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
return Err(TraceStoreError::NotFound(name.to_string()))
}
Err(e) => return Err(TraceStoreError::Io(e)),
};
let index: Index = serde_json::from_str(&index_text).map_err(|source| {
TraceStoreError::Parse { file: index_path.display().to_string(), source }
})?;
let mut taps = Vec::with_capacity(index.taps.len());
for tap_name in &index.taps {
let path = run_dir.join(format!("{tap_name}.json"));
let text = fs::read_to_string(&path)?;
let tap: ColumnarTrace = serde_json::from_str(&text).map_err(|source| {
TraceStoreError::Parse { file: path.display().to_string(), source }
})?;
taps.push(tap);
}
Ok(RunTraces { manifest: index.manifest, taps })
}
/// Classify a name by its on-disk shape (the read-side of total name
/// resolution). Top-level `index.json` -> `Run`; else ≥1 immediate subdir with
/// `index.json` (depth-1, the campaign nominee layout) OR ≥1 immediate subdir
/// whose OWN immediate subdir has `index.json` (depth-2, the #224
/// sweep/walk-forward per-cell member fan-out) -> `Family`; else `NotFound`.
pub fn name_kind(&self, name: &str) -> NameKind {
let run_dir = self.dir.join(name);
if run_dir.join("index.json").is_file() {
return NameKind::Run;
}
if has_family_member(&run_dir) {
return NameKind::Family;
}
NameKind::NotFound
}
/// Read every member of a family. Two on-disk shapes are resolved: the
/// depth-1 campaign layout (`<name>/<key>/index.json`, `key` = the immediate
/// subdir name) and the depth-2 #224 sweep/walk-forward fan-out
/// (`<name>/<cell>/<member>/index.json`, `key` = `<cell>/<member>` so members
/// stay distinguishable across cells). Read via `read("<name>/<key>")`,
/// collected and **sorted by key** (deterministic, C1). An absent/empty
/// family reads as an empty Vec (treat-as-absent); a malformed member
/// propagates its `Parse` error.
pub fn read_family(&self, name: &str) -> Result<Vec<FamilyMember>, TraceStoreError> {
let run_dir = self.dir.join(name);
let entries = match fs::read_dir(&run_dir) {
Ok(e) => e,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()),
Err(e) => return Err(TraceStoreError::Io(e)),
};
let mut members = Vec::new();
for entry in entries {
let entry = entry?;
let path = entry.path();
if path.join("index.json").is_file() {
// depth-1: the immediate subdir IS a member.
let key = match entry.file_name().into_string() {
Ok(k) => k,
Err(_) => continue, // non-UTF8 dir name: skip
};
let traces = self.read(&format!("{name}/{key}"))?;
members.push(FamilyMember { key, traces });
continue;
}
// depth-2: the immediate subdir is a cell; its OWN immediate subdirs
// with `index.json` are the members.
let cell_name = match entry.file_name().into_string() {
Ok(k) => k,
Err(_) => continue, // non-UTF8 dir name: skip
};
let cell_entries = match fs::read_dir(&path) {
Ok(e) => e,
Err(_) => continue, // stray file or unreadable: skip
};
for cell_entry in cell_entries.flatten() {
if !cell_entry.path().join("index.json").is_file() {
continue; // stray file or non-member subdir
}
let member_name = match cell_entry.file_name().into_string() {
Ok(k) => k,
Err(_) => continue, // non-UTF8 dir name: skip
};
let key = format!("{cell_name}/{member_name}");
let traces = self.read(&format!("{name}/{key}"))?;
members.push(FamilyMember { key, traces });
}
}
members.sort_by(|a, b| a.key.cmp(&b.key));
Ok(members)
}
/// The write-guard: `Ok` means a write of `intent` may proceed at this name,
/// **not** that the name is unused. It refuses only cross-kind reuse (a name
/// already a run when writing a family, or vice versa), so the one ambiguous
/// on-disk state — a name used by both a run and a family — is unreachable. A
/// same-kind overwrite and a fresh name both return `Ok` (the name may be very
/// much in use). Called once per command before any write.
pub fn ensure_name_free(
&self,
name: &str,
intent: WriteKind,
) -> Result<(), TraceStoreError> {
match (intent, self.name_kind(name)) {
(WriteKind::Run, NameKind::Family) => {
Err(TraceStoreError::NameTaken { name: name.to_string(), existing: NameKind::Family })
}
(WriteKind::Family, NameKind::Run) => {
Err(TraceStoreError::NameTaken { name: name.to_string(), existing: NameKind::Run })
}
_ => Ok(()),
}
}
}
/// What can go wrong reading or writing the trace store.
#[derive(Debug)]
pub enum TraceStoreError {
/// An I/O error reading or writing a trace file.
Io(std::io::Error),
/// A trace/index file did not parse (the offending file named).
Parse { file: String, source: serde_json::Error },
/// No recorded run of this name (no `index.json` under `traces/<name>/`).
NotFound(String),
/// A trace name already used by the other kind (run vs family) — refuse reuse.
NameTaken { name: String, existing: NameKind },
}
impl fmt::Display for TraceStoreError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
TraceStoreError::Io(e) => write!(f, "trace store i/o: {e}"),
TraceStoreError::Parse { file, source } => {
write!(f, "trace store parse error in {file}: {source}")
}
TraceStoreError::NotFound(name) => {
write!(f, "no recorded run '{name}' under runs/traces")
}
TraceStoreError::NameTaken { name, existing } => {
let kind = match existing {
NameKind::Run => "run",
NameKind::Family => "family",
NameKind::NotFound => "name",
};
write!(f, "'{name}' already used as a {kind}; pick another --trace name")
}
}
}
}
impl std::error::Error for TraceStoreError {}
impl From<std::io::Error> for TraceStoreError {
fn from(e: std::io::Error) -> Self {
TraceStoreError::Io(e)
}
}
#[cfg(test)]
mod tests {
use super::*;
use aura_core::{Scalar, ScalarKind, Timestamp};
fn temp_traces_root(name: &str) -> PathBuf {
let dir =
std::env::temp_dir().join(format!("aura-tracestore-{}-{}", std::process::id(), name));
let _ = fs::remove_dir_all(&dir);
fs::create_dir_all(&dir).expect("create temp traces root");
dir
}
fn sample_manifest() -> RunManifest {
RunManifest {
commit: "c".to_string(),
params: vec![("p".to_string(), Scalar::f64(1.0))],
window: (Timestamp(1), Timestamp(5)),
seed: 0,
broker: "sim-optimal(pip_size=0.0001)".to_string(),
selection: None,
instrument: None,
topology_hash: None,
project: None,
}
}
fn sample_taps() -> Vec<ColumnarTrace> {
vec![
ColumnarTrace::from_rows(
"equity",
&[ScalarKind::F64],
&[(Timestamp(1), vec![Scalar::f64(0.0)]), (Timestamp(2), vec![Scalar::f64(0.5)])],
),
ColumnarTrace::from_rows(
"exposure",
&[ScalarKind::F64],
&[(Timestamp(1), vec![Scalar::f64(1.0)]), (Timestamp(2), vec![Scalar::f64(-1.0)])],
),
]
}
#[test]
fn write_then_read_round_trips() {
let root = temp_traces_root("roundtrip");
let store = TraceStore::open(&root);
let manifest = sample_manifest();
let taps = sample_taps();
store.write("demo", &manifest, &taps).expect("write");
let back = store.read("demo").expect("read");
assert_eq!(back.manifest, manifest);
assert_eq!(back.taps, taps);
assert!(root.join("traces/demo/index.json").exists());
assert!(root.join("traces/demo/equity.json").exists());
assert!(root.join("traces/demo/exposure.json").exists());
let _ = fs::remove_dir_all(&root);
}
#[test]
fn read_missing_run_is_not_found() {
let root = temp_traces_root("missing");
let store = TraceStore::open(&root);
match store.read("nope") {
Err(TraceStoreError::NotFound(name)) => assert_eq!(name, "nope"),
other => panic!("expected NotFound, got {other:?}"),
}
let _ = fs::remove_dir_all(&root);
}
#[test]
fn corrupt_tap_file_is_a_parse_error_naming_the_file() {
let root = temp_traces_root("corrupt");
let store = TraceStore::open(&root);
store.write("demo", &sample_manifest(), &sample_taps()).expect("write");
fs::write(root.join("traces/demo/equity.json"), "not json").expect("clobber");
match store.read("demo") {
Err(TraceStoreError::Parse { file, .. }) => assert!(file.ends_with("equity.json")),
other => panic!("expected Parse, got {other:?}"),
}
let _ = fs::remove_dir_all(&root);
}
#[test]
fn name_kind_classifies_run_family_and_absent() {
let root = temp_traces_root("namekind");
let store = TraceStore::open(&root);
store.write("solo", &sample_manifest(), &sample_taps()).expect("write solo");
store.write("fam/m1", &sample_manifest(), &sample_taps()).expect("write m1");
assert_eq!(store.name_kind("solo"), NameKind::Run);
assert_eq!(store.name_kind("fam"), NameKind::Family);
assert_eq!(store.name_kind("ghost"), NameKind::NotFound);
let _ = fs::remove_dir_all(&root);
}
#[test]
fn read_family_returns_members_sorted_ignoring_strays() {
let root = temp_traces_root("readfamily");
let store = TraceStore::open(&root);
// members written out of order; read_family must sort by key.
store.write("fam/b", &sample_manifest(), &sample_taps()).expect("write b");
store.write("fam/a", &sample_manifest(), &sample_taps()).expect("write a");
// a stray file directly under the family dir is ignored (no index.json there).
fs::write(root.join("traces/fam/stray.txt"), "x").expect("stray");
let members = store.read_family("fam").expect("read_family");
let keys: Vec<&str> = members.iter().map(|m| m.key.as_str()).collect();
assert_eq!(keys, vec!["a", "b"]);
assert_eq!(members[0].traces.taps.len(), 2);
// an absent family reads as empty (treat-as-absent).
assert!(store.read_family("ghost").expect("ghost").is_empty());
let _ = fs::remove_dir_all(&root);
}
#[test]
fn name_kind_classifies_depth_two_sweep_fanout_as_family() {
let root = temp_traces_root("namekind-depth2");
let store = TraceStore::open(&root);
// the #224 sweep/walk-forward fan-out: <name>/<cell>/<member>/index.json,
// two levels below the family name — no top-level or depth-1 index.json.
store.write("fam2/cellA/m1", &sample_manifest(), &sample_taps()).expect("write m1");
assert_eq!(store.name_kind("fam2"), NameKind::Family);
let _ = fs::remove_dir_all(&root);
}
#[test]
fn read_family_resolves_depth_two_sweep_fanout_with_cell_member_keys() {
let root = temp_traces_root("readfamily-depth2");
let store = TraceStore::open(&root);
// members written out of order across two cells; read_family must sort by
// the combined "<cell>/<member>" key, keeping members distinguishable
// across cells.
store.write("fam2/cellB/m1", &sample_manifest(), &sample_taps()).expect("write cellB/m1");
store.write("fam2/cellA/m2", &sample_manifest(), &sample_taps()).expect("write cellA/m2");
store.write("fam2/cellA/m1", &sample_manifest(), &sample_taps()).expect("write cellA/m1");
// a stray file directly under a cell dir is ignored (no index.json there).
fs::write(root.join("traces/fam2/cellA/stray.txt"), "x").expect("stray");
let members = store.read_family("fam2").expect("read_family");
let keys: Vec<&str> = members.iter().map(|m| m.key.as_str()).collect();
assert_eq!(keys, vec!["cellA/m1", "cellA/m2", "cellB/m1"]);
assert_eq!(members[0].traces.taps.len(), 2);
let _ = fs::remove_dir_all(&root);
}
#[test]
fn ensure_name_free_refuses_cross_kind_reuse() {
let root = temp_traces_root("guard");
let store = TraceStore::open(&root);
store.write("asrun", &sample_manifest(), &sample_taps()).expect("write run");
store.write("asfam/m1", &sample_manifest(), &sample_taps()).expect("write fam");
assert!(matches!(
store.ensure_name_free("asrun", WriteKind::Family),
Err(TraceStoreError::NameTaken { .. })
));
assert!(matches!(
store.ensure_name_free("asfam", WriteKind::Run),
Err(TraceStoreError::NameTaken { .. })
));
// same-kind re-run and a fresh name are fine.
assert!(store.ensure_name_free("asrun", WriteKind::Run).is_ok());
assert!(store.ensure_name_free("asfam", WriteKind::Family).is_ok());
assert!(store.ensure_name_free("fresh", WriteKind::Run).is_ok());
let _ = fs::remove_dir_all(&root);
}
}