Files
Aura/docs/specs/0001-core-streaming-substrate.md
T
Brummel b622a74571 spec: cycle 0001 core streaming substrate
The first walking-skeleton cycle: the aura-core data substrate everything
streams through — the four scalar base types (Scalar/ScalarKind/Timestamp,
C7), the fixed-capacity financial-indexed SoA Column with its zero-copy
Window (C8), the type-erased AnyColumn edge with an edge-time kind check
(C7), and the per-column run_count freshness primitive (C5). Sharpens the
RustAst reference: no boxed Value, no VecDeque, no RefCell/Rc, no per-event
allocation. Node trait, Ctx, firing policies, and the sim loop are deferred
to subsequent cycles.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-03 11:30:01 +02:00

14 KiB

Core Streaming Substrate — Design Spec

Date: 2026-06-03 Status: Draft — awaiting user spec review Authors: orchestrator + Claude

Goal

Deliver the foundation everything in aura streams through: the four scalar base types, the bounded-lookback Structure-of-Arrays column with financial-style indexing, the zero-copy read window a node sees, the type-erased edge carrier, and the per-column run-count that the freshness model is built on. This is the first cycle of the walking-skeleton milestone (ingest → one signal → deterministic backtest → position table → sim-optimal broker → pip-equity); it is the substrate the subsequent Node/Ctx/engine-loop cycles build on.

The reference is RustAst's src/ast/rtl/series and streams — a working implementation of this substrate. aura keeps its shape (newest-at-0 indexing, a per-series push counter for freshness, SoA composites) and sharpens it: no boxed Value, no VecDeque, no RefCell/Rc, no per-event allocation.

This cycle implements only the data substrate. Out of scope (subsequent cycles): the Node trait, the evaluation Ctx, firing policies A/B, the deterministic sim loop / freshness-gated recompute, the ingestion boundary / k-way merge, and the named OHLCV bundle. Those are named here only to fix this cycle's boundary; this cycle delivers the primitives they will consume.

Architecture

aura-core gains four small modules, no dependencies added:

  • scalar — the closed four-type scalar set: the Timestamp newtype, the Scalar Copy-POD carrier, and the ScalarKind tag (C7).
  • columnColumn<T>: a fixed-capacity ring buffer, pre-sized to a lookback limit at construction, financial-indexed (index 0 = newest), carrying a monotonic run_count (C5/C8). Window<'_, T> is its zero-copy read view.
  • anyAnyColumn: the type-erased edge carrier, an enum over the four Column<T> kinds. The kind check is paid at the edge (wiring time); the hot loop pushes into the concrete Column<T> monomorphically (C7).
  • errorKindMismatch, the wiring-time guard error returned when a Scalar of the wrong kind is pushed into an AnyColumn.

The whole substrate is interior-mutability-free. A sim is single-threaded and non-concurrent (C1), so the engine owns columns by &mut for the push side and hands out &Column / Window (shared borrow) on the read side — RustAst's RefCell-on-every-series disappears.

Concrete code shapes

North-star slice (the consumer this substrate exists to serve)

A future aura-std SMA node — not built this cycle — reads its input as an indexed, newest-at-0 window. This is the minimal honest slice proving the substrate's read API fits the node use case it exists for:

// Illustrative future consumer (NOT delivered this cycle): the read shape an
// indicator node wants. `length` is the node's pre-sized lookback (C8).
fn sma(window: &Window<'_, f64>, length: usize) -> Option<f64> {
    if window.len() < length {
        return None; // not yet warmed up (C8: None = not-yet-warmed-up)
    }
    let sum: f64 = (0..length).map(|k| window[k]).sum(); // [0] newest, [length-1] oldest
    Some(sum / length as f64)
}

The point of the cycle is that Window makes this code possible with no lookback bookkeeping in the node: the engine pre-sized the buffer, the node only reads by index.

Delivered code — the scalar set (scalar.rs)

/// Canonical engine time: epoch-nanoseconds UTC. Newtype over i64 (C7).
/// `Default` (epoch 0) is required so `Column<Timestamp>` can pre-size its ring.
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub struct Timestamp(pub i64);

/// The four scalar base types, type-erased for an edge. Copy POD — no `dyn Any`,
/// no heap, no Rc (C7 forbids all three on the hot path).
#[derive(Clone, Copy, Debug, PartialEq)]
pub enum Scalar {
    I64(i64),
    F64(f64),
    Bool(bool),
    Ts(Timestamp),
}

/// The kind tag of a scalar / column, used for the edge-time type check.
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
pub enum ScalarKind {
    I64,
    F64,
    Bool,
    Timestamp,
}

impl Scalar {
    pub fn kind(self) -> ScalarKind { /* match */ }
}

// Total, infallible widenings into the carrier:
impl From<i64> for Scalar { /* ... */ }
impl From<f64> for Scalar { /* ... */ }
impl From<bool> for Scalar { /* ... */ }
impl From<Timestamp> for Scalar { /* ... */ }

