Files
Aura/crates/aura-std/src/recorder.rs
T
Brummel cd3d1ca9ed refactor(aura-core): split Scalar into a tag-free Cell + ScalarKind
Motivation
----------
`Scalar` was a tagged enum (I64/F64/Bool/Ts), so every scalar value
physically carried its own kind tag. But the kind is already known from
the schema/port/column the value flows through (C7: the type is a
property of the column, not of the value — the hot path is already
columnar `Column<T>`, and `AnyColumn::get` *reconstructs* the tag from
the column on the way out). The per-value tag was therefore redundant
with the kind the surrounding context already holds.

That redundancy had three costs:

  * It baked an implicit `match` (a branch) into every function that read
    a Scalar payload — even where the caller statically knew the type.
    The tag could never be exploited away.
  * Size: a tagged enum is tag + payload = 16 bytes (f64/i64 alignment),
    twice the 8 bytes the value needs. A `Column<Scalar>` would be double
    the memory and half the cache utilisation.
  * It is the shared root of several downstream papercuts we keep hitting
    — the lossy f64 manifest field, the `unreachable!` panic on a
    non-numeric param, the serde-tag question — all symptoms of "the type
    is baked into the value".

Change
------
Introduce `Cell`: a type-erased 64-bit word (`struct Cell(u64)`) that is
not readable without external type context. It is constructed per base
type (`from_i64/from_f64/from_bool/from_ts`) and read only by naming the
type at the call site (`i64()/f64()/bool()/ts()`) — each a branch-free
bit-cast. The hot path resolves the kind once at the boundary (from the
schema) and then reads natively, with no per-value branch. `Cell` knows
nothing of `Scalar` or `ScalarKind`; the dependency is strictly one-way,
and it lives in its own `cell.rs` (more is planned on top of it).

`Scalar` becomes `struct { kind: ScalarKind, cell: Cell }` — the
self-describing form for the dynamic boundaries (builder binding,
serialization, rendering), built on top of `Cell`. Its `as_*` accessors
now `debug_assert` the kind and return the native value (free in
release); calling the wrong accessor is a caller bug, not a checked
`Option`. The variant constructors `Scalar::I64(..)` become associated
fns `Scalar::i64(..)`.

`PartialEq` is hand-written (not derived) to preserve the former enum's
value semantics: kinds must match, then native payloads compare, so f64
keeps IEEE-754 behaviour (`NaN != NaN`, `+0.0 == -0.0`) and a kind
mismatch is never equal even when the raw words coincide. A fixture
(`scalar_eq_is_value_not_bitwise`) pins exactly the cases where bit- and
value-equality diverge, so it can't silently regress. `Cell`'s own
`Eq`/`Hash` stay bitwise — correct for a raw word.

The change is behaviour-preserving: Scalar's observable behaviour is
identical to the pre-Cell enum (the value-equality fixture proves it);
only the internal representation changed. The ~440 call sites across the
workspace are a mechanical constructor rename plus ~12 destructuring
sites (match-arms / `let`-patterns) rewritten to `kind()` + `as_*`.

Verified: cargo build --workspace --all-targets, cargo clippy --workspace
--all-targets -- -D warnings, cargo test --workspace — all green.
2026-06-16 12:12:52 +02:00

153 lines
6.2 KiB
Rust

