From e5d80959e0a192e7332e686c452fb884a0447464 Mon Sep 17 00:00:00 2001 From: Brummel Date: Thu, 4 Jun 2026 14:59:13 +0200 Subject: [PATCH] =?UTF-8?q?test:=20M=C3=97N=C3=97K=20stress=20matrix=20?= =?UTF-8?q?=E2=80=94=20node=20fan-out,=20deep=20chains,=20wide=20layers?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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) 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 --- crates/aura-engine/src/harness.rs | 380 ++++++++++++++++++++++++++++++ 1 file changed, 380 insertions(+) diff --git a/crates/aura-engine/src/harness.rs b/crates/aura-engine/src/harness.rs index 39e03e0..064eea1 100644 --- a/crates/aura-engine/src/harness.rs +++ b/crates/aura-engine/src/harness.rs @@ -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)> = ra.try_iter().collect(); + let s2: Vec<(Timestamp, Vec)> = rb.try_iter().collect(); + let s3: Vec<(Timestamp, Vec)> = 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::>(), expected); + assert_eq!(rb2.try_iter().collect::>(), expected); + assert_eq!(rc2.try_iter().collect::>(), 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)> = rr.try_iter().collect(); + let sub: Vec<(Timestamp, Vec)> = 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::>(), raw_expected); // deterministic + assert_eq!(rs2.try_iter().collect::>(), 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)> = rany.try_iter().collect(); + let bar: Vec<(Timestamp, Vec)> = 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::>(), any_expected); // deterministic + assert_eq!(rbar2.try_iter().collect::>(), 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)> = 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::>(), 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)> = ra.try_iter().collect(); + let w3: Vec<(Timestamp, Vec)> = rb.try_iter().collect(); + let w4: Vec<(Timestamp, Vec)> = 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::>(), e2); // deterministic + assert_eq!(rb2.try_iter().collect::>(), e3); + assert_eq!(rc2.try_iter().collect::>(), 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)> = rr.try_iter().collect(); + let sub: Vec<(Timestamp, Vec)> = rs.try_iter().collect(); + let i64s: Vec<(Timestamp, Vec)> = 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::>(), raw_expected); // deterministic + assert_eq!(rs2.try_iter().collect::>(), sub_expected); + assert_eq!(ri2.try_iter().collect::>(), i64_expected); + } }