eec129ee51
The engine's ingestion input was the only type encoding an eager-dataflow
assumption: `Harness::run(Vec<Vec<(Timestamp, Scalar)>>)` — one fully
materialized stream per source. At realistic research scale (20y, 3 tick
streams, ~7.9e9 ticks) that layout is ~190 GB resident, non-residable. Yet
the merge loop never needs the whole stream: it only ever peeks the head
timestamp of each live source and pops one record. So the eager Vec is
gratuitous.
This lands the first half of the Source seam (plan Tasks 1-2):
- A `Source` trait (`peek(&self) -> Option<Timestamp>`,
`next(&mut) -> Option<(Timestamp, Scalar)>`), object-safe so the merge
holds `Vec<Box<dyn Source>>`. The peek/next split is the producer contract
the k-way merge drives.
- `VecSource`, an adapter over the old `Vec<(Timestamp, Scalar)>` — the
cycle-0011 eager shape, named. It IS the old behaviour.
- `Harness::run` re-typed to `Vec<Box<dyn Source>>`; the merge loop's index
arithmetic becomes peek/next. C4 tie-by-source-index is preserved (the
strictly-`<` replace + source-order scan is byte-for-byte the old
semantics). The forward+eval body is unchanged.
- All 53 call sites threaded VecSource-wrapped in the same change, so the
workspace compiles atomically and run output is byte-identical (C1).
Behaviour-preserving is the load-bearing guarantee: the full suite is green
and unchanged (the two-run C1 determinism pairs — real_bars r1/r2, blueprint
sweep-equality, harness h/h2 — stay byte-identical); aura-engine unit tests
129 -> 132 (+3 VecSource unit tests). No residency change yet: VecSource holds
the same Vec as before. The lazy data-server-backed streaming source and its
end-to-end residency proof are the second half (plan Tasks 3-4).
refs #71
405 lines
16 KiB
Rust
405 lines
16 KiB
Rust
//! 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 JSON (C14): the schema is tiny,
|
|
//! closed, and flat. `to_json` renders via serde (the report types derive it,
|
|
//! cycle 0029) — the same encoder the run registry uses, so a record's stdout
|
|
//! and on-disk shapes coincide.
|
|
|
|
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, serde::Serialize, serde::Deserialize)]
|
|
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, serde::Serialize, serde::Deserialize)]
|
|
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, serde::Serialize, serde::Deserialize)]
|
|
pub struct RunReport {
|
|
pub manifest: RunManifest,
|
|
pub metrics: RunMetrics,
|
|
}
|
|
|
|
impl RunReport {
|
|
/// Render the canonical, machine-readable JSON (C14) via serde — the same
|
|
/// encoder the run registry uses on disk, so a record's stdout shape and its
|
|
/// `runs.jsonl` shape are byte-identical. `params` is an array of
|
|
/// `[name, value]` pairs; finite `f64` fields carry a fractional part
|
|
/// (`2.0`). Consumers must parse values as numbers, never key off the token
|
|
/// shape.
|
|
pub fn to_json(&self) -> String {
|
|
serde_json::to_string(self).expect("a finite RunReport always serializes")
|
|
}
|
|
}
|
|
|
|
/// 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
|
|
}
|
|
}
|
|
|
|
/// 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, FlatGraph, Harness, SourceSpec, Target, VecSource};
|
|
use aura_core::{Firing, NodeSchema, PortSpec, ScalarKind};
|
|
use aura_std::{Exposure, Recorder, SimBroker, Sma, Sub};
|
|
use std::sync::mpsc;
|
|
|
|
/// The declared signature of a `Recorder` over one f64 column (the sink shape
|
|
/// the two-sink harness uses).
|
|
fn f64_recorder_sig() -> NodeSchema {
|
|
NodeSchema {
|
|
inputs: vec![PortSpec { kind: ScalarKind::F64, firing: Firing::Any, name: "in".into() }],
|
|
output: vec![],
|
|
params: vec![],
|
|
}
|
|
}
|
|
|
|
/// 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()
|
|
}
|
|
|
|
/// 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(FlatGraph {
|
|
nodes: 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
|
|
],
|
|
signatures: vec![
|
|
Sma::builder().schema().clone(),
|
|
Sma::builder().schema().clone(),
|
|
Sub::builder().schema().clone(),
|
|
Exposure::builder().schema().clone(),
|
|
SimBroker::builder(0.0001).schema().clone(),
|
|
f64_recorder_sig(),
|
|
f64_recorder_sig(),
|
|
],
|
|
sources: 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
|
|
],
|
|
}],
|
|
edges: 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![Box::new(VecSource::new(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.0],["sma_slow",4.0],["exposure_scale",1.0]],"window":[1,6],"seed":0,"broker":"sim-optimal(pip_size=1.0)"},"metrics":{"total_pips":12.0,"max_drawdown":1.0,"exposure_sign_flips":1}}"#,
|
|
);
|
|
}
|
|
|
|
#[test]
|
|
fn to_json_equals_serde_disk_shape() {
|
|
// the same RunReport value the canonical-form test builds.
|
|
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 },
|
|
};
|
|
// stdout (to_json) and disk (serde_json::to_string) are now the same bytes.
|
|
assert_eq!(report.to_json(), serde_json::to_string(&report).unwrap());
|
|
}
|
|
|
|
#[test]
|
|
fn runreport_serde_round_trips() {
|
|
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 },
|
|
};
|
|
let json = serde_json::to_string(&report).expect("serialize RunReport");
|
|
// window is a 2-element [from, to] array (Timestamp newtype is transparent)
|
|
assert!(json.contains("\"window\":[1,6]"), "window shape: {json}");
|
|
let back: RunReport = serde_json::from_str(&json).expect("deserialize RunReport");
|
|
assert_eq!(back, report);
|
|
}
|
|
}
|