feat(aura-ingest): data-server M1 ingestion at the one merge boundary
aura's first real data source (closes #7, Walking-skeleton milestone). A new aura-ingest workspace crate transposes data-server's AoS M1Parsed records into SoA base columns (C7), normalizes Unix-ms to canonical epoch-ns at the one ingestion boundary (C3), and feeds the existing k-way merge a real close-price stream so a backtest runs over real bars, deterministically (C1). Surface: - unix_ms_to_epoch_ns(time_ms) -> Timestamp (= ms * 1_000_000), the single C3 unit normalization. - M1Columns: the OHLCV bar as a bundle of base columns (open/high/low/close/ spread: f64, volume: i64, ts: epoch-ns) — SoA, C7. - transpose_m1(&[M1Parsed]) -> M1Columns: pure AoS->SoA (C1). - M1Columns::close_stream() -> Vec<(Timestamp, Scalar)>: the price input the SMA-cross sample strategy consumes. - load_m1_window(server, symbol, from_ms, to_ms): drains data-server's chronological chunks, transposes once at the boundary. Design (per spec 0011): - Eager materialization, not a lazy/shared Source abstraction: C12's cross-sim Arc<[T]> sharing has no consumer until the multi-sim orchestration cycle; deferred, the transpose logic carries over. - aura-ingest is the external-dependency firewall. data-server transitively pulls chrono + regex + zip (+ their trees) into Cargo.lock; isolating it in this one crate keeps aura-core/std/engine zero-external-dependency and the frozen deploy artifact (C6 replays recorded streams, never re-ingests) clean. cargo test --workspace now needs a populated cargo cache (one Gitea fetch, done). - Scope: boundary + tests only, no aura run CLI arg surface (M1 first; tick and the CLI wiring are follow-ups). Tests: 6 hermetic unit tests on hand-built M1Parsed (normalization, field-wise transpose, purity, close_stream order, empty edges) + one gated integration test (tests/real_bars.rs) that runs the cycle-0007 sample harness over real AAPL.US 2006-08 close bars and asserts finite metrics + two-run bit-identical JSON (C1); it skips cleanly where /mnt/tickdata is absent. On this machine it ran the real path. Workspace gates green: test (86), clippy -D warnings, doc -D warnings. Minor: transient unused-import warnings in the new lib.rs across the additive build steps resolved once load_m1_window consumed the imports; final -D warnings gate clean.
This commit is contained in:
@@ -0,0 +1,18 @@
|
||||
[package]
|
||||
name = "aura-ingest"
|
||||
edition.workspace = true
|
||||
version.workspace = true
|
||||
license.workspace = true
|
||||
publish.workspace = true
|
||||
|
||||
[dependencies]
|
||||
aura-core = { path = "../aura-core" }
|
||||
# data-server: aura's first real data source AND its one external dependency
|
||||
# (transitively chrono + regex + zip). Isolated in this crate — the firewall
|
||||
# that keeps aura-core/std/engine zero-external-dep (spec §Architecture).
|
||||
data-server = { git = "http://192.168.178.103:3000/Brummel/data-server.git", branch = "main" }
|
||||
|
||||
[dev-dependencies]
|
||||
# the integration test bootstraps a sample harness + folds a RunReport
|
||||
aura-engine = { path = "../aura-engine" }
|
||||
aura-std = { path = "../aura-std" }
|
||||
@@ -0,0 +1,190 @@
|
||||
//! aura's first real data source: the C3/C7 ingestion boundary.
|
||||
//!
|
||||
//! Transposes [`data_server`]'s Array-of-Structs `M1Parsed` records into aura's
|
||||
//! Structure-of-Arrays base columns (C7) and normalizes data-server's
|
||||
//! Unix-millisecond time to aura's canonical epoch-nanosecond [`Timestamp`] at
|
||||
//! this one boundary (C3). A transposed [`M1Columns`] exposes its close column
|
||||
//! as the engine source-stream shape (`Vec<(Timestamp, Scalar)>`) the SMA-cross
|
||||
//! sample strategy consumes.
|
||||
//!
|
||||
//! This crate is the workspace's **external-dependency firewall**: it is the
|
||||
//! only crate that links `data-server` (and its transitive `chrono`/`regex`/
|
||||
//! `zip`), so `aura-core`/`aura-std`/`aura-engine` stay zero-external-dependency.
|
||||
|
||||
use aura_core::{Scalar, Timestamp};
|
||||
use data_server::records::M1Parsed;
|
||||
use data_server::DataServer;
|
||||
use std::sync::Arc;
|
||||
|
||||
/// Normalize data-server's Unix-millisecond time to aura's canonical epoch-ns
|
||||
/// [`Timestamp`] — the single unit normalization of C3, performed at the one
|
||||
/// ingestion boundary and nowhere else. The `i64` epoch-ns range covers all
|
||||
/// market data through year ~2262, well beyond any real file.
|
||||
pub fn unix_ms_to_epoch_ns(time_ms: i64) -> Timestamp {
|
||||
Timestamp(time_ms * 1_000_000)
|
||||
}
|
||||
|
||||
/// One M1 window transposed Array-of-Structs → Structure-of-Arrays (C7): the
|
||||
/// OHLCV bar as a bundle of base columns, time already normalized to epoch-ns.
|
||||
/// `volume` is the one `i64` column; the price/spread columns are `f64`. All
|
||||
/// columns share an index: `ts[i]` is the timestamp of `close[i]`, etc.
|
||||
#[derive(Clone, Debug, PartialEq)]
|
||||
pub struct M1Columns {
|
||||
pub ts: Vec<Timestamp>,
|
||||
pub open: Vec<f64>,
|
||||
pub high: Vec<f64>,
|
||||
pub low: Vec<f64>,
|
||||
pub close: Vec<f64>,
|
||||
pub spread: Vec<f64>,
|
||||
pub volume: Vec<i64>,
|
||||
}
|
||||
|
||||
impl M1Columns {
|
||||
/// Pre-size all columns for `n` bars.
|
||||
fn with_capacity(n: usize) -> Self {
|
||||
Self {
|
||||
ts: Vec::with_capacity(n),
|
||||
open: Vec::with_capacity(n),
|
||||
high: Vec::with_capacity(n),
|
||||
low: Vec::with_capacity(n),
|
||||
close: Vec::with_capacity(n),
|
||||
spread: Vec::with_capacity(n),
|
||||
volume: Vec::with_capacity(n),
|
||||
}
|
||||
}
|
||||
|
||||
/// The close column as an engine source stream: each normalized timestamp
|
||||
/// zipped with its close as a [`Scalar::F64`]. Ascending in timestamp iff
|
||||
/// the bars were (the C3 ingestion precondition `Harness::run` relies on).
|
||||
/// This is the price input the SMA-cross sample strategy consumes; a project
|
||||
/// that wants another field reads the public columns and zips its own.
|
||||
pub fn close_stream(&self) -> Vec<(Timestamp, Scalar)> {
|
||||
self.ts
|
||||
.iter()
|
||||
.zip(&self.close)
|
||||
.map(|(&t, &v)| (t, Scalar::F64(v)))
|
||||
.collect()
|
||||
}
|
||||
}
|
||||
|
||||
/// Transpose data-server's AoS M1 records into aura's SoA columns (C7),
|
||||
/// normalizing time at the boundary (C3). Pure: identical bars yield identical
|
||||
/// columns (C1) — it reads no clock and no external state.
|
||||
pub fn transpose_m1(bars: &[M1Parsed]) -> M1Columns {
|
||||
let mut c = M1Columns::with_capacity(bars.len());
|
||||
for b in bars {
|
||||
c.ts.push(unix_ms_to_epoch_ns(b.time_ms));
|
||||
c.open.push(b.open);
|
||||
c.high.push(b.high);
|
||||
c.low.push(b.low);
|
||||
c.close.push(b.close);
|
||||
c.spread.push(b.spread);
|
||||
c.volume.push(b.volume);
|
||||
}
|
||||
c
|
||||
}
|
||||
|
||||
/// Drain a data-server M1 window into transposed SoA columns. Reads the
|
||||
/// chronological chunk iterator to exhaustion into one AoS buffer, then
|
||||
/// transposes once at the boundary (C3 — one merge point, not per chunk).
|
||||
/// `None` if the symbol has no data in `[from_ms, to_ms]` (data-server's own
|
||||
/// `None`); the inclusive Unix-ms window bounds are data-server's contract.
|
||||
pub fn load_m1_window(
|
||||
server: &Arc<DataServer>,
|
||||
symbol: &str,
|
||||
from_ms: i64,
|
||||
to_ms: i64,
|
||||
) -> Option<M1Columns> {
|
||||
let mut it = server.stream_m1_windowed(symbol, Some(from_ms), Some(to_ms))?;
|
||||
let mut bars: Vec<M1Parsed> = Vec::new();
|
||||
// M1Parsed is Copy; &Arc<[M1Parsed]> derefs to &[M1Parsed] in arg position.
|
||||
while let Some(chunk) = it.next_chunk() {
|
||||
bars.extend_from_slice(&chunk);
|
||||
}
|
||||
Some(transpose_m1(&bars))
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn unix_ms_to_epoch_ns_scales_ms_to_ns() {
|
||||
// data-server's own epoch fixture: 2017-03-01 00:00 UTC.
|
||||
assert_eq!(
|
||||
unix_ms_to_epoch_ns(1_488_326_400_000),
|
||||
Timestamp(1_488_326_400_000_000_000)
|
||||
);
|
||||
// epoch maps to epoch.
|
||||
assert_eq!(unix_ms_to_epoch_ns(0), Timestamp(0));
|
||||
}
|
||||
|
||||
/// A hand-built M1 bar with distinct per-field values so a transpose test
|
||||
/// can tell the columns apart.
|
||||
fn full_bar(time_ms: i64) -> M1Parsed {
|
||||
M1Parsed {
|
||||
time_ms,
|
||||
open: 1.0,
|
||||
high: 2.0,
|
||||
low: 0.5,
|
||||
close: 1.5,
|
||||
spread: 0.1,
|
||||
volume: 100,
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn transpose_m1_maps_every_field_to_its_column() {
|
||||
let bars = [full_bar(1_000), full_bar(2_000)];
|
||||
let c = transpose_m1(&bars);
|
||||
// ts normalized ms -> ns; each column equal length to the input.
|
||||
assert_eq!(c.ts, vec![Timestamp(1_000_000_000), Timestamp(2_000_000_000)]);
|
||||
assert_eq!(c.open, vec![1.0, 1.0]);
|
||||
assert_eq!(c.high, vec![2.0, 2.0]);
|
||||
assert_eq!(c.low, vec![0.5, 0.5]);
|
||||
assert_eq!(c.close, vec![1.5, 1.5]);
|
||||
assert_eq!(c.spread, vec![0.1, 0.1]);
|
||||
// volume is the one i64 column.
|
||||
assert_eq!(c.volume, vec![100_i64, 100_i64]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn transpose_m1_is_pure() {
|
||||
let bars = [full_bar(1_000), full_bar(2_000), full_bar(3_000)];
|
||||
assert_eq!(transpose_m1(&bars), transpose_m1(&bars));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn transpose_m1_empty_input_yields_empty_columns() {
|
||||
let c = transpose_m1(&[]);
|
||||
assert!(c.ts.is_empty());
|
||||
assert!(c.close.is_empty());
|
||||
assert!(c.volume.is_empty());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn close_stream_zips_ts_with_close_in_order() {
|
||||
// distinct close per bar so order is observable; ts ascending.
|
||||
let bars = [
|
||||
M1Parsed { time_ms: 1, open: 0.0, high: 0.0, low: 0.0, close: 10.0, spread: 0.0, volume: 0 },
|
||||
M1Parsed { time_ms: 2, open: 0.0, high: 0.0, low: 0.0, close: 20.0, spread: 0.0, volume: 0 },
|
||||
M1Parsed { time_ms: 3, open: 0.0, high: 0.0, low: 0.0, close: 30.0, spread: 0.0, volume: 0 },
|
||||
];
|
||||
let stream = transpose_m1(&bars).close_stream();
|
||||
assert_eq!(
|
||||
stream,
|
||||
vec![
|
||||
(Timestamp(1_000_000), Scalar::F64(10.0)),
|
||||
(Timestamp(2_000_000), Scalar::F64(20.0)),
|
||||
(Timestamp(3_000_000), Scalar::F64(30.0)),
|
||||
]
|
||||
);
|
||||
// ascending in timestamp (the merge precondition).
|
||||
assert!(stream.windows(2).all(|w| w[0].0 < w[1].0));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn close_stream_empty_on_empty_columns() {
|
||||
assert!(transpose_m1(&[]).close_stream().is_empty());
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,107 @@
|
||||
//! Gated integration test: a real data-server M1 close stream driven through
|
||||
//! the cycle-0007 signal-quality sample harness (SMA-cross → Exposure →
|
||||
//! SimBroker), folded into a `RunReport`. Skips with a note where the local
|
||||
//! Pepperstone data directory is absent, so `cargo test --workspace` stays green
|
||||
//! anywhere; it exercises the real ingestion path where data exists.
|
||||
|
||||
use std::sync::mpsc;
|
||||
use std::sync::Arc;
|
||||
|
||||
use aura_core::{Firing, Scalar, ScalarKind, Timestamp};
|
||||
use aura_engine::{
|
||||
f64_field, summarize, Edge, Harness, RunManifest, RunReport, SourceSpec, Target,
|
||||
};
|
||||
use aura_ingest::load_m1_window;
|
||||
use aura_std::{Exposure, Recorder, SimBroker, Sma, Sub};
|
||||
use data_server::{DataServer, DEFAULT_DATA_PATH};
|
||||
|
||||
/// Bootstrap the cycle-0007 two-sink signal-quality harness (mirrors
|
||||
/// `aura-engine`'s `report::tests::build_two_sink_harness`, with the shipped
|
||||
/// `aura_std::Recorder`), run it on `prices`, and fold the recorded equity +
|
||||
/// exposure into a `RunReport` whose window is the first/last real bar ts.
|
||||
fn run_sample_over(prices: Vec<(Timestamp, Scalar)>) -> RunReport {
|
||||
let (tx_eq, rx_eq) = mpsc::channel();
|
||||
let (tx_ex, rx_ex) = mpsc::channel();
|
||||
let mut h = Harness::bootstrap(
|
||||
vec![
|
||||
Box::new(Sma::new(2)), // 0
|
||||
Box::new(Sma::new(4)), // 1
|
||||
Box::new(Sub::new()), // 2
|
||||
Box::new(Exposure::new(0.5)), // 3
|
||||
Box::new(SimBroker::new(0.0001)), // 4
|
||||
Box::new(Recorder::new(&[ScalarKind::F64], Firing::Any, tx_eq)), // 5 equity sink
|
||||
Box::new(Recorder::new(&[ScalarKind::F64], Firing::Any, tx_ex)), // 6 exposure sink
|
||||
],
|
||||
vec![SourceSpec {
|
||||
kind: ScalarKind::F64,
|
||||
targets: vec![
|
||||
Target { node: 0, slot: 0 },
|
||||
Target { node: 1, slot: 0 },
|
||||
Target { node: 4, slot: 1 }, // price into the broker
|
||||
],
|
||||
}],
|
||||
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 }, // equity -> sink 5
|
||||
Edge { from: 3, to: 6, slot: 0, from_field: 0 }, // exposure -> sink 6
|
||||
],
|
||||
)
|
||||
.expect("valid signal-quality DAG");
|
||||
|
||||
let window = (
|
||||
prices.first().map(|&(t, _)| t).unwrap_or(Timestamp(0)),
|
||||
prices.last().map(|&(t, _)| t).unwrap_or(Timestamp(0)),
|
||||
);
|
||||
h.run(vec![prices]);
|
||||
|
||||
let eq_rows: Vec<(Timestamp, Vec<Scalar>)> = rx_eq.try_iter().collect();
|
||||
let ex_rows: Vec<(Timestamp, Vec<Scalar>)> = rx_ex.try_iter().collect();
|
||||
let equity = f64_field(&eq_rows, 0);
|
||||
let exposure = f64_field(&ex_rows, 0);
|
||||
let metrics = summarize(&equity, &exposure);
|
||||
|
||||
RunReport {
|
||||
manifest: RunManifest {
|
||||
commit: "real-bars-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,
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn sample_strategy_runs_over_real_m1_bars_deterministically() {
|
||||
let server = Arc::new(DataServer::new(DEFAULT_DATA_PATH));
|
||||
if !server.has_symbol("AAPL.US") {
|
||||
eprintln!("skip: no local data at {DEFAULT_DATA_PATH} (symbol AAPL.US absent)");
|
||||
return; // hermetic elsewhere; exercises the real path where files exist
|
||||
}
|
||||
// 2006-08 in inclusive Unix-ms (data-server skips files outside the window).
|
||||
let (from_ms, to_ms) = (1_154_390_400_000_i64, 1_157_068_799_999_i64);
|
||||
|
||||
let load = || {
|
||||
load_m1_window(&server, "AAPL.US", from_ms, to_ms)
|
||||
.expect("AAPL.US has data in the 2006-08 window")
|
||||
.close_stream()
|
||||
};
|
||||
|
||||
let prices = load();
|
||||
assert!(!prices.is_empty(), "window resolved to zero bars");
|
||||
|
||||
let r1 = run_sample_over(prices);
|
||||
assert!(r1.metrics.total_pips.is_finite(), "a backtest ran over real bars");
|
||||
|
||||
// same window -> bit-identical report (C1).
|
||||
let r2 = run_sample_over(load());
|
||||
assert_eq!(r1.to_json(), r2.to_json());
|
||||
}
|
||||
Reference in New Issue
Block a user