From 5228542bfda00d876e2aa8c4c50b39d3fa9d16dc Mon Sep 17 00:00:00 2001 From: Brummel Date: Wed, 3 Jun 2026 15:55:20 +0200 Subject: [PATCH] feat: node output is a record (composite stream), not a single scalar Cycle 0005 (BLOCKER #1): a node's output generalizes from one scalar to a record of K >= 1 named base-scalar columns, with a scalar being the degenerate K = 1 case. A node keeps exactly one output port; its payload is now an ordered bundle of base columns (a composite stream, C7). This unblocks every multi-column producer the engine needs next -- OHLCV bars and, later, the C10 position-event table -- none of which could exist while a node emitted at most one column. Revises C8, sharpens C7 in the design ledger. Contract (aura-core): - `NodeSchema.output: ScalarKind` -> `Vec` (FieldSpec = { name: &'static str, kind }; the field position is what an Edge binds, the name is metadata for later sinks/playground, C18). - `Node::eval -> Option` -> `-> Option<&[Scalar]>`: a node fills a buffer it owns (sized once at construction) and returns a borrowed K-field row; `None` still means filter/not-warmed. This is the C7-faithful representation chosen with the user: zero per-cycle heap allocation on the forward path (rejected: Option> allocates per fire; an inline fixed [Scalar; N] bakes a width cap; engine-owned output columns invert the eval model). Scalar output is the degenerate 1-field record -- no separate scalar path. Field-wise binding (aura-engine): - `Edge` gains `from_field: usize` -- which producer column this edge forwards. Consuming a whole record is N edges; there is no "bind whole record" mechanism. - bootstrap kind-checks per field (`from.output.get(from_field)` -> BadIndex if out of range; `field.kind != slot.kind` -> KindMismatch). - the run loop copies a producer's returned row into one reused scratch buffer (resolving the borrow; no per-cycle alloc once warm) and scatters scratch[from_field] into each out-edge slot. `NodeBox.out_len` (set at bootstrap) backs a debug_assert on the returned row width. `run` returns `Vec>>` (the observed node's full row per cycle -- a materialization-surface alloc, not the inter-node hot path). - the K fields of one record are co-fresh by construction (one eval, one timestamp), so C6 is untouched. aura-std: Sma/Sub migrate to the 1-field degenerate case (hold a [Scalar; 1] buffer; behaviour unchanged). Sub drops #[derive(Default)] (Scalar has no Default) for a manual Default. Proof (aura-engine tests, +5): an Ohlcv 5-field bundler fed by five timestamp-aligned barrier sources emits one complete bar per timestamp (ohlcv_bundles_five_field_record); a downstream Sub binds high(1)-low(2) field-wise, observed end-to-end + a determinism re-run (edge_binds_single_field_high_minus_low); a second consumer binds close(3)-open(0) on the same record (distinct_edges_read_distinct_fields); must-fail: bootstrap_rejects_from_field_out_of_range (BadIndex), bootstrap_rejects_per_field_kind_mismatch (a TwoField [f64,i64] producer, the i64 field into an f64 slot). All prior 0003/0004 tests adapted to from_field + the Vec return, expected vectors unchanged (scalar = 1-field record). Gates re-run by the orchestrator (not just the agent): cargo build/test --workspace --all-targets (36: aura-core 19 + aura-std 3 + aura-engine 14, 0 failed) / clippy --workspace --all-targets -D warnings clean; surface-purity grep (no dyn-Any / Rc / RefCell) clean. Ledger C8/C7 edits verified (Option<&[Scalar]> present, "one series per node" gone, cycle-0005 realization note added). The glossary composite/node record-reality pass is the audit-time follow-up. closes #1 refs walking-skeleton Co-Authored-By: Claude Opus 4.8 (1M context) --- crates/aura-core/src/lib.rs | 4 +- crates/aura-core/src/node.rs | 31 ++- crates/aura-engine/src/harness.rs | 360 ++++++++++++++++++++++++++---- crates/aura-std/src/sma.rs | 22 +- crates/aura-std/src/sub.rs | 24 +- docs/design/INDEX.md | 25 ++- 6 files changed, 384 insertions(+), 82 deletions(-) diff --git a/crates/aura-core/src/lib.rs b/crates/aura-core/src/lib.rs index cb6b9b6..dd79add 100644 --- a/crates/aura-core/src/lib.rs +++ b/crates/aura-core/src/lib.rs @@ -18,7 +18,7 @@ //! //! - [`Node`] — the `schema`/`eval` contract every node implements (C8), with //! [`NodeSchema`] / [`InputSpec`] declaring inputs (kind + lookback) and the -//! single output kind; +//! output record ([`FieldSpec`] columns; length 1 = scalar); //! - [`Ctx`] — the per-`eval` read-side: zero-copy, financial-indexed [`Window`] //! access into each input (closing the cycle-0001 read-side gap on //! [`AnyColumn`]). @@ -39,5 +39,5 @@ pub use any::AnyColumn; pub use column::{Column, Window}; pub use ctx::Ctx; pub use error::KindMismatch; -pub use node::{Firing, InputSpec, Node, NodeSchema}; +pub use node::{FieldSpec, Firing, InputSpec, Node, NodeSchema}; pub use scalar::{Scalar, ScalarKind, Timestamp}; diff --git a/crates/aura-core/src/node.rs b/crates/aura-core/src/node.rs index f723d87..3b89497 100644 --- a/crates/aura-core/src/node.rs +++ b/crates/aura-core/src/node.rs @@ -1,6 +1,7 @@ //! The node contract (C8): the interface every node implements. A node declares //! its inputs (each with its scalar kind, lookback depth, and firing policy, C6) -//! and its output kind via `schema`, and computes one cycle's output via `eval`. +//! and its output record (1..K base columns, C7) via `schema`, and computes one +//! cycle's row via `eval`. //! Tunable params (C12/C19) are deliberately not part of the schema yet — see //! spec 0002's "Out of scope". @@ -29,21 +30,33 @@ pub struct InputSpec { pub firing: Firing, } -/// A node's declared interface: its inputs (in order) and its single output -/// kind. Built once at wiring, never on the hot path — the `Vec` is fine here. +/// One declared output column of a node's record: its name (metadata for sinks / +/// the playground, C18) and its scalar kind. The position of a `FieldSpec` in +/// `NodeSchema.output` is what an `Edge` binds (`Edge::from_field`); the name is +/// not load-bearing for wiring. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct FieldSpec { + pub name: &'static str, + pub kind: ScalarKind, +} + +/// A node's declared interface: its inputs (in order) and its output record — an +/// ordered list of named base columns; length 1 is a scalar (the degenerate +/// case). Built once at wiring, never on the hot path — the `Vec`s are fine here. #[derive(Clone, Debug, PartialEq, Eq)] pub struct NodeSchema { pub inputs: Vec, - pub output: ScalarKind, + pub output: Vec, } -/// The universal composable dataflow unit (C8): at most one output, a producer -/// or transformer. `schema` declares the interface; `eval` computes one cycle's -/// output (`None` = filter / not-yet-warmed-up). `&mut self` because a node may -/// keep its own derived state. +/// The universal composable dataflow unit (C8): one output port carrying a record +/// of 1..K base columns, a producer or transformer. `schema` declares the +/// interface; `eval` computes one cycle's row, returning a borrowed slice into a +/// buffer the node owns (`None` = filter / not-yet-warmed-up). `&mut self` because +/// a node may keep its own derived state (and its output buffer). pub trait Node { fn schema(&self) -> NodeSchema; - fn eval(&mut self, ctx: Ctx<'_>) -> Option; + fn eval(&mut self, ctx: Ctx<'_>) -> Option<&[Scalar]>; } #[cfg(test)] diff --git a/crates/aura-engine/src/harness.rs b/crates/aura-engine/src/harness.rs index e2b5814..25e16fc 100644 --- a/crates/aura-engine/src/harness.rs +++ b/crates/aura-engine/src/harness.rs @@ -5,7 +5,9 @@ //! the flat, monomorphized sharpening of RustAst's reference-counted, //! interior-mutable observer push graph: no reference counting, no interior //! mutability, no per-cycle allocation (C1/C7). Each node owns its input columns -//! (the cycle-0002 shape), so `Ctx` is unchanged. +//! (the cycle-0002 shape), so `Ctx` is unchanged. A node's `eval` returns a +//! borrowed record (`Option<&[Scalar]>`); each out-edge forwards one field of it +//! (`Edge::from_field`) into a consumer slot, so the K fields are co-fresh (C6). //! //! Firing (C5/C6) gates re-evaluation. Two read clocks drive it, both stamped per //! input slot on every push: `fresh_at` (the `cycle_id` of the last push — @@ -20,12 +22,15 @@ use aura_core::{AnyColumn, Ctx, Firing, Node, Scalar, ScalarKind, Timestamp}; -/// A producer-output -> consumer-input-slot forwarding edge. +/// Forwards one field (`from_field`) of a producer's output record into a +/// consumer's input slot. Consuming a whole record is N such edges, one per field +/// (there is no "bind whole record" mechanism). #[derive(Clone, Copy, Debug, PartialEq, Eq)] pub struct Edge { pub from: usize, pub to: usize, pub slot: usize, + pub from_field: usize, } /// An input slot the source value is forwarded into each cycle. @@ -74,6 +79,7 @@ struct NodeBox { inputs: Vec, firing: Vec, slots: Vec, + out_len: usize, } /// A bootstrapped, frozen root graph instance plus its deterministic run loop. @@ -133,7 +139,8 @@ impl Harness { .iter() .map(|_| SlotState { fresh_at: 0, last_ts: Timestamp(i64::MIN) }) .collect(); - boxes.push(NodeBox { node: nd, inputs, firing, slots }); + let out_len = schema.output.len(); + boxes.push(NodeBox { node: nd, inputs, firing, slots, out_len }); } // source targets: each source's value must match each of its target slots' kind @@ -150,15 +157,16 @@ impl Harness { } } - // edges: indices in range, producer output kind == consumer slot kind + // edges: indices in range, producer output field kind == consumer slot kind let mut out_edges: Vec> = vec![Vec::new(); n]; for &e in &edges { let from = schemas.get(e.from).ok_or(BootstrapError::BadIndex)?; let to = schemas.get(e.to).ok_or(BootstrapError::BadIndex)?; + let field = from.output.get(e.from_field).ok_or(BootstrapError::BadIndex)?; let slot = to.inputs.get(e.slot).ok_or(BootstrapError::BadIndex)?; - if from.output != slot.kind { + if field.kind != slot.kind { return Err(BootstrapError::KindMismatch { - producer: from.output, + producer: field.kind, consumer: slot.kind, }); } @@ -202,7 +210,7 @@ impl Harness { /// and produced output, `None` when it held or filtered). One stream per /// source, each ascending in timestamp (C3 ingestion precondition). Allocates /// nothing per cycle beyond the output vector. - pub fn run(&mut self, streams: Vec>) -> Vec> { + pub fn run(&mut self, streams: Vec>) -> Vec>> { assert_eq!( streams.len(), self.sources.len(), @@ -219,6 +227,7 @@ impl Harness { let mut cursor: Vec = vec![0; streams.len()]; let mut cycle_id: u64 = 0; let mut out = Vec::new(); + let mut scratch: Vec = Vec::new(); loop { // pick the live source head with the smallest (timestamp, source index) @@ -251,8 +260,9 @@ impl Harness { } // evaluate in topological order; gate by firing; forward Some outputs - let mut observed = None; + let mut observed: Option> = None; 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) @@ -260,17 +270,22 @@ impl Harness { if !fired { continue; // hold: no eval, no push } - let result = { + let result: Option<&[Scalar]> = { let nb = &mut nodes[nidx]; nb.node.eval(Ctx::new(&nb.inputs)) }; if nidx == observe { - observed = result; + observed = result.map(|row| row.to_vec()); } - if let Some(v) = result { + 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(v).expect("edge kind checked at wiring"); + 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 }; } } @@ -334,7 +349,7 @@ mod tests { // InputSpec / NodeSchema are not imported by harness.rs production code (it // only reads `nd.schema()` fields, never naming the types), so they are not // brought in by `use super::*` — the fixtures construct them, so import here. - use aura_core::{InputSpec, NodeSchema}; + use aura_core::{FieldSpec, InputSpec, NodeSchema}; use aura_std::{Sma, Sub}; /// Build an f64 source stream from (timestamp, value) points. @@ -347,7 +362,9 @@ mod tests { /// Mode A as-of join: a 2-input f64 sum that fires whenever either input is /// fresh, holding the other. Warm-up returns None until both have a value. - struct AsOfSum; + struct AsOfSum { + out: [Scalar; 1], + } impl Node for AsOfSum { fn schema(&self) -> NodeSchema { NodeSchema { @@ -355,22 +372,25 @@ mod tests { InputSpec { kind: ScalarKind::F64, lookback: 1, firing: Firing::Any }, InputSpec { kind: ScalarKind::F64, lookback: 1, firing: Firing::Any }, ], - output: ScalarKind::F64, + output: vec![FieldSpec { name: "value", kind: ScalarKind::F64 }], } } - fn eval(&mut self, ctx: Ctx<'_>) -> Option { + fn eval(&mut self, ctx: Ctx<'_>) -> Option<&[Scalar]> { let a = ctx.f64_in(0); let b = ctx.f64_in(1); if a.is_empty() || b.is_empty() { return None; } - Some(Scalar::F64(a[0] + b[0])) + self.out[0] = Scalar::F64(a[0] + b[0]); + Some(&self.out) } } /// Mode B barrier join: a 2-input f64 sum that fires only when both inputs /// share the current cycle timestamp (both warm by construction when it fires). - struct BarrierSum; + struct BarrierSum { + out: [Scalar; 1], + } impl Node for BarrierSum { fn schema(&self) -> NodeSchema { NodeSchema { @@ -378,18 +398,21 @@ mod tests { InputSpec { kind: ScalarKind::F64, lookback: 1, firing: Firing::Barrier(0) }, InputSpec { kind: ScalarKind::F64, lookback: 1, firing: Firing::Barrier(0) }, ], - output: ScalarKind::F64, + output: vec![FieldSpec { name: "value", kind: ScalarKind::F64 }], } } - fn eval(&mut self, ctx: Ctx<'_>) -> Option { - Some(Scalar::F64(ctx.f64_in(0)[0] + ctx.f64_in(1)[0])) + fn eval(&mut self, ctx: Ctx<'_>) -> Option<&[Scalar]> { + self.out[0] = Scalar::F64(ctx.f64_in(0)[0] + ctx.f64_in(1)[0]); + Some(&self.out) } } /// Mixed A+B: barrier pair (inputs 0,1 in group 0) plus an as-of input /// (input 2). Fires when the pair completes (holding input 2) OR when input 2 /// ticks (holding the pair) — the OR-combine. - struct MixedSum; + struct MixedSum { + out: [Scalar; 1], + } impl Node for MixedSum { fn schema(&self) -> NodeSchema { NodeSchema { @@ -398,17 +421,80 @@ mod tests { InputSpec { kind: ScalarKind::F64, lookback: 1, firing: Firing::Barrier(0) }, InputSpec { kind: ScalarKind::F64, lookback: 1, firing: Firing::Any }, ], - output: ScalarKind::F64, + output: vec![FieldSpec { name: "value", kind: ScalarKind::F64 }], } } - fn eval(&mut self, ctx: Ctx<'_>) -> Option { + fn eval(&mut self, ctx: Ctx<'_>) -> Option<&[Scalar]> { let a = ctx.f64_in(0); let b = ctx.f64_in(1); let c = ctx.f64_in(2); if a.is_empty() || b.is_empty() || c.is_empty() { return None; } - Some(Scalar::F64(a[0] + b[0] + c[0])) + self.out[0] = Scalar::F64(a[0] + b[0] + c[0]); + Some(&self.out) + } + } + + /// A neutral multi-field producer: five f64 inputs bundled into one 5-field + /// record. No trading-domain logic — it proves the K > 1 output mechanism in + /// isolation. The five inputs are a Barrier(0) group, so the node emits one + /// complete bar only when all five share the cycle's timestamp. + struct Ohlcv { + out: [Scalar; 5], + } + impl Node for Ohlcv { + fn schema(&self) -> NodeSchema { + NodeSchema { + inputs: vec![ + InputSpec { kind: ScalarKind::F64, lookback: 1, firing: Firing::Barrier(0) }, + InputSpec { kind: ScalarKind::F64, lookback: 1, firing: Firing::Barrier(0) }, + InputSpec { kind: ScalarKind::F64, lookback: 1, firing: Firing::Barrier(0) }, + InputSpec { kind: ScalarKind::F64, lookback: 1, firing: Firing::Barrier(0) }, + InputSpec { kind: ScalarKind::F64, lookback: 1, firing: Firing::Barrier(0) }, + ], + output: vec![ + FieldSpec { name: "open", kind: ScalarKind::F64 }, + FieldSpec { name: "high", kind: ScalarKind::F64 }, + FieldSpec { name: "low", kind: ScalarKind::F64 }, + FieldSpec { name: "close", kind: ScalarKind::F64 }, + FieldSpec { name: "volume", kind: ScalarKind::F64 }, + ], + } + } + fn eval(&mut self, ctx: Ctx<'_>) -> Option<&[Scalar]> { + for i in 0..5 { + let w = ctx.f64_in(i); + if w.is_empty() { + return None; // not yet warmed + } + self.out[i] = Scalar::F64(w[0]); + } + Some(&self.out) // one 5-field record, all fields co-fresh + } + } + + /// A producer whose output record mixes kinds: field 0 is f64, field 1 is i64. + /// Used only to prove the bootstrap kind check is per-field (field 0 would bind + /// into an f64 slot; field 1 would not). Its `eval` never runs in these tests — + /// bootstrap rejects the wiring first. + struct TwoField { + out: [Scalar; 2], + } + impl Node for TwoField { + fn schema(&self) -> NodeSchema { + NodeSchema { + inputs: vec![InputSpec { kind: ScalarKind::F64, lookback: 1, firing: Firing::Any }], + output: vec![ + FieldSpec { name: "f", kind: ScalarKind::F64 }, + FieldSpec { name: "i", kind: ScalarKind::I64 }, + ], + } + } + fn eval(&mut self, _ctx: Ctx<'_>) -> Option<&[Scalar]> { + self.out[0] = Scalar::F64(0.0); + self.out[1] = Scalar::I64(0); + Some(&self.out) } } @@ -428,9 +514,9 @@ mod tests { vec![ None, None, - Some(Scalar::F64(2.0)), - Some(Scalar::F64(3.0)), - Some(Scalar::F64(4.0)), + Some(vec![Scalar::F64(2.0)]), + Some(vec![Scalar::F64(3.0)]), + Some(vec![Scalar::F64(4.0)]), ] ); } @@ -446,7 +532,7 @@ mod tests { kind: ScalarKind::F64, targets: vec![Target { node: 0, slot: 0 }, Target { node: 1, slot: 0 }], }], - vec![Edge { from: 0, to: 2, slot: 0 }, Edge { from: 1, to: 2, slot: 1 }], + vec![Edge { from: 0, to: 2, slot: 0, from_field: 0 }, Edge { from: 1, to: 2, slot: 1, from_field: 0 }], 2, ) .expect("valid DAG") @@ -462,9 +548,9 @@ mod tests { None, None, None, - Some(Scalar::F64(2.0)), - Some(Scalar::F64(2.0)), - Some(Scalar::F64(2.0)), + Some(vec![Scalar::F64(2.0)]), + Some(vec![Scalar::F64(2.0)]), + Some(vec![Scalar::F64(2.0)]), ] ); @@ -478,7 +564,7 @@ mod tests { // source 0 ticks t=1,2,3,4; source 1 ticks t=2,4 (slower); both inputs Any. let build = || { Harness::bootstrap( - vec![Box::new(AsOfSum)], + vec![Box::new(AsOfSum { out: [Scalar::F64(0.0)] })], vec![ SourceSpec { kind: ScalarKind::F64, targets: vec![Target { node: 0, slot: 0 }] }, SourceSpec { kind: ScalarKind::F64, targets: vec![Target { node: 0, slot: 1 }] }, @@ -499,10 +585,10 @@ mod tests { vec![ None, None, - Some(Scalar::F64(120.0)), - Some(Scalar::F64(130.0)), - Some(Scalar::F64(140.0)), - Some(Scalar::F64(240.0)), + Some(vec![Scalar::F64(120.0)]), + Some(vec![Scalar::F64(130.0)]), + Some(vec![Scalar::F64(140.0)]), + Some(vec![Scalar::F64(240.0)]), ] ); @@ -515,7 +601,7 @@ mod tests { // identical wiring to mode A, but both inputs are Barrier(0). let build = || { Harness::bootstrap( - vec![Box::new(BarrierSum)], + vec![Box::new(BarrierSum { out: [Scalar::F64(0.0)] })], vec![ SourceSpec { kind: ScalarKind::F64, targets: vec![Target { node: 0, slot: 0 }] }, SourceSpec { kind: ScalarKind::F64, targets: vec![Target { node: 0, slot: 1 }] }, @@ -536,10 +622,10 @@ mod tests { vec![ None, None, - Some(Scalar::F64(120.0)), + Some(vec![Scalar::F64(120.0)]), None, None, - Some(Scalar::F64(240.0)), + Some(vec![Scalar::F64(240.0)]), ] ); @@ -557,12 +643,12 @@ mod tests { // distinct from the multi-source barrier above. let build = || { Harness::bootstrap( - vec![Box::new(Sma::new(2)), Box::new(Sma::new(4)), Box::new(BarrierSum)], + vec![Box::new(Sma::new(2)), Box::new(Sma::new(4)), Box::new(BarrierSum { out: [Scalar::F64(0.0)] })], 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 }, Edge { from: 1, to: 2, slot: 1 }], + vec![Edge { from: 0, to: 2, slot: 0, from_field: 0 }, Edge { from: 1, to: 2, slot: 1, from_field: 0 }], 2, ) .expect("valid DAG") @@ -579,9 +665,9 @@ mod tests { None, None, None, - Some(Scalar::F64(28.0)), - Some(Scalar::F64(32.0)), - Some(Scalar::F64(36.0)), + Some(vec![Scalar::F64(28.0)]), + Some(vec![Scalar::F64(32.0)]), + Some(vec![Scalar::F64(36.0)]), ] ); @@ -594,7 +680,7 @@ mod tests { fn mixed_a_and_b_or_combine_on_one_node() { // in0,in1 = barrier group 0 (sources 0,1); in2 = as-of (source 2). let mut h = Harness::bootstrap( - vec![Box::new(MixedSum)], + vec![Box::new(MixedSum { out: [Scalar::F64(0.0)] })], vec![ SourceSpec { kind: ScalarKind::F64, targets: vec![Target { node: 0, slot: 0 }] }, SourceSpec { kind: ScalarKind::F64, targets: vec![Target { node: 0, slot: 1 }] }, @@ -618,8 +704,8 @@ mod tests { vec![ None, None, - Some(Scalar::F64(221.0)), - Some(Scalar::F64(223.0)), + Some(vec![Scalar::F64(221.0)]), + Some(vec![Scalar::F64(223.0)]), None, ] ); @@ -631,7 +717,7 @@ mod tests { let err = Harness::bootstrap( vec![Box::new(Sma::new(1)), Box::new(Sma::new(1))], vec![], - vec![Edge { from: 0, to: 1, slot: 0 }, Edge { from: 1, to: 0, slot: 0 }], + vec![Edge { from: 0, to: 1, slot: 0, from_field: 0 }, Edge { from: 1, to: 0, slot: 0, from_field: 0 }], 0, ) .unwrap_err(); @@ -666,4 +752,180 @@ mod tests { .unwrap_err(); assert_eq!(err, BootstrapError::BadIndex); } + + /// Build five timestamp-aligned f64 sources feeding Ohlcv's five barrier slots. + fn ohlcv_streams() -> Vec> { + vec![ + f64_stream(&[(1, 10.0), (2, 20.0)]), // open + f64_stream(&[(1, 15.0), (2, 25.0)]), // high + f64_stream(&[(1, 8.0), (2, 19.0)]), // low + f64_stream(&[(1, 12.0), (2, 22.0)]), // close + f64_stream(&[(1, 100.0), (2, 200.0)]), // volume + ] + } + + fn ohlcv_sources() -> Vec { + (0..5) + .map(|slot| SourceSpec { + kind: ScalarKind::F64, + targets: vec![Target { node: 0, slot }], + }) + .collect() + } + + #[test] + fn ohlcv_bundles_five_field_record() { + // node 0 = Ohlcv; five sources feed O/H/L/C/V; observe node 0. The barrier + // fires once all five share the timestamp, so each bar appears on the fifth + // cycle of its timestamp (the four partial cycles hold -> None). + let mut h = Harness::bootstrap( + vec![Box::new(Ohlcv { out: [Scalar::F64(0.0); 5] })], + ohlcv_sources(), + vec![], + 0, + ) + .expect("valid"); + let out = h.run(ohlcv_streams()); + assert_eq!( + out, + vec![ + None, + None, + None, + None, + Some(vec![ + Scalar::F64(10.0), + Scalar::F64(15.0), + Scalar::F64(8.0), + Scalar::F64(12.0), + Scalar::F64(100.0), + ]), + None, + None, + None, + None, + Some(vec![ + Scalar::F64(20.0), + Scalar::F64(25.0), + Scalar::F64(19.0), + Scalar::F64(22.0), + Scalar::F64(200.0), + ]), + ] + ); + } + + #[test] + fn edge_binds_single_field_high_minus_low() { + // nodes [Ohlcv (0), Sub (1)]; Sub binds field 1 (high) and field 2 (low) + // of the Ohlcv record -> high - low == the bar range. Observing Sub proves + // from_field routes the right columns (not field 0), and that the two bound + // fields are co-fresh (Sub's Any inputs both fire in the bar's cycle). + let build = || { + Harness::bootstrap( + vec![ + Box::new(Ohlcv { out: [Scalar::F64(0.0); 5] }), + Box::new(Sub::new()), + ], + 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 + ], + 1, + ) + .expect("valid DAG") + }; + let mut h = build(); + let out = h.run(ohlcv_streams()); + // bar1: 15 - 8 = 7; bar2: 25 - 19 = 6 (each on the bar's fifth cycle). + assert_eq!( + out, + vec![ + None, + None, + None, + None, + Some(vec![Scalar::F64(7.0)]), + None, + None, + None, + None, + Some(vec![Scalar::F64(6.0)]), + ] + ); + + // determinism (C1): a second identical run is bit-identical + let mut h2 = build(); + assert_eq!(h2.run(ohlcv_streams()), out); + } + + #[test] + fn distinct_edges_read_distinct_fields() { + // Same Ohlcv, a different consumer: Sub binds field 3 (close) and field 0 + // (open) -> close - open. Proves two different edges on one record read two + // different fields (3 and 0, neither of them the high/low pair above). + let mut h = Harness::bootstrap( + vec![ + Box::new(Ohlcv { out: [Scalar::F64(0.0); 5] }), + Box::new(Sub::new()), + ], + 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 + ], + 1, + ) + .expect("valid DAG"); + let out = h.run(ohlcv_streams()); + // bar1: 12 - 10 = 2; bar2: 22 - 20 = 2. + assert_eq!( + out, + vec![ + None, + None, + None, + None, + Some(vec![Scalar::F64(2.0)]), + None, + None, + None, + None, + Some(vec![Scalar::F64(2.0)]), + ] + ); + } + + #[test] + fn bootstrap_rejects_from_field_out_of_range() { + // Sma(0) has a 1-field output (index 0 only); an edge reading field 9 is + // out of range -> BadIndex (caught before any kind check). + 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 }], + 0, + ) + .unwrap_err(); + assert_eq!(err, BootstrapError::BadIndex); + } + + #[test] + fn bootstrap_rejects_per_field_kind_mismatch() { + // TwoField(0) output: field 0 f64, field 1 i64. Binding field 1 (i64) into + // Sma(1)'s f64 input slot is a per-field kind mismatch -> KindMismatch. (The + // mismatch is field-specific: from_field 0 would have matched.) + 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 }], + 1, + ) + .unwrap_err(); + assert_eq!( + err, + BootstrapError::KindMismatch { producer: ScalarKind::I64, consumer: ScalarKind::F64 } + ); + } } diff --git a/crates/aura-std/src/sma.rs b/crates/aura-std/src/sma.rs index 8a98f8f..b7ed80f 100644 --- a/crates/aura-std/src/sma.rs +++ b/crates/aura-std/src/sma.rs @@ -3,18 +3,19 @@ //! `Node` contract is authorable from a downstream crate and evaluable with no //! engine present (the test drives it by hand, as the sim loop later will). -use aura_core::{Ctx, Firing, InputSpec, Node, NodeSchema, Scalar, ScalarKind}; +use aura_core::{Ctx, FieldSpec, Firing, InputSpec, Node, NodeSchema, Scalar, ScalarKind}; /// Simple moving average over the last `length` values of one f64 input. pub struct Sma { length: usize, + out: [Scalar; 1], } impl Sma { /// Build an SMA of window `length` (must be >= 1). pub fn new(length: usize) -> Self { assert!(length >= 1, "SMA length must be >= 1"); - Self { length } + Self { length, out: [Scalar::F64(0.0)] } } } @@ -26,11 +27,11 @@ impl Node for Sma { lookback: self.length, firing: Firing::Any, }], - output: ScalarKind::F64, + output: vec![FieldSpec { name: "value", kind: ScalarKind::F64 }], } } - fn eval(&mut self, ctx: Ctx<'_>) -> Option { + fn eval(&mut self, ctx: Ctx<'_>) -> Option<&[Scalar]> { let w = ctx.f64_in(0); if w.len() < self.length { return None; // not yet warmed up @@ -39,7 +40,8 @@ impl Node for Sma { for k in 0..self.length { sum += w[k]; // index 0 = newest (financial indexing) } - Some(Scalar::F64(sum / self.length as f64)) + self.out[0] = Scalar::F64(sum / self.length as f64); + Some(&self.out) } } @@ -65,7 +67,11 @@ mod tests { for (v, want) in feed.iter().zip(expect) { inputs[0].push(Scalar::F64(*v)).unwrap(); - assert_eq!(sma.eval(Ctx::new(&inputs)), want.map(Scalar::F64)); + let got = sma.eval(Ctx::new(&inputs)); + match want { + None => assert_eq!(got, None), + Some(m) => assert_eq!(got, Some([Scalar::F64(m)].as_slice())), + } } } @@ -75,9 +81,9 @@ mod tests { let mut inputs = vec![AnyColumn::with_capacity(ScalarKind::F64, 1)]; inputs[0].push(Scalar::F64(7.0)).unwrap(); - assert_eq!(sma.eval(Ctx::new(&inputs)), Some(Scalar::F64(7.0))); + assert_eq!(sma.eval(Ctx::new(&inputs)), Some([Scalar::F64(7.0)].as_slice())); inputs[0].push(Scalar::F64(9.0)).unwrap(); - assert_eq!(sma.eval(Ctx::new(&inputs)), Some(Scalar::F64(9.0))); + assert_eq!(sma.eval(Ctx::new(&inputs)), Some([Scalar::F64(9.0)].as_slice())); } } diff --git a/crates/aura-std/src/sub.rs b/crates/aura-std/src/sub.rs index b4135cd..c013e58 100644 --- a/crates/aura-std/src/sub.rs +++ b/crates/aura-std/src/sub.rs @@ -3,17 +3,24 @@ //! real fan-out + join to run (two SMAs joining into one node), exercising //! multi-input `Ctx` access inside a running graph. -use aura_core::{Ctx, Firing, InputSpec, Node, NodeSchema, Scalar, ScalarKind}; +use aura_core::{Ctx, FieldSpec, Firing, InputSpec, Node, NodeSchema, Scalar, ScalarKind}; /// Two-input f64 difference: input 0 minus input 1. Emits `None` until both /// inputs have a value. -#[derive(Default)] -pub struct Sub; +pub struct Sub { + out: [Scalar; 1], +} impl Sub { /// Build a `Sub` node. pub fn new() -> Self { - Self + Self { out: [Scalar::F64(0.0)] } + } +} + +impl Default for Sub { + fn default() -> Self { + Self::new() } } @@ -24,17 +31,18 @@ impl Node for Sub { InputSpec { kind: ScalarKind::F64, lookback: 1, firing: Firing::Any }, InputSpec { kind: ScalarKind::F64, lookback: 1, firing: Firing::Any }, ], - output: ScalarKind::F64, + output: vec![FieldSpec { name: "value", kind: ScalarKind::F64 }], } } - fn eval(&mut self, ctx: Ctx<'_>) -> Option { + fn eval(&mut self, ctx: Ctx<'_>) -> Option<&[Scalar]> { let a = ctx.f64_in(0); let b = ctx.f64_in(1); if a.is_empty() || b.is_empty() { return None; } - Some(Scalar::F64(a[0] - b[0])) + self.out[0] = Scalar::F64(a[0] - b[0]); + Some(&self.out) } } @@ -57,6 +65,6 @@ mod tests { // both present -> a - b inputs[1].push(Scalar::F64(4.0)).unwrap(); - assert_eq!(sub.eval(Ctx::new(&inputs)), Some(Scalar::F64(6.0))); + assert_eq!(sub.eval(Ctx::new(&inputs)), Some([Scalar::F64(6.0)].as_slice())); } } diff --git a/docs/design/INDEX.md b/docs/design/INDEX.md index b0775e2..41ea1f0 100644 --- a/docs/design/INDEX.md +++ b/docs/design/INDEX.md @@ -168,7 +168,10 @@ in a cycle carries that cycle's timestamp). ### C7 — Four scalar base types, streamed as SoA **Guarantee.** Only `i64`, `f64`, `bool`, `timestamp` (newtype over i64, epoch-ns UTC) are streamed, as columnar Structure-of-Arrays. Composite streams -(OHLCV) are bundles of base columns. Edges are type-erased to these four kinds; +(OHLCV) are bundles of base columns — this is the **node-output model** too: a +node emits a record of 1..K base columns (C8), each forwarded field-wise to a +consumer slot; the bundle is structural, never a fifth scalar type. Edges are +type-erased to these four kinds; the type check is paid once at wiring/sim-start, then the topology is frozen per sim → direct dispatch, no per-event allocation. **Forbids.** Streaming non-scalars (String, Records, tables, calendars) — those @@ -182,19 +185,29 @@ Type-erasure at the edge is also forced by the cdylib boundary (C13). **Guarantee.** A node implements `schema()` (declares each input's scalar type, required lookback depth, and firing group, **and the node's own tunable parameters — typed, with ranges**, which aggregate into the blueprint's -param-space the optimizer sweeps, C12/C19/C20) + `eval(ctx) -> Option`. The +param-space the optimizer sweeps, C12/C19/C20) + `eval(ctx) -> Option<&[Scalar]>`. The engine provides read-only, zero-copy windows into each input's SoA ring buffer (`ctx.f64_in(x)[k]`, sized at wiring); a node may *additionally* keep its own mutable series for derived/intermediate state. `None`/Void return = filter / not-yet-warmed-up. A node is a **producer, a consumer, or both**: a -producer/transformer exposes **at most one** output (one series per node); a -**pure consumer (sink)** — chart, equity, logger — has **no** output. Sources -are pure producers; sinks are pure consumers. +producer/transformer exposes **one output port**, whose payload is a **record of +1..K base-scalar columns** (a scalar is the degenerate K=1 record; an `eval` +returns a borrowed row, one value per column); a **pure consumer (sink)** — chart, +equity, logger — has **no** output. Sources are pure producers; sinks are pure +consumers. **Forbids.** A node sizing/growing its input lookback at runtime; more than one -output per node (model as multiple nodes); copy-on-read of input history. +output **port** per node; a fifth scalar type or a heterogeneous output payload +(a record is a bundle of base columns, C7); copy-on-read of input history. **Why.** Engine-provided windows mean LLM-authored code cannot mis-manage lookback bookkeeping, and history passes through zero-copy. Fixed, pre-sized buffers suit deterministic, pre-dimensioned sims (no realloc in the hot loop). +**Realization (cycle 0005).** `NodeSchema.output` is a `Vec` (named base +columns; length 1 = scalar). Binding is **field-wise only**: `Edge::from_field` +selects one producer column per edge; consuming a whole record is N edges (no +"bind whole record" mechanism). The K fields of one record are **co-fresh by +construction** (one `eval`, one timestamp), so C6 is untouched. `eval` returns +`Option<&[Scalar]>` — a borrowed row into a node-owned buffer — so the forward +path allocates nothing per cycle (C7). ### C9 — Fractal, acyclic composition **Guarantee.** A composite is itself a `Node` that wires a sub-graph and exposes