//! `Recorder` — a reusable recording sink (the glossary *sink* role, C8/C22):
//! a pure consumer that, each fired cycle, sends `(ctx.now(), row)` — the newest
//! value of each declared input column — to an out-of-graph `mpsc` destination it
//! holds. It produces nothing (`output: vec![]`), so it is a leaf in the DAG. The
//! `mpsc::Sender` keeps the engine's purity invariant (C7): the node carries no
//! `Rc`/`RefCell` interior mutability, only an owned channel handle. Supports all
//! four base scalar kinds so any column can be persisted; returns `None` (filters)
//! until every input column is warm.
use aura_core::{Ctx, Firing, Node, NodeSchema, PortSpec, PrimitiveBuilder, Scalar, ScalarKind, Timestamp};
use std::sync::mpsc::Sender;
/// A recording sink over `kinds.len()` input columns. Each fired cycle it reads
/// the newest value of every column and sends the row to `tx`; it returns `None`
/// (records, forwards nothing) and `None` during warm-up until all columns have a
/// value.
pub struct Recorder {
kinds: Vec<ScalarKind>,
tx: Sender<(Timestamp, Vec<Scalar>)>,
}
impl Recorder {
/// A recorder over one input column per entry in `kinds`, each with the given
/// `firing` policy, sending recorded `(timestamp, row)` pairs to `tx`. The
/// `firing` policy is a property of the declared signature (carried by the
/// `PrimitiveBuilder`'s `PortSpec`s), not of the built sink, so it is part of
/// the construction contract but not stored on the instance.
pub fn new(kinds: &[ScalarKind], _firing: Firing, tx: Sender<(Timestamp, Vec<Scalar>)>) -> Self {
Self { kinds: kinds.to_vec(), tx }
}
/// The param-generic recipe for a blueprint primitive. The channel + kinds + firing
/// are non-param construction args (captured), not tunable params; the node
/// declares none. The input ports (one per kind) are threaded into the schema
/// statically; `tx`/`kinds` are cloned for the build closure.
pub fn builder(
kinds: Vec<ScalarKind>,
firing: Firing,
tx: Sender<(Timestamp, Vec<Scalar>)>,
) -> PrimitiveBuilder {
let inputs = kinds
.iter()
.enumerate()
.map(|(i, &kind)| PortSpec { kind, firing, name: format!("col[{i}]") })
.collect();
let build_kinds = kinds.clone();
PrimitiveBuilder::new(
"Recorder",
NodeSchema { inputs, output: vec![], params: vec![] }, // sink: empty output (C8)
move |_| Box::new(Recorder::new(&build_kinds, firing, tx.clone())),
)
}
}
impl Node for Recorder {
fn lookbacks(&self) -> Vec<usize> {
vec![1; self.kinds.len()]
}
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() {
// newest of each column by kind; `?` returns None (warm-up) if cold.
let scalar = match kind {
ScalarKind::F64 => Scalar::f64(ctx.f64_in(i).get(0)?),
ScalarKind::I64 => Scalar::i64(ctx.i64_in(i).get(0)?),
ScalarKind::Bool => Scalar::bool(ctx.bool_in(i).get(0)?),
ScalarKind::Timestamp => Scalar::ts(ctx.ts_in(i).get(0)?),
};
row.push(scalar);
}
let _ = self.tx.send((ctx.now(), row));
None
}
fn label(&self) -> String {
"Recorder".to_string()
}
}
#[cfg(test)]
mod tests {
use super::*;
use aura_core::{AnyColumn, Timestamp};
use std::sync::mpsc;
#[test]
fn recorder_captures_f64_stream_after_warmup() {
let (tx, rx) = mpsc::channel();
let mut rec = Recorder::new(&[ScalarKind::F64], Firing::Any, tx);
// a sink declares no output (C8) — asserted on the param-generic builder
let (tx_b, _rx_b) = mpsc::channel();
assert!(
Recorder::builder(vec![ScalarKind::F64], Firing::Any, tx_b)
.schema()
.output
.is_empty(),
"a sink declares no output (C8)"
);
// size the one f64 input column from the node's lookback, as bootstrap would.
let mut inputs = vec![AnyColumn::with_capacity(ScalarKind::F64, rec.lookbacks()[0])];
// cold: returns None and records nothing.
assert_eq!(rec.eval(Ctx::new(&inputs, Timestamp(1))), None);
assert!(rx.try_recv().is_err());
// warm: returns None (pure consumer) but records (now, [F64(newest)]).
for (t, v) in [(2_i64, 10.0_f64), (3, 20.0), (4, 30.0)] {
inputs[0].push(Scalar::f64(v)).unwrap();
assert_eq!(rec.eval(Ctx::new(&inputs, Timestamp(t))), None);
}
let rows: Vec<(Timestamp, Vec<Scalar>)> = rx.try_iter().collect();
assert_eq!(
rows,
vec![
(Timestamp(2), vec![Scalar::f64(10.0)]),
(Timestamp(3), vec![Scalar::f64(20.0)]),
(Timestamp(4), vec![Scalar::f64(30.0)]),
]
);
}
#[test]
fn recorder_is_none_until_all_columns_warm() {
let (tx, rx) = mpsc::channel();
let mut rec = Recorder::new(&[ScalarKind::F64, ScalarKind::F64], Firing::Any, tx);
let mut inputs = vec![
AnyColumn::with_capacity(ScalarKind::F64, 1),
AnyColumn::with_capacity(ScalarKind::F64, 1),
];
// only column 0 present -> None, nothing recorded.
inputs[0].push(Scalar::f64(1.0)).unwrap();
assert_eq!(rec.eval(Ctx::new(&inputs, Timestamp(1))), None);
assert!(rx.try_recv().is_err());
// both present -> records the full row (still returns None).
inputs[1].push(Scalar::f64(2.0)).unwrap();
assert_eq!(rec.eval(Ctx::new(&inputs, Timestamp(2))), None);
let rows: Vec<(Timestamp, Vec<Scalar>)> = rx.try_iter().collect();
assert_eq!(rows, vec![(Timestamp(2), vec![Scalar::f64(1.0), Scalar::f64(2.0)])]);
}
#[test]
fn input_slots_are_named_col_index() {
let (tx, _rx) = std::sync::mpsc::channel();
let r = Recorder::builder(vec![ScalarKind::F64, ScalarKind::F64], Firing::Any, tx);
let names: Vec<String> = r.schema().inputs.iter().map(|p| p.name.clone()).collect();
assert_eq!(names, ["col[0]", "col[1]"]);
}
}