From 9337a8585db60b88622c26273d44ef31637d20c8 Mon Sep 17 00:00:00 2001 From: Brummel Date: Mon, 15 Jun 2026 10:13:17 +0200 Subject: [PATCH] plan: 0041 source ingestion seam Decomposes the boss-signed spec into four tasks: (1) the `Source` trait + `VecSource` adapter in aura-engine; (2) the atomic `Harness::run` re-type to `Vec>` with all 53 call sites threaded behaviour-preserving; (3) the streaming `M1FieldSource` over a data-server window; (4) a gated streaming e2e + the measured residency predicate. refs #71 --- docs/plans/0041-source-ingestion-seam.md | 871 +++++++++++++++++++++++ 1 file changed, 871 insertions(+) create mode 100644 docs/plans/0041-source-ingestion-seam.md diff --git a/docs/plans/0041-source-ingestion-seam.md b/docs/plans/0041-source-ingestion-seam.md new file mode 100644 index 0000000..66ed675 --- /dev/null +++ b/docs/plans/0041-source-ingestion-seam.md @@ -0,0 +1,871 @@ +# Source ingestion seam — Implementation Plan + +> **Parent spec:** `docs/specs/0041-source-ingestion-seam.md` +> +> **For agentic workers:** REQUIRED SUB-SKILL: use the `implement` skill to +> run this plan. Steps use `- [ ]` checkboxes for tracking. + +**Goal:** Re-type the engine's ingestion input from a materialized +`Vec>` into a `Source` producer trait the k-way merge +drives by `peek`/`next`, and prove it against a real, lazily-streamed +`data-server` window — closing the cycle-0011 eager-materialization gap. + +**Architecture:** `aura-engine` gains a `Source` trait + a `VecSource` adapter; +`Harness::run` is re-typed to `Vec>` with the merge loop driven +by `peek`/`next`, and every existing call site is wrapped in `VecSource` +(behaviour-preserving, byte-identical). `aura-ingest` gains a streaming +`M1FieldSource` over a `data-server` M1 window (per-pull decode from a borrowed +`Arc` chunk, lazy `next_chunk` refill, `ms→epoch-ns` at the seam), proven by a +gated end-to-end backtest plus a measured residency predicate. + +**Tech Stack:** `aura-engine` (`harness.rs`, `lib.rs`), `aura-ingest` +(`lib.rs`, `tests/`), `data-server` (`stream_m1_windowed`, `SymbolChunkIter`, +`next_chunk`, `loader::CHUNK_SIZE`, `M1Parsed`). + +--- + +**Files this plan creates or modifies:** + +- Modify: `crates/aura-engine/src/harness.rs` — add `Source` trait + `VecSource` + (after `SourceSpec`, ~line 52); re-type `Harness::run` (line 220) + merge loop + (220-258); re-type test helper `ohlcv_streams()` (line 900); wrap 31 call sites. +- Modify: `crates/aura-engine/src/lib.rs:25-30,47` — refresh module doc; add + `Source`, `VecSource` to the crate-root re-export. +- Modify: `crates/aura-engine/src/blueprint.rs` — wrap 10 call sites. +- Modify: `crates/aura-engine/src/sweep.rs:263` — wrap 1 call site. +- Modify: `crates/aura-engine/src/report.rs:223` — wrap 1 call site. +- Modify: `crates/aura-cli/src/main.rs` — wrap 4 call sites. +- Modify: `crates/aura-ingest/src/lib.rs` — add `M1Field` + `decode` + `M1FieldSource` + (near `load_m1_window`, ~line 110); add `data_server::SymbolChunkIter` import. +- Modify: `crates/aura-ingest/tests/real_bars.rs:72` — wrap 1 call site. +- Create: `crates/aura-ingest/tests/streaming_seam.rs` — gated streaming e2e + + residency-drive + sharing tests. +- Test: `harness.rs` `#[cfg(test)] mod tests` — `VecSource` unit tests. +- Test: `aura-ingest/src/lib.rs` `#[cfg(test)] mod tests` — `decode` unit tests. + +--- + +## Task 1: `Source` trait + `VecSource` adapter (aura-engine) + +**Files:** +- Modify: `crates/aura-engine/src/harness.rs` (after `SourceSpec`, ~line 52) +- Modify: `crates/aura-engine/src/lib.rs:25-30,47` +- Test: `crates/aura-engine/src/harness.rs` (`#[cfg(test)] mod tests`) + +This task adds the new types only; `Harness::run` is NOT re-typed here (that is +Task 2). `Source`/`VecSource` compile and their tests pass standalone. + +- [ ] **Step 1: Add the `Source` trait + `VecSource` to `harness.rs`** + +Insert immediately after the `SourceSpec` struct (after its closing `}` at +`harness.rs:52`, before the `FlatGraph` doc comment at line 54): + +```rust +/// The ingestion-boundary producer the k-way merge drives (C3). A source yields +/// timestamped scalars in ascending timestamp order (the C3 ingestion +/// precondition). Object-safe: the merge holds `Vec>`. +pub trait Source { + /// The timestamp of the head record, without consuming it. `None` = the + /// source is exhausted. Cheap and side-effect-free: it reads a pre-decoded + /// head, never advancing the underlying stream. + fn peek(&self) -> Option; + + /// Pop the head `(timestamp, scalar)` and advance. `None` = exhausted. + /// After `next` returns `Some`, `peek` reflects the new head. + fn next(&mut self) -> Option<(Timestamp, Scalar)>; +} + +/// A `Source` over a fully materialized stream — the cycle-0011 eager shape, +/// named. Every pre-seam call site wraps its `Vec<(Timestamp, Scalar)>` in this, +/// so run output is byte-identical (C1) and residency is unchanged. +pub struct VecSource { + stream: Vec<(Timestamp, Scalar)>, + cursor: usize, +} + +impl VecSource { + pub fn new(stream: Vec<(Timestamp, Scalar)>) -> Self { + Self { stream, cursor: 0 } + } +} + +impl Source for VecSource { + fn peek(&self) -> Option { + self.stream.get(self.cursor).map(|&(t, _)| t) + } + fn next(&mut self) -> Option<(Timestamp, Scalar)> { + let item = self.stream.get(self.cursor).copied()?; + self.cursor += 1; + Some(item) + } +} +``` + +- [ ] **Step 2: Add `VecSource` unit tests** + +Append inside the existing `#[cfg(test)] mod tests { ... }` block in +`harness.rs` (the module already has `use super::*;`): + +```rust +#[test] +fn vec_source_peek_is_non_consuming_and_idempotent() { + let mut s = VecSource::new(vec![(Timestamp(1), Scalar::F64(10.0)), (Timestamp(2), Scalar::F64(20.0))]); + // peek twice: same head, no advance. + assert_eq!(s.peek(), Some(Timestamp(1))); + assert_eq!(s.peek(), Some(Timestamp(1))); + // next pops it; peek now reflects the new head. + assert_eq!(Source::next(&mut s), Some((Timestamp(1), Scalar::F64(10.0)))); + assert_eq!(s.peek(), Some(Timestamp(2))); +} + +#[test] +fn vec_source_exhaustion_yields_none() { + let mut s = VecSource::new(vec![(Timestamp(5), Scalar::I64(7))]); + assert_eq!(Source::next(&mut s), Some((Timestamp(5), Scalar::I64(7)))); + assert_eq!(s.peek(), None); + assert_eq!(Source::next(&mut s), None); +} + +#[test] +fn vec_source_empty_peeks_none() { + let mut s = VecSource::new(vec![]); + assert_eq!(s.peek(), None); + assert_eq!(Source::next(&mut s), None); +} +``` + +(`Source::next(&mut s)` is spelled fully-qualified so it never collides with +`Iterator::next` in scope.) + +- [ ] **Step 3: Re-export the new types and refresh the module doc** + +In `crates/aura-engine/src/lib.rs`, change the re-export at line 47 from: + +```rust +pub use harness::{BootstrapError, Edge, FlatGraph, Harness, SourceSpec, Target}; +``` + +to: + +```rust +pub use harness::{BootstrapError, Edge, FlatGraph, Harness, Source, SourceSpec, Target, VecSource}; +``` + +And in the module doc, change the "Still to come" line at lines 25-26 from: + +```rust +//! Still to come (subsequent cycles): the `Source` trait + data-server ingestion +//! and source-native time normalization (C3/C11), the broker-independent +``` + +to: + +```rust +//! Still to come (subsequent cycles): the broker-independent +``` + +- [ ] **Step 4: Build and run the `VecSource` tests** + +Run: `cargo test -p aura-engine vec_source` +Expected: PASS (3 tests: `vec_source_peek_is_non_consuming_and_idempotent`, +`vec_source_exhaustion_yields_none`, `vec_source_empty_peeks_none`). + +- [ ] **Step 5: Lint gate** + +Run: `cargo clippy -p aura-engine --all-targets -- -D warnings` +Expected: clean (no warnings; the pub re-exported `Source`/`VecSource` are not +dead code). + +--- + +## Task 2: Re-type `Harness::run` + thread all callers (atomic compile unit) + +**Files:** +- Modify: `crates/aura-engine/src/harness.rs:220-258` (signature + merge loop), + `:900` (`ohlcv_streams` helper), 31 call sites +- Modify: `crates/aura-engine/src/blueprint.rs` (10 sites) +- Modify: `crates/aura-engine/src/sweep.rs:263` (1 site) +- Modify: `crates/aura-engine/src/report.rs:223` (1 site) +- Modify: `crates/aura-cli/src/main.rs` (4 sites) +- Modify: `crates/aura-ingest/tests/real_bars.rs:72` (1 site) + +**Why one task:** re-typing `run`'s signature makes the whole workspace fail to +compile until every one of the 53 callers is threaded. The signature change and +all callers must land together; the build gate at the end is unsatisfiable +otherwise (self-review item 7). The 4 `h.run(ohlcv_streams())` sites are threaded +by re-typing the helper, not the call site. + +**The mechanical wrap rule (applies to every `vec![...]` call-site literal):** +each top-level element `X` of the `vec![...]` passed to `.run(` is replaced by +`Box::new(VecSource::new(X))`. The expected parameter type `Vec>` +drives the unsizing coercion, so no `as Box` is needed. +Example: `h.run(vec![prices])` → `h.run(vec![Box::new(VecSource::new(prices))])`; +`h.run(vec![s0, s1])` → `h.run(vec![Box::new(VecSource::new(s0)), Box::new(VecSource::new(s1))])`. + +- [ ] **Step 1: Re-type the `run` signature + merge loop in `harness.rs`** + +Replace the current signature + pick scan (`harness.rs:220-257`). Current: + +```rust + pub fn run(&mut self, streams: Vec>) { + assert_eq!( + streams.len(), + self.sources.len(), + "run: one stream per source required (got {} streams for {} sources)", + streams.len(), + self.sources.len() + ); + + // disjoint field borrows so the topo walk can read topo/out_edges/sources + // while mutating nodes + let Harness { nodes, topo, out_edges, sources } = self; + + let mut cursor: Vec = vec![0; streams.len()]; + let mut cycle_id: u64 = 0; + let mut scratch: Vec = Vec::new(); + + loop { + // pick the live source head with the smallest (timestamp, source index) + let mut pick: Option = None; + for (s, stream) in streams.iter().enumerate() { + if cursor[s] < stream.len() { + match pick { + None => pick = Some(s), + Some(p) => { + if stream[cursor[s]].0 < streams[p][cursor[p]].0 { + pick = Some(s); + } + } + } + } + } + let s = match pick { + Some(s) => s, + None => break, // all streams exhausted + }; + let (ts, value) = streams[s][cursor[s]]; + cursor[s] += 1; + cycle_id += 1; +``` + +Replace with (sources arg is `mut`; the forward+eval body below line 257 is +UNCHANGED): + +```rust + pub fn run(&mut self, mut sources_in: Vec>) { + assert_eq!( + sources_in.len(), + self.sources.len(), + "run: one source per SourceSpec required (got {} sources for {} specs)", + sources_in.len(), + self.sources.len() + ); + + // disjoint field borrows so the topo walk can read topo/out_edges/sources + // while mutating nodes + let Harness { nodes, topo, out_edges, sources } = self; + + let mut cycle_id: u64 = 0; + let mut scratch: Vec = Vec::new(); + + loop { + // pick the live source head with the smallest (timestamp, source index). + // strictly-`<` replace + source-order scan preserves C4 tie-breaking. + let mut pick: Option<(usize, Timestamp)> = None; + for (s, src) in sources_in.iter().enumerate() { + if let Some(ts) = src.peek() { + match pick { + None => pick = Some((s, ts)), + Some((_, best)) if ts < best => pick = Some((s, ts)), + _ => {} + } + } + } + let s = match pick { + Some((s, _)) => s, + None => break, // all sources exhausted + }; + // `&mut *sources_in[s]`: deref the Box to `dyn Source` then re-borrow, + // so UFCS resolves Self = dyn Source (Box does not itself + // impl Source). Fully-qualified to never collide with Iterator::next. + let (ts, value) = + Source::next(&mut *sources_in[s]).expect("peeked Some ⇒ next is Some"); + cycle_id += 1; +``` + +(The local binding is named `sources_in` to avoid shadowing the `sources` +destructured from `self` — the `SourceSpec` wiring used in the forward step at +the unchanged line `for t in sources[s].targets.iter()`.) + +- [ ] **Step 2: Verify the forward/eval body still references `sources` (the wiring)** + +Confirm the unchanged forwarding line (was `harness.rs:261`) still reads: + +```rust + for t in sources[s].targets.iter() { +``` + +`sources` here is the `Vec` destructured from `self` — unchanged. No +edit; this step is a read-check that the rename in Step 1 did not collide. + +- [ ] **Step 3: Re-type the `ohlcv_streams()` test helper in `harness.rs:900`** + +Replace (lines 899-908): + +```rust + /// Build five timestamp-aligned f64 sources feeding Ohlcv's five barrier slots. + fn ohlcv_streams() -> Vec> { + vec![ + f64_stream(&[(1, 10.0), (2, 20.0)]), // open + f64_stream(&[(1, 15.0), (2, 25.0)]), // high + f64_stream(&[(1, 8.0), (2, 19.0)]), // low + f64_stream(&[(1, 12.0), (2, 22.0)]), // close + f64_stream(&[(1, 100.0), (2, 200.0)]), // volume + ] + } +``` + +with (returns boxed sources, so the four `h.run(ohlcv_streams())` callers need +no change): + +```rust + /// Build five timestamp-aligned f64 sources feeding Ohlcv's five barrier slots. + fn ohlcv_streams() -> Vec> { + vec![ + Box::new(VecSource::new(f64_stream(&[(1, 10.0), (2, 20.0)]))), // open + Box::new(VecSource::new(f64_stream(&[(1, 15.0), (2, 25.0)]))), // high + Box::new(VecSource::new(f64_stream(&[(1, 8.0), (2, 19.0)]))), // low + Box::new(VecSource::new(f64_stream(&[(1, 12.0), (2, 22.0)]))), // close + Box::new(VecSource::new(f64_stream(&[(1, 100.0), (2, 200.0)]))), // volume + ] + } +``` + +- [ ] **Step 4: Wrap the 31 remaining `harness.rs` call sites** + +Apply the wrap rule to each site below (the 4 `h.run(ohlcv_streams())` at +`:948/:1000/:1013/:1043` are NOT in this list — they are handled by Step 3). +Each `vec![...]` element becomes `Box::new(VecSource::new())`: + +```text +:594 h.run(vec![f64_stream(&[(1,1.0),(2,2.0),(3,3.0),(4,4.0),(5,5.0)])]); +:642 h.run(vec![prices.clone()]); :657 h2.run(vec![prices]); +:689 h.run(vec![s0.clone(), s1.clone()]); :704 h2.run(vec![s0, s1]); +:735 h.run(vec![s0.clone(), s1.clone()]); :748 h2.run(vec![s0, s1]); +:789 h.run(vec![prices.clone()]); :804 h2.run(vec![prices]); +:834 h.run(vec![s0, s1, s2]); +:1122 h.run(vec![f64_stream(&[(1,10.0),(2,12.0),(3,14.0),(4,16.0),(5,18.0)])]); +:1172 h.run(vec![ f64_stream(..)×5 ]); (wrap each of the 5 elements) +:1219 h.run(vec![ <4 raw inner-vec literals> ]); (wrap each of the 4 — see Step 5) +:1257 h.run(vec![f64_stream(&[(1,10.0),(2,20.0),(3,30.0)])]); +:1292 a.run(vec![prices.clone()]); :1297 b.run(vec![prices]); +:1327 h.run(vec![s0, s1]); :1362 h.run(vec![s0, s1]); +:1447 h.run(vec![prices.clone()]); :1466 h2.run(vec![prices]); +:1512 h.run(vec![prices.clone()]); :1535 h2.run(vec![prices]); +:1579 h.run(vec![a.clone(), b.clone()]); :1600 h2.run(vec![a, b]); +:1640 h.run(vec![prices.clone()]); :1651 h2.run(vec![prices]); +:1701 h.run(vec![prices.clone()]); :1728 h2.run(vec![prices]); +:1788 h.run(vec![prices.clone(), counts.clone()]); :1817 h2.run(vec![prices, counts]); +:1861 h.run(vec![f64_stream(&[...])]); (multi-line; wrap the single element) +:1926 h.run(vec![f64_stream(&[...])]); (multi-line; wrap the single element) +``` + +(Line numbers drift as edits are applied; match on the call expression, not the +absolute line. After this step, `git grep -n 'run(vec!' crates/aura-engine/src/harness.rs` +should return zero un-wrapped `vec![` — every element is now `Box::new(VecSource::new(...))`.) + +- [ ] **Step 5: Wrap the raw inner-vec literal at `harness.rs:1219`** + +Replace: + +```rust + h.run(vec![ + vec![(Timestamp(1), Scalar::I64(7))], + vec![(Timestamp(2), Scalar::F64(1.5))], + vec![(Timestamp(3), Scalar::Bool(true))], + vec![(Timestamp(4), Scalar::Ts(Timestamp(99)))], + ]); +``` + +with: + +```rust + h.run(vec![ + Box::new(VecSource::new(vec![(Timestamp(1), Scalar::I64(7))])), + Box::new(VecSource::new(vec![(Timestamp(2), Scalar::F64(1.5))])), + Box::new(VecSource::new(vec![(Timestamp(3), Scalar::Bool(true))])), + Box::new(VecSource::new(vec![(Timestamp(4), Scalar::Ts(Timestamp(99)))])), + ]); +``` + +- [ ] **Step 6: Wrap the 10 `blueprint.rs` call sites** + +Apply the wrap rule to each (match on the expression): + +```text +:1023 named.run(vec![synthetic_prices()]); :1031 positional.run(vec![synthetic_prices()]); +:1742 flat.run(vec![prices.clone()]); :1749 composed.run(vec![prices]); +:1812 h.run(vec![prices]); :1829 a.run(vec![prices.clone()]); +:1835 b.run(vec![prices]); :1873 a.run(vec![prices.clone()]); +:1877 b.run(vec![prices]); :2395 h.run(vec![prices]); +``` + +Each → `.run(vec![Box::new(VecSource::new())]);`. + +- [ ] **Step 7: Wrap the `sweep.rs` and `report.rs` call sites** + +`crates/aura-engine/src/sweep.rs:263`: +`h.run(vec![synthetic_prices()]);` → `h.run(vec![Box::new(VecSource::new(synthetic_prices()))]);` + +`crates/aura-engine/src/report.rs:223` (multi-line `f64_stream` element): wrap the +single element → `h.run(vec![Box::new(VecSource::new(f64_stream(&[...])))]);` + +- [ ] **Step 8: Wrap the 4 `aura-cli/src/main.rs` call sites** + +```text +:140 h.run(vec![prices]); :283 h.run(vec![prices]); +:462 h.run(vec![prices]); :527 h.run(vec![showcase_prices()]); +``` + +Each → `h.run(vec![Box::new(VecSource::new())]);`. `main.rs` already +imports from `aura_engine`; ensure `VecSource` is in the import list (the file +imports `Harness` etc. from `aura_engine::{...}` — add `VecSource`). + +- [ ] **Step 9: Wrap the `real_bars.rs` call site** + +`crates/aura-ingest/tests/real_bars.rs:72`: +`h.run(vec![prices]);` → `h.run(vec![Box::new(VecSource::new(prices))]);` +Add `VecSource` to the `use aura_engine::{...}` import list at `real_bars.rs:11`. + +- [ ] **Step 10: Workspace build gate** + +Run: `cargo build --workspace --tests` +Expected: 0 errors (every one of the 53 callers threaded; the only `fn run` in +the workspace is `Harness::run`, so there is no method-resolution collision). + +- [ ] **Step 11: Behaviour-preserving test gate** + +Run: `cargo test --workspace` +Expected: all green — byte-identical to pre-seam. The two-run C1 determinism +pairs (`real_bars.rs` `r1`/`r2` `to_json`, `blueprint.rs` sweep-equality, +`harness.rs` `h`/`h2` pairs) stay green; no test output changes. + +- [ ] **Step 12: Lint gate** + +Run: `cargo clippy --workspace --all-targets -- -D warnings` +Expected: clean. + +--- + +## Task 3: `M1Field` + `M1FieldSource` streaming source (aura-ingest) + +**Files:** +- Modify: `crates/aura-ingest/src/lib.rs` (imports at 15-18; new items near + `load_m1_window` at ~110) +- Test: `crates/aura-ingest/src/lib.rs` (`#[cfg(test)] mod tests`) + +`decode` is a **free function** of `(field, bar)` (the spec calls it "a pure +function of `(bar, field)`"), so the hermetic decode test needs no iterator. + +- [ ] **Step 0: Move `aura-engine` to `[dependencies]`** + +`M1FieldSource` lives in `src/lib.rs` and implements `aura_engine::Source`, so +`aura-engine` must be a regular dependency — it is currently `[dev-dependencies]` +only (`Cargo.toml:17`). In `crates/aura-ingest/Cargo.toml`, add under +`[dependencies]` (after the `data-server` line): + +```toml +aura-engine = { path = "../aura-engine" } +``` + +and REMOVE the now-duplicate `aura-engine = { path = "../aura-engine" }` line +from `[dev-dependencies]` (cargo warns on a dependency present in both). Leave +`aura-std` in `[dev-dependencies]` — only the tests use it. Regular deps are +visible to tests, so `real_bars.rs` still resolves `aura_engine`. + +Run: `cargo tree -p aura-ingest -e normal -i aura-engine 2>&1 | head -2` +Expected: `aura-engine` now appears as a normal (non-dev) dependency. + +- [ ] **Step 1: Add the `SymbolChunkIter` import** + +In `crates/aura-ingest/src/lib.rs`, change line 17: + +```rust +use data_server::DataServer; +``` + +to: + +```rust +use data_server::{DataServer, SymbolChunkIter}; +``` + +- [ ] **Step 2: Add `M1Field`, `decode`, and `M1FieldSource`** + +Insert after `load_m1_window` (after its closing `}` at `lib.rs:123`): + +```rust +/// Which base column of an M1 bar a source streams (C7: a composite window is a +/// bundle of base columns; one `Source` per consumed field). +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum M1Field { + Open, + High, + Low, + Close, + Spread, + Volume, +} + +/// Project one M1 bar into `(normalized ts, scalar)` for `field`. Pure: a +/// function of `(field, bar)` only — no iterator, no clock — so it is testable +/// in isolation. Time is normalized ms→epoch-ns at this one seam (C3). +fn decode(field: M1Field, bar: &M1Parsed) -> (Timestamp, Scalar) { + let value = match field { + M1Field::Open => Scalar::F64(bar.open), + M1Field::High => Scalar::F64(bar.high), + M1Field::Low => Scalar::F64(bar.low), + M1Field::Close => Scalar::F64(bar.close), + M1Field::Spread => Scalar::F64(bar.spread), + M1Field::Volume => Scalar::I64(bar.volume), + }; + (unix_ms_to_epoch_ns(bar.time_ms), value) +} + +/// A streaming [`Source`](aura_engine::Source) over a data-server M1 window. Holds +/// at most one `Arc` chunk (a zero-copy clone of the cache's chunk) and a cursor; +/// constructs each `Scalar` per-pull; refills via `next_chunk()` when the chunk +/// drains. Resident footprint is O(one chunk), independent of window length. +pub struct M1FieldSource { + iter: SymbolChunkIter, + chunk: Option>, + pos: usize, + field: M1Field, + head: Option<(Timestamp, Scalar)>, +} + +impl M1FieldSource { + /// Open over `[from_ms, to_ms]` (inclusive Unix-ms, data-server's contract). + /// `None` when no archived file overlaps the window (unknown symbol / + /// far-future window) — propagating `stream_m1_windowed`'s file-level Option. + /// A window that overlaps a file but holds zero matching bars yields a source + /// whose first `peek` is `None` (immediately exhausted), not `None` here. + pub fn open( + server: &Arc, + symbol: &str, + from_ms: Option, + to_ms: Option, + field: M1Field, + ) -> Option { + let iter = server.stream_m1_windowed(symbol, from_ms, to_ms)?; + let mut s = Self { iter, chunk: None, pos: 0, field, head: None }; + s.advance(); + Some(s) + } + + /// Decode the head at the cursor, refilling chunks as needed. Sets + /// `self.head = None` at exhaustion. + fn advance(&mut self) { + loop { + if self.chunk.is_none() { + self.chunk = self.iter.next_chunk(); + self.pos = 0; + if self.chunk.is_none() { + self.head = None; + return; + } + } + let chunk = self.chunk.as_ref().unwrap(); + if self.pos < chunk.len() { + self.head = Some(decode(self.field, &chunk[self.pos])); + return; + } + self.chunk = None; // current chunk drained — loop to refill + } + } + + /// Records resident in this source right now: the current chunk's length (or + /// 0 at exhaustion). The residency probe — bounded by one chunk length by + /// construction (the source holds at most one `Arc` chunk and no accumulating + /// field), so it can never grow with window length. + pub fn resident_records(&self) -> usize { + self.chunk.as_ref().map_or(0, |c| c.len()) + } +} + +impl aura_engine::Source for M1FieldSource { + fn peek(&self) -> Option { + self.head.map(|(t, _)| t) + } + fn next(&mut self) -> Option<(Timestamp, Scalar)> { + let item = self.head?; + self.pos += 1; + self.advance(); + Some(item) + } +} +``` + +(`aura-engine` is made a regular dependency in Step 0, so `aura_engine::Source` +resolves from `src/lib.rs`.) + +- [ ] **Step 4: Add hermetic `decode` unit tests** + +Append inside the existing `#[cfg(test)] mod tests` block in +`aura-ingest/src/lib.rs` (it already has `use super::*;` and a `full_bar` +helper at `lib.rs:142`): + +```rust +#[test] +fn decode_projects_each_field_to_its_scalar_kind() { + // full_bar: open 1.0, high 2.0, low 0.5, close 1.5, spread 0.1, volume 100. + let bar = full_bar(1_000); + assert_eq!(decode(M1Field::Open, &bar), (Timestamp(1_000_000_000), Scalar::F64(1.0))); + assert_eq!(decode(M1Field::High, &bar), (Timestamp(1_000_000_000), Scalar::F64(2.0))); + assert_eq!(decode(M1Field::Low, &bar), (Timestamp(1_000_000_000), Scalar::F64(0.5))); + assert_eq!(decode(M1Field::Close, &bar), (Timestamp(1_000_000_000), Scalar::F64(1.5))); + assert_eq!(decode(M1Field::Spread, &bar), (Timestamp(1_000_000_000), Scalar::F64(0.1))); + // volume is the one i64 column. + assert_eq!(decode(M1Field::Volume, &bar), (Timestamp(1_000_000_000), Scalar::I64(100))); +} +``` + +- [ ] **Step 5: Build + run the decode tests + lint** + +Run: `cargo test -p aura-ingest decode_projects` +Expected: PASS (`decode_projects_each_field_to_its_scalar_kind`). + +Run: `cargo clippy -p aura-ingest --all-targets -- -D warnings` +Expected: clean. + +--- + +## Task 4: Gated streaming e2e + residency proof (aura-ingest) + +**Files:** +- Create: `crates/aura-ingest/tests/streaming_seam.rs` + +Mirrors the gated pattern of `real_bars.rs` (skip where local data is absent). +Three tests: the streaming backtest e2e, the residency-drive predicate, and the +N-field sharing property. + +- [ ] **Step 1: Write the gated streaming test file** + +Create `crates/aura-ingest/tests/streaming_seam.rs`: + +```rust +//! Gated integration test for the Source ingestion seam: a real data-server M1 +//! close stream driven LAZILY through the signal-quality sample harness via +//! `M1FieldSource`, plus the residency predicate (peak resident records ≤ one +//! chunk, independent of window length) and the N-field sharing property. Skips +//! where the local Pepperstone data directory is absent. + +use std::sync::mpsc; +use std::sync::Arc; + +use aura_core::{Firing, NodeSchema, PortSpec, Scalar, ScalarKind, Timestamp}; +use aura_engine::{ + f64_field, summarize, Edge, FlatGraph, Harness, RunManifest, RunReport, Source, SourceSpec, + Target, VecSource, +}; +use aura_ingest::{M1Field, M1FieldSource}; +use aura_std::{Exposure, Recorder, SimBroker, Sma, Sub}; +use data_server::loader::CHUNK_SIZE; +use data_server::{DataServer, DEFAULT_DATA_PATH}; + +const SYMBOL: &str = "AAPL.US"; +// 2006-08 in inclusive Unix-ms (a full month of M1 ⇒ many > CHUNK_SIZE bars). +const FROM_MS: i64 = 1_154_390_400_000; +const TO_MS: i64 = 1_157_068_799_999; + +/// Bootstrap the cycle-0007 two-sink signal-quality harness and run it on a +/// single boxed source, folding the recorded equity + exposure into a RunReport. +fn run_sample(window: (Timestamp, Timestamp), source: Box) -> RunReport { + let (tx_eq, rx_eq) = mpsc::channel(); + let (tx_ex, rx_ex) = mpsc::channel(); + let f64_recorder_sig = || NodeSchema { + inputs: vec![PortSpec { kind: ScalarKind::F64, firing: Firing::Any, name: "in".into() }], + output: vec![], + params: vec![], + }; + let mut h = Harness::bootstrap(FlatGraph { + nodes: vec![ + Box::new(Sma::new(2)), + Box::new(Sma::new(4)), + Box::new(Sub::new()), + Box::new(Exposure::new(0.5)), + Box::new(SimBroker::new(0.0001)), + Box::new(Recorder::new(&[ScalarKind::F64], Firing::Any, tx_eq)), + Box::new(Recorder::new(&[ScalarKind::F64], Firing::Any, tx_ex)), + ], + signatures: vec![ + Sma::builder().schema().clone(), + Sma::builder().schema().clone(), + Sub::builder().schema().clone(), + Exposure::builder().schema().clone(), + SimBroker::builder(0.0001).schema().clone(), + f64_recorder_sig(), + f64_recorder_sig(), + ], + sources: vec![SourceSpec { + kind: ScalarKind::F64, + targets: vec![ + Target { node: 0, slot: 0 }, + Target { node: 1, slot: 0 }, + Target { node: 4, slot: 1 }, + ], + }], + edges: vec![ + Edge { from: 0, to: 2, slot: 0, from_field: 0 }, + Edge { from: 1, to: 2, slot: 1, from_field: 0 }, + Edge { from: 2, to: 3, slot: 0, from_field: 0 }, + Edge { from: 3, to: 4, slot: 0, from_field: 0 }, + Edge { from: 4, to: 5, slot: 0, from_field: 0 }, + Edge { from: 3, to: 6, slot: 0, from_field: 0 }, + ], + }) + .expect("valid signal-quality DAG"); + + h.run(vec![source]); + + let eq_rows: Vec<(Timestamp, Vec)> = rx_eq.try_iter().collect(); + let ex_rows: Vec<(Timestamp, Vec)> = rx_ex.try_iter().collect(); + let equity = f64_field(&eq_rows, 0); + let exposure = f64_field(&ex_rows, 0); + RunReport { + manifest: RunManifest { + commit: "streaming-seam-test".to_string(), + params: vec![ + ("sma_fast".to_string(), 2.0), + ("sma_slow".to_string(), 4.0), + ("exposure_scale".to_string(), 0.5), + ], + window, + seed: 0, + broker: "sim-optimal(pip_size=0.0001)".to_string(), + }, + metrics: summarize(&equity, &exposure), + } +} + +fn skip_if_no_data(server: &Arc) -> bool { + if !server.has_symbol(SYMBOL) { + eprintln!("skip: no local data at {DEFAULT_DATA_PATH} (symbol {SYMBOL} absent)"); + return true; + } + false +} + +#[test] +fn streaming_close_source_backtests_end_to_end_deterministically() { + let server = Arc::new(DataServer::new(DEFAULT_DATA_PATH)); + if skip_if_no_data(&server) { + return; + } + let open_close = || { + M1FieldSource::open(&server, SYMBOL, Some(FROM_MS), Some(TO_MS), M1Field::Close) + .expect("AAPL.US has data in the 2006-08 window") + }; + + // window bounds for the manifest: drain a source once to read first/last ts. + let mut probe = open_close(); + let first = probe.peek().expect("non-empty window"); + let mut last = first; + while let Some((t, _)) = Source::next(&mut probe) { + last = t; + } + let window = (first, last); + + let r1 = run_sample(window, Box::new(open_close())); + assert!(r1.metrics.total_pips.is_finite(), "a backtest ran over the streamed window"); + + // same window streamed again ⇒ bit-identical report (C1). + let r2 = run_sample(window, Box::new(open_close())); + assert_eq!(r1.to_json(), r2.to_json()); +} + +#[test] +fn residency_is_bounded_by_one_chunk_independent_of_window_length() { + let server = Arc::new(DataServer::new(DEFAULT_DATA_PATH)); + if skip_if_no_data(&server) { + return; + } + let mut src = M1FieldSource::open(&server, SYMBOL, Some(FROM_MS), Some(TO_MS), M1Field::Close) + .expect("AAPL.US has data in the 2006-08 window"); + + // Drive the source the way `run` does (peek to pick, next to pop), sampling + // resident records at every pull. + let mut peak = 0usize; + let mut total = 0usize; + while src.peek().is_some() { + peak = peak.max(src.resident_records()); + Source::next(&mut src); + total += 1; + } + + // multi-chunk window: more records than a single chunk holds, so the bound + // below is non-vacuous (a O(window) source would have peaked at `total`). + assert!(total > CHUNK_SIZE, "window must span multiple chunks (got {total} records)"); + // O(one chunk), NOT O(window): the per-pull ceiling never exceeds one chunk, + // regardless of how many chunks the window spans. + assert!(peak <= CHUNK_SIZE, "resident records {peak} exceeded one chunk ({CHUNK_SIZE})"); + assert!(peak > 0, "a non-empty window resided at least one record"); +} + +#[test] +fn two_field_sources_share_one_window() { + let server = Arc::new(DataServer::new(DEFAULT_DATA_PATH)); + if skip_if_no_data(&server) { + return; + } + // close + volume over the same window: each is an independent Source pulling + // its field; both stream to exhaustion (the N-sources-share-the-ts-axis shape). + let mut close = M1FieldSource::open(&server, SYMBOL, Some(FROM_MS), Some(TO_MS), M1Field::Close) + .expect("close source"); + let mut volume = M1FieldSource::open(&server, SYMBOL, Some(FROM_MS), Some(TO_MS), M1Field::Volume) + .expect("volume source"); + + let mut n_close = 0usize; + while let Some((_, v)) = Source::next(&mut close) { + assert!(matches!(v, Scalar::F64(_)), "close is an f64 column"); + n_close += 1; + } + let mut n_volume = 0usize; + while let Some((_, v)) = Source::next(&mut volume) { + assert!(matches!(v, Scalar::I64(_)), "volume is an i64 column"); + n_volume += 1; + } + assert_eq!(n_close, n_volume, "both fields share the same ts axis ⇒ same count"); + assert!(n_close > 0); +} +``` + +(The `VecSource` import is retained for parity even though this file uses only +`M1FieldSource`; drop it if clippy flags it unused — see Step 3.) + +- [ ] **Step 2: Run the gated streaming tests** + +Run: `cargo test -p aura-ingest --test streaming_seam` +Expected (where local data is present): PASS, 3 tests +(`streaming_close_source_backtests_end_to_end_deterministically`, +`residency_is_bounded_by_one_chunk_independent_of_window_length`, +`two_field_sources_share_one_window`). +Expected (where local data is absent): all 3 print `skip:` and pass (hermetic). + +- [ ] **Step 3: Lint + unused-import scrub** + +Run: `cargo clippy -p aura-ingest --all-targets -- -D warnings` +Expected: clean. If `VecSource` (or any import) is flagged unused in +`streaming_seam.rs`, remove it from the `use aura_engine::{...}` line. + +- [ ] **Step 4: Full workspace gate** + +Run: `cargo test --workspace` +Expected: all green. + +Run: `cargo clippy --workspace --all-targets -- -D warnings` +Expected: clean. + +Run: `cargo doc --workspace --no-deps 2>&1` +Expected: clean (the new `pub` items — `Source`, `VecSource`, `M1Field`, +`M1FieldSource` — carry doc comments; no missing-doc or broken-intra-doc-link).