feat: recording is a node role, not a type (multi-sink substrate)
Replace the engine's single observe: usize recording affordance with recording-by-node, so one run records many streams. 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 pure-consumer node returns None, and a node may record AND return a forwarded output in the same eval (the C8 "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. Engine surface shrinks: Ctx gains now: Timestamp + now() (C2-causal, the present cycle's timestamp); Harness loses the observe field, its observe >= n bootstrap check, and the per-cycle observed-row collection; bootstrap drops its 4th param; run returns (). Recorded streams are now sparse and timestamped (a record per fired cycle) instead of the dense Vec<Option<row>>. The Ctx::new signature change touched 9 call sites across three crates (not the 3 the spec estimated) — aura-std's node tests and the engine run loop were threaded too. The engine test suite migrated to a test-local Recorder fixture whose read-back is an mpsc channel, never Rc/RefCell, keeping aura-engine/src purity-clean (C7). Eight new proof tests cover the multi-sink headline, producer-and-sink, mixed-kind recording, all-fields tap, both recorder firing modes, determinism, and recorder-edge kind rejection. C8/C22 gain cycle-0006 realization notes. Gates: workspace test 45 green (core 20, std 3, engine 22), clippy -D warnings clean, purity grep clean (only a comment names Rc/RefCell). closes #2
This commit is contained in:
@@ -6,16 +6,25 @@
|
||||
use crate::{AnyColumn, Timestamp, Window};
|
||||
|
||||
/// Read-only, zero-copy view of a node's inputs for one `eval`, in schema
|
||||
/// order. `Copy` because it is just a borrow of the input slice.
|
||||
/// 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) for one `eval`.
|
||||
pub fn new(inputs: &'a [AnyColumn]) -> Self {
|
||||
Self { inputs }
|
||||
/// 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
|
||||
}
|
||||
|
||||
/// Zero-copy `f64` window into input `i` (index 0 = newest). Panics if input
|
||||
@@ -61,7 +70,7 @@ impl<'a> Ctx<'a> {
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::{Scalar, ScalarKind};
|
||||
use crate::{Scalar, ScalarKind, Timestamp};
|
||||
|
||||
#[test]
|
||||
fn ctx_hands_financial_indexed_windows() {
|
||||
@@ -69,7 +78,7 @@ mod tests {
|
||||
for v in [10.0_f64, 20.0, 30.0] {
|
||||
inputs[0].push(Scalar::F64(v)).unwrap();
|
||||
}
|
||||
let ctx = Ctx::new(&inputs);
|
||||
let ctx = Ctx::new(&inputs, Timestamp(0));
|
||||
let w = ctx.f64_in(0);
|
||||
assert_eq!(w.len(), 3);
|
||||
assert_eq!(w[0], 30.0); // newest
|
||||
@@ -84,7 +93,7 @@ mod tests {
|
||||
];
|
||||
inputs[0].push(Scalar::F64(1.5)).unwrap();
|
||||
inputs[1].push(Scalar::I64(42)).unwrap();
|
||||
let ctx = Ctx::new(&inputs);
|
||||
let ctx = Ctx::new(&inputs, Timestamp(0));
|
||||
assert_eq!(ctx.f64_in(0)[0], 1.5);
|
||||
assert_eq!(ctx.i64_in(1)[0], 42);
|
||||
}
|
||||
@@ -94,7 +103,14 @@ mod tests {
|
||||
fn ctx_panics_on_kind_mismatch() {
|
||||
let mut inputs = vec![AnyColumn::with_capacity(ScalarKind::I64, 2)];
|
||||
inputs[0].push(Scalar::I64(7)).unwrap();
|
||||
let ctx = Ctx::new(&inputs);
|
||||
let ctx = Ctx::new(&inputs, Timestamp(0));
|
||||
let _ = ctx.f64_in(0); // wrong kind → panic
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn ctx_now_returns_cycle_timestamp() {
|
||||
let inputs: Vec<AnyColumn> = vec![];
|
||||
let ctx = Ctx::new(&inputs, Timestamp(42));
|
||||
assert_eq!(ctx.now(), Timestamp(42));
|
||||
}
|
||||
}
|
||||
|
||||
+567
-179
File diff suppressed because it is too large
Load Diff
@@ -48,7 +48,7 @@ impl Node for Sma {
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use aura_core::AnyColumn;
|
||||
use aura_core::{AnyColumn, Timestamp};
|
||||
|
||||
#[test]
|
||||
fn sma_warms_up_then_tracks_the_window_mean() {
|
||||
@@ -67,7 +67,7 @@ mod tests {
|
||||
|
||||
for (v, want) in feed.iter().zip(expect) {
|
||||
inputs[0].push(Scalar::F64(*v)).unwrap();
|
||||
let got = sma.eval(Ctx::new(&inputs));
|
||||
let got = sma.eval(Ctx::new(&inputs, Timestamp(0)));
|
||||
match want {
|
||||
None => assert_eq!(got, None),
|
||||
Some(m) => assert_eq!(got, Some([Scalar::F64(m)].as_slice())),
|
||||
@@ -81,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)].as_slice()));
|
||||
assert_eq!(sma.eval(Ctx::new(&inputs, Timestamp(0))), 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)].as_slice()));
|
||||
assert_eq!(sma.eval(Ctx::new(&inputs, Timestamp(0))), Some([Scalar::F64(9.0)].as_slice()));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -49,7 +49,7 @@ impl Node for Sub {
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use aura_core::AnyColumn;
|
||||
use aura_core::{AnyColumn, Timestamp};
|
||||
|
||||
#[test]
|
||||
fn sub_is_difference_once_both_inputs_present() {
|
||||
@@ -61,10 +61,10 @@ mod tests {
|
||||
|
||||
// only input 0 present -> None
|
||||
inputs[0].push(Scalar::F64(10.0)).unwrap();
|
||||
assert_eq!(sub.eval(Ctx::new(&inputs)), None);
|
||||
assert_eq!(sub.eval(Ctx::new(&inputs, Timestamp(0))), None);
|
||||
|
||||
// both present -> a - b
|
||||
inputs[1].push(Scalar::F64(4.0)).unwrap();
|
||||
assert_eq!(sub.eval(Ctx::new(&inputs)), Some([Scalar::F64(6.0)].as_slice()));
|
||||
assert_eq!(sub.eval(Ctx::new(&inputs, Timestamp(0))), Some([Scalar::F64(6.0)].as_slice()));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -208,6 +208,16 @@ selects one producer column per edge; consuming a whole record is N edges (no
|
||||
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).
|
||||
**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).
|
||||
|
||||
### C9 — Fractal, acyclic composition
|
||||
**Guarantee.** A composite is itself a `Node` that wires a sub-graph and exposes
|
||||
@@ -493,6 +503,15 @@ records. Making sinks the one recording-and-observability mechanism keeps "what
|
||||
can I see?" answerable by "what did I instrument?", and keeps the engine
|
||||
UI-agnostic (C14). Live param tuning (runtime values, no topology change —
|
||||
C12 / C19 Fork A) gives the interactive feel without a wiring DSL.
|
||||
**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.
|
||||
|
||||
---
|
||||
|
||||
|
||||
Reference in New Issue
Block a user