Files
Aura/crates/aura-engine/src/report.rs
T
Brummel b91c89f964 feat(broker-foundation): InstrumentSpec deploy metadata + PositionEvent schema
The realistic-broker milestone's (C10 A-side) two foundation value types — no
broker, no derivation yet, fully tested.

#113 — InstrumentSpec gains six deploy-grade fields (instrument_id,
contract_size, pip_value_per_lot, min_lot, lot_step, quote_currency) and the
vetted table is populated for GER40/FRA40/EURUSD/GBPUSD/USDCAD. The struct stays
Copy (quote_currency is &'static str); the _ => None refuse-don't-guess arm is the
permanent floor. This is the milestone's permanent authored floor (recorded-metadata
tier is forward work, #124).

#114 — PositionAction { Buy, Sell, Close } (closed enum, serde-encoded as a bare
i64 0/1/2 via From/TryFrom) + PositionEvent in aura-engine beside RunMetrics,
faithful to the C10 column spec: event_ts, action, position_id, instrument_id,
unsigned volume; no open_ts; direction is the action (no signed-volume trick).
Re-exported from the crate root.

Derived fork decisions (recorded on #113): quote_currency &'static str; action as
ordinal i64; stable instrument_id + overridable vetted values.

Verified by the orchestrator: cargo test --workspace = 458 passed / 0 failed;
clippy --all-targets -D warnings clean. An adversarial 4-lens review (ledger
faithfulness, downstream soundness, serde/wire robustness, test honesty) returned
3 SOUND + 1 CONCERN: quote_currency &'static str cannot hold a runtime-sourced
currency at the #124 resolver — not a blocker (#115/#116 read only the authored
floor), logged on #124 as a deliberate deferred design choice. The implementer's
deviation from the plan snippet (vec![..] -> [..]) avoids clippy::useless_vec;
semantics identical.

Two tester-authored consumer-boundary E2E tests added beyond the inline units:
instrument_spec_deploy_metadata.rs (public resolution seam) and
position_event_table.rs (persisted bare-int wire shape + reversal table).

closes #113 #114
2026-06-22 19:51:17 +02:00

