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
90 lines
2.9 KiB
Rust
90 lines
2.9 KiB
Rust
//! `Sma` — simple moving average over the last `length` values of one f64
|
|
//! input. The walking skeleton's first worked node: it proves the `aura-core`
|
|
//! `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, 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, out: [Scalar::F64(0.0)] }
|
|
}
|
|
}
|
|
|
|
impl Node for Sma {
|
|
fn schema(&self) -> NodeSchema {
|
|
NodeSchema {
|
|
inputs: vec![InputSpec {
|
|
kind: ScalarKind::F64,
|
|
lookback: self.length,
|
|
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.len() < self.length {
|
|
return None; // not yet warmed up
|
|
}
|
|
let mut sum = 0.0;
|
|
for k in 0..self.length {
|
|
sum += w[k]; // index 0 = newest (financial indexing)
|
|
}
|
|
self.out[0] = Scalar::F64(sum / self.length as f64);
|
|
Some(&self.out)
|
|
}
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::*;
|
|
use aura_core::{AnyColumn, Timestamp};
|
|
|
|
#[test]
|
|
fn sma_warms_up_then_tracks_the_window_mean() {
|
|
let mut sma = Sma::new(3);
|
|
let schema = sma.schema();
|
|
|
|
// size the input column from the schema, as the engine will at wiring
|
|
let mut inputs = vec![AnyColumn::with_capacity(
|
|
schema.inputs[0].kind,
|
|
schema.inputs[0].lookback,
|
|
)];
|
|
|
|
let feed = [1.0_f64, 2.0, 3.0, 4.0, 5.0];
|
|
// means of [1,2,3], [2,3,4], [3,4,5] once warmed up
|
|
let expect = [None, None, Some(2.0), Some(3.0), Some(4.0)];
|
|
|
|
for (v, want) in feed.iter().zip(expect) {
|
|
inputs[0].push(Scalar::F64(*v)).unwrap();
|
|
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())),
|
|
}
|
|
}
|
|
}
|
|
|
|
#[test]
|
|
fn sma_length_one_is_identity() {
|
|
let mut sma = Sma::new(1);
|
|
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, Timestamp(0))), Some([Scalar::F64(7.0)].as_slice()));
|
|
|
|
inputs[0].push(Scalar::F64(9.0)).unwrap();
|
|
assert_eq!(sma.eval(Ctx::new(&inputs, Timestamp(0))), Some([Scalar::F64(9.0)].as_slice()));
|
|
}
|
|
}
|