feat(aura-engine): run metrics + manifest report surface

Cycle 0009 (Walking skeleton milestone). A run's two C18 "from day one"
artefacts — a reproducible manifest and summary metrics — land as a new,
pure-additive `report` module in aura-engine:

- summarize(equity, exposure) -> RunMetrics: a post-run pure reduction over a
  run's recorded pip-equity + exposure streams. total_pips (last cumulative
  value), max_drawdown (worst running-peak-to-trough), exposure_sign_flips (a
  turnover proxy: adjacent normalized-sign changes, flat distinct from
  long/short via a three-way sign0). Total and pure: empty -> zeros, identical
  inputs -> identical metrics (C1/C12).
- RunManifest: the reproducible descriptor (commit, params as name->value
  pairs, data-window, seed, broker label). Caller-supplied — the engine cannot
  introspect a git commit or seed. params is the honest precursor to the
  deferred typed param-space (node.rs; 0008 audit).
- RunReport { manifest, metrics } + to_json(): canonical, hand-rolled JSON for
  the structured C14 face.
- f64_field(rows, field): bridges a recording sink's Vec<Scalar> rows to
  summarize; panics on a kind/width mismatch (a wiring bug), like the engine's
  other "checked at wiring" violations.

Why post-run, not a node: a node emits one record per eval with no terminal
eval (C8), so an end-of-run reduction has no home as a node; the World drains
its cycle-0006 recording sinks after Harness::run and folds them here. Matches
C18's "store manifests + metrics, re-derive results on demand".

Why hand-rolled JSON, not serde: the workspace is deliberately zero-dependency
and the schema is tiny/closed/flat; serde is reconsidered when the run-registry
milestone grows the schema. (User-approved at spec review.)

Scope: the surface + a determinism end-to-end test reusing the cycle-0007
two-sink harness; the `aura run` subcommand that prints the report is #8. No
engine/Harness/node-contract change; workspace stays zero-dependency.

Verified: cargo test --workspace green (aura-engine 40, incl. 10 new report
tests); clippy --all-targets -D warnings clean; RUSTDOCFLAGS=-D warnings cargo
doc clean.

