test: M×N×K stress matrix — node fan-out, deep chains, wide layers

The milestone's green end-to-end gate (#3): prove the closed
bootstrap + run substrate carries arbitrary node-level
fan-out / fan-in / depth / width DAGs deterministically (C1) and computes
every recorded stream correctly. Closes the two coverage holes the engine
tests left open after cycle 0006 — node fan-out (one PRODUCING node read
by several consumers; prior coverage was source fan-out only) and
deep-chain / wide-layer topologies.

Six tests appended to harness.rs's test module, composing the existing
Sma / Sub / Recorder / BarrierSum fixtures (no new fixtures, no aura-std
surface — C9):
- node_fan_out_identical_taps_record_identical_streams
- node_fan_out_divergent_consumers_each_compute_their_own
- node_fan_out_under_mixed_firing_each_consumer_records_per_policy
- deep_transform_chain_propagates_end_to_end
- wide_parallel_layer_multi_sink_records_each_stream
- milestone_end_to_end_mixed_dag_records_every_stream_deterministically

Each asserts exact (Timestamp, Vec<Scalar>) recorded values (hand-computed
from the SMA/Sub + firing/warm-up semantics) AND determinism via a second
fresh-harness drain proving bit-identical streams. 28 engine tests green
(22 + 6), clippy -D warnings clean. No engine misbehaviour surfaced.

