1100a60c76
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
117 lines
4.0 KiB
Rust
117 lines
4.0 KiB
Rust
//! The evaluation context (C8): the read-side window access a node sees in
|
|
//! `eval`. The engine sizes and types each input from the node's `schema` at
|
|
//! wiring, so the typed accessors below treat a kind mismatch as an engine bug
|
|
//! (panic), not a user-facing error.
|
|
|
|
use crate::{AnyColumn, Timestamp, Window};
|
|
|
|
/// 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
|
|
}
|
|
|
|
/// Zero-copy `f64` window into input `i` (index 0 = newest). Panics if input
|
|
/// `i` is not an `f64` edge — a wiring bug, never reachable from a correctly
|
|
/// wired graph.
|
|
pub fn f64_in(&self, i: usize) -> Window<'a, f64> {
|
|
let inputs: &'a [AnyColumn] = self.inputs;
|
|
inputs[i]
|
|
.as_f64()
|
|
.expect("input kind mismatch (checked at wiring) — engine bug")
|
|
.window()
|
|
}
|
|
|
|
/// Zero-copy `i64` window into input `i` (index 0 = newest). See `f64_in`.
|
|
pub fn i64_in(&self, i: usize) -> Window<'a, i64> {
|
|
let inputs: &'a [AnyColumn] = self.inputs;
|
|
inputs[i]
|
|
.as_i64()
|
|
.expect("input kind mismatch (checked at wiring) — engine bug")
|
|
.window()
|
|
}
|
|
|
|
/// Zero-copy `bool` window into input `i` (index 0 = newest). See `f64_in`.
|
|
pub fn bool_in(&self, i: usize) -> Window<'a, bool> {
|
|
let inputs: &'a [AnyColumn] = self.inputs;
|
|
inputs[i]
|
|
.as_bool()
|
|
.expect("input kind mismatch (checked at wiring) — engine bug")
|
|
.window()
|
|
}
|
|
|
|
/// Zero-copy `timestamp` window into input `i` (index 0 = newest). See
|
|
/// `f64_in`.
|
|
pub fn ts_in(&self, i: usize) -> Window<'a, Timestamp> {
|
|
let inputs: &'a [AnyColumn] = self.inputs;
|
|
inputs[i]
|
|
.as_ts()
|
|
.expect("input kind mismatch (checked at wiring) — engine bug")
|
|
.window()
|
|
}
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::*;
|
|
use crate::{Scalar, ScalarKind, Timestamp};
|
|
|
|
#[test]
|
|
fn ctx_hands_financial_indexed_windows() {
|
|
let mut inputs = vec![AnyColumn::with_capacity(ScalarKind::F64, 4)];
|
|
for v in [10.0_f64, 20.0, 30.0] {
|
|
inputs[0].push(Scalar::F64(v)).unwrap();
|
|
}
|
|
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
|
|
assert_eq!(w[2], 10.0); // oldest
|
|
}
|
|
|
|
#[test]
|
|
fn ctx_addresses_multiple_inputs() {
|
|
let mut inputs = vec![
|
|
AnyColumn::with_capacity(ScalarKind::F64, 2),
|
|
AnyColumn::with_capacity(ScalarKind::I64, 2),
|
|
];
|
|
inputs[0].push(Scalar::F64(1.5)).unwrap();
|
|
inputs[1].push(Scalar::I64(42)).unwrap();
|
|
let ctx = Ctx::new(&inputs, Timestamp(0));
|
|
assert_eq!(ctx.f64_in(0)[0], 1.5);
|
|
assert_eq!(ctx.i64_in(1)[0], 42);
|
|
}
|
|
|
|
#[test]
|
|
#[should_panic(expected = "engine bug")]
|
|
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, 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));
|
|
}
|
|
}
|