closes #6
This commit is contained in:
2026-06-04 18:57:14 +02:00
parent cb66708d26
commit 80d7bfbbf0
2 changed files with 437 additions and 3 deletions
+15 -3
View File
@@ -13,14 +13,26 @@
//! - [`BootstrapError`] — wiring faults caught once, at bootstrap (kind
//! mismatch, bad index, directed cycle).
//!
//! Delivered in cycle 0009 — the run report surface:
//!
//! - [`RunMetrics`] / [`RunManifest`] / [`RunReport`] — the `(manifest, metrics)`
//! pair C18 mandates per run, with [`RunReport::to_json`] for the structured
//! C14 face;
//! - [`summarize`] — the post-run pure reduction over a run's recorded
//! pip-equity + exposure streams into [`RunMetrics`]; [`f64_field`] bridges a
//! recording sink's `Vec<Scalar>` rows to it.
//!
//! Still to come (subsequent cycles): the `Source` trait + data-server ingestion
//! and source-native time normalization (C3/C11), the broker-independent
//! position-event output and downstream broker nodes (C10), and the atomic sim
//! unit `(topology + params + data-window + seed) -> metrics` that the sweep /
//! optimize / walk-forward / Monte-Carlo axes orchestrate.
//! position-event output and downstream broker nodes (C10), and the param
//! injection + orchestration axes of the atomic sim unit
//! (`(topology + params + data-window + seed)`, swept by optimize / walk-forward
//! / Monte-Carlo) — its `-> metrics` reduction now ships via [`summarize`].
//!
//! Visualization is never here: it is a downstream consumer node on the streams.
mod harness;
mod report;
pub use harness::{BootstrapError, Edge, Harness, SourceSpec, Target};
pub use report::{f64_field, summarize, RunManifest, RunMetrics, RunReport};
+422
View File
@@ -0,0 +1,422 @@
//! Run summary metrics + the reproducible run manifest (C18 / C12): the
//! `(manifest, metrics)` pair a run produces "from day one". The metrics are a
//! **post-run pure reduction** over a run's recorded streams — a node cannot
//! reduce end-of-run (C8 caps a node at one record per `eval`, with no terminal
//! `eval`), so the World drains its recording sinks after [`Harness::run`](crate::Harness::run)
//! and folds them here. Output is canonical, hand-rolled JSON (C14): the schema
//! is tiny, closed, and flat, so the deliberately zero-dependency workspace
//! stays zero-dependency.
use aura_core::{Scalar, Timestamp};
/// Summary metrics reduced from a run's recorded streams — the `-> metrics`
/// half of C12's atomic sim unit. Pure function of the recorded streams.
#[derive(Clone, Debug, PartialEq)]
pub struct RunMetrics {
/// Final cumulative pip equity — the last value of the (cumulative)
/// pip-equity curve. `0.0` if the curve is empty.
pub total_pips: f64,
/// Largest peak-to-trough drop on the cumulative pip curve:
/// `max_t (running_peak(t) - equity(t))`, always `>= 0.0` (`0.0` if the
/// curve is monotonic non-decreasing or empty).
pub max_drawdown: f64,
/// Count of adjacent recorded exposure samples whose sign differs (a zero
/// exposure normalizes to sign `0`, so flat is distinct from long/short).
/// A turnover proxy: it counts long<->short reversals *and* transitions
/// into/out of flat — the plain sign-change count over the exposure series.
pub exposure_sign_flips: u64,
}
/// The reproducible run descriptor (C18). **Caller-supplied**: the engine
/// cannot introspect a git commit, an RNG seed, or a broker label — the World
/// that bootstraps and runs the harness fills these in.
#[derive(Clone, Debug, PartialEq)]
pub struct RunManifest {
/// Node/engine identity: the git commit of the frozen artifact (C18 —
/// commit = identity; the frozen bot *is* a commit).
pub commit: String,
/// The bound tuning params as ordered `name -> value` pairs — the precursor
/// to the eventual typed param-space (deferred; see spec Non-goals).
pub params: Vec<(String, f64)>,
/// The data-window: inclusive `(from, to)` epoch-ns bounds (C12).
pub window: (Timestamp, Timestamp),
/// The RNG seed (C12 seed-as-input). `0` for a seed-free synthetic run.
pub seed: u64,
/// The broker profile label, e.g. `"sim-optimal(pip_size=0.0001)"`.
pub broker: String,
}
/// A run's full structured result: the descriptor plus the metrics it
/// reproduces. The durable run record of C18 ("stores manifests + metrics,
/// re-derives full results on demand").
#[derive(Clone, Debug, PartialEq)]
pub struct RunReport {
pub manifest: RunManifest,
pub metrics: RunMetrics,
}
impl RunReport {
/// Render canonical, machine-readable JSON (C14). Hand-rolled — the schema
/// is tiny, closed, and flat, so the deliberately zero-dependency workspace
/// stays so. Field order is fixed; `f64` uses the round-trippable `{}`
/// shortest form (finite values only — pip equity and exposure are finite
/// by construction); `params` renders as a JSON object in insertion order.
pub fn to_json(&self) -> String {
let m = &self.manifest;
let mut params = String::from("{");
for (i, (name, value)) in m.params.iter().enumerate() {
if i > 0 {
params.push(',');
}
params.push_str(&json_str(name));
params.push(':');
params.push_str(&value.to_string());
}
params.push('}');
format!(
"{{\"manifest\":{{\"commit\":{commit},\"params\":{params},\"window\":[{from},{to}],\"seed\":{seed},\"broker\":{broker}}},\"metrics\":{{\"total_pips\":{pips},\"max_drawdown\":{dd},\"exposure_sign_flips\":{flips}}}}}",
commit = json_str(&m.commit),
from = m.window.0.0,
to = m.window.1.0,
seed = m.seed,
broker = json_str(&m.broker),
pips = self.metrics.total_pips,
dd = self.metrics.max_drawdown,
flips = self.metrics.exposure_sign_flips,
)
}
}
/// Reduce a run's recorded pip-equity + exposure streams into summary metrics.
/// Pure — identical inputs yield identical metrics (C1/C12). Timestamps are
/// carried in the input to match exactly what a sink records; the reduction
/// itself is value-only (it does not read the timestamps).
pub fn summarize(
equity: &[(Timestamp, f64)],
exposure: &[(Timestamp, f64)],
) -> RunMetrics {
// total pips: the last cumulative equity value (0.0 if empty).
let total_pips = equity.last().map(|&(_, v)| v).unwrap_or(0.0);
// max drawdown: the largest running-peak-minus-value, always >= 0.0.
let mut peak = f64::NEG_INFINITY;
let mut max_drawdown = 0.0_f64;
for &(_, v) in equity {
if v > peak {
peak = v;
}
let dd = peak - v;
if dd > max_drawdown {
max_drawdown = dd;
}
}
// exposure sign-flips: adjacent samples whose normalized sign differs.
let mut exposure_sign_flips = 0u64;
let mut prev: Option<f64> = None;
for &(_, v) in exposure {
let s = sign0(v);
if let Some(p) = prev
&& s != p
{
exposure_sign_flips += 1;
}
prev = Some(s);
}
RunMetrics { total_pips, max_drawdown, exposure_sign_flips }
}
/// Three-way sign: `-1.0` / `0.0` / `+1.0`. Unlike `f64::signum` (which returns
/// `+1.0` for `+0.0`), a zero exposure maps to `0.0` so flat is distinct from
/// long/short in the sign-flip count.
fn sign0(v: f64) -> f64 {
if v > 0.0 {
1.0
} else if v < 0.0 {
-1.0
} else {
0.0
}
}
/// Minimal JSON string rendering: wrap in quotes, escape `"` and `\`. Our
/// caller-supplied labels (commit hash, param names, broker label) contain
/// neither control characters nor other JSON-significant bytes, so these two
/// escapes are sufficient.
fn json_str(s: &str) -> String {
let mut out = String::with_capacity(s.len() + 2);
out.push('"');
for c in s.chars() {
match c {
'"' => out.push_str("\\\""),
'\\' => out.push_str("\\\\"),
_ => out.push(c),
}
}
out.push('"');
out
}
/// Bridge a recording sink's recorded `(ts, row)` stream to [`summarize`]:
/// extract one `f64` field of each row into `(ts, f64)` samples. Panics if a
/// row has no such field or the field is not an `f64` scalar — a wiring bug (a
/// sink's declared kinds are fixed at bootstrap, so a correctly-wired
/// equity/exposure sink always yields `f64` at field 0), surfaced like the
/// engine's other "checked at wiring" contract violations rather than silently
/// dropped.
pub fn f64_field(rows: &[(Timestamp, Vec<Scalar>)], field: usize) -> Vec<(Timestamp, f64)> {
rows.iter()
.map(|(ts, row)| {
let Some(&scalar) = row.get(field) else {
panic!("f64_field: row has no field {field} (row width {})", row.len());
};
let Scalar::F64(v) = scalar else {
panic!("f64_field: field {field} is not an f64 scalar: {scalar:?}");
};
(*ts, v)
})
.collect()
}
#[cfg(test)]
mod tests {
use super::*;
use crate::{Edge, Harness, SourceSpec, Target};
use aura_core::{Ctx, Firing, InputSpec, Node, NodeSchema, ScalarKind};
use aura_std::{Exposure, SimBroker, Sma, Sub};
use std::sync::mpsc;
/// Build an f64 source stream from (timestamp, value) points (mirrors the
/// harness.rs test helper; the e2e test needs its own copy — the harness
/// test module's is private to that module).
fn f64_stream(points: &[(i64, f64)]) -> Vec<(Timestamp, Scalar)> {
points.iter().map(|&(t, v)| (Timestamp(t), Scalar::F64(v))).collect()
}
/// A recording sink that sends `(now, row)` out of the graph each fired
/// cycle and returns `None` (pure consumer, C8). Re-declared here because
/// the identically-shaped fixture in `harness.rs` is private to that test
/// module.
struct Recorder {
kinds: Vec<ScalarKind>,
firing: Firing,
tx: mpsc::Sender<(Timestamp, Vec<Scalar>)>,
}
impl Recorder {
fn new(
kinds: &[ScalarKind],
firing: Firing,
tx: mpsc::Sender<(Timestamp, Vec<Scalar>)>,
) -> Self {
Self { kinds: kinds.to_vec(), firing, tx }
}
}
impl Node for Recorder {
fn schema(&self) -> NodeSchema {
NodeSchema {
inputs: self
.kinds
.iter()
.map(|&kind| InputSpec { kind, lookback: 1, firing: self.firing })
.collect(),
output: vec![],
}
}
fn eval(&mut self, ctx: Ctx<'_>) -> Option<&[Scalar]> {
// this test records only f64 streams (one f64 column per sink)
let mut row = Vec::with_capacity(self.kinds.len());
for i in 0..self.kinds.len() {
let w = ctx.f64_in(i);
if w.is_empty() {
return None;
}
row.push(Scalar::F64(w[0]));
}
let _ = self.tx.send((ctx.now(), row));
None
}
}
/// Bootstrap the cycle-0007 signal-quality harness with TWO sinks: one on
/// the SimBroker equity output (node 4 -> node 5) and one on the Exposure
/// output (node 3 -> node 6). Returns the harness plus the two receivers.
#[allow(clippy::type_complexity)]
fn build_two_sink_harness() -> (
Harness,
mpsc::Receiver<(Timestamp, Vec<Scalar>)>,
mpsc::Receiver<(Timestamp, Vec<Scalar>)>,
) {
let (tx_eq, rx_eq) = mpsc::channel();
let (tx_ex, rx_ex) = mpsc::channel();
let h = Harness::bootstrap(
vec![
Box::new(Sma::new(2)), // 0
Box::new(Sma::new(4)), // 1
Box::new(Sub::new()), // 2
Box::new(Exposure::new(0.5)), // 3
Box::new(SimBroker::new(0.0001)), // 4
Box::new(Recorder::new(&[ScalarKind::F64], Firing::Any, tx_eq)), // 5 equity sink
Box::new(Recorder::new(&[ScalarKind::F64], Firing::Any, tx_ex)), // 6 exposure sink
],
vec![SourceSpec {
kind: ScalarKind::F64,
targets: vec![
Target { node: 0, slot: 0 },
Target { node: 1, slot: 0 },
Target { node: 4, slot: 1 }, // price into the broker
],
}],
vec![
Edge { from: 0, to: 2, slot: 0, from_field: 0 },
Edge { from: 1, to: 2, slot: 1, from_field: 0 },
Edge { from: 2, to: 3, slot: 0, from_field: 0 },
Edge { from: 3, to: 4, slot: 0, from_field: 0 },
Edge { from: 4, to: 5, slot: 0, from_field: 0 }, // equity -> sink 5
Edge { from: 3, to: 6, slot: 0, from_field: 0 }, // exposure -> sink 6
],
)
.expect("valid signal-quality DAG");
(h, rx_eq, rx_ex)
}
fn run_once() -> RunReport {
let (mut h, rx_eq, rx_ex) = build_two_sink_harness();
h.run(vec![f64_stream(&[
(1, 1.0000),
(2, 1.0010),
(3, 1.0025),
(4, 1.0020),
(5, 1.0040),
])]);
let eq_rows: Vec<(Timestamp, Vec<Scalar>)> = rx_eq.try_iter().collect();
let ex_rows: Vec<(Timestamp, Vec<Scalar>)> = rx_ex.try_iter().collect();
let equity = f64_field(&eq_rows, 0);
let exposure = f64_field(&ex_rows, 0);
let metrics = summarize(&equity, &exposure);
RunReport {
manifest: RunManifest {
commit: "test-commit".to_string(),
params: vec![
("sma_fast".to_string(), 2.0),
("sma_slow".to_string(), 4.0),
("exposure_scale".to_string(), 0.5),
],
window: (Timestamp(1), Timestamp(5)),
seed: 0,
broker: "sim-optimal(pip_size=0.0001)".to_string(),
},
metrics,
}
}
#[test]
fn report_is_deterministic_end_to_end() {
let r1 = run_once();
let r2 = run_once();
// a run actually emitted metrics over a non-empty pip curve
assert!(r1.metrics.total_pips.is_finite());
// same manifest -> same metrics (C1/C12): two runs are bit-identical
assert_eq!(r1.metrics, r2.metrics);
assert_eq!(r1.to_json(), r2.to_json());
}
fn samples(values: &[f64]) -> Vec<(Timestamp, f64)> {
values
.iter()
.enumerate()
.map(|(i, &v)| (Timestamp(i as i64 + 1), v))
.collect()
}
#[test]
fn summarize_total_pips_is_last_cumulative_value() {
let equity = samples(&[0.0, 5.0, 4.0, 12.0]);
let m = summarize(&equity, &[]);
assert_eq!(m.total_pips, 12.0);
}
#[test]
fn summarize_is_zero_on_empty_streams() {
let m = summarize(&[], &[]);
assert_eq!(m.total_pips, 0.0);
assert_eq!(m.max_drawdown, 0.0);
assert_eq!(m.exposure_sign_flips, 0);
}
#[test]
fn summarize_max_drawdown_is_worst_peak_to_trough() {
// peak 10 then trough 5 (drop 5), recovers to 8; worst drop is 5,
// not the final drop (10 -> 8 = 2).
let equity = samples(&[0.0, 10.0, 5.0, 8.0]);
let m = summarize(&equity, &[]);
assert_eq!(m.max_drawdown, 5.0);
}
#[test]
fn summarize_max_drawdown_zero_on_monotonic_curve() {
let equity = samples(&[0.0, 1.0, 2.0, 3.0]);
let m = summarize(&equity, &[]);
assert_eq!(m.max_drawdown, 0.0);
}
#[test]
fn summarize_sign_flips_counts_signum_changes() {
// signum series: + + - 0 - -> flips at +->-, -->0, 0->- = 3.
let exposure = samples(&[0.5, 0.5, -0.5, 0.0, -0.5]);
let m = summarize(&[], &exposure);
assert_eq!(m.exposure_sign_flips, 3);
}
#[test]
fn summarize_sign_flips_zero_on_constant_sign() {
let exposure = samples(&[0.2, 0.5, 1.0, 0.7]);
let m = summarize(&[], &exposure);
assert_eq!(m.exposure_sign_flips, 0);
}
#[test]
fn f64_field_projects_the_named_field() {
let rows = vec![
(Timestamp(1), vec![Scalar::F64(1.5), Scalar::I64(9)]),
(Timestamp(2), vec![Scalar::F64(2.5), Scalar::I64(8)]),
];
assert_eq!(
f64_field(&rows, 0),
vec![(Timestamp(1), 1.5), (Timestamp(2), 2.5)],
);
}
#[test]
#[should_panic(expected = "not an f64 scalar")]
fn f64_field_panics_on_kind_mismatch() {
let rows = vec![(Timestamp(1), vec![Scalar::I64(7)])];
let _ = f64_field(&rows, 0);
}
#[test]
fn to_json_renders_the_canonical_form() {
let report = RunReport {
manifest: RunManifest {
commit: "abc123".to_string(),
params: vec![
("sma_fast".to_string(), 2.0),
("sma_slow".to_string(), 4.0),
("exposure_scale".to_string(), 1.0),
],
window: (Timestamp(1), Timestamp(6)),
seed: 0,
broker: "sim-optimal(pip_size=1.0)".to_string(),
},
metrics: RunMetrics {
total_pips: 12.0,
max_drawdown: 1.0,
exposure_sign_flips: 1,
},
};
assert_eq!(
report.to_json(),
r#"{"manifest":{"commit":"abc123","params":{"sma_fast":2,"sma_slow":4,"exposure_scale":1},"window":[1,6],"seed":0,"broker":"sim-optimal(pip_size=1.0)"},"metrics":{"total_pips":12,"max_drawdown":1,"exposure_sign_flips":1}}"#,
);
}
}