closes #3
This commit is contained in:
2026-06-04 14:59:13 +02:00
parent 3a02fbf383
commit e5d80959e0
+380
View File
@@ -1316,4 +1316,384 @@ mod tests {
BootstrapError::KindMismatch { producer: ScalarKind::I64, consumer: ScalarKind::F64 }
);
}
// --- #3 stress matrix: M-producer x N-consumer x K-sink DAGs ---
// The closed `Harness::bootstrap` + `run` substrate must carry arbitrary
// node-level fan-out / fan-in / depth / width deterministically (C1) and
// compute every recorded stream correctly. Earlier cycles proved *source*
// fan-out, single chains and firing in isolation; these close the node
// fan-out / deep-chain / wide-layer holes #3 names. No trading domain.
#[test]
fn node_fan_out_identical_taps_record_identical_streams() {
// PROPERTY: one PRODUCING node read by several consumers via distinct
// edges from the same `from` node feeds every consumer the identical,
// correct stream — node fan-out (not source fan-out: a single SMA(3)
// output is forwarded down three edges to three recorders).
let build = |t1, t2, t3| {
Harness::bootstrap(
vec![
Box::new(Sma::new(3)),
Box::new(Recorder::new(&[ScalarKind::F64], Firing::Any, t1)),
Box::new(Recorder::new(&[ScalarKind::F64], Firing::Any, t2)),
Box::new(Recorder::new(&[ScalarKind::F64], Firing::Any, t3)),
],
vec![SourceSpec {
kind: ScalarKind::F64,
targets: vec![Target { node: 0, slot: 0 }],
}],
vec![
Edge { from: 0, to: 1, slot: 0, from_field: 0 },
Edge { from: 0, to: 2, slot: 0, from_field: 0 },
Edge { from: 0, to: 3, slot: 0, from_field: 0 },
],
)
.expect("valid fan-out DAG")
};
let prices = f64_stream(&[(1, 2.0), (2, 4.0), (3, 6.0), (4, 8.0), (5, 10.0)]);
let (a1, ra) = mpsc::channel();
let (b1, rb) = mpsc::channel();
let (c1, rc) = mpsc::channel();
let mut h = build(a1, b1, c1);
h.run(vec![prices.clone()]);
let s1: Vec<(Timestamp, Vec<Scalar>)> = ra.try_iter().collect();
let s2: Vec<(Timestamp, Vec<Scalar>)> = rb.try_iter().collect();
let s3: Vec<(Timestamp, Vec<Scalar>)> = rc.try_iter().collect();
// SMA(3) warms at cycle 3: mean(2,4,6)=4, mean(4,6,8)=6, mean(6,8,10)=8.
let expected = vec![
(Timestamp(3), vec![Scalar::F64(4.0)]),
(Timestamp(4), vec![Scalar::F64(6.0)]),
(Timestamp(5), vec![Scalar::F64(8.0)]),
];
assert_eq!(s1, expected);
assert_eq!(s2, expected); // every tap sees the identical shared stream
assert_eq!(s3, expected);
// determinism (C1): a fresh harness drains bit-identical.
let (a2, ra2) = mpsc::channel();
let (b2, rb2) = mpsc::channel();
let (c2, rc2) = mpsc::channel();
let mut h2 = build(a2, b2, c2);
h2.run(vec![prices]);
assert_eq!(ra2.try_iter().collect::<Vec<_>>(), expected);
assert_eq!(rb2.try_iter().collect::<Vec<_>>(), expected);
assert_eq!(rc2.try_iter().collect::<Vec<_>>(), expected);
}
#[test]
fn node_fan_out_divergent_consumers_each_compute_their_own() {
// PROPERTY: one producer (SMA(2)) read by consumers that do DIFFERENT
// things — recorded raw by one sink AND fed as input0 of a Sub whose
// other input is a second producer (SMA(4)) — yields both the raw tap
// and the downstream-combined stream, each independently correct.
let build = |t_raw, t_sub| {
Harness::bootstrap(
vec![
Box::new(Sma::new(2)), // 0: shared producer
Box::new(Sma::new(4)), // 1: second producer
Box::new(Sub::new()), // 2: SMA(2) - SMA(4)
Box::new(Recorder::new(&[ScalarKind::F64], Firing::Any, t_raw)), // 3: raw tap of 0
Box::new(Recorder::new(&[ScalarKind::F64], Firing::Any, t_sub)), // 4: tap of Sub
],
vec![SourceSpec {
kind: ScalarKind::F64,
targets: vec![Target { node: 0, slot: 0 }, Target { node: 1, slot: 0 }],
}],
vec![
Edge { from: 0, to: 3, slot: 0, from_field: 0 }, // SMA(2) -> raw recorder
Edge { from: 0, to: 2, slot: 0, from_field: 0 }, // SMA(2) -> Sub.in0
Edge { from: 1, to: 2, slot: 1, from_field: 0 }, // SMA(4) -> Sub.in1
Edge { from: 2, to: 4, slot: 0, from_field: 0 }, // Sub -> Sub recorder
],
)
.expect("valid divergent-fan-out DAG")
};
let prices = f64_stream(&[(1, 10.0), (2, 12.0), (3, 14.0), (4, 16.0), (5, 18.0), (6, 20.0)]);
let (tr, rr) = mpsc::channel();
let (ts, rs) = mpsc::channel();
let mut h = build(tr, ts);
h.run(vec![prices.clone()]);
let raw: Vec<(Timestamp, Vec<Scalar>)> = rr.try_iter().collect();
let sub: Vec<(Timestamp, Vec<Scalar>)> = rs.try_iter().collect();
// SMA(2): 11,13,15,17,19 at cycles 2..6 (raw tap).
let raw_expected = vec![
(Timestamp(2), vec![Scalar::F64(11.0)]),
(Timestamp(3), vec![Scalar::F64(13.0)]),
(Timestamp(4), vec![Scalar::F64(15.0)]),
(Timestamp(5), vec![Scalar::F64(17.0)]),
(Timestamp(6), vec![Scalar::F64(19.0)]),
];
// Sub warms once SMA(4) is warm (cycle 4): 15-13, 17-15, 19-17 -> 2.
let sub_expected = vec![
(Timestamp(4), vec![Scalar::F64(2.0)]),
(Timestamp(5), vec![Scalar::F64(2.0)]),
(Timestamp(6), vec![Scalar::F64(2.0)]),
];
assert_eq!(raw, raw_expected);
assert_eq!(sub, sub_expected);
let (tr2, rr2) = mpsc::channel();
let (ts2, rs2) = mpsc::channel();
let mut h2 = build(tr2, ts2);
h2.run(vec![prices]);
assert_eq!(rr2.try_iter().collect::<Vec<_>>(), raw_expected); // deterministic
assert_eq!(rs2.try_iter().collect::<Vec<_>>(), sub_expected);
}
#[test]
fn node_fan_out_under_mixed_firing_each_consumer_records_per_policy() {
// PROPERTY: one producer (SMA(2)) read simultaneously by an as-of
// (Firing::Any) consumer and a Firing::Barrier consumer — each records
// per ITS OWN firing policy off the SAME shared upstream value. The Any
// tap fires on every SMA push; the BarrierSum fires only on the cycles
// where the held SMA output and a second source coincide on a timestamp.
let build = |t_any, t_bar| {
Harness::bootstrap(
vec![
Box::new(Sma::new(2)), // 0: shared producer (src A)
Box::new(Recorder::new(&[ScalarKind::F64], Firing::Any, t_any)), // 1: as-of tap of 0
Box::new(BarrierSum { out: [Scalar::F64(0.0)] }), // 2: SMA(2) + src B (barrier)
Box::new(Recorder::new(&[ScalarKind::F64], Firing::Any, t_bar)), // 3: tap of barrier
],
vec![
SourceSpec { kind: ScalarKind::F64, targets: vec![Target { node: 0, slot: 0 }] }, // A
SourceSpec { kind: ScalarKind::F64, targets: vec![Target { node: 2, slot: 1 }] }, // B
],
vec![
Edge { from: 0, to: 1, slot: 0, from_field: 0 }, // SMA(2) -> as-of recorder
Edge { from: 0, to: 2, slot: 0, from_field: 0 }, // SMA(2) -> barrier.in0
Edge { from: 2, to: 3, slot: 0, from_field: 0 }, // barrier -> recorder
],
)
.expect("valid mixed-firing fan-out DAG")
};
let a = f64_stream(&[(1, 10.0), (2, 12.0), (3, 14.0), (4, 16.0)]); // drives SMA(2)
let b = f64_stream(&[(2, 100.0), (4, 200.0)]); // barrier.in1
let (ta, rany) = mpsc::channel();
let (tb, rbar) = mpsc::channel();
let mut h = build(ta, tb);
h.run(vec![a.clone(), b.clone()]);
let any: Vec<(Timestamp, Vec<Scalar>)> = rany.try_iter().collect();
let bar: Vec<(Timestamp, Vec<Scalar>)> = rbar.try_iter().collect();
// As-of tap records every SMA(2) fire: 11@t2, 13@t3, 15@t4.
let any_expected = vec![
(Timestamp(2), vec![Scalar::F64(11.0)]),
(Timestamp(3), vec![Scalar::F64(13.0)]),
(Timestamp(4), vec![Scalar::F64(15.0)]),
];
// Barrier fires only where held SMA output and src B share a timestamp:
// t=2 (SMA held=11 + B=100 = 111), t=4 (SMA held=15 + B=200 = 215).
let bar_expected = vec![
(Timestamp(2), vec![Scalar::F64(111.0)]),
(Timestamp(4), vec![Scalar::F64(215.0)]),
];
assert_eq!(any, any_expected);
assert_eq!(bar, bar_expected);
let (ta2, rany2) = mpsc::channel();
let (tb2, rbar2) = mpsc::channel();
let mut h2 = build(ta2, tb2);
h2.run(vec![a, b]);
assert_eq!(rany2.try_iter().collect::<Vec<_>>(), any_expected); // deterministic
assert_eq!(rbar2.try_iter().collect::<Vec<_>>(), bar_expected);
}
#[test]
fn deep_transform_chain_propagates_end_to_end() {
// PROPERTY: a linear chain of several transform nodes (SMA(2) -> SMA(2)
// -> SMA(2) -> recorder) propagates values correctly through depth, with
// each stage's warm-up delaying the tail — closes "deep chains untested".
let build = |tx| {
Harness::bootstrap(
vec![
Box::new(Sma::new(2)), // 0
Box::new(Sma::new(2)), // 1
Box::new(Sma::new(2)), // 2
Box::new(Recorder::new(&[ScalarKind::F64], Firing::Any, tx)), // 3 tail
],
vec![SourceSpec {
kind: ScalarKind::F64,
targets: vec![Target { node: 0, slot: 0 }],
}],
vec![
Edge { from: 0, to: 1, slot: 0, from_field: 0 },
Edge { from: 1, to: 2, slot: 0, from_field: 0 },
Edge { from: 2, to: 3, slot: 0, from_field: 0 },
],
)
.expect("valid deep chain")
};
let prices = f64_stream(&[(1, 2.0), (2, 4.0), (3, 6.0), (4, 8.0), (5, 10.0)]);
let (tx, rx) = mpsc::channel();
let mut h = build(tx);
h.run(vec![prices.clone()]);
let out: Vec<(Timestamp, Vec<Scalar>)> = rx.try_iter().collect();
// stage1 (c2..c5): 3,5,7,9. stage2 (c3..): 4,6,8. stage3 (c4..): 5,7.
let expected = vec![
(Timestamp(4), vec![Scalar::F64(5.0)]),
(Timestamp(5), vec![Scalar::F64(7.0)]),
];
assert_eq!(out, expected);
let (tx2, rx2) = mpsc::channel();
let mut h2 = build(tx2);
h2.run(vec![prices]);
assert_eq!(rx2.try_iter().collect::<Vec<_>>(), expected); // deterministic
}
#[test]
fn wide_parallel_layer_multi_sink_records_each_stream() {
// PROPERTY: one source fanned to several PARALLEL producers of different
// params (SMA(2), SMA(3), SMA(4)), each recorded by its own sink in one
// run, yields each parallel stream correctly — closes "wide layers
// untested" and exercises multi-sink at width.
let build = |t2, t3, t4| {
Harness::bootstrap(
vec![
Box::new(Sma::new(2)), // 0
Box::new(Sma::new(3)), // 1
Box::new(Sma::new(4)), // 2
Box::new(Recorder::new(&[ScalarKind::F64], Firing::Any, t2)), // 3
Box::new(Recorder::new(&[ScalarKind::F64], Firing::Any, t3)), // 4
Box::new(Recorder::new(&[ScalarKind::F64], Firing::Any, t4)), // 5
],
vec![SourceSpec {
kind: ScalarKind::F64,
targets: vec![
Target { node: 0, slot: 0 },
Target { node: 1, slot: 0 },
Target { node: 2, slot: 0 },
],
}],
vec![
Edge { from: 0, to: 3, slot: 0, from_field: 0 },
Edge { from: 1, to: 4, slot: 0, from_field: 0 },
Edge { from: 2, to: 5, slot: 0, from_field: 0 },
],
)
.expect("valid wide layer")
};
let prices = f64_stream(&[(1, 10.0), (2, 12.0), (3, 14.0), (4, 16.0), (5, 18.0)]);
let (a, ra) = mpsc::channel();
let (b, rb) = mpsc::channel();
let (c, rc) = mpsc::channel();
let mut h = build(a, b, c);
h.run(vec![prices.clone()]);
let w2: Vec<(Timestamp, Vec<Scalar>)> = ra.try_iter().collect();
let w3: Vec<(Timestamp, Vec<Scalar>)> = rb.try_iter().collect();
let w4: Vec<(Timestamp, Vec<Scalar>)> = rc.try_iter().collect();
let e2 = vec![
(Timestamp(2), vec![Scalar::F64(11.0)]),
(Timestamp(3), vec![Scalar::F64(13.0)]),
(Timestamp(4), vec![Scalar::F64(15.0)]),
(Timestamp(5), vec![Scalar::F64(17.0)]),
];
let e3 = vec![
(Timestamp(3), vec![Scalar::F64(12.0)]),
(Timestamp(4), vec![Scalar::F64(14.0)]),
(Timestamp(5), vec![Scalar::F64(16.0)]),
];
let e4 = vec![
(Timestamp(4), vec![Scalar::F64(13.0)]),
(Timestamp(5), vec![Scalar::F64(15.0)]),
];
assert_eq!(w2, e2);
assert_eq!(w3, e3);
assert_eq!(w4, e4);
let (a2, ra2) = mpsc::channel();
let (b2, rb2) = mpsc::channel();
let (c2, rc2) = mpsc::channel();
let mut h2 = build(a2, b2, c2);
h2.run(vec![prices]);
assert_eq!(ra2.try_iter().collect::<Vec<_>>(), e2); // deterministic
assert_eq!(rb2.try_iter().collect::<Vec<_>>(), e3);
assert_eq!(rc2.try_iter().collect::<Vec<_>>(), e4);
}
#[test]
fn milestone_end_to_end_mixed_dag_records_every_stream_deterministically() {
// PROPERTY (#3 headline): a single richer DAG combining node fan-out
// (SMA(2) -> raw recorder AND Sub) + source fan-out (source -> SMA(2),
// SMA(4)) + fan-in (Sub) + multiple sinks of MIXED scalar kinds (two f64
// streams + one i64 stream) records every stream correctly AND is fully
// deterministic — the "pure compute substrate carries arbitrary
// M-producer x N-consumer x K-sink DAGs" gate.
let build = |t_raw, t_sub, t_i64| {
Harness::bootstrap(
vec![
Box::new(Sma::new(2)), // 0: shared f64 producer
Box::new(Sma::new(4)), // 1: second f64 producer
Box::new(Sub::new()), // 2: SMA(2) - SMA(4)
Box::new(Recorder::new(&[ScalarKind::F64], Firing::Any, t_raw)), // 3
Box::new(Recorder::new(&[ScalarKind::F64], Firing::Any, t_sub)), // 4
Box::new(Recorder::new(&[ScalarKind::I64], Firing::Any, t_i64)), // 5: i64 sink
],
vec![
SourceSpec {
kind: ScalarKind::F64,
targets: vec![Target { node: 0, slot: 0 }, Target { node: 1, slot: 0 }],
},
SourceSpec { kind: ScalarKind::I64, targets: vec![Target { node: 5, slot: 0 }] },
],
vec![
Edge { from: 0, to: 3, slot: 0, from_field: 0 }, // SMA(2) raw
Edge { from: 0, to: 2, slot: 0, from_field: 0 }, // SMA(2) -> Sub.in0
Edge { from: 1, to: 2, slot: 1, from_field: 0 }, // SMA(4) -> Sub.in1
Edge { from: 2, to: 4, slot: 0, from_field: 0 }, // Sub -> recorder
],
)
.expect("valid milestone DAG")
};
let prices =
f64_stream(&[(1, 10.0), (2, 12.0), (3, 14.0), (4, 16.0), (5, 18.0), (6, 20.0)]);
let counts: Vec<(Timestamp, Scalar)> = vec![
(Timestamp(1), Scalar::I64(7)),
(Timestamp(2), Scalar::I64(8)),
(Timestamp(3), Scalar::I64(9)),
];
let (tr, rr) = mpsc::channel();
let (ts, rs) = mpsc::channel();
let (ti, ri) = mpsc::channel();
let mut h = build(tr, ts, ti);
h.run(vec![prices.clone(), counts.clone()]);
let raw: Vec<(Timestamp, Vec<Scalar>)> = rr.try_iter().collect();
let sub: Vec<(Timestamp, Vec<Scalar>)> = rs.try_iter().collect();
let i64s: Vec<(Timestamp, Vec<Scalar>)> = ri.try_iter().collect();
let raw_expected = vec![
(Timestamp(2), vec![Scalar::F64(11.0)]),
(Timestamp(3), vec![Scalar::F64(13.0)]),
(Timestamp(4), vec![Scalar::F64(15.0)]),
(Timestamp(5), vec![Scalar::F64(17.0)]),
(Timestamp(6), vec![Scalar::F64(19.0)]),
];
let sub_expected = vec![
(Timestamp(4), vec![Scalar::F64(2.0)]),
(Timestamp(5), vec![Scalar::F64(2.0)]),
(Timestamp(6), vec![Scalar::F64(2.0)]),
];
let i64_expected = vec![
(Timestamp(1), vec![Scalar::I64(7)]),
(Timestamp(2), vec![Scalar::I64(8)]),
(Timestamp(3), vec![Scalar::I64(9)]),
];
assert_eq!(raw, raw_expected);
assert_eq!(sub, sub_expected);
assert_eq!(i64s, i64_expected);
let (tr2, rr2) = mpsc::channel();
let (ts2, rs2) = mpsc::channel();
let (ti2, ri2) = mpsc::channel();
let mut h2 = build(tr2, ts2, ti2);
h2.run(vec![prices, counts]);
assert_eq!(rr2.try_iter().collect::<Vec<_>>(), raw_expected); // deterministic
assert_eq!(rs2.try_iter().collect::<Vec<_>>(), sub_expected);
assert_eq!(ri2.try_iter().collect::<Vec<_>>(), i64_expected);
}
}