774 lines
31 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, ScalarKind, Timestamp};
use std::collections::HashMap;
/// 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 three position-event actions (C10). Direction IS the action; volume is
/// unsigned. Serde-encoded as its i64 mapping (`Buy=0, Sell=1, Close=2`) so the
/// persisted/columnar form stays C7-scalar and ledger-faithful (`action: i64`).
#[derive(Clone, Copy, Debug, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[serde(into = "i64", try_from = "i64")]
pub enum PositionAction {
Buy,
Sell,
Close,
}
impl From<PositionAction> for i64 {
fn from(a: PositionAction) -> i64 {
match a {
PositionAction::Buy => 0,
PositionAction::Sell => 1,
PositionAction::Close => 2,
}
}
}
impl TryFrom<i64> for PositionAction {
type Error = String;
fn try_from(v: i64) -> Result<Self, Self::Error> {
match v {
0 => Ok(PositionAction::Buy),
1 => Ok(PositionAction::Sell),
2 => Ok(PositionAction::Close),
other => Err(format!("invalid PositionAction i64: {other}")),
}
}
}
/// One row of C10's derived position-event table — the broker-independent audit
/// view (the first difference of the exposure state). A post-run value type
/// (sibling of [`RunMetrics`]), NOT a per-`eval` node output (C8). Multiple events
/// may share one `event_ts` (a reversal: Close then open, close-before-open). No
/// `open_ts` — a position's open time is its opening event's `event_ts`.
#[derive(Clone, Copy, Debug, PartialEq, serde::Serialize, serde::Deserialize)]
pub struct PositionEvent {
pub event_ts: Timestamp,
pub action: PositionAction,
/// Monotonic, assigned at open. A Close references an existing `position_id`.
pub position_id: i64,
pub instrument_id: i64,
/// Lots, unsigned. A partial close carries its own (smaller) volume.
pub volume: f64,
}
/// 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. Each value is a
/// self-describing [`Scalar`], so the param's kind (an `i64` length vs an
/// `f64` scale) survives into the record instead of collapsing to `f64`.
pub params: Vec<(String, Scalar)>,
/// 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 where `value` is a self-describing tagged scalar
/// (serde's externally-tagged enum: `{"I64": 10}` for a length, `{"F64": 2.5}`
/// for a scale). Consumers parse the tagged object, never a bare number.
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());
};
if scalar.kind() != ScalarKind::F64 {
panic!("f64_field: field {field} is not an f64 scalar: {scalar:?}");
}
(*ts, scalar.as_f64())
})
.collect()
}
/// One spine row joined with each side stream's row recorded at the same
/// timestamp. `sides` is parallel to the `sides` argument of [`join_on_ts`]; an
/// entry is `None` where that side did not fire at this spine timestamp.
#[derive(Clone, Debug, PartialEq)]
pub struct JoinedRow {
pub ts: Timestamp,
pub spine: Vec<Scalar>,
pub sides: Vec<Option<Vec<Scalar>>>,
}
/// Join recording-sink tap streams on their recorded timestamp (C8/C18: a post-run
/// reduction over recorded sink output; C3: NOT an in-graph join).
///
/// `spine` defines the row set — exactly one [`JoinedRow`] per spine entry, in
/// spine order. Each side stream is looked up by timestamp: `Some(row)` where it
/// fired at that timestamp, `None` where it did not. The helper does not interpret
/// a row's columns (it returns each whole); the caller maps `None` to whatever
/// default its column means.
///
/// Precondition (C1): each stream has at most one row per timestamp — a sink fires
/// at most once per cycle and cycles have unique timestamps. A duplicate timestamp
/// within one stream resolves last-write-wins. A side row whose timestamp is absent
/// from the spine is dropped (the spine defines the rows).
pub fn join_on_ts(
spine: &[(Timestamp, Vec<Scalar>)],
sides: &[&[(Timestamp, Vec<Scalar>)]],
) -> Vec<JoinedRow> {
let side_maps: Vec<HashMap<i64, &Vec<Scalar>>> = sides
.iter()
.map(|s| s.iter().map(|(t, row)| (t.0, row)).collect())
.collect();
spine
.iter()
.map(|(ts, row)| JoinedRow {
ts: *ts,
spine: row.clone(),
sides: side_maps.iter().map(|m| m.get(&ts.0).map(|r| (*r).clone())).collect(),
})
.collect()
}
/// One drained recorder tap in columnar (SoA) form: parallel arrays, one per
/// recorded column, plus the shared recorded-timestamp axis. Chart-ready (a column
/// is a series of numbers) and kind-tagged (C7) — values are coerced to f64 for
/// plotting; the base type survives in `kinds`. Pure: identical rows always encode
/// to the same value (C1).
#[derive(Clone, Debug, PartialEq, serde::Serialize, serde::Deserialize)]
pub struct ColumnarTrace {
pub tap: String,
pub kinds: Vec<String>,
pub ts: Vec<i64>,
pub columns: Vec<Vec<f64>>,
}
impl ColumnarTrace {
/// Transpose a drained tap's `(ts, row)` pairs into columns. An empty `rows`
/// yields empty `ts` and one empty column per kind. Each cell is coerced to f64:
/// f64 as-is, i64 as f64, bool 1.0/0.0, timestamp epoch as f64. Panics if any
/// row's width disagrees with `kinds.len()` — a wiring bug (a sink's column
/// count is fixed at bootstrap), surfaced as a named panic like
/// [`f64_field`] rather than a bare index-out-of-bounds (wider) or silent
/// ragged columns (narrower).
pub fn from_rows(tap: &str, kinds: &[ScalarKind], rows: &[(Timestamp, Vec<Scalar>)]) -> Self {
let ts: Vec<i64> = rows.iter().map(|(t, _)| t.0).collect();
let mut columns: Vec<Vec<f64>> = vec![Vec::with_capacity(rows.len()); kinds.len()];
for (_, row) in rows {
if row.len() != kinds.len() {
panic!(
"from_rows: row width {} disagrees with kinds.len() {} (tap {tap:?})",
row.len(),
kinds.len(),
);
}
for (c, scalar) in row.iter().enumerate() {
columns[c].push(scalar_to_f64(*scalar));
}
}
ColumnarTrace {
tap: tap.to_string(),
kinds: kinds.iter().map(kind_tag).collect(),
ts,
columns,
}
}
/// Inverse for the serve/align path: rebuild `(ts, row)` pairs with each cell as
/// `Scalar::f64(columns[c][r])` — uniformly f64 (the on-disk store is f64
/// columns; the base type survives only in `kinds`). A true value round-trip for
/// f64 taps; keeps the serve-side `as_f64` projection total.
pub fn to_rows(&self) -> Vec<(Timestamp, Vec<Scalar>)> {
self.ts
.iter()
.enumerate()
.map(|(r, &t)| {
let row = self.columns.iter().map(|col| Scalar::f64(col[r])).collect();
(Timestamp(t), row)
})
.collect()
}
}
/// The four-kind tag string for a column (the on-disk `kinds` form; `ScalarKind`
/// derives no serde, so tags are plain strings).
fn kind_tag(kind: &ScalarKind) -> String {
match kind {
ScalarKind::F64 => "F64",
ScalarKind::I64 => "I64",
ScalarKind::Bool => "Bool",
ScalarKind::Timestamp => "Timestamp",
}
.to_string()
}
/// Coerce a recorded scalar to the f64 a chart plots: f64 as-is, i64 as f64,
/// bool 1.0/0.0, timestamp epoch as f64.
fn scalar_to_f64(s: Scalar) -> f64 {
match s.kind() {
ScalarKind::F64 => s.as_f64(),
ScalarKind::I64 => s.as_i64() as f64,
ScalarKind::Bool => {
if s.as_bool() {
1.0
} else {
0.0
}
}
ScalarKind::Timestamp => s.as_ts().0 as f64,
}
}
#[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;
#[test]
fn position_action_round_trips_through_i64() {
for a in [PositionAction::Buy, PositionAction::Sell, PositionAction::Close] {
let n: i64 = a.into();
assert_eq!(PositionAction::try_from(n), Ok(a));
}
assert_eq!(i64::from(PositionAction::Buy), 0);
assert_eq!(i64::from(PositionAction::Sell), 1);
assert_eq!(i64::from(PositionAction::Close), 2);
}
#[test]
fn position_action_rejects_out_of_range_i64() {
assert!(PositionAction::try_from(3).is_err());
assert!(PositionAction::try_from(-1).is_err());
}
#[test]
fn position_event_serde_round_trips_with_bare_int_action() {
let ev = PositionEvent {
event_ts: Timestamp(42),
action: PositionAction::Sell,
position_id: 7,
instrument_id: 3,
volume: 0.5,
};
let json = serde_json::to_string(&ev).expect("serialize");
// action encodes as a bare integer (C7 scalar shape), not a tagged enum
assert!(json.contains("\"action\":1"), "action not bare-int encoded: {json}");
let back: PositionEvent = serde_json::from_str(&json).expect("deserialize");
assert_eq!(back, ev);
}
#[test]
fn position_events_may_share_one_event_ts_on_reversal() {
// a stop-and-reverse: Close then open at the SAME event_ts (close-before-open).
let ts = Timestamp(100);
let close = PositionEvent {
event_ts: ts, action: PositionAction::Close,
position_id: 1, instrument_id: 3, volume: 0.5,
};
let open = PositionEvent {
event_ts: ts, action: PositionAction::Sell,
position_id: 2, instrument_id: 3, volume: 0.5,
};
let table = [close, open];
assert_eq!(table[0].event_ts, table[1].event_ts);
assert_eq!(table[0].action, PositionAction::Close);
assert_eq!(table[1].action, PositionAction::Sell);
}
/// 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(), Scalar::i64(2)),
("sma_slow".to_string(), Scalar::i64(4)),
("exposure_scale".to_string(), Scalar::f64(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(), Scalar::i64(2)),
("sma_slow".to_string(), Scalar::i64(4)),
("exposure_scale".to_string(), Scalar::f64(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",{"I64":2}],["sma_slow",{"I64":4}],["exposure_scale",{"F64":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(), Scalar::i64(2)),
("sma_slow".to_string(), Scalar::i64(4)),
("exposure_scale".to_string(), Scalar::f64(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(), Scalar::i64(2)),
("sma_slow".to_string(), Scalar::i64(4)),
("exposure_scale".to_string(), Scalar::f64(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);
}
#[test]
fn join_on_ts_aligns_streams_of_different_cardinality() {
// spine fires every bar; side A is one row shorter (no ts 10, like cold
// Delay(1) on the first bar); side B fires on a subset (only ts 20, 40,
// like a Session filter before the open).
let spine = vec![
(Timestamp(10), vec![Scalar::f64(1.0)]),
(Timestamp(20), vec![Scalar::f64(2.0)]),
(Timestamp(30), vec![Scalar::f64(3.0)]),
(Timestamp(40), vec![Scalar::f64(4.0)]),
];
let side_a = vec![
(Timestamp(20), vec![Scalar::bool(true)]),
(Timestamp(30), vec![Scalar::bool(false)]),
(Timestamp(40), vec![Scalar::bool(true)]),
];
let side_b = vec![
(Timestamp(20), vec![Scalar::i64(0)]),
(Timestamp(40), vec![Scalar::i64(2)]),
];
let joined = join_on_ts(&spine, &[&side_a, &side_b]);
// one row per spine entry, in spine order
assert_eq!(joined.len(), 4);
assert_eq!(
joined.iter().map(|j| j.ts).collect::<Vec<_>>(),
vec![Timestamp(10), Timestamp(20), Timestamp(30), Timestamp(40)]
);
// ts 10: spine present, both sides absent (the zip-by-index misalignment case)
assert_eq!(joined[0].spine, vec![Scalar::f64(1.0)]);
assert_eq!(joined[0].sides[0], None);
assert_eq!(joined[0].sides[1], None);
// ts 20: both sides present and aligned to THIS ts
assert_eq!(joined[1].sides[0], Some(vec![Scalar::bool(true)]));
assert_eq!(joined[1].sides[1], Some(vec![Scalar::i64(0)]));
// ts 30: side A present, side B absent
assert_eq!(joined[2].sides[0], Some(vec![Scalar::bool(false)]));
assert_eq!(joined[2].sides[1], None);
// ts 40: both present
assert_eq!(joined[3].sides[0], Some(vec![Scalar::bool(true)]));
assert_eq!(joined[3].sides[1], Some(vec![Scalar::i64(2)]));
}
#[test]
fn join_on_ts_drops_side_rows_absent_from_spine() {
// a side row whose ts is not in the spine is dropped — the spine defines
// the row set.
let spine = vec![(Timestamp(10), vec![Scalar::f64(1.0)])];
let side = vec![
(Timestamp(10), vec![Scalar::i64(7)]),
(Timestamp(99), vec![Scalar::i64(8)]), // ts 99 absent from spine -> dropped
];
let joined = join_on_ts(&spine, &[&side]);
assert_eq!(joined.len(), 1);
assert_eq!(joined[0].ts, Timestamp(10));
assert_eq!(joined[0].sides[0], Some(vec![Scalar::i64(7)]));
}
#[test]
fn columnar_trace_round_trips_f64_rows() {
let rows = vec![
(Timestamp(2), vec![Scalar::f64(10.0)]),
(Timestamp(3), vec![Scalar::f64(20.0)]),
(Timestamp(4), vec![Scalar::f64(30.0)]),
];
let ct = ColumnarTrace::from_rows("equity", &[ScalarKind::F64], &rows);
assert_eq!(ct.tap, "equity");
assert_eq!(ct.kinds, vec!["F64".to_string()]);
assert_eq!(ct.ts, vec![2, 3, 4]);
assert_eq!(ct.columns, vec![vec![10.0, 20.0, 30.0]]);
// to_rows is the inverse for f64 taps
assert_eq!(ct.to_rows(), rows);
}
#[test]
fn columnar_trace_empty_rows_yields_empty_columns_sized_to_kinds() {
let ct = ColumnarTrace::from_rows("exposure", &[ScalarKind::F64], &[]);
assert_eq!(ct.ts, Vec::<i64>::new());
assert_eq!(ct.columns, vec![Vec::<f64>::new()]);
assert!(ct.to_rows().is_empty());
}
#[test]
fn columnar_trace_coerces_non_f64_kinds_to_f64() {
let rows = vec![
(Timestamp(1), vec![Scalar::i64(7), Scalar::bool(true)]),
(Timestamp(2), vec![Scalar::i64(9), Scalar::bool(false)]),
];
let ct = ColumnarTrace::from_rows("mix", &[ScalarKind::I64, ScalarKind::Bool], &rows);
assert_eq!(ct.kinds, vec!["I64".to_string(), "Bool".to_string()]);
assert_eq!(ct.columns, vec![vec![7.0, 9.0], vec![1.0, 0.0]]);
}
/// The fourth coercion arm: a Timestamp-kind column tags as "Timestamp" and
/// its epoch-ns survives the f64 projection (`scalar_to_f64` reads `as_ts().0`),
/// closing the kind for which `from_rows` would otherwise be untested.
#[test]
fn columnar_trace_coerces_timestamp_kind_to_epoch_f64() {
let rows = vec![
(Timestamp(1), vec![Scalar::ts(Timestamp(1_700_000_000_000_000_000))]),
(Timestamp(2), vec![Scalar::ts(Timestamp(1_700_000_000_000_000_060))]),
];
let ct = ColumnarTrace::from_rows("clock", &[ScalarKind::Timestamp], &rows);
assert_eq!(ct.kinds, vec!["Timestamp".to_string()]);
assert_eq!(
ct.columns,
vec![vec![1_700_000_000_000_000_000.0, 1_700_000_000_000_000_060.0]],
);
}
/// A row wider than the declared kinds is a wiring bug (a sink's column count
/// is fixed at bootstrap), surfaced as a named panic like `f64_field` — not a
/// bare index-out-of-bounds.
#[test]
#[should_panic(expected = "row width 2 disagrees with kinds.len() 1")]
fn from_rows_panics_on_row_wider_than_kinds() {
let rows = vec![(Timestamp(1), vec![Scalar::f64(1.0), Scalar::f64(2.0)])];
let _ = ColumnarTrace::from_rows("wide", &[ScalarKind::F64], &rows);
}
/// Symmetric to the wide case: a row narrower than the declared kinds is the
/// same wiring-bug class and carries the same named panic.
#[test]
#[should_panic(expected = "row width 1 disagrees with kinds.len() 2")]
fn from_rows_panics_on_row_narrower_than_kinds() {
let rows = vec![(Timestamp(1), vec![Scalar::f64(1.0)])];
let _ = ColumnarTrace::from_rows("narrow", &[ScalarKind::F64, ScalarKind::F64], &rows);
}
}