From 7b8e8010dd79715f1ef3704eba10d7036921382c Mon Sep 17 00:00:00 2001 From: Brummel Date: Thu, 4 Jun 2026 14:28:14 +0200 Subject: [PATCH] plan: 0006 sink recording MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Six tasks: (1) aura-core Ctx::now(); (2) thread the Ctx::new signature change through aura-std's five test call sites; (3) aura-engine production shrink (remove observe, run -> (), thread the run-loop Ctx); (4) migrate the engine test suite to a test-local Recorder fixture + drained channel (the finish-threading task — the test module is one compile unit, so it restores the full cargo test); (5) eight new node-recording proof tests; (6) ledger C8/C22 realization notes + the ship gate (workspace test, clippy, purity grep). Recon found the spec undercounts Ctx::new call sites (9, not 3): the aura-std and aura-engine run-loop sites also break on the signature change and are threaded in tasks 2-3. refs #2 --- docs/plans/0006-sink-recording.md | 1374 +++++++++++++++++++++++++++++ 1 file changed, 1374 insertions(+) create mode 100644 docs/plans/0006-sink-recording.md diff --git a/docs/plans/0006-sink-recording.md b/docs/plans/0006-sink-recording.md new file mode 100644 index 0000000..5ee6c51 --- /dev/null +++ b/docs/plans/0006-sink-recording.md @@ -0,0 +1,1374 @@ +# Sink recording — recording is a node role, not a type — Implementation Plan + +> **Parent spec:** `docs/specs/0006-sink-recording.md` +> +> **For agentic workers:** REQUIRED SUB-SKILL: use the `implement` skill to run +> this plan. Steps use `- [ ]` checkboxes for tracking. + +**Goal:** Replace the engine's single `observe: usize` recording affordance with +recording-by-node, so one run records many streams; recording is an out-of-graph +`eval` side effect to a destination the node holds, with no `Sink` type. + +**Architecture:** Three production changes — `Ctx` gains `now: Timestamp` + +`now()` (aura-core); `Harness` loses `observe` (field, bootstrap param + check, +per-cycle collection) and `run` returns `()` constructing `Ctx::new(&inputs, ts)` +(aura-engine). A test-local `Recorder` fixture (holding an `mpsc::Sender`) proves +multi-stream recording; existing firing/DAG tests migrate from the dense +`Vec>` to a sparse, timestamped drained channel. Ledger C8/C22 gain +cycle-0006 realization notes. + +**Tech Stack:** Rust workspace (aura-core ← aura-std ← aura-engine); +`std::sync::mpsc` for the test read-back; the existing per-field kind check (0005) +covers recorder edges. + +--- + +## Files this plan creates or modifies + +- Modify: `crates/aura-core/src/ctx.rs:10-19` — `Ctx` gains `now: Timestamp` + + `now()`; `Ctx::new` takes `now`; 3 in-crate test call sites updated; new unit + test `ctx_now_returns_cycle_timestamp`. +- Modify: `crates/aura-std/src/sma.rs:51,70,84,87` — 3 `Ctx::new` test call sites + thread a timestamp; test-module import gains `Timestamp`. +- Modify: `crates/aura-std/src/sub.rs:52,64,68` — 2 `Ctx::new` test call sites + thread a timestamp; test-module import gains `Timestamp`. +- Modify: `crates/aura-engine/src/harness.rs` — production: remove `observe` + (field `:91`, Debug `:105`, bootstrap param `:119` + check `:122-124`, + construction `:204`, destructure `:224-225`, loop bookkeeping `:263,:277-279, + :293,:295`), change `run` to `()` (`:213`), thread `Ctx::new(&nb.inputs, ts)` + (`:275`), fix doc-comments (`:8,:59,:208-212`). Tests: add `Recorder` + + `TapForward` fixtures; migrate the 14 bootstrap/run tests; add 8 new proof + tests. +- Modify: `docs/design/INDEX.md:210,495` — C8 + C22 cycle-0006 realization notes. + +**Decision recorded (orchestrator):** the spec lists both a migrated +`ohlcv_bundles_five_field_record` and a new `recorder_taps_all_fields_of_a_record`; +they prove overlapping mechanics. Both are kept and differentiated: the migrated +test records **two** bars (focus: barrier timing — one record per bar, on the +fifth cycle of each timestamp), the new test records **one** bar (focus: a 5-input +recorder taps all five fields via five field-wise edges, 0005 — asserts the row +has five fields). No spec-named test is dropped. + +**Out of scope (stays #3):** the exhaustive multi-producer × multi-consumer × +multi-sink stress matrix. This plan ships the substrate (#2) only. + +--- + +## Task 1: `Ctx::now()` (aura-core) + +**Files:** +- Modify: `crates/aura-core/src/ctx.rs` + +- [ ] **Step 1: Add the `now` field, the `now()` accessor, and update `Ctx::new`** + +Replace the struct + `impl` opening (`ctx.rs:8-19`): + +```rust +/// Read-only, zero-copy view of a node's inputs for one `eval`, in schema +/// order, plus the cycle's timestamp (C4). `Copy` because it is just a borrow of +/// the input slice plus a `Copy` timestamp. +#[derive(Clone, Copy)] +pub struct Ctx<'a> { + inputs: &'a [AnyColumn], + now: Timestamp, +} + +impl<'a> Ctx<'a> { + /// Wrap the per-input columns (in schema-declared order) and the cycle + /// timestamp for one `eval`. + pub fn new(inputs: &'a [AnyColumn], now: Timestamp) -> Self { + Self { inputs, now } + } + + /// The current cycle's timestamp (C4). Causal — the present cycle's + /// timestamp, never the future (C2) — so reading it introduces no look-ahead. + pub fn now(&self) -> Timestamp { + self.now + } +``` + +(`Timestamp` is already imported at `ctx.rs:6`.) + +- [ ] **Step 2: Thread the timestamp through the 3 in-crate test call sites** + +In `ctx.rs` `#[cfg(test)] mod tests`, the three `Ctx::new(&inputs)` calls become +`Ctx::new(&inputs, Timestamp(0))`: + +- `ctx.rs:72` (in `ctx_hands_financial_indexed_windows`): +```rust + let ctx = Ctx::new(&inputs, Timestamp(0)); +``` +- `ctx.rs:87` (in `ctx_addresses_multiple_inputs`): +```rust + let ctx = Ctx::new(&inputs, Timestamp(0)); +``` +- `ctx.rs:97` (in `ctx_panics_on_kind_mismatch`): +```rust + let ctx = Ctx::new(&inputs, Timestamp(0)); +``` + +The test-module import (`ctx.rs:64`, `use crate::{Scalar, ScalarKind};`) gains +`Timestamp`: +```rust + use crate::{Scalar, ScalarKind, Timestamp}; +``` + +- [ ] **Step 3: Add the `ctx_now_returns_cycle_timestamp` unit test** + +Append inside `#[cfg(test)] mod tests` (after `ctx_panics_on_kind_mismatch`, before +the closing `}` at `ctx.rs:100`): + +```rust + #[test] + fn ctx_now_returns_cycle_timestamp() { + let inputs: Vec = vec![]; + let ctx = Ctx::new(&inputs, Timestamp(42)); + assert_eq!(ctx.now(), Timestamp(42)); + } +``` + +- [ ] **Step 4: Verify aura-core compiles and its tests pass** + +Run: `cargo test -p aura-core` +Expected: PASS — `test result: ok. 20 passed; 0 failed` (the prior 19 + the new +`ctx_now_returns_cycle_timestamp`). The downstream crates are not compiled by a +`-p aura-core` run, so their still-old `Ctx::new` calls do not break this gate. + +--- + +## Task 2: thread `Ctx::new` through aura-std (caller-threading forced by Task 1) + +**Files:** +- Modify: `crates/aura-std/src/sma.rs` +- Modify: `crates/aura-std/src/sub.rs` + +The `Ctx::new` signature change in Task 1 breaks every aura-std test call site. +This task threads all five (compile-driven enumeration: `sma.rs:70,84,87`; +`sub.rs:64,68`). Mechanical — the node tests do not assert on `now()`, so a +`Timestamp(0)` placeholder is correct. + +- [ ] **Step 1: Update the three `Ctx::new` call sites in `sma.rs`** + +Add `Timestamp` to the test-module import (`sma.rs:51`, `use aura_core::AnyColumn;`): +```rust + use aura_core::{AnyColumn, Timestamp}; +``` + +`sma.rs:70` (in `sma_warms_up_then_tracks_the_window_mean`): +```rust + let got = sma.eval(Ctx::new(&inputs, Timestamp(0))); +``` +`sma.rs:84` (in `sma_length_one_is_identity`): +```rust + assert_eq!(sma.eval(Ctx::new(&inputs, Timestamp(0))), Some([Scalar::F64(7.0)].as_slice())); +``` +`sma.rs:87` (same test): +```rust + assert_eq!(sma.eval(Ctx::new(&inputs, Timestamp(0))), Some([Scalar::F64(9.0)].as_slice())); +``` + +- [ ] **Step 2: Update the two `Ctx::new` call sites in `sub.rs`** + +Add `Timestamp` to the test-module import (`sub.rs:52`, `use aura_core::AnyColumn;`): +```rust + use aura_core::{AnyColumn, Timestamp}; +``` + +`sub.rs:64` (in `sub_is_difference_once_both_inputs_present`): +```rust + assert_eq!(sub.eval(Ctx::new(&inputs, Timestamp(0))), None); +``` +`sub.rs:68` (same test): +```rust + assert_eq!(sub.eval(Ctx::new(&inputs, Timestamp(0))), Some([Scalar::F64(6.0)].as_slice())); +``` + +- [ ] **Step 3: Verify aura-std compiles and its tests pass** + +Run: `cargo test -p aura-std` +Expected: PASS — `test result: ok. 3 passed; 0 failed`. aura-engine (downstream) +is not compiled by this gate, so its still-old API does not break it. + +--- + +## Task 3: shrink the engine surface — remove `observe`, `run -> ()` (aura-engine production) + +**Files:** +- Modify: `crates/aura-engine/src/harness.rs` (production code only; the + `#[cfg(test)] mod tests` is migrated in Task 4) + +This task changes three production surfaces (`observe` removal, `run` return type, +the run-loop `Ctx::new`) and the doc-comments that reference them. The +`#[cfg(test)]` test module is **left broken on purpose** — it still calls the old +4-arg `bootstrap` and binds `run`'s return — and is restored in Task 4. The gate +here is therefore a **production-only build** (`cargo build`, which does not +compile `#[cfg(test)]` modules), per the planner's compile-gate-ordering rule. + +- [ ] **Step 1: Fix the module doc and the `BadIndex` doc** + +`harness.rs:8` — the module doc currently reads (exact substring to replace): +```rust +//! (the cycle-0002 shape), so `Ctx` is unchanged. A node's `eval` returns a +``` +Replace that one line with: +```rust +//! (the cycle-0002 shape), so `Ctx` borrows them read-only and additionally +//! carries the cycle timestamp (`ctx.now()`, C4). A node's `eval` returns a +``` + +`harness.rs:59` — drop the observe clause from the `BadIndex` doc: +```rust + /// A node or slot index in an edge or target is out of range. + BadIndex, +``` + +- [ ] **Step 2: Remove the `observe` field and its `Debug` line** + +`harness.rs:86-92` — the struct loses `observe`: +```rust +/// A bootstrapped, frozen root graph instance plus its deterministic run loop. +pub struct Harness { + nodes: Vec, + topo: Vec, + out_edges: Vec>, + sources: Vec, +} +``` + +`harness.rs:100-106` — the `Debug` impl loses the `observe` field line: +```rust + f.debug_struct("Harness") + .field("nodes", &self.nodes.len()) + .field("topo", &self.topo) + .field("out_edges", &self.out_edges) + .field("sources", &self.sources) + .finish() +``` + +- [ ] **Step 3: Drop the `observe` parameter, its bounds check, and its construction** + +`harness.rs:115-124` — `bootstrap` loses its fourth parameter and the +`observe >= n` check: +```rust + pub fn bootstrap( + nodes: Vec>, + sources: Vec, + edges: Vec, + ) -> Result { + let n = nodes.len(); + + let schemas: Vec<_> = nodes.iter().map(|nd| nd.schema()).collect(); +``` +(The `let n = nodes.len();` line stays — `n` is still used to size `out_edges` and +`indeg`. Only the `if observe >= n { ... }` block is removed.) + +`harness.rs:199-205` — the construction loses `observe`: +```rust + Ok(Harness { + nodes: boxes, + topo, + out_edges, + sources, + }) +``` + +- [ ] **Step 4: Change `run` to return `()`, drop the observe bookkeeping, thread `ts` into `Ctx`** + +Replace the **entire** `run` function (`harness.rs:208-296`, from the `/// Drive +the sources` doc-comment through the closing `}` of `run`) with the block below. +Changes vs. the original: return type `()`; doc rewritten; `observe` dropped from +the destructure and the `let observe = *observe;` line gone; `let mut out` gone; +`let mut observed` gone; `Ctx::new(&nb.inputs, ts)`; the `if nidx == observe` +branch gone; `out.push(observed)` and the trailing `out` return gone. + +```rust + /// Drive the sources, k-way-merged in timestamp order (ties by source index, + /// C4). One stream per source, each ascending in timestamp (C3 ingestion + /// precondition). Recording is a node-side concern: a recording node pushes + /// its record to a destination it holds (out of graph) inside `eval`; the + /// engine only routes in-graph edges and is oblivious to the side effect. + /// Allocates nothing per cycle beyond the reused scratch buffer. + pub fn run(&mut self, streams: Vec>) { + assert_eq!( + streams.len(), + self.sources.len(), + "run: one stream per source required (got {} streams for {} sources)", + streams.len(), + self.sources.len() + ); + + // disjoint field borrows so the topo walk can read topo/out_edges/sources + // while mutating nodes + let Harness { nodes, topo, out_edges, sources } = self; + + let mut cursor: Vec = vec![0; streams.len()]; + let mut cycle_id: u64 = 0; + let mut scratch: Vec = Vec::new(); + + loop { + // pick the live source head with the smallest (timestamp, source index) + let mut pick: Option = None; + for (s, stream) in streams.iter().enumerate() { + if cursor[s] < stream.len() { + match pick { + None => pick = Some(s), + Some(p) => { + if stream[cursor[s]].0 < streams[p][cursor[p]].0 { + pick = Some(s); + } + } + } + } + } + let s = match pick { + Some(s) => s, + None => break, // all streams exhausted + }; + let (ts, value) = streams[s][cursor[s]]; + cursor[s] += 1; + cycle_id += 1; + + // forward the source value into its target slots, stamping freshness + for t in sources[s].targets.iter() { + let nb = &mut nodes[t.node]; + nb.inputs[t.slot].push(value).expect("source kind checked at wiring"); + nb.slots[t.slot] = SlotState { fresh_at: cycle_id, last_ts: ts }; + } + + // evaluate in topological order; gate by firing; forward Some outputs + for &nidx in topo.iter() { + let out_len = nodes[nidx].out_len; + let fired = { + let nb = &nodes[nidx]; + fires(&nb.firing, &nb.slots, cycle_id, ts) + }; + if !fired { + continue; // hold: no eval, no push + } + let result: Option<&[Scalar]> = { + let nb = &mut nodes[nidx]; + nb.node.eval(Ctx::new(&nb.inputs, ts)) + }; + if let Some(row) = result { + debug_assert_eq!(row.len(), out_len, "node returned a row of the wrong width"); + scratch.clear(); + scratch.extend_from_slice(row); + for e in out_edges[nidx].iter() { + let nb = &mut nodes[e.to]; + nb.inputs[e.slot] + .push(scratch[e.from_field]) + .expect("edge kind checked at wiring"); + nb.slots[e.slot] = SlotState { fresh_at: cycle_id, last_ts: ts }; + } + } + } + } + } + +- [ ] **Step 5: Verify the production library builds (test module intentionally still broken)** + +Run: `cargo build --workspace` +Expected: PASS — `Finished` with no errors. `cargo build` does not compile +`#[cfg(test)]` modules, so the not-yet-migrated `harness.rs` tests do not break +this gate. aura-core and aura-std production are already on the new API (Tasks 1-2); +no production caller of `run`/`bootstrap` exists outside the test module +(verified: aura-cli does not reference them). + +--- + +## Task 4: migrate the engine test suite onto the recording API (aura-engine tests) + +**Files:** +- Modify: `crates/aura-engine/src/harness.rs` (`#[cfg(test)] mod tests` only) + +This is the "finish-threading" task: the `harness.rs` test module is one +compilation unit, so it compiles only once **all** 14 `bootstrap` call sites drop +their 4th argument and all run-binding sites move to the drained channel. The gate +is therefore the full `cargo test -p aura-engine`. Behaviour is preserved: a +recorded stream is exactly the old `Some(row)` entries, now sparse and tagged with +each firing cycle's timestamp. + +- [ ] **Step 1: Add the test-module `mpsc` import and the `Recorder` fixture** + +At the top of `#[cfg(test)] mod tests` (`harness.rs:348-353`), add the `mpsc` +import after the existing `use` lines: +```rust + use std::sync::mpsc; +``` + +Add the `Recorder` fixture alongside the other fixtures (after `TwoField`, before +the first `#[test]` at `harness.rs:501`): +```rust + /// A recording node (test-local fixture; stands in for a downstream author's + /// chart/registry sink). It declares typed input slots and holds an + /// `mpsc::Sender`; on every fired cycle it reads the newest of each input plus + /// `ctx.now()`, sends the timestamped record out of the graph, and returns + /// `None` (pure consumer — C8). Read-back is via the channel, never `Rc`/ + /// `RefCell`, so `aura-engine/src` stays free of the interior-mutability the + /// purity invariant (C7) forbids. + struct Recorder { + kinds: Vec, + firing: Firing, + tx: mpsc::Sender<(Timestamp, Vec)>, + } + impl Recorder { + fn new( + kinds: &[ScalarKind], + firing: Firing, + tx: mpsc::Sender<(Timestamp, Vec)>, + ) -> 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![], // pure sink: no output port + } + } + fn eval(&mut self, ctx: Ctx<'_>) -> Option<&[Scalar]> { + let mut row = Vec::with_capacity(self.kinds.len()); + for (i, &kind) in self.kinds.iter().enumerate() { + let v = match kind { + ScalarKind::I64 => { + let w = ctx.i64_in(i); + if w.is_empty() { + return None; // not yet warmed + } + Scalar::I64(w[0]) + } + ScalarKind::F64 => { + let w = ctx.f64_in(i); + if w.is_empty() { + return None; + } + Scalar::F64(w[0]) + } + ScalarKind::Bool => { + let w = ctx.bool_in(i); + if w.is_empty() { + return None; + } + Scalar::Bool(w[0]) + } + ScalarKind::Timestamp => { + let w = ctx.ts_in(i); + if w.is_empty() { + return None; + } + Scalar::Ts(w[0]) + } + }; + row.push(v); + } + let _ = self.tx.send((ctx.now(), row)); // out-of-graph side effect + None // records, forwards nothing + } + } +``` + +- [ ] **Step 2: Migrate `chain_source_sma_runs`** + +Replace `harness.rs:501-522`: +```rust + #[test] + fn chain_source_sma_runs() { + // node 0 = SMA(3); source -> SMA(3).in0; node 1 = Recorder taps node 0. + let (tx, rx) = mpsc::channel(); + let mut h = Harness::bootstrap( + vec![ + Box::new(Sma::new(3)), + Box::new(Recorder::new(&[ScalarKind::F64], Firing::Any, tx)), + ], + vec![SourceSpec { kind: ScalarKind::F64, targets: vec![Target { node: 0, slot: 0 }] }], + vec![Edge { from: 0, to: 1, slot: 0, from_field: 0 }], + ) + .expect("valid"); + h.run(vec![f64_stream(&[(1, 1.0), (2, 2.0), (3, 3.0), (4, 4.0), (5, 5.0)])]); + let got: Vec<(Timestamp, Vec)> = rx.try_iter().collect(); + // SMA(3) warms at cycle 3; the recorder captures only fired cycles, each + // tagged with the cycle's timestamp (sparse — no None hold-rows). + assert_eq!( + got, + vec![ + (Timestamp(3), vec![Scalar::F64(2.0)]), + (Timestamp(4), vec![Scalar::F64(3.0)]), + (Timestamp(5), vec![Scalar::F64(4.0)]), + ] + ); + } +``` + +- [ ] **Step 3: Migrate `fan_out_join_dag_runs_deterministically`** + +Replace `harness.rs:524-560`: +```rust + #[test] + fn fan_out_join_dag_runs_deterministically() { + // 0 = SMA(2), 1 = SMA(4), 2 = Sub; source fans into both SMAs; SMAs join + // into Sub; node 3 = Recorder taps Sub — the 0003 baseline on the new API. + let build = |tx| { + Harness::bootstrap( + vec![ + Box::new(Sma::new(2)), + Box::new(Sma::new(4)), + Box::new(Sub::new()), + Box::new(Recorder::new(&[ScalarKind::F64], Firing::Any, tx)), + ], + vec![SourceSpec { + kind: ScalarKind::F64, + targets: vec![Target { node: 0, slot: 0 }, Target { node: 1, slot: 0 }], + }], + 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 }, + ], + ) + .expect("valid DAG") + }; + let prices = f64_stream(&[(1, 10.0), (2, 12.0), (3, 14.0), (4, 16.0), (5, 18.0), (6, 20.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(); + // Sub fires once SMA(4) is warm (cycle 4): 15-13, 17-15, 19-17 -> 2. + assert_eq!( + out, + vec![ + (Timestamp(4), vec![Scalar::F64(2.0)]), + (Timestamp(5), vec![Scalar::F64(2.0)]), + (Timestamp(6), vec![Scalar::F64(2.0)]), + ] + ); + + // determinism (C1): a second identical run drains a bit-identical stream. + let (tx2, rx2) = mpsc::channel(); + let mut h2 = build(tx2); + h2.run(vec![prices]); + let out2: Vec<(Timestamp, Vec)> = rx2.try_iter().collect(); + assert_eq!(out2, out); + } +``` + +- [ ] **Step 4: Migrate `mode_a_as_of_fires_on_any_fresh_and_holds`** + +Replace `harness.rs:562-597`: +```rust + #[test] + fn mode_a_as_of_fires_on_any_fresh_and_holds() { + // AsOfSum @0; node 1 = Recorder taps it. source 0 ticks t=1..4; source 1 + // ticks t=2,4 (slower); both AsOfSum inputs Any. + let build = |tx| { + Harness::bootstrap( + vec![ + Box::new(AsOfSum { out: [Scalar::F64(0.0)] }), + Box::new(Recorder::new(&[ScalarKind::F64], Firing::Any, tx)), + ], + vec![ + SourceSpec { kind: ScalarKind::F64, targets: vec![Target { node: 0, slot: 0 }] }, + SourceSpec { kind: ScalarKind::F64, targets: vec![Target { node: 0, slot: 1 }] }, + ], + vec![Edge { from: 0, to: 1, slot: 0, from_field: 0 }], + ) + .expect("valid") + }; + let s0 = f64_stream(&[(1, 10.0), (2, 20.0), (3, 30.0), (4, 40.0)]); + let s1 = f64_stream(&[(2, 100.0), (4, 200.0)]); + + let (tx, rx) = mpsc::channel(); + let mut h = build(tx); + h.run(vec![s0.clone(), s1.clone()]); + let out: Vec<(Timestamp, Vec)> = rx.try_iter().collect(); + // holds s1=100 across t=3 and the t=4 s0-cycle; emits on every tick once warm. + assert_eq!( + out, + vec![ + (Timestamp(2), vec![Scalar::F64(120.0)]), + (Timestamp(3), vec![Scalar::F64(130.0)]), + (Timestamp(4), vec![Scalar::F64(140.0)]), + (Timestamp(4), vec![Scalar::F64(240.0)]), + ] + ); + + let (tx2, rx2) = mpsc::channel(); + let mut h2 = build(tx2); + h2.run(vec![s0, s1]); + let out2: Vec<(Timestamp, Vec)> = rx2.try_iter().collect(); + assert_eq!(out2, out); // deterministic + } +``` + +- [ ] **Step 5: Migrate `mode_b_barrier_fires_only_on_timestamp_coincidence`** + +Replace `harness.rs:599-634`: +```rust + #[test] + fn mode_b_barrier_fires_only_on_timestamp_coincidence() { + // identical wiring to mode A, but both BarrierSum inputs are Barrier(0). + let build = |tx| { + Harness::bootstrap( + vec![ + Box::new(BarrierSum { out: [Scalar::F64(0.0)] }), + Box::new(Recorder::new(&[ScalarKind::F64], Firing::Any, tx)), + ], + vec![ + SourceSpec { kind: ScalarKind::F64, targets: vec![Target { node: 0, slot: 0 }] }, + SourceSpec { kind: ScalarKind::F64, targets: vec![Target { node: 0, slot: 1 }] }, + ], + vec![Edge { from: 0, to: 1, slot: 0, from_field: 0 }], + ) + .expect("valid") + }; + let s0 = f64_stream(&[(1, 10.0), (2, 20.0), (3, 30.0), (4, 40.0)]); + let s1 = f64_stream(&[(2, 100.0), (4, 200.0)]); + + let (tx, rx) = mpsc::channel(); + let mut h = build(tx); + h.run(vec![s0.clone(), s1.clone()]); + let out: Vec<(Timestamp, Vec)> = rx.try_iter().collect(); + // records ONLY at t=2 and t=4 where both inputs share the timestamp. + assert_eq!( + out, + vec![ + (Timestamp(2), vec![Scalar::F64(120.0)]), + (Timestamp(4), vec![Scalar::F64(240.0)]), + ] + ); + + let (tx2, rx2) = mpsc::channel(); + let mut h2 = build(tx2); + h2.run(vec![s0, s1]); + let out2: Vec<(Timestamp, Vec)> = rx2.try_iter().collect(); + assert_eq!(out2, out); // deterministic + } +``` + +- [ ] **Step 6: Migrate `within_source_diamond_rejoin_barrier_fires`** + +Replace `harness.rs:636-677`: +```rust + #[test] + fn within_source_diamond_rejoin_barrier_fires() { + // One source fans out through SMA(2), SMA(4) that rejoin at a Barrier(0) + // node; node 3 = Recorder taps the barrier. Every push in a cycle carries + // that cycle's timestamp, so once both SMAs warm and emit in the same + // cycle, both barrier inputs share the timestamp and the barrier fires. + let build = |tx| { + Harness::bootstrap( + vec![ + Box::new(Sma::new(2)), + Box::new(Sma::new(4)), + Box::new(BarrierSum { out: [Scalar::F64(0.0)] }), + Box::new(Recorder::new(&[ScalarKind::F64], Firing::Any, tx)), + ], + vec![SourceSpec { + kind: ScalarKind::F64, + targets: vec![Target { node: 0, slot: 0 }, Target { node: 1, slot: 0 }], + }], + 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 }, + ], + ) + .expect("valid DAG") + }; + let prices = f64_stream(&[(1, 10.0), (2, 12.0), (3, 14.0), (4, 16.0), (5, 18.0), (6, 20.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(); + // SMA(4) warms at cycle 4; from then both paths emit each cycle at the same + // timestamp, so the barrier fires: SMA(2)+SMA(4) = 15+13, 17+15, 19+17. + assert_eq!( + out, + vec![ + (Timestamp(4), vec![Scalar::F64(28.0)]), + (Timestamp(5), vec![Scalar::F64(32.0)]), + (Timestamp(6), vec![Scalar::F64(36.0)]), + ] + ); + + let (tx2, rx2) = mpsc::channel(); + let mut h2 = build(tx2); + h2.run(vec![prices]); + let out2: Vec<(Timestamp, Vec)> = rx2.try_iter().collect(); + assert_eq!(out2, out); + } +``` + +- [ ] **Step 7: Migrate `mixed_a_and_b_or_combine_on_one_node`** + +Replace `harness.rs:679-712`: +```rust + #[test] + fn mixed_a_and_b_or_combine_on_one_node() { + // MixedSum @0 (in0,in1 barrier group 0; in2 as-of); node 1 = Recorder taps it. + let (tx, rx) = mpsc::channel(); + let mut h = Harness::bootstrap( + vec![ + Box::new(MixedSum { out: [Scalar::F64(0.0)] }), + Box::new(Recorder::new(&[ScalarKind::F64], Firing::Any, tx)), + ], + vec![ + SourceSpec { kind: ScalarKind::F64, targets: vec![Target { node: 0, slot: 0 }] }, + SourceSpec { kind: ScalarKind::F64, targets: vec![Target { node: 0, slot: 1 }] }, + SourceSpec { kind: ScalarKind::F64, targets: vec![Target { node: 0, slot: 2 }] }, + ], + vec![Edge { from: 0, to: 1, slot: 0, from_field: 0 }], + ) + .expect("valid"); + let s0 = f64_stream(&[(2, 20.0), (5, 50.0)]); // in0 (barrier) + let s1 = f64_stream(&[(2, 200.0)]); // in1 (barrier) + let s2 = f64_stream(&[(1, 1.0), (3, 3.0)]); // in2 (as-of) + + h.run(vec![s0, s1, s2]); + let out: Vec<(Timestamp, Vec)> = rx.try_iter().collect(); + // c3: barrier pair completes at t=2, holds c=1 -> 221. c4: as-of input + // ticks at t=3, holds the pair -> 223. c1 filters; c2,c5 hold (no record). + assert_eq!( + out, + vec![ + (Timestamp(2), vec![Scalar::F64(221.0)]), + (Timestamp(3), vec![Scalar::F64(223.0)]), + ] + ); + } +``` + +- [ ] **Step 8: Drop the `observe` argument from the four unchanged `bootstrap_rejects_*` tests** + +These tests call `bootstrap(...).unwrap_err()` and do not run; the only change is +removing the trailing 4th argument. + +`bootstrap_rejects_a_cycle` (`harness.rs:717-722`) — remove the `0,` at `:721`: +```rust + let err = Harness::bootstrap( + vec![Box::new(Sma::new(1)), Box::new(Sma::new(1))], + vec![], + vec![Edge { from: 0, to: 1, slot: 0, from_field: 0 }, Edge { from: 1, to: 0, slot: 0, from_field: 0 }], + ) + .unwrap_err(); +``` + +`bootstrap_rejects_a_kind_mismatch` (`harness.rs:730-735`) — remove the `0,` at `:735`: +```rust + let err = Harness::bootstrap( + vec![Box::new(Sma::new(1))], + vec![SourceSpec { kind: ScalarKind::I64, targets: vec![Target { node: 0, slot: 0 }] }], + vec![], + ) + .unwrap_err(); +``` + +`bootstrap_rejects_from_field_out_of_range` (`harness.rs:904-909`) — remove the +`0,` at `:908`: +```rust + let err = Harness::bootstrap( + vec![Box::new(Sma::new(1)), Box::new(Sma::new(1))], + vec![SourceSpec { kind: ScalarKind::F64, targets: vec![Target { node: 0, slot: 0 }] }], + vec![Edge { from: 0, to: 1, slot: 0, from_field: 9 }], + ) + .unwrap_err(); +``` + +`bootstrap_rejects_per_field_kind_mismatch` (`harness.rs:919-924`) — remove the +`1,` at `:923`: +```rust + let err = Harness::bootstrap( + vec![Box::new(TwoField { out: [Scalar::F64(0.0), Scalar::I64(0)] }), Box::new(Sma::new(1))], + vec![SourceSpec { kind: ScalarKind::F64, targets: vec![Target { node: 0, slot: 0 }] }], + vec![Edge { from: 0, to: 1, slot: 0, from_field: 1 }], + ) + .unwrap_err(); +``` + +- [ ] **Step 9: Repurpose `bootstrap_rejects_a_bad_index` (the `observe 5` trigger no longer exists)** + +Replace `harness.rs:743-754`: +```rust + #[test] + fn bootstrap_rejects_a_bad_index() { + // an edge target node (9) that does not exist -> BadIndex. (The old trigger + // — an out-of-range observe index — is gone with `observe`; BadIndex itself + // is unchanged, only the path that reaches it.) + let err = Harness::bootstrap( + vec![Box::new(Sma::new(1)), Box::new(Sma::new(1))], + vec![SourceSpec { kind: ScalarKind::F64, targets: vec![Target { node: 0, slot: 0 }] }], + vec![Edge { from: 0, to: 9, slot: 0, from_field: 0 }], + ) + .unwrap_err(); + assert_eq!(err, BootstrapError::BadIndex); + } +``` + +- [ ] **Step 10: Migrate `ohlcv_bundles_five_field_record` (two bars, all five fields)** + +Replace `harness.rs:776-816`: +```rust + #[test] + fn ohlcv_bundles_five_field_record() { + // node 0 = Ohlcv; five sources feed O/H/L/C/V; node 1 = a 5-input Recorder + // taps all five fields via five edges. The barrier fires once all five share + // the timestamp, so each bar is recorded once, on the fifth cycle of its ts. + let (tx, rx) = mpsc::channel(); + let mut h = Harness::bootstrap( + vec![ + Box::new(Ohlcv { out: [Scalar::F64(0.0); 5] }), + Box::new(Recorder::new( + &[ScalarKind::F64, ScalarKind::F64, ScalarKind::F64, ScalarKind::F64, ScalarKind::F64], + Firing::Any, + tx, + )), + ], + ohlcv_sources(), + vec![ + Edge { from: 0, to: 1, slot: 0, from_field: 0 }, // open + Edge { from: 0, to: 1, slot: 1, from_field: 1 }, // high + Edge { from: 0, to: 1, slot: 2, from_field: 2 }, // low + Edge { from: 0, to: 1, slot: 3, from_field: 3 }, // close + Edge { from: 0, to: 1, slot: 4, from_field: 4 }, // volume + ], + ) + .expect("valid"); + h.run(ohlcv_streams()); + let out: Vec<(Timestamp, Vec)> = rx.try_iter().collect(); + assert_eq!( + out, + vec![ + (Timestamp(1), vec![ + Scalar::F64(10.0), + Scalar::F64(15.0), + Scalar::F64(8.0), + Scalar::F64(12.0), + Scalar::F64(100.0), + ]), + (Timestamp(2), vec![ + Scalar::F64(20.0), + Scalar::F64(25.0), + Scalar::F64(19.0), + Scalar::F64(22.0), + Scalar::F64(200.0), + ]), + ] + ); + } +``` + +- [ ] **Step 11: Migrate `edge_binds_single_field_high_minus_low`** + +Replace `harness.rs:818-861`: +```rust + #[test] + fn edge_binds_single_field_high_minus_low() { + // [Ohlcv (0), Sub (1), Recorder (2)]; Sub binds high (field 1) and low + // (field 2) of the Ohlcv record -> high - low; the Recorder taps Sub. + // Proves from_field routes the right columns (not field 0) and the two + // bound fields are co-fresh (Sub's Any inputs both fire in the bar's cycle). + let build = |tx| { + Harness::bootstrap( + vec![ + Box::new(Ohlcv { out: [Scalar::F64(0.0); 5] }), + Box::new(Sub::new()), + Box::new(Recorder::new(&[ScalarKind::F64], Firing::Any, tx)), + ], + ohlcv_sources(), + vec![ + Edge { from: 0, to: 1, slot: 0, from_field: 1 }, // high + Edge { from: 0, to: 1, slot: 1, from_field: 2 }, // low + Edge { from: 1, to: 2, slot: 0, from_field: 0 }, // Sub -> Recorder + ], + ) + .expect("valid DAG") + }; + let (tx, rx) = mpsc::channel(); + let mut h = build(tx); + h.run(ohlcv_streams()); + let out: Vec<(Timestamp, Vec)> = rx.try_iter().collect(); + // bar1: 15 - 8 = 7; bar2: 25 - 19 = 6 (each on the bar's fifth cycle). + assert_eq!( + out, + vec![ + (Timestamp(1), vec![Scalar::F64(7.0)]), + (Timestamp(2), vec![Scalar::F64(6.0)]), + ] + ); + + let (tx2, rx2) = mpsc::channel(); + let mut h2 = build(tx2); + h2.run(ohlcv_streams()); + let out2: Vec<(Timestamp, Vec)> = rx2.try_iter().collect(); + assert_eq!(out2, out); + } +``` + +- [ ] **Step 12: Migrate `distinct_edges_read_distinct_fields`** + +Replace `harness.rs:863-898`: +```rust + #[test] + fn distinct_edges_read_distinct_fields() { + // Same Ohlcv, a different consumer: Sub binds close (field 3) and open + // (field 0) -> close - open; the Recorder taps Sub. Proves two edges on one + // record read two different fields (3 and 0, not the high/low pair above). + let (tx, rx) = mpsc::channel(); + let mut h = Harness::bootstrap( + vec![ + Box::new(Ohlcv { out: [Scalar::F64(0.0); 5] }), + Box::new(Sub::new()), + Box::new(Recorder::new(&[ScalarKind::F64], Firing::Any, tx)), + ], + ohlcv_sources(), + vec![ + Edge { from: 0, to: 1, slot: 0, from_field: 3 }, // close + Edge { from: 0, to: 1, slot: 1, from_field: 0 }, // open + Edge { from: 1, to: 2, slot: 0, from_field: 0 }, // Sub -> Recorder + ], + ) + .expect("valid DAG"); + h.run(ohlcv_streams()); + let out: Vec<(Timestamp, Vec)> = rx.try_iter().collect(); + // bar1: 12 - 10 = 2; bar2: 22 - 20 = 2. + assert_eq!( + out, + vec![ + (Timestamp(1), vec![Scalar::F64(2.0)]), + (Timestamp(2), vec![Scalar::F64(2.0)]), + ] + ); + } +``` + +- [ ] **Step 13: Verify the migrated suite compiles and is green** + +Run: `cargo test -p aura-engine` +Expected: PASS — `test result: ok. 14 passed; 0 failed` (the 14 pre-existing +tests, now on the recording API; behaviour preserved). This confirms the test +module compiles again (every `bootstrap`/`run` call site migrated). + +--- + +## Task 5: new proof tests for node-recording (aura-engine tests) + +**Files:** +- Modify: `crates/aura-engine/src/harness.rs` (`#[cfg(test)] mod tests` only) + +These are the #2 deliverable. They are additive — they compile against the +now-migrated API and the `Recorder` fixture from Task 4. + +- [ ] **Step 1: Add the `TapForward` fixture (producer-and-sink in one node)** + +Add after the `Recorder` fixture (before the first `#[test]`): +```rust + /// A node that records AND forwards: it sends `(now, value)` out of the graph + /// (sink side effect) and returns its value as a one-field output the engine + /// forwards downstream (producer). Proves the C8 "both" role. + struct TapForward { + out: [Scalar; 1], + tx: mpsc::Sender<(Timestamp, Vec)>, + } + impl Node for TapForward { + fn schema(&self) -> NodeSchema { + NodeSchema { + inputs: vec![InputSpec { kind: ScalarKind::F64, lookback: 1, firing: Firing::Any }], + output: vec![FieldSpec { name: "value", kind: ScalarKind::F64 }], + } + } + fn eval(&mut self, ctx: Ctx<'_>) -> Option<&[Scalar]> { + let w = ctx.f64_in(0); + if w.is_empty() { + return None; + } + let v = w[0]; + let _ = self.tx.send((ctx.now(), vec![Scalar::F64(v)])); // sink side effect + self.out[0] = Scalar::F64(v); + Some(&self.out) // producer output: engine forwards it + } + } +``` + +- [ ] **Step 2: Add `multi_sink_records_distinct_interior_streams` (the headline)** + +Append at the end of `#[cfg(test)] mod tests` (before the module's closing `}`): +```rust + #[test] + fn multi_sink_records_distinct_interior_streams() { + // Two recorders tap SMA(2) and SMA(4) in ONE run -> one run records many + // streams (the #2 headline). Each drained stream is individually correct. + let (tx_fast, rx_fast) = mpsc::channel(); + let (tx_slow, rx_slow) = mpsc::channel(); + let mut h = Harness::bootstrap( + vec![ + Box::new(Sma::new(2)), + Box::new(Sma::new(4)), + Box::new(Recorder::new(&[ScalarKind::F64], Firing::Any, tx_fast)), + Box::new(Recorder::new(&[ScalarKind::F64], Firing::Any, tx_slow)), + ], + vec![SourceSpec { + kind: ScalarKind::F64, + targets: vec![Target { node: 0, slot: 0 }, Target { node: 1, slot: 0 }], + }], + vec![ + Edge { from: 0, to: 2, slot: 0, from_field: 0 }, // SMA(2) -> recorder fast + Edge { from: 1, to: 3, slot: 0, from_field: 0 }, // SMA(4) -> recorder slow + ], + ) + .expect("valid DAG"); + h.run(vec![f64_stream(&[(1, 10.0), (2, 12.0), (3, 14.0), (4, 16.0), (5, 18.0)])]); + let fast: Vec<(Timestamp, Vec)> = rx_fast.try_iter().collect(); + let slow: Vec<(Timestamp, Vec)> = rx_slow.try_iter().collect(); + // SMA(2) warms at cycle 2, SMA(4) at cycle 4 — two different-rate streams. + assert_eq!( + fast, + 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)]), + ] + ); + assert_eq!( + slow, + vec![ + (Timestamp(4), vec![Scalar::F64(13.0)]), + (Timestamp(5), vec![Scalar::F64(15.0)]), + ] + ); + } +``` + +- [ ] **Step 3: Add `recorder_taps_all_fields_of_a_record` (one bar, five edges)** + +Append: +```rust + #[test] + fn recorder_taps_all_fields_of_a_record() { + // A 5-input Recorder taps all five OHLCV fields via five field-wise edges + // (0005: N edges, no whole-record bind); its recorded row is the whole bar. + let (tx, rx) = mpsc::channel(); + let mut h = Harness::bootstrap( + vec![ + Box::new(Ohlcv { out: [Scalar::F64(0.0); 5] }), + Box::new(Recorder::new( + &[ScalarKind::F64, ScalarKind::F64, ScalarKind::F64, ScalarKind::F64, ScalarKind::F64], + Firing::Any, + tx, + )), + ], + ohlcv_sources(), + vec![ + Edge { from: 0, to: 1, slot: 0, from_field: 0 }, + Edge { from: 0, to: 1, slot: 1, from_field: 1 }, + Edge { from: 0, to: 1, slot: 2, from_field: 2 }, + Edge { from: 0, to: 1, slot: 3, from_field: 3 }, + Edge { from: 0, to: 1, slot: 4, from_field: 4 }, + ], + ) + .expect("valid"); + h.run(vec![ + f64_stream(&[(1, 10.0)]), + f64_stream(&[(1, 15.0)]), + f64_stream(&[(1, 8.0)]), + f64_stream(&[(1, 12.0)]), + f64_stream(&[(1, 100.0)]), + ]); + let out: Vec<(Timestamp, Vec)> = rx.try_iter().collect(); + assert_eq!(out.len(), 1); + assert_eq!(out[0].1.len(), 5); // all five fields recorded as one row + assert_eq!( + out, + vec![(Timestamp(1), vec![ + Scalar::F64(10.0), + Scalar::F64(15.0), + Scalar::F64(8.0), + Scalar::F64(12.0), + Scalar::F64(100.0), + ])] + ); + } +``` + +- [ ] **Step 4: Add `recorder_records_mixed_scalar_kinds`** + +Append: +```rust + #[test] + fn recorder_records_mixed_scalar_kinds() { + // A recorder with i64 + f64 + bool + timestamp inputs records a four-field + // mixed-kind row -> recording is not f64-only. Four sources tick once each + // at t=1,2,3,4; only on cycle 4 are all slots warm, so it records once, + // holding the earlier-ticked values. + let (tx, rx) = mpsc::channel(); + let mut h = Harness::bootstrap( + vec![Box::new(Recorder::new( + &[ScalarKind::I64, ScalarKind::F64, ScalarKind::Bool, ScalarKind::Timestamp], + Firing::Any, + tx, + ))], + vec![ + SourceSpec { kind: ScalarKind::I64, targets: vec![Target { node: 0, slot: 0 }] }, + SourceSpec { kind: ScalarKind::F64, targets: vec![Target { node: 0, slot: 1 }] }, + SourceSpec { kind: ScalarKind::Bool, targets: vec![Target { node: 0, slot: 2 }] }, + SourceSpec { kind: ScalarKind::Timestamp, targets: vec![Target { node: 0, slot: 3 }] }, + ], + vec![], + ) + .expect("valid"); + h.run(vec![ + vec![(Timestamp(1), Scalar::I64(7))], + vec![(Timestamp(2), Scalar::F64(1.5))], + vec![(Timestamp(3), Scalar::Bool(true))], + vec![(Timestamp(4), Scalar::Ts(Timestamp(99)))], + ]); + let out: Vec<(Timestamp, Vec)> = rx.try_iter().collect(); + assert_eq!( + out, + vec![(Timestamp(4), vec![ + Scalar::I64(7), + Scalar::F64(1.5), + Scalar::Bool(true), + Scalar::Ts(Timestamp(99)), + ])] + ); + } +``` + +- [ ] **Step 5: Add `node_is_producer_and_sink_at_once`** + +Append: +```rust + #[test] + fn node_is_producer_and_sink_at_once() { + // TapForward records its input AND forwards it downstream; a second + // Recorder taps the forwarded output. Both channels see the same stream -> + // one node is producer and sink at once (C8 "both"). + let (tx_tap, rx_tap) = mpsc::channel(); + let (tx_down, rx_down) = mpsc::channel(); + let mut h = Harness::bootstrap( + vec![ + Box::new(TapForward { out: [Scalar::F64(0.0)], tx: tx_tap }), + Box::new(Recorder::new(&[ScalarKind::F64], Firing::Any, tx_down)), + ], + vec![SourceSpec { kind: ScalarKind::F64, targets: vec![Target { node: 0, slot: 0 }] }], + vec![Edge { from: 0, to: 1, slot: 0, from_field: 0 }], + ) + .expect("valid"); + h.run(vec![f64_stream(&[(1, 10.0), (2, 20.0), (3, 30.0)])]); + let tapped: Vec<(Timestamp, Vec)> = rx_tap.try_iter().collect(); + let downstream: Vec<(Timestamp, Vec)> = rx_down.try_iter().collect(); + let expected = vec![ + (Timestamp(1), vec![Scalar::F64(10.0)]), + (Timestamp(2), vec![Scalar::F64(20.0)]), + (Timestamp(3), vec![Scalar::F64(30.0)]), + ]; + assert_eq!(tapped, expected); // it recorded (sink side effect) + assert_eq!(downstream, expected); // and forwarded (producer output) + } +``` + +- [ ] **Step 6: Add `recording_is_deterministic`** + +Append: +```rust + #[test] + fn recording_is_deterministic() { + // Two fresh harnesses, two channels, identical input -> bit-identical + // recorded streams (C1). + let build = |tx| { + Harness::bootstrap( + vec![ + Box::new(Sma::new(3)), + Box::new(Recorder::new(&[ScalarKind::F64], Firing::Any, tx)), + ], + vec![SourceSpec { kind: ScalarKind::F64, targets: vec![Target { node: 0, slot: 0 }] }], + vec![Edge { from: 0, to: 1, slot: 0, from_field: 0 }], + ) + .expect("valid") + }; + let prices = f64_stream(&[(1, 1.0), (2, 2.0), (3, 3.0), (4, 4.0), (5, 5.0)]); + + let (tx_a, rx_a) = mpsc::channel(); + let mut a = build(tx_a); + a.run(vec![prices.clone()]); + let run_a: Vec<(Timestamp, Vec)> = rx_a.try_iter().collect(); + + let (tx_b, rx_b) = mpsc::channel(); + let mut b = build(tx_b); + b.run(vec![prices]); + let run_b: Vec<(Timestamp, Vec)> = rx_b.try_iter().collect(); + + assert_eq!(run_a, run_b); + assert!(!run_a.is_empty()); // and it actually recorded something + } +``` + +- [ ] **Step 7: Add the two recorder firing-mode tests** + +Append: +```rust + #[test] + fn recorder_barrier_firing_records_only_on_coincidence() { + // A 2-input Barrier(0) recorder records only on cycles where both inputs + // share the timestamp — the recorder's OWN firing policy gates recording. + let (tx, rx) = mpsc::channel(); + let mut h = Harness::bootstrap( + vec![Box::new(Recorder::new( + &[ScalarKind::F64, ScalarKind::F64], + Firing::Barrier(0), + tx, + ))], + vec![ + SourceSpec { kind: ScalarKind::F64, targets: vec![Target { node: 0, slot: 0 }] }, + SourceSpec { kind: ScalarKind::F64, targets: vec![Target { node: 0, slot: 1 }] }, + ], + vec![], + ) + .expect("valid"); + let s0 = f64_stream(&[(1, 10.0), (2, 20.0), (3, 30.0), (4, 40.0)]); + let s1 = f64_stream(&[(2, 100.0), (4, 200.0)]); + h.run(vec![s0, s1]); + let out: Vec<(Timestamp, Vec)> = rx.try_iter().collect(); + // records ONLY at t=2 and t=4 (both inputs coincide); holds otherwise. + assert_eq!( + out, + vec![ + (Timestamp(2), vec![Scalar::F64(20.0), Scalar::F64(100.0)]), + (Timestamp(4), vec![Scalar::F64(40.0), Scalar::F64(200.0)]), + ] + ); + } + + #[test] + fn recorder_any_firing_records_on_each_fresh() { + // A 2-input Any recorder records on any-fresh once both are warm (as-of), + // holding the stale input. + let (tx, rx) = mpsc::channel(); + let mut h = Harness::bootstrap( + vec![Box::new(Recorder::new( + &[ScalarKind::F64, ScalarKind::F64], + Firing::Any, + tx, + ))], + vec![ + SourceSpec { kind: ScalarKind::F64, targets: vec![Target { node: 0, slot: 0 }] }, + SourceSpec { kind: ScalarKind::F64, targets: vec![Target { node: 0, slot: 1 }] }, + ], + vec![], + ) + .expect("valid"); + let s0 = f64_stream(&[(1, 10.0), (2, 20.0), (3, 30.0), (4, 40.0)]); + let s1 = f64_stream(&[(2, 100.0), (4, 200.0)]); + h.run(vec![s0, s1]); + let out: Vec<(Timestamp, Vec)> = rx.try_iter().collect(); + // from t=2 on, records every cycle holding the stale input; two cycles fall + // on t=4 (the s0 tick then the s1 tick). + assert_eq!( + out, + vec![ + (Timestamp(2), vec![Scalar::F64(20.0), Scalar::F64(100.0)]), + (Timestamp(3), vec![Scalar::F64(30.0), Scalar::F64(100.0)]), + (Timestamp(4), vec![Scalar::F64(40.0), Scalar::F64(100.0)]), + (Timestamp(4), vec![Scalar::F64(40.0), Scalar::F64(200.0)]), + ] + ); + } +``` + +- [ ] **Step 8: Add `bootstrap_rejects_kind_mismatched_recorder_edge`** + +Append: +```rust + #[test] + fn bootstrap_rejects_kind_mismatched_recorder_edge() { + // TwoField output: field 0 f64, field 1 i64. Binding field 1 (i64) into a + // Recorder's f64 input slot is a per-field kind mismatch -> KindMismatch + // (0005's check already covers recorder edges; recording adds no new hole). + let (tx, _rx) = mpsc::channel(); + let err = Harness::bootstrap( + vec![ + Box::new(TwoField { out: [Scalar::F64(0.0), Scalar::I64(0)] }), + Box::new(Recorder::new(&[ScalarKind::F64], Firing::Any, tx)), + ], + vec![SourceSpec { kind: ScalarKind::F64, targets: vec![Target { node: 0, slot: 0 }] }], + vec![Edge { from: 0, to: 1, slot: 0, from_field: 1 }], + ) + .unwrap_err(); + assert_eq!( + err, + BootstrapError::KindMismatch { producer: ScalarKind::I64, consumer: ScalarKind::F64 } + ); + } +``` + +- [ ] **Step 9: Verify the full engine suite is green** + +Run: `cargo test -p aura-engine` +Expected: PASS — `test result: ok. 22 passed; 0 failed` (14 migrated + 8 new). + +--- + +## Task 6: ledger realization notes + final ship gate + +**Files:** +- Modify: `docs/design/INDEX.md` + +- [ ] **Step 1: Append the C8 cycle-0006 realization note** + +After the existing `**Realization (cycle 0005).**` paragraph in C8 (ends at +`INDEX.md:210`), insert: +```markdown +**Realization (cycle 0006).** The pure-consumer (sink) half of this contract is +now realized at the substrate: **recording is a node role, not a type.** A +recording node reads its typed input windows + `ctx.now()` in `eval` and pushes +the record to a destination it holds as a field (a channel, a chart handle) — an +**out-of-graph side effect**. There is no `Sink` type, trait, or engine flag: a +node that only records returns `None` (pure consumer), and a node may record +**and** return an output the engine forwards in the same `eval` (the "both" +case). In-graph routing stays engine-owned data (the edge table); the escape out +of the graph is the node's own side effect — and that boundary is the +determinism / graph-as-data boundary (C1/C7). +``` + +- [ ] **Step 2: Append the C22 cycle-0006 realization note** + +After the C22 `**Why.**` paragraph (ends at `INDEX.md:495`, before the `---` at +`:497`), insert: +```markdown +**Realization (cycle 0006).** Sinks-as-recording-mechanism is realized at the +substrate level: a recorded trace is exactly what a recording node pushed out of +the graph (no engine recording registry; the constructing World holds each +recording node's destination). The engine's single `observe: usize` affordance is +removed — `Harness::run` returns `()` and recording is a node-side concern, so one +run records *many* streams (one per recording node) instead of exactly one row. +Recorded streams are sparse and timestamped (a record per fired cycle, tagged +`ctx.now()`), matching a trace of timestamped events (C18). No new contract; the +`Harness` API change (observe removed, `run -> ()`) is recorded here. +``` + +- [ ] **Step 3: Final ship gate — full suite, clippy, purity grep** + +Run: `cargo test --workspace` +Expected: PASS — `0 failed` across all crates (aura-core 20, aura-std 3, +aura-engine 22). + +Run: `cargo clippy --workspace --all-targets -- -D warnings` +Expected: PASS — `Finished` with no warnings (compiles every test target too). + +Run: `git grep -nE 'dyn Any|Rc<|RefCell' crates/aura-engine/src; echo "exit=$status"` +Expected: no matching lines — the purity grep finds nothing in engine source (the +`Recorder` read-back uses `mpsc`, not interior mutability). Output is just +`exit=1` (fish: grep's no-match exit code), with no preceding match lines. + +Run: `git grep -nE '\bobserve\b' crates/aura-engine/src` +Expected: no matches that name the removed field — the only acceptable residue is +none (the module-doc "observer push" phrase at `harness.rs:6` uses "observer", not +"observe", and is RustAst-contrast prose left intact). If any `observe` field +reference remains, it is a missed deletion from Task 3.