From 2ffcd9f03dfeea397be21f86aab58b976a5305dc Mon Sep 17 00:00:00 2001 From: Brummel Date: Thu, 18 Jun 2026 12:26:18 +0200 Subject: [PATCH] =?UTF-8?q?plan:=20real-data=20source-open=20seam=20?= =?UTF-8?q?=E2=80=94=20execution=20plan?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Bite-sized, placeholder-free projection of spec 0052 for the implement skill: Task 1 the additive aura-ingest surface + hermetic round-trip test; Task 2 a gated behaviour-preservation test (open_ohlc Timestamp path == ms-path open); Task 3 the atomic compile-driven consumer migration (shared retype + 6 GER40 consumers); Task 4 the compiler-invisible compare.rs + acceptance grep gates. refs #80, #81, #92 --- docs/plans/0052-real-data-source-open-seam.md | 773 ++++++++++++++++++ 1 file changed, 773 insertions(+) create mode 100644 docs/plans/0052-real-data-source-open-seam.md diff --git a/docs/plans/0052-real-data-source-open-seam.md b/docs/plans/0052-real-data-source-open-seam.md new file mode 100644 index 0000000..3819adc --- /dev/null +++ b/docs/plans/0052-real-data-source-open-seam.md @@ -0,0 +1,773 @@ +# Real-data source-open seam — Implementation Plan + +> **Parent spec:** `docs/specs/0052-real-data-source-open-seam.md` +> +> **For agentic workers:** REQUIRED SUB-SKILL: use the `implement` skill to run +> this plan. Steps use `- [ ]` checkboxes for tracking. + +**Goal:** Consolidate the `aura-ingest` source-open surface so a real OHLC bundle +builds from `aura-ingest` alone in the engine's native epoch-ns `Timestamp` +currency, removing the per-consumer `ns_to_ms` hand-divide (#80), the duplicated +order-sensitive OHLC opener (#92), and the direct `data_server` reach (#81) — +behaviour-preserving (byte-identical `RunReport`s). + +**Architecture:** Four additive items in `crates/aura-ingest/src/lib.rs` (private +`epoch_ns_to_unix_ms`, `M1FieldSource::open_window`, free `open_ohlc`, +`default_data_server` + `pub use data_server::{DataServer, DEFAULT_DATA_PATH}`) +own the one ms↔ns crossing. Then the consumer layer (the shared +`breakout_real.rs` helpers + every GER40 example/test) migrates to the +Timestamp-native surface. The shared helper file is `#[path] mod`-included into +every consumer target, so deleting `open_ohlc_sources` / retyping +`utc_month_window_ms` + `report_from_trace` breaks all six compile-driven +consumers at once — they migrate in one atomic task. `ger40_breakout_compare.rs` +opens OHLC by its own hand-loop and calls neither retyped helper, so it survives +the build unmigrated (the compiler-invisible site) and is migrated by inspection +in its own task. + +**Tech Stack:** Rust, `aura-ingest` crate (lib + examples + integration tests), +`cargo build/test/clippy --workspace`. Gated tests skip-with-note when the local +Pepperstone archive is absent. + +--- + +## Files this plan creates or modifies + +- Modify: `crates/aura-ingest/src/lib.rs` — the four additive items + round-trip test (Task 1). +- Create: `crates/aura-ingest/tests/open_ohlc_seam.rs` — gated behaviour-preservation test (Task 2). +- Modify: `crates/aura-ingest/examples/shared/breakout_real.rs` — remove `open_ohlc_sources`, retype `utc_month_window_ms`→`utc_month_window` and `report_from_trace` (Task 3). +- Modify: `crates/aura-ingest/examples/ger40_breakout_real.rs` (Task 3). +- Modify: `crates/aura-ingest/examples/ger40_breakout_sweep.rs` (Task 3). +- Modify: `crates/aura-ingest/examples/ger40_breakout_walkforward.rs` — also delete local `ns_to_ms` (Task 3). +- Modify: `crates/aura-ingest/tests/ger40_breakout_world.rs` — also delete local `ns_to_ms`, two call sites (Task 3). +- Modify: `crates/aura-ingest/tests/ger40_breakout_real.rs` (Task 3). +- Modify: `crates/aura-ingest/tests/ger40_breakout_blueprint.rs` — two call sites (Task 3). +- Modify: `crates/aura-ingest/examples/ger40_breakout_compare.rs` — compiler-invisible, by inspection (Task 4). + +**Out of scope (do NOT touch):** the `WindowRoller` origin/end constructions +(`ger40_breakout_walkforward.rs:100-105`, `tests/ger40_breakout_world.rs:165-168`) +build the roll span from chrono via the public `unix_ms_to_epoch_ns` — the +sanctioned ms→ns *entry* of wall-clock into engine currency, not the forbidden +consumer-side ns→ms divide. They already produce `Timestamp` and are not +perturbed by the retype. The gated single-field ms-path tests (`real_bars.rs`, +`unbounded_window.rs`, `ger40_archive_utc_dst.rs`, `streaming_seam.rs`) keep their +`data_server` import (out of scope per spec AC1). + +--- + +## Task 1: `aura-ingest` library surface — the four additive items + +**Files:** +- Modify: `crates/aura-ingest/src/lib.rs` + +This task is purely additive: the existing ms-based `open` and +`unix_ms_to_epoch_ns` are untouched, so nothing else in the workspace breaks. + +- [ ] **Step 1: Reconcile the `data_server` import to add the re-export** + +In `crates/aura-ingest/src/lib.rs`, replace the import line 25: + +```rust +use data_server::{DataServer, SymbolChunkIter}; +``` + +with (keeps `SymbolChunkIter` for the struct field; `DataServer` now arrives via a +public re-export so internal code still sees it): + +```rust +use data_server::SymbolChunkIter; + +/// Re-export of the data-server archive entry points (#81): a real-data source +/// builds from `aura-ingest` alone — a consumer never names the external +/// `data_server` crate directly. +pub use data_server::{DataServer, DEFAULT_DATA_PATH}; +``` + +- [ ] **Step 2: Add the seam-private inverse `epoch_ns_to_unix_ms`** + +Immediately after `unix_ms_to_epoch_ns` (right after its closing brace at +`lib.rs:34`), add: + +```rust +/// Inverse of [`unix_ms_to_epoch_ns`]: project aura's canonical epoch-ns +/// [`Timestamp`] back to data-server's Unix-millisecond time. **Private** — the +/// seam owns the ms↔ns convention (C3); a consumer threads `Timestamp`s and never +/// converts. Floor division by 1e6 ns/ms (archived M1 data has no sub-ms instant). +fn epoch_ns_to_unix_ms(ts: Timestamp) -> i64 { + ts.0 / 1_000_000 +} +``` + +- [ ] **Step 3: Add `M1FieldSource::open_window`** + +In the `impl M1FieldSource` block, immediately after `open` (after its closing +brace at `lib.rs:217`, before the `advance` doc comment), add: + +```rust + /// Open over the `[from, to]` window in aura's native epoch-ns [`Timestamp`] + /// currency — the engine-side mirror of [`open`](Self::open), which takes + /// data-server's Unix-ms. Each bound is mapped through the seam-private + /// `epoch_ns_to_unix_ms` and delegated to `open`, so the ms↔ns crossing + /// happens once, here, and a consumer never divides. `None` bounds preserve + /// open-ended windows exactly as `open` does. + pub fn open_window( + server: &Arc, + symbol: &str, + from: Option, + to: Option, + field: M1Field, + ) -> Option { + Self::open( + server, + symbol, + from.map(epoch_ns_to_unix_ms), + to.map(epoch_ns_to_unix_ms), + field, + ) + } +``` + +- [ ] **Step 4: Add the free fns `open_ohlc` and `default_data_server`** + +Immediately after the `impl aura_engine::Source for M1FieldSource` block closes +(after `lib.rs:272`, before the `#[cfg(test)]` line at `lib.rs:274`), add: + +```rust +/// Open the four real OHLC [`M1FieldSource`]s for `symbol` over the closed +/// epoch-ns window `[from, to]`, in the FIXED order open, high, low, close — the +/// C4 merge tie-break order a resampler's `Barrier(0)` group depends on (#92). +/// This is the single vetted home of that order; consumers never spell it out. +/// All four sources share the one `Arc` (one `FileCache`), so a +/// window's bars are parsed once and reused across the four field decodes, and +/// the same `Arc` flows across the disjoint sims of a sweep / walk-forward (C12). +/// +/// Returns `None` if any field has no archived file overlapping the window (the +/// caller skips cleanly), propagating each [`open_window`](M1FieldSource::open_window)'s +/// file-level `Option`. A window that overlaps a file but holds zero bars yields +/// four sources whose first `peek` is `None`, not a `None` here. +pub fn open_ohlc( + server: &Arc, + symbol: &str, + from: Timestamp, + to: Timestamp, +) -> Option>> { + let open = M1FieldSource::open_window(server, symbol, Some(from), Some(to), M1Field::Open)?; + let high = M1FieldSource::open_window(server, symbol, Some(from), Some(to), M1Field::High)?; + let low = M1FieldSource::open_window(server, symbol, Some(from), Some(to), M1Field::Low)?; + let close = M1FieldSource::open_window(server, symbol, Some(from), Some(to), M1Field::Close)?; + Some(vec![Box::new(open), Box::new(high), Box::new(low), Box::new(close)]) +} + +/// Construct the default data-server over the local archive at +/// [`DEFAULT_DATA_PATH`], wrapped in the `Arc` the streaming sources share (#81). +/// Build it once and clone the `Arc` across a family's sims — one cache, never +/// one server per field (C12). +pub fn default_data_server() -> Arc { + Arc::new(DataServer::new(DEFAULT_DATA_PATH)) +} +``` + +- [ ] **Step 5: Add the hermetic round-trip test** + +In `mod tests`, immediately after `unix_ms_to_epoch_ns_scales_ms_to_ns` (after its +closing brace at `lib.rs:287`), add: + +```rust + #[test] + fn epoch_ns_to_unix_ms_inverts_unix_ms_to_epoch_ns() { + // The seam-owned inverse round-trips the forward normalization for every + // ms instant, so a Timestamp window bound fed back to data-server's ms + // contract recovers the exact ms (C3: one currency crossing, owned here). + for ms in [0_i64, 1, 1_488_326_400_000, 1_727_000_000_000] { + assert_eq!(epoch_ns_to_unix_ms(unix_ms_to_epoch_ns(ms)), ms); + } + } +``` + +- [ ] **Step 6: Build and run the lib tests** + +Run: `cargo build -p aura-ingest` +Expected: clean build, 0 errors. + +Run: `cargo test -p aura-ingest --lib` +Expected: `test result: ok.` — includes +`epoch_ns_to_unix_ms_inverts_unix_ms_to_epoch_ns ... ok` and +`unix_ms_to_epoch_ns_scales_ms_to_ns ... ok`. + +--- + +## Task 2: gated behaviour-preservation test (AC5) + +**Files:** +- Create: `crates/aura-ingest/tests/open_ohlc_seam.rs` + +Proves the consolidated Timestamp opener reproduces the untouched ms-path `open` +bit-for-bit. Depends only on Task 1; written against the unchanged `build_harness` +/ `drain_trace` shared helpers and the ms-path `open`, so it compiles and passes +both before and after the Task 3 consumer migration. + +- [ ] **Step 1: Write the test file** + +Create `crates/aura-ingest/tests/open_ohlc_seam.rs` with exactly: + +```rust +//! Gated behaviour-preservation proof for the consolidated OHLC opener (#80/#92): +//! the library `open_ohlc` over an epoch-ns `Timestamp` window produces a recorded +//! `(held, equity)` series BIT-IDENTICAL to opening the four `M1FieldSource`s by +//! hand over the equivalent Unix-ms window via the untouched ms-path `open`. Same +//! instants, same bars, same numbers — the migration changes the surface, not the +//! behaviour. +//! +//! Mirrors the other gated GER40 tests: skips with a note where the local +//! Pepperstone archive is absent, so `cargo test --workspace` stays green +//! anywhere; exercises the real path where the files exist. + +use std::sync::Arc; + +use aura_engine::{Harness, Source}; +use chrono::TimeZone; +use data_server::{DataServer, DEFAULT_DATA_PATH}; + +use aura_ingest::{open_ohlc, M1Field, M1FieldSource}; + +#[path = "../examples/shared/breakout_real.rs"] +mod breakout_real; +use breakout_real::*; + +/// Run a fresh harness over `sources` and drain the recorded `(held, equity)` +/// series (the same fold the determinism tests use). +fn series_from(mut h: Harness, taps: &Taps, sources: Vec>) -> Vec<(f64, f64)> { + h.run(sources); + drop(h); + let trace = drain_trace(taps); + trace.iter().map(|b| (b.held, b.equity)).collect() +} + +/// Property: the Timestamp `open_ohlc` path and the hand-rolled ms-path `open` +/// (×4, in O/H/L/C order) over the equivalent window yield bit-identical recorded +/// series — behaviour-preserving consolidation (AC5). +#[test] +fn open_ohlc_timestamp_path_matches_ms_path_open() { + let server = Arc::new(DataServer::new(DEFAULT_DATA_PATH)); + if !server.has_symbol(SYMBOL) { + eprintln!("skip: no local data at {DEFAULT_DATA_PATH} (symbol {SYMBOL} absent)"); + return; + } + + // September 2024 (inclusive UTC), the window the other gated tests pin. + let from_ms = chrono::Utc.with_ymd_and_hms(2024, 9, 1, 0, 0, 0).unwrap().timestamp_millis(); + let to_ms = chrono::Utc.with_ymd_and_hms(2024, 9, 30, 23, 59, 59).unwrap().timestamp_millis(); + let from_ts = aura_ingest::unix_ms_to_epoch_ns(from_ms); + let to_ts = aura_ingest::unix_ms_to_epoch_ns(to_ms); + + // Path A: the consolidated Timestamp opener. + let (h_a, taps_a) = build_harness(); + let sources_a = open_ohlc(&server, SYMBOL, from_ts, to_ts).expect("window overlaps GER40 data"); + let series_a = series_from(h_a, &taps_a, sources_a); + assert!(!series_a.is_empty(), "window resolved to zero bars"); + + // Path B: the untouched ms-path open, four fields by hand in O/H/L/C order. + let fields = [M1Field::Open, M1Field::High, M1Field::Low, M1Field::Close]; + let mut sources_b: Vec> = Vec::with_capacity(4); + for &f in &fields { + let s = M1FieldSource::open(&server, SYMBOL, Some(from_ms), Some(to_ms), f) + .expect("window overlaps GER40 data"); + sources_b.push(Box::new(s)); + } + let (h_b, taps_b) = build_harness(); + let series_b = series_from(h_b, &taps_b, sources_b); + + assert_eq!( + series_a, series_b, + "open_ohlc (Timestamp) must reproduce the ms-path open bit-identically (behaviour-preserving, #80/#92)", + ); +} +``` + +- [ ] **Step 2: Build and run the new test target** + +Run: `cargo test -p aura-ingest --test open_ohlc_seam` +Expected: `test result: ok. 1 passed` (the test runs where the archive exists; it +prints a `skip:` note and still passes where the archive is absent). + +--- + +## Task 3: atomic consumer migration (shared retype + the six compile-driven consumers) + +**Files:** +- Modify: `crates/aura-ingest/examples/shared/breakout_real.rs` +- Modify: `crates/aura-ingest/examples/ger40_breakout_real.rs` +- Modify: `crates/aura-ingest/examples/ger40_breakout_sweep.rs` +- Modify: `crates/aura-ingest/examples/ger40_breakout_walkforward.rs` +- Modify: `crates/aura-ingest/tests/ger40_breakout_world.rs` +- Modify: `crates/aura-ingest/tests/ger40_breakout_real.rs` +- Modify: `crates/aura-ingest/tests/ger40_breakout_blueprint.rs` + +**Atomicity note for the executor:** `breakout_real.rs` is `#[path] mod`-included +into all six consumer targets. Steps 1-4 (the shared retype) make every consumer +fail to compile until Steps 5-10 thread them. Do NOT build between edit steps — +the build gate is Steps 11-13, after all seven files are edited. `compare.rs` is +deliberately NOT touched here (it calls neither retyped helper and still +compiles); it is Task 4. + +**The single migration shape applied at every call site:** +- `open_ohlc_sources(server, from_ms, to_ms)` → `open_ohlc(server, SYMBOL, from, to)` (now `Timestamp` bounds). +- `utc_month_window_ms(y, m)` → `utc_month_window(y, m)` (returns `(Timestamp, Timestamp)`). +- a local helper's `from_ms: i64, to_ms: i64` params → `from: Timestamp, to: Timestamp`; its `RunManifest.window` `(aura_ingest::unix_ms_to_epoch_ns(from_ms), …(to_ms))` → `(from, to)`. +- `Arc::new(DataServer::new(DEFAULT_DATA_PATH))` → `default_data_server()`. +- `use data_server::{DataServer, DEFAULT_DATA_PATH};` → an `aura_ingest` re-export import; `ns_to_ms(w.is.0)` etc. → `w.is.0` straight through. + +- [ ] **Step 1: shared `breakout_real.rs` — prune now-unused imports** + +Delete line 40 (`use std::sync::Arc;` — only `open_ohlc_sources` used it) and +line 50 (`use data_server::DataServer;` — same). In the `aura_engine` import +(lines 43-46) remove `Source` (only `open_ohlc_sources`'s return type used it). +The import becomes: + +```rust +use aura_engine::{ + summarize, Composite, Edge, FlatGraph, GraphBuilder, Harness, RunManifest, RunReport, + SourceSpec, Target, +}; +``` + +- [ ] **Step 2: shared `breakout_real.rs` — remove `open_ohlc_sources`** + +Delete the whole helper including its doc comment — lines 309-328 (from +`/// Open the four real OHLC …` through the closing `}` of `open_ohlc_sources`). + +- [ ] **Step 3: shared `breakout_real.rs` — retype `utc_month_window_ms` → `utc_month_window`** + +Replace the helper at lines 330-346 with: + +```rust +/// Build the inclusive window for the whole calendar month `(year, month)` in +/// UTC: `[first instant of the 1st, last ms of the last day]`, in aura's native +/// epoch-ns [`Timestamp`] currency. UTC keeps the boundary trivial to reason +/// about; the Session node still indexes bars in Berlin wall-clock inside the +/// graph. The SAME UTC instants as before, typed as `Timestamp` (each bound +/// wrapped through the one `unix_ms_to_epoch_ns` seam), so the consumer layer is +/// Timestamp-native and no caller round-trips ms→Timestamp→ms. +pub fn utc_month_window(year: i32, month: u32) -> (Timestamp, Timestamp) { + let from = chrono::Utc + .with_ymd_and_hms(year, month, 1, 0, 0, 0) + .unwrap() + .timestamp_millis(); + // first instant of the next month, minus one ms => last ms of this month. + let (ny, nm) = if month == 12 { (year + 1, 1) } else { (year, month + 1) }; + let next = chrono::Utc + .with_ymd_and_hms(ny, nm, 1, 0, 0, 0) + .unwrap() + .timestamp_millis(); + ( + aura_ingest::unix_ms_to_epoch_ns(from), + aura_ingest::unix_ms_to_epoch_ns(next - 1), + ) +} +``` + +- [ ] **Step 4: shared `breakout_real.rs` — retype `report_from_trace`** + +In the helper at lines 407-431: change the doc line 406 from +`/// the requested `[from_ms, to_ms]` normalized to epoch-ns.` to +`/// the requested `[from, to]` epoch-ns window.`; change the signature and the +`window:` field. The signature line 407 becomes: + +```rust +pub fn report_from_trace(trace: &[BarTrace], from: Timestamp, to: Timestamp) -> RunReport { +``` + +and the `window:` field (lines 422-425) becomes: + +```rust + window: (from, to), +``` + +- [ ] **Step 5: `examples/ger40_breakout_real.rs`** + +Delete line 23 (`use std::sync::Arc;`) and replace line 27 +(`use data_server::{DataServer, DEFAULT_DATA_PATH};`) with: + +```rust +use aura_ingest::{default_data_server, open_ohlc, DEFAULT_DATA_PATH}; +``` + +Line 40: `let server = Arc::new(DataServer::new(DEFAULT_DATA_PATH));` → +`let server = default_data_server();` + +Line 49: `let (from_ms, to_ms) = utc_month_window_ms(YEAR, MONTH);` → +`let (from, to) = utc_month_window(YEAR, MONTH);` + +Line 51: `let Some(sources) = open_ohlc_sources(&server, from_ms, to_ms) else {` → +`let Some(sources) = open_ohlc(&server, SYMBOL, from, to) else {` + +Line 75: `let report = report_from_trace(&trace, from_ms, to_ms);` → +`let report = report_from_trace(&trace, from, to);` + +- [ ] **Step 6: `examples/ger40_breakout_sweep.rs`** + +Replace line 21 (`use data_server::{DataServer, DEFAULT_DATA_PATH};`) with (keep +`use std::sync::Arc;` at line 16 — `Arc::clone` is still used at line 110): + +```rust +use aura_ingest::{default_data_server, open_ohlc, DataServer, DEFAULT_DATA_PATH}; +``` + +Retype `run_point` (line 36) — params `from_ms: i64, to_ms: i64` → +`from: Timestamp, to: Timestamp`: + +```rust +fn run_point(server: &Arc, point: &[Cell], from: Timestamp, to: Timestamp) -> RunReport { +``` + +Line 42: `open_ohlc_sources(server, from_ms, to_ms)` → +`open_ohlc(server, SYMBOL, from, to)` + +Lines 59-62 (the `window:` field) → `window: (from, to),` + +Line 71: `let server = Arc::new(DataServer::new(DEFAULT_DATA_PATH));` → +`let server = default_data_server();` + +Lines 78-79: +`let (from_ms, _) = utc_month_window_ms(2024, 1);` → +`let (from, _) = utc_month_window(2024, 1);` and +`let (_, to_ms) = utc_month_window_ms(2024, 12);` → +`let (_, to) = utc_month_window(2024, 12);` + +Line 112: `run_point(&server_for_closure, pt, from_ms, to_ms)` → +`run_point(&server_for_closure, pt, from, to)` + +- [ ] **Step 7: `examples/ger40_breakout_walkforward.rs`** + +Replace line 27 (`use data_server::{DataServer, DEFAULT_DATA_PATH};`) with (keep +`use std::sync::Arc;` — `Arc::clone` at line 128 still used): + +```rust +use aura_ingest::{default_data_server, open_ohlc, DataServer, DEFAULT_DATA_PATH}; +``` + +Delete the local `ns_to_ms` helper and its doc comment — lines 38-43 (from +`/// ns → ms: …` through the closing `}`). + +Retype `run_point` (line 48) — `from_ms: i64, to_ms: i64` → +`from: Timestamp, to: Timestamp`: + +```rust +fn run_point( + server: &Arc, + point: &[Cell], + from: Timestamp, + to: Timestamp, +) -> (RunReport, Vec<(Timestamp, f64)>) { +``` + +Line 59: `open_ohlc_sources(server, from_ms, to_ms)` → +`open_ohlc(server, SYMBOL, from, to)` + +Lines 76-79 (the `window:` field) → `window: (from, to),` + +Line 89: `let server = Arc::new(DataServer::new(DEFAULT_DATA_PATH));` → +`let server = default_data_server();` + +**Leave lines 100-105 unchanged** (the `WindowRoller` origin/end via +`aura_ingest::unix_ms_to_epoch_ns` — out of scope, the sanctioned wall-clock entry). + +Line 146: delete `let (is_from, is_to) = (ns_to_ms(w.is.0), ns_to_ms(w.is.1));` + +Line 148: `run_point(&server_for_closure, pt, is_from, is_to).0` → +`run_point(&server_for_closure, pt, w.is.0, w.is.1).0` + +Line 164: delete `let (oos_from, oos_to) = (ns_to_ms(w.oos.0), ns_to_ms(w.oos.1));` + +Lines 165-166: `run_point(&server_for_closure, &chosen, oos_from, oos_to)` → +`run_point(&server_for_closure, &chosen, w.oos.0, w.oos.1)` + +- [ ] **Step 8: `tests/ger40_breakout_world.rs`** + +Replace line 30 (`use data_server::{DataServer, DEFAULT_DATA_PATH};`) with (keep +`use std::sync::Arc;` — `Arc::clone` at lines 118/186 still used): + +```rust +use aura_ingest::{default_data_server, open_ohlc, DataServer, DEFAULT_DATA_PATH}; +``` + +Delete the local `ns_to_ms` helper and its doc comment — lines 36-41. + +Retype `run_point` (line 49) — `from_ms: i64, to_ms: i64` → +`from: Timestamp, to: Timestamp`: + +```rust +fn run_point(server: &Arc, point: &[Cell], from: Timestamp, to: Timestamp) -> RunReport { +``` + +Line 54: `open_ohlc_sources(server, from_ms, to_ms)` → +`open_ohlc(server, SYMBOL, from, to)` + +Lines 70-71 (the `window:` field) → `window: (from, to),` + +Line 94: `let server = Arc::new(DataServer::new(DEFAULT_DATA_PATH));` → +`let server = default_data_server();` + +Line 100: `let (from_ms, to_ms) = utc_month_window_ms(2024, 9);` → +`let (from, to) = utc_month_window(2024, 9);` + +Line 120: `run_point(&server_for_closure, pt, from_ms, to_ms)` → +`run_point(&server_for_closure, pt, from, to)` + +Line 155: `let server = Arc::new(DataServer::new(DEFAULT_DATA_PATH));` → +`let server = default_data_server();` + +**Leave lines 165-168 unchanged** (the `WindowRoller` origin/end — out of scope). + +Line 204: delete `let (is_from, is_to) = (ns_to_ms(w.is.0), ns_to_ms(w.is.1));` + +Line 206: `run_point(&server_for_closure, pt, is_from, is_to)` → +`run_point(&server_for_closure, pt, w.is.0, w.is.1)` + +Line 222: delete `let (oos_from, oos_to) = (ns_to_ms(w.oos.0), ns_to_ms(w.oos.1));` + +Line 223: `let oos_report = run_point(&server_for_closure, &chosen, oos_from, oos_to);` → +`let oos_report = run_point(&server_for_closure, &chosen, w.oos.0, w.oos.1);` + +Line 228: `let sources = open_ohlc_sources(&server_for_closure, oos_from, oos_to)` → +`let sources = open_ohlc(&server_for_closure, SYMBOL, w.oos.0, w.oos.1)` + +- [ ] **Step 9: `tests/ger40_breakout_real.rs`** + +Replace line 15 (`use data_server::{DataServer, DEFAULT_DATA_PATH};`) with (keep +`use std::sync::Arc;` — `Arc` is still named in `run_real`'s +signature; add `Timestamp`, which this file does not yet import): + +```rust +use aura_core::Timestamp; +use aura_ingest::{default_data_server, open_ohlc, DataServer, DEFAULT_DATA_PATH}; +``` + +Retype `run_real` (line 28) — `from_ms: i64, to_ms: i64` → +`from: Timestamp, to: Timestamp`: + +```rust +fn run_real( + server: &Arc, + from: Timestamp, + to: Timestamp, +) -> (aura_engine::RunReport, Vec<(f64, f64)>) { +``` + +Line 34: `open_ohlc_sources(server, from_ms, to_ms)` → +`open_ohlc(server, SYMBOL, from, to)` + +Line 42: `let report = report_from_trace(&trace, from_ms, to_ms);` → +`let report = report_from_trace(&trace, from, to);` + +Line 59: `let server = Arc::new(DataServer::new(DEFAULT_DATA_PATH));` → +`let server = default_data_server();` + +Line 65: `let (from_ms, to_ms) = utc_month_window_ms(YEAR, MONTH);` → +`let (from, to) = utc_month_window(YEAR, MONTH);` + +Line 67: `let (r1, s1) = run_real(&server, from_ms, to_ms);` → +`let (r1, s1) = run_real(&server, from, to);` + +Line 85: `let (r2, s2) = run_real(&server, from_ms, to_ms);` → +`let (r2, s2) = run_real(&server, from, to);` + +- [ ] **Step 10: `tests/ger40_breakout_blueprint.rs`** + +Replace line 21 (`use data_server::{DataServer, DEFAULT_DATA_PATH};`) with (keep +`use std::sync::Arc;` — `Arc` named in the run helpers; extend the +`aura_core` import at line 19 to add `Timestamp`): + +```rust +use aura_ingest::{default_data_server, open_ohlc, DataServer, DEFAULT_DATA_PATH}; +``` + +and change line 19 `use aura_core::Scalar;` to `use aura_core::{Scalar, Timestamp};`. + +Retype `run_blueprint` (line 75) and `run_flatgraph` (line 91) — both +`from_ms: i64, to_ms: i64` → `from: Timestamp, to: Timestamp`: + +```rust +fn run_blueprint(server: &Arc, from: Timestamp, to: Timestamp) -> Vec<(f64, f64)> { +``` +```rust +fn run_flatgraph(server: &Arc, from: Timestamp, to: Timestamp) -> Vec<(f64, f64)> { +``` + +Line 83: `open_ohlc_sources(server, from_ms, to_ms)` → +`open_ohlc(server, SYMBOL, from, to)` + +Line 94: `open_ohlc_sources(server, from_ms, to_ms)` → +`open_ohlc(server, SYMBOL, from, to)` + +Line 111: `let server = Arc::new(DataServer::new(DEFAULT_DATA_PATH));` → +`let server = default_data_server();` + +Line 117: `let (from_ms, to_ms) = utc_month_window_ms(YEAR, MONTH);` → +`let (from, to) = utc_month_window(YEAR, MONTH);` + +Lines 119, 123, 130: `run_blueprint(&server, from_ms, to_ms)` / +`run_flatgraph(&server, from_ms, to_ms)` → `(&server, from, to)` in all three. + +- [ ] **Step 11: workspace build gate** + +Run: `cargo build --workspace --all-targets` +Expected: clean build, 0 errors (every compile-driven consumer threaded; +`compare.rs` still compiles with its old hand-loop — it is migrated in Task 4). + +- [ ] **Step 12: workspace test gate** + +Run: `cargo test --workspace` +Expected: `test result: ok` across all crates, 0 failed. The gated GER40 tests +(`ger40_breakout_real`, `ger40_breakout_world`, `ger40_breakout_blueprint`, +`open_ohlc_seam`) run where the archive exists and print a `skip:` note otherwise — +either way no failures, and the C1/C23 determinism assertions prove the migrated +path is byte-identical. + +- [ ] **Step 13: clippy gate** + +Run: `cargo clippy --workspace --all-targets -- -D warnings` +Expected: 0 warnings (no unused-import warnings — the pruned `Arc` / `data_server` +/ `Source` imports of Step 1 and Steps 5-10 are accounted for). + +--- + +## Task 4: `compare.rs` by-inspection migration + acceptance verification + +**Files:** +- Modify: `crates/aura-ingest/examples/ger40_breakout_compare.rs` + +The compiler-invisible site: it opens OHLC by its own hand-loop and builds its +window inline via chrono, so nothing in Task 3 flags it. Migrated here by +inspection, then the acceptance greps verify the whole migration. + +- [ ] **Step 1: imports** + +Delete line 23 only if `Arc` becomes unused — it does NOT (`run_cell`'s signature +names `&Arc`), so **keep** `use std::sync::Arc;`. Delete line 27 +(`use chrono::TimeZone;` — its only use is the inline window removed below). +Replace line 29 (`use data_server::{DataServer, DEFAULT_DATA_PATH};`) and line 31 +(`use aura_ingest::{M1Field, M1FieldSource};`) with a single import: + +```rust +use aura_ingest::{default_data_server, open_ohlc, DataServer, DEFAULT_DATA_PATH}; +``` + +In the `aura_engine` import (line 26) remove `Source` (only the hand-loop used it): + +```rust +use aura_engine::{summarize, RunMetrics}; +``` + +- [ ] **Step 2: `run_cell` signature + replace the hand-loop with `open_ohlc`** + +Retype `run_cell` (line 55) — `from_ms: i64, to_ms: i64` → +`from: Timestamp, to: Timestamp`: + +```rust +fn run_cell( + server: &Arc, + symbol: &str, + bar_period: i64, + from: Timestamp, + to: Timestamp, +) -> Option<(RunMetrics, u32)> { +``` + +Update the doc line 51 from `/// `[from_ms, to_ms]` of real `symbol` OHLC` to +`/// `[from, to]` (epoch-ns) of real `symbol` OHLC`. + +Replace the hand-spelled loop (lines 74-80): + +```rust + let fields = [M1Field::Open, M1Field::High, M1Field::Low, M1Field::Close]; + let mut srcs: Vec> = Vec::with_capacity(4); + for &f in &fields { + let s = M1FieldSource::open(server, symbol, Some(from_ms), Some(to_ms), f)?; + srcs.push(Box::new(s)); + } + h.run(srcs); +``` + +with: + +```rust + let srcs = open_ohlc(server, symbol, from, to)?; + h.run(srcs); +``` + +- [ ] **Step 3: replace the inline chrono window with `utc_month_window`** + +Replace the inline window in `main` (lines 113-121): + +```rust + // Full year 2024 for body. + let from_ms = chrono::Utc + .with_ymd_and_hms(2024, 1, 1, 0, 0, 0) + .unwrap() + .timestamp_millis(); + let to_ms = chrono::Utc + .with_ymd_and_hms(2024, 12, 31, 23, 59, 59) + .unwrap() + .timestamp_millis(); +``` + +with (the same full-year UTC span via the shared Timestamp helper — the +end-of-year boundary widens by <1s and captures no extra minute bar): + +```rust + // Full year 2024 (the same UTC span, via the shared Timestamp helper). + let (from, _) = utc_month_window(2024, 1); + let (_, to) = utc_month_window(2024, 12); +``` + +- [ ] **Step 4: `default_data_server` + threaded call sites** + +Line 107: `let server = Arc::new(DataServer::new(DEFAULT_DATA_PATH));` → +`let server = default_data_server();` + +Line 138: `run_cell(&server, symbol, BAR_MINUTES, from_ms, to_ms)` → +`run_cell(&server, symbol, BAR_MINUTES, from, to)` + +Line 159: `run_cell(&server, SYMBOL, period, from_ms, to_ms)` → +`run_cell(&server, SYMBOL, period, from, to)` + +- [ ] **Step 5: build + clippy gate** + +Run: `cargo build --workspace --all-targets` +Expected: clean build, 0 errors. + +Run: `cargo clippy --workspace --all-targets -- -D warnings` +Expected: 0 warnings. + +- [ ] **Step 6: acceptance grep — `open_ohlc_sources` fully gone (AC6)** + +Run: `grep -rn 'open_ohlc_sources' crates/aura-ingest` +Expected: no output (exit status 1) — the helper, every call, and every prose +mention are removed. + +- [ ] **Step 7: acceptance grep — no consumer-side `ns_to_ms` (AC2)** + +Run: `grep -rn 'ns_to_ms' crates/aura-ingest/examples crates/aura-ingest/tests` +Expected: no output (exit status 1) — both local `ns_to_ms` divides are deleted. +(The seam-internal `epoch_ns_to_unix_ms` lives in `src/lib.rs`, outside this scope, +and does not contain the substring `ns_to_ms`.) + +- [ ] **Step 8: acceptance grep — no `data_server` import in any migrated consumer (AC1)** + +Run: `grep -ln 'use data_server' crates/aura-ingest/examples/ger40_breakout_*.rs crates/aura-ingest/tests/ger40_breakout_*.rs` +Expected: no output (exit status 1) — every GER40 OHLC consumer builds from +`aura-ingest` alone. (The out-of-scope single-field ms-path tests `real_bars.rs`, +`unbounded_window.rs`, `ger40_archive_utc_dst.rs`, `streaming_seam.rs` keep theirs; +they are not matched by this glob.) + +- [ ] **Step 9: acceptance grep — `epoch_ns_to_unix_ms` is private (AC4)** + +Run: `grep -n 'fn epoch_ns_to_unix_ms' crates/aura-ingest/src/lib.rs` +Expected: exactly one line, `fn epoch_ns_to_unix_ms(ts: Timestamp) -> i64 {` — +**not** prefixed with `pub`. (A `grep -n 'pub fn epoch_ns_to_unix_ms'` returns +nothing.) + +- [ ] **Step 10: final full-suite gate** + +Run: `cargo test --workspace` +Expected: `test result: ok` across all crates, 0 failed.