Delivered code — the column and its window (column.rs)

/// A fixed-capacity, bounded-lookback SoA column. Pre-sized at construction
/// (C8: no realloc in the hot loop); financial-indexed (index 0 = newest);
/// carries a monotonic push counter for the freshness model (C5).
pub struct Column<T> {
    buf: Box<[T]>,   // capacity == lookback; never reallocated
    head: usize,     // ring write cursor (next slot)
    len: usize,      // valid elements, saturates at capacity
    run_count: u64,  // monotonic total pushes — RustAst's total_count (C5)
}

impl<T: Copy + Default> Column<T> {
    /// `lookback` is the window depth the engine sized for this edge (C8).
    pub fn with_capacity(lookback: usize) -> Self { /* Box<[T]> of len lookback */ }

    /// Append the newest value. O(1), no allocation. Bumps `run_count`.
    /// On overflow the oldest value is overwritten (true ring).
    pub fn push(&mut self, v: T) { /* buf[head]=v; head=(head+1)%cap; len=min(len+1,cap); run_count+=1 */ }

    /// Financial index: 0 = newest, `len()-1` = oldest. None if `k >= len`
    /// (C8: None == not-yet-warmed-up / out of range). Zero-copy.
    pub fn get(&self, k: usize) -> Option<T> { /* buf[(head - 1 - k) mod cap] */ }

    pub fn len(&self) -> usize { self.len }
    pub fn capacity(&self) -> usize { self.buf.len() }
    pub fn run_count(&self) -> u64 { self.run_count }

    /// A zero-copy read view handed to a reading node.
    pub fn window(&self) -> Window<'_, T> { Window { col: self } }
}

/// Read-only, zero-copy view into a `Column` — what a node sees (C8).
pub struct Window<'a, T> {
    col: &'a Column<T>,
}

impl<'a, T: Copy + Default> Window<'a, T> {
    pub fn get(&self, k: usize) -> Option<T> { self.col.get(k) }
    pub fn len(&self) -> usize { self.col.len() }
    pub fn run_count(&self) -> u64 { self.col.run_count() }
}

/// `window[k]` — newest-at-0 indexed access; panics on out-of-range, matching
/// slice `Index` semantics. Nodes that may be cold use `get` instead.
impl<'a, T: Copy + Default> core::ops::Index<usize> for Window<'a, T> {
    type Output = T;
    fn index(&self, k: usize) -> &T { /* &buf[(head-1-k) mod cap], bounds-checked */ }
}

Delivered code — the type-erased edge (any.rs)

/// A type-erased edge: the four concrete columns behind one enum (C7). The kind
/// check happens here, at wiring time; the hot loop grabs the concrete column
/// once (`as_f64_mut`) and pushes monomorphically — no per-event dispatch.
pub enum AnyColumn {
    I64(Column<i64>),
    F64(Column<f64>),
    Bool(Column<bool>),
    Ts(Column<Timestamp>),
}

impl AnyColumn {
    pub fn with_capacity(kind: ScalarKind, lookback: usize) -> Self { /* match kind */ }

    pub fn kind(&self) -> ScalarKind { /* match self */ }
    pub fn len(&self) -> usize { /* match self */ }
    pub fn run_count(&self) -> u64 { /* match self */ }

    /// Edge-time typed push. Returns `KindMismatch` if `v.kind() != self.kind()`
    /// — the C7 type-erasure guard. The hot path avoids this by taking the
    /// concrete `Column<T>` once via `as_*_mut`.
    pub fn push(&mut self, v: Scalar) -> Result<(), KindMismatch> { /* match (self, v) */ }

    /// Read one type-erased value (0 = newest). None if cold / out of range.
    pub fn get(&self, k: usize) -> Option<Scalar> { /* match self -> Scalar::… */ }

    /// Concrete-column accessors for the monomorphic hot path.
    pub fn as_f64_mut(&mut self) -> Option<&mut Column<f64>> { /* … */ }
    pub fn as_i64_mut(&mut self) -> Option<&mut Column<i64>> { /* … */ }
    pub fn as_bool_mut(&mut self) -> Option<&mut Column<bool>> { /* … */ }
    pub fn as_ts_mut(&mut self) -> Option<&mut Column<Timestamp>> { /* … */ }
}

Delivered code — the must-fail guard fixture (error.rs + a test)

The C7 type-erasure contract is load-bearing: a wrong-kind push must be rejected, never silently coerced. The test whose correct behaviour is rejection:

#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub struct KindMismatch {
    pub expected: ScalarKind,
    pub got: ScalarKind,
}

#[test]
fn wrong_kind_push_is_rejected() {
    let mut edge = AnyColumn::with_capacity(ScalarKind::I64, 8);
    let err = edge.push(Scalar::F64(1.5)).unwrap_err(); // pushing f64 into an i64 edge
    assert_eq!(err, KindMismatch { expected: ScalarKind::I64, got: ScalarKind::F64 });
    assert_eq!(edge.len(), 0);       // nothing was stored
    assert_eq!(edge.run_count(), 0); // and the counter did not move
}

Implementation shape (secondary — before → after)

aura-core/src/lib.rs today is doc-comment-only (no pub items). After this cycle it declares the four modules and re-exports the surface:

// before: lib.rs is documentation only, defines no types.

// after:
mod scalar;
mod column;
mod any;
mod error;

pub use scalar::{Scalar, ScalarKind, Timestamp};
pub use column::{Column, Window};
pub use any::AnyColumn;
pub use error::KindMismatch;

Components

  • Timestampi64 epoch-ns newtype, fully ordered (the ingestion boundary will normalize source-native units into it in a later cycle; here it is just the type).
  • Scalar / ScalarKind — the closed carrier + tag. Scalar is Copy; From impls give ergonomic widening; kind() projects to the tag.
  • Column<T> — the ring. Owns a pre-sized Box<[T]>; push is O(1) and allocation-free; get/Index implement newest-at-0; run_count is the freshness primitive.
  • Window<'_, T> — the borrow a reader gets; carries no storage, only a reference; Index for the warm path, get for the possibly-cold path.
  • AnyColumn — the edge; type check at the boundary, concrete-column accessors for the hot loop.
  • KindMismatch — the guard error.

Data flow

  1. The engine (later cycle) sizes an edge: AnyColumn::with_capacity(kind, lookback) — buffer allocated once, never again (C8).
  2. A producer pushes its newest output for the cycle. Wiring-time / dynamic path: edge.push(scalar)? (kind-checked). Hot path: the engine took edge.as_f64_mut() once at sim start and calls col.push(x) monomorphically each cycle — no per-event dispatch, no allocation (C7).
  3. run_count increments on every push; the freshness model (later cycle) compares an input's run_count against the value it last saw to decide fire-vs-hold (C5). This cycle delivers the counter, not the comparison loop.
  4. A reading node receives col.window() and reads by index, newest at 0 (C8).

There is no merge, no clock, and no recompute logic in this cycle — those are the engine's, downstream. The substrate is pure data structures + their invariants.

Error handling

  • Wrong-kind pushErr(KindMismatch { expected, got }); the column is left untouched (len and run_count unchanged). This is the only fallible operation in the substrate.
  • Cold / out-of-range readget returns None (C8: not-yet-warmed-up); Index panics like a slice (the warm-path ergonomic, used only when the node has already checked len).
  • No panics on the push/normal-read path; no allocation after construction; no interior mutability to poison.

Testing strategy

Unit tests live beside each module; all run under cargo test --workspace.

  • Column ring semantics: push then get(0) is newest; get(len-1) is oldest; pushing more than capacity keeps exactly the newest capacity values and drops the oldest (wraparound correctness); len saturates at capacity; get(k>=len) == None.
  • run_count is the freshness primitive: monotonic across pushes (including after the ring wraps — it counts pushes, not stored len); never moved by a read.
  • Window is a zero-copy view: values read through the window equal the column's; taking a window does not change the column's capacity (no realloc); Index and get agree on the warm range.
  • Scalar / ScalarKind: kind() is correct for each variant; From widenings round-trip through get on an AnyColumn.
  • AnyColumn type erasure: same-kind push stores and is read back as the right-kind Scalar; run_count/len delegate; the concrete as_*_mut accessor returns Some only for the matching kind.
  • The must-fail guard (wrong_kind_push_is_rejected, shown above): a wrong-kind push is rejected and leaves the column untouched.

Acceptance criteria

Applying aura's feature-acceptance lens (CLAUDE.md):

  1. The intended audience naturally reaches for it. The node author (LLM via builder APIs, later cycles) reads inputs as window[k] newest-at-0 — the north-star slice above is exactly the SMA an author writes, with zero lookback bookkeeping. The substrate is the thing every node, source, and sink stands on.
  2. It measurably improves correctness / removes redundancy. Engine-sized, pre-allocated windows mean node code cannot mis-manage lookback (C8 rationale); the closed Scalar set + AnyColumn removes RustAst's boxed Value and per-series RefCell entirely.
  3. It reintroduces no failure class the core constraints exist to kill. No dyn Any, no per-event heap allocation, no interior mutability, fixed pre-sized buffers (C1/C7/C8 all upheld); the wrong-kind push is structurally rejected, not coerced.

Concretely, the cycle ships when:

  • cargo build --workspace, cargo test --workspace, and cargo clippy --workspace --all-targets -- -D warnings are all green.
  • Every test in the Testing-strategy list exists and passes, including the wrong_kind_push_is_rejected must-fail guard.
  • aura-core exports Scalar, ScalarKind, Timestamp, Column, Window, AnyColumn, KindMismatch and nothing leaks RefCell/Rc/dyn Any on the substrate surface.