From 1d14ebc1cdfb9e05b87249c5d46b07403e2b4f7f Mon Sep 17 00:00:00 2001 From: Brummel Date: Sat, 11 Jul 2026 03:37:14 +0200 Subject: [PATCH] feat(cli,engine): multi-column opening on every real-data path (#231 task 4) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The six Close-only open sites route through the resolved binding's column set via open_columns (probe_window, open_real_source, run_sources/windowed_sources, the campaign member open, the trace re-run open) — a strategy declaring high/low/close roles gets exactly those columns, merged in the canonical C4 order. Synthetic data refuses any binding beyond {close} with the honest --real remedy (the walk generates a close series only); the sweep/mc family builders ride their established Err contract, run/walkforward/reproduce refuse exit-1. Proof at two layers: an engine e2e composes an OHLC blueprint from data and runs it over four inline VecSources to hand-computed rows (the merge-order pin), and an archive-gated CLI e2e sweeps a high/low blueprint over GER40 opening the High and Low columns. Verified: engine OHLC e2e green, full workspace suite green, clippy -D warnings clean. refs #231 --- crates/aura-cli/src/campaign_run.rs | 42 ++--- crates/aura-cli/src/main.rs | 158 +++++++++++++----- crates/aura-cli/tests/cli_run.rs | 156 +++++++++++++++++ .../tests/ohlc_from_blueprint_data.rs | 137 +++++++++++++++ 4 files changed, 428 insertions(+), 65 deletions(-) create mode 100644 crates/aura-engine/tests/ohlc_from_blueprint_data.rs diff --git a/crates/aura-cli/src/campaign_run.rs b/crates/aura-cli/src/campaign_run.rs index 43e369c..4e0e3cf 100644 --- a/crates/aura-cli/src/campaign_run.rs +++ b/crates/aura-cli/src/campaign_run.rs @@ -24,7 +24,7 @@ use aura_engine::{ blueprint_from_json, f64_field, summarize, summarize_r, ColumnarTrace, FamilySelection, RunReport, }; -use aura_ingest::{instrument_geometry, unix_ms_to_epoch_ns, M1Field, M1FieldSource}; +use aura_ingest::{instrument_geometry, open_columns_window, unix_ms_to_epoch_ns}; use aura_registry::{CampaignRunRecord, WriteKind}; use aura_research::{ campaign_to_json, content_id_of, parse_campaign, parse_process, validate_campaign, @@ -263,12 +263,14 @@ impl MemberRunner for CliMemberRunner<'_> { })?; let from = unix_ms_to_epoch_ns(window_ms.0); let to = unix_ms_to_epoch_ns(window_ms.1); - let source = match M1FieldSource::open_window( + // One source per resolved binding column, canonical order (the same + // order `wrap_r` declares the roles in — the shared plan). + let sources = match open_columns_window( &self.server, &cell.instrument, Some(from), Some(to), - M1Field::Close, + &binding.columns(), ) { Some(s) => s, // No archived file overlaps the window at all. @@ -279,10 +281,11 @@ impl MemberRunner for CliMemberRunner<'_> { }); } }; - // A window that overlaps a file but holds zero matching bars yields a - // source whose first peek is None (open_window's documented contract) - // — the same no-data condition. - if aura_engine::Source::peek(&source).is_none() { + // A window that overlaps a file but holds zero matching bars yields + // sources whose first peek is None (open_window's documented contract) + // — the same no-data condition (all columns decode the same bars, so + // probing the first source covers the set). + if aura_engine::Source::peek(sources[0].as_ref()).is_none() { return Err(MemberFault::NoData { instrument: cell.instrument.clone(), window_ms, @@ -298,7 +301,7 @@ impl MemberRunner for CliMemberRunner<'_> { signal, &point, &space, - vec![Box::new(source)], + sources, (from, to), 0, geo.pip_size, @@ -823,23 +826,23 @@ pub(crate) fn persist_campaign_traces( cell_rec.instrument, from.0, to.0 ) }; - let source = M1FieldSource::open_window( - server, - &cell_rec.instrument, - Some(from), - Some(to), - M1Field::Close, - ) - .ok_or_else(no_data)?; - if aura_engine::Source::peek(&source).is_none() { - return Err(no_data()); - } let signal = blueprint_from_json(blueprint_json, &|t| env.resolve(t)) .expect("stored blueprint passed the referential gate; reload is infallible"); // Same name-default binding the member ran under (campaign // overrides thread through in the DataSection task). let binding = crate::binding::resolve_binding(&cell_rec.strategy, signal.input_roles(), &BTreeMap::new())?; + let sources = open_columns_window( + server, + &cell_rec.instrument, + Some(from), + Some(to), + &binding.columns(), + ) + .ok_or_else(no_data)?; + if aura_engine::Source::peek(sources[0].as_ref()).is_none() { + return Err(no_data()); + } // The cell's OWN regime (#219/#212), not the hardcoded default: the // recorded member ran under `cell_rec.regime`, bound through the same // `stop_rule_for_regime` helper `CliMemberRunner::run_member` uses, so @@ -853,7 +856,6 @@ pub(crate) fn persist_campaign_traces( let mut h = crate::wrap_r(signal, tx_eq, tx_ex, tx_r, tx_req, stop, false, geo.pip_size, &binding) .bootstrap_with_cells(&point) .expect("the member's point re-bootstraps (it already ran this realization)"); - let sources: Vec> = vec![Box::new(source)]; h.run(sources); // Drain ALL FOUR channels (`run_signal_r` leaves req undrained; the // trace path must not): eq/ex/req feed the taps, r feeds the metrics. diff --git a/crates/aura-cli/src/main.rs b/crates/aura-cli/src/main.rs index 35dc6c8..d3f4dfd 100644 --- a/crates/aura-cli/src/main.rs +++ b/crates/aura-cli/src/main.rs @@ -516,11 +516,15 @@ fn pip_or_refuse( } } -/// Probe the full data window: open a single-pass probe `M1FieldSource`, drain it -/// for the first/last timestamp, and return `(first, last)`. Refuses (via -/// `no_real_data`) when the symbol/window yields no source or no bars. Shared by -/// `open_real_source` (which needs the manifest window from a probe separate from -/// the run source) and `DataSource::full_window`. +/// Probe the full data window: open a single-pass probe source through the +/// shared opener, drain it for the first/last timestamp, and return +/// `(first, last)`. Refuses (via `no_real_data`) when the symbol/window yields +/// no source or no bars. Shared by `open_real_source` (which needs the manifest +/// window from a probe separate from the run sources) and +/// `DataSource::full_window`. The probe drains the CLOSE column regardless of +/// the strategy's binding: the bounds are field-independent (every archived bar +/// carries all six fields at one timestamp), and file-level absence is +/// field-independent too. fn probe_window( server: &std::sync::Arc, symbol: &str, @@ -528,28 +532,32 @@ fn probe_window( to_ms: Option, env: &project::Env, ) -> (Timestamp, Timestamp) { - let mut probe = aura_ingest::M1FieldSource::open(server, symbol, from_ms, to_ms, aura_ingest::M1Field::Close) - .unwrap_or_else(|| no_real_data(symbol, env)); - let first = aura_engine::Source::peek(&probe).unwrap_or_else(|| no_real_data(symbol, env)); + let mut probe = aura_ingest::open_columns(server, symbol, from_ms, to_ms, &[aura_ingest::M1Field::Close]) + .unwrap_or_else(|| no_real_data(symbol, env)) + .pop() + .expect("open_columns yields one source per requested field"); + let first = aura_engine::Source::peek(probe.as_ref()).unwrap_or_else(|| no_real_data(symbol, env)); let mut last = first; - while let Some((t, _)) = aura_engine::Source::next(&mut probe) { + while let Some((t, _)) = aura_engine::Source::next(&mut *probe) { last = t; } (first, last) } -/// Open a real M1-close source for a recorded symbol over an optional window, returning -/// the run source paired with its manifest `window` and per-instrument `pip_size`. +/// Open the real M1 sources for a recorded symbol over an optional window — one +/// source per resolved binding column, in canonical order — returning the run +/// sources paired with the manifest `window` and the per-instrument `pip_size`. /// Single home of the real-source construction the single-run handlers share — the /// sidecar-pip lookup, the `DataServer` `has_symbol` refusal, the probe-window pass, and -/// the run-source `open` (each refusal an stderr + exit 1). Pre-data refusals keep the +/// the run-source open (each refusal an stderr + exit 1). Pre-data refusals keep the /// pip honest by construction. Used by `resolve_run_data`. fn open_real_source( symbol: &str, from_ms: Option, to_ms: Option, env: &project::Env, -) -> (Box, (Timestamp, Timestamp), f64) { + fields: &[aura_ingest::M1Field], +) -> (Vec>, (Timestamp, Timestamp), f64) { // Per-instrument pip from the recorded sidecar; resolved BEFORE bar-data access // so an instrument with no geometry refuses without touching the archive. let server = std::sync::Arc::new(data_server::DataServer::new(env.data_path())); @@ -559,12 +567,9 @@ fn open_real_source( } // Manifest window: drain a separate probe (single-pass Source) for first/last ts. let window = probe_window(&server, symbol, from_ms, to_ms, env); - let source: Box = - match aura_ingest::M1FieldSource::open(&server, symbol, from_ms, to_ms, aura_ingest::M1Field::Close) { - Some(s) => Box::new(s), - None => no_real_data(symbol, env), - }; - (source, window, pip) + let sources = aura_ingest::open_columns(&server, symbol, from_ms, to_ms, fields) + .unwrap_or_else(|| no_real_data(symbol, env)); + (sources, window, pip) } impl DataSource { @@ -625,28 +630,31 @@ impl DataSource { } } - /// A fresh full-window source per member (single-pass). Synthetic: showcase. - fn run_sources(&self, env: &project::Env) -> Vec> { + /// A fresh full-window source set per member (single-pass): the synthetic + /// showcase close stream, or one real source per resolved binding column + /// in canonical order (callers guard the synthetic arm to `{close}`). + fn run_sources(&self, env: &project::Env, fields: &[aura_ingest::M1Field]) -> Vec> { match self { DataSource::Synthetic => vec![Box::new(VecSource::new(showcase_prices()))], - DataSource::Real { server, symbol, from_ms, to_ms, .. } => vec![Box::new( - aura_ingest::M1FieldSource::open(server, symbol, *from_ms, *to_ms, aura_ingest::M1Field::Close) - .unwrap_or_else(|| no_real_data(symbol, env)), - )], + DataSource::Real { server, symbol, from_ms, to_ms, .. } => { + aura_ingest::open_columns(server, symbol, *from_ms, *to_ms, fields) + .unwrap_or_else(|| no_real_data(symbol, env)) + } } } - /// A fresh windowed source for an IS/OOS sub-window (walk-forward). Synthetic: - /// `walkforward_window_source`. Real: `open_window` (ns-native `Timestamp`). + /// A fresh windowed source set for an IS/OOS sub-window (walk-forward). + /// Synthetic: `walkforward_window_source`. Real: one source per resolved + /// binding column over the ns-native window. fn windowed_sources( - &self, from: Timestamp, to: Timestamp, env: &project::Env, + &self, from: Timestamp, to: Timestamp, env: &project::Env, fields: &[aura_ingest::M1Field], ) -> Vec> { match self { DataSource::Synthetic => vec![Box::new(walkforward_window_source(from, to))], - DataSource::Real { server, symbol, .. } => vec![Box::new( - aura_ingest::M1FieldSource::open_window(server, symbol, Some(from), Some(to), aura_ingest::M1Field::Close) - .unwrap_or_else(|| no_real_data(symbol, env)), - )], + DataSource::Real { server, symbol, .. } => { + aura_ingest::open_columns_window(server, symbol, Some(from), Some(to), fields) + .unwrap_or_else(|| no_real_data(symbol, env)) + } } } @@ -1118,6 +1126,10 @@ fn reproduce_family_in( eprintln!("aura: {m}"); std::process::exit(1); }); + if matches!(data, DataSource::Synthetic) && !binding.close_only() { + eprintln!("aura: {}", binding::synthetic_refusal(&hash, &binding)); + std::process::exit(1); + } let space = wrap_r(reload(), tx_eq, tx_ex, tx_r, tx_req, stop, true, SYNTHETIC_PIP_SIZE, &binding).param_space(); let point = point_from_params(&space, &stored.manifest.params); // A MonteCarlo member carries no tuning params (the params-join is empty), so its @@ -1150,7 +1162,7 @@ fn reproduce_family_in( // stored window bounds; the winner params come from the shared // manifest->cells recovery below (as Sweep members do). let (from, to) = stored.manifest.window; - let s = data.windowed_sources(from, to, env); + let s = data.windowed_sources(from, to, env, &binding.columns()); let w = window_of(&s).expect("non-empty OOS window"); (s, w) } @@ -1163,9 +1175,9 @@ fn reproduce_family_in( _ => match data { DataSource::Real { .. } => { let (from, to) = stored.manifest.window; - (data.windowed_sources(from, to, env), (from, to)) + (data.windowed_sources(from, to, env, &binding.columns()), (from, to)) } - DataSource::Synthetic => (data.run_sources(env), data.full_window(env)), + DataSource::Synthetic => (data.run_sources(env, &binding.columns()), data.full_window(env)), }, }; let rerun = run_blueprint_member( @@ -1374,13 +1386,15 @@ fn wrap_r( } /// Resolve a `RunData` selector to the `(sources, window, pip_size)` triple the -/// r-sma run paths feed to the harness: the built-in synthetic R stream, -/// or a lazily-streamed real M1 close source (with its sidecar pip + probed window). +/// run paths feed to the harness: the built-in synthetic R stream (close-only — +/// the caller guards the binding shape), or the lazily-streamed real sources of +/// the binding's resolved columns (with the sidecar pip + probed window). /// Single definition used by `run_signal_r`. #[allow(clippy::type_complexity)] fn resolve_run_data( data: &RunData, env: &project::Env, + binding: &binding::ResolvedBinding, ) -> ( Vec>, (Timestamp, Timestamp), @@ -1394,8 +1408,7 @@ fn resolve_run_data( (sources, window, SYNTHETIC_PIP_SIZE) } RunData::Real { symbol, from, to } => { - let (source, window, pip_size) = open_real_source(symbol, *from, *to, env); - (vec![source], window, pip_size) + open_real_source(symbol, *from, *to, env, &binding.columns()) } } } @@ -1416,6 +1429,10 @@ fn run_signal_r( eprintln!("aura: {m}"); std::process::exit(1); }); + if matches!(data, RunData::Synthetic) && !binding.close_only() { + eprintln!("aura: {}", binding::synthetic_refusal(signal.name(), &binding)); + std::process::exit(1); + } let names: Vec = signal .param_space() .iter() @@ -1427,7 +1444,7 @@ fn run_signal_r( // The req tap (r_equity recorder) is wired but not persisted on this path; keep the // receiver alive so the sink's sends do not fail, but do not drain it. let (tx_req, _rx_req) = mpsc::channel(); - let (sources, window, pip_size) = resolve_run_data(&data, env); + let (sources, window, pip_size) = resolve_run_data(&data, env, &binding); let wrapped = wrap_r(signal, tx_eq, tx_ex, tx_r, tx_req, StopRule::Vol { length: R_SMA_STOP_LENGTH, k: R_SMA_STOP_K }, false, pip_size, &binding); let flat = wrapped .compile_with_params(params) @@ -1569,6 +1586,9 @@ fn blueprint_sweep_family( // Strict binding resolution (name defaults — the verb path carries no // campaign overrides): the family's open plan and wrap plan in one value. let binding = binding::resolve_binding(probe_signal.name(), probe_signal.input_roles(), &BTreeMap::new())?; + if matches!(data, DataSource::Synthetic) && !binding.close_only() { + return Err(binding::synthetic_refusal(probe_signal.name(), &binding)); + } let pip = data.pip_size(); let window = data.full_window(env); // a single throwaway floated build, only to resolve param_space (borrow) then seed @@ -1601,7 +1621,7 @@ fn blueprint_sweep_family( .sweep(|point| { // fresh per-member graph (Composite is !Clone, reload per member) run through // the shared reduce-mode member path — the same fn reproduction re-runs. - run_blueprint_member(reload(doc), point, &space, data.run_sources(env), window, 0, pip, &topo, env, StopRule::Vol { length: R_SMA_STOP_LENGTH, k: R_SMA_STOP_K }, &binding) + run_blueprint_member(reload(doc), point, &space, data.run_sources(env, &binding.columns()), window, 0, pip, &topo, env, StopRule::Vol { length: R_SMA_STOP_LENGTH, k: R_SMA_STOP_K }, &binding) }) // render the sweep terminal's BindError to a message (the fn's String error contract), // so `UnknownKnob("nope")` still surfaces verbatim at the IO wrapper. @@ -1632,7 +1652,7 @@ fn blueprint_sweep_over( binder = binder.axis(n, vals.clone()); } binder.sweep_with_lattice(|point| { - let sources = data.windowed_sources(from, to, env); + let sources = data.windowed_sources(from, to, env, &binding.columns()); let window = window_of(&sources).expect("non-empty in-sample window"); run_blueprint_member(reload(doc), point, &space, sources, window, 0, pip, &topo, env, StopRule::Vol { length: R_SMA_STOP_LENGTH, k: R_SMA_STOP_K }, binding) }) @@ -1650,7 +1670,7 @@ fn run_oos_blueprint( let reload = blueprint_from_json(doc, &|t| env.resolve(t)) .expect("doc parse-validated at the dispatch boundary; reload is infallible"); let pip = data.pip_size(); - let sources = data.windowed_sources(from, to, env); + let sources = data.windowed_sources(from, to, env, &binding.columns()); let window = window_of(&sources).expect("non-empty out-of-sample window"); let report = run_blueprint_member(reload, params, space, sources, window, 0, pip, topo, env, StopRule::Vol { length: R_SMA_STOP_LENGTH, k: R_SMA_STOP_K }, binding); (Vec::new(), report) @@ -1685,6 +1705,10 @@ fn blueprint_walkforward_family( eprintln!("aura: {m}"); std::process::exit(1); }); + if matches!(data, DataSource::Synthetic) && !binding.close_only() { + eprintln!("aura: {}", binding::synthetic_refusal(probe_signal.name(), &binding)); + std::process::exit(1); + } // Validate the `--axis` grid ONCE at the dispatch boundary, mirroring `aura sweep` // (which resolves its axes a single time before any member runs). `walk_forward` fans // the per-window closure out across the windows in parallel, so a `BindError` raised @@ -1746,6 +1770,11 @@ fn blueprint_mc_family( // Strict binding resolution (name defaults — mc's synthetic family binds // no campaign overrides); the exit-free Err contract of this builder. let binding = binding::resolve_binding(probe_signal.name(), probe_signal.input_roles(), &BTreeMap::new())?; + if !binding.close_only() { + // MC draws ALWAYS run the seeded synthetic close walk (real-data mc + // routes through the campaign sugar and never reaches this builder). + return Err(binding::synthetic_refusal(probe_signal.name(), &binding)); + } let pip = data.pip_size(); // probe the wrapped param_space (the same probe the sweep resolves against); // MC needs it empty. `blueprint_axis_probe` is the single source of that wrap. @@ -3871,13 +3900,52 @@ mod tests { let env = project::Env::std(); let d = DataSource::Synthetic; assert_eq!(d.pip_size(), SYNTHETIC_PIP_SIZE); - assert!(!d.run_sources(&env).is_empty()); + assert!(!d.run_sources(&env, &[aura_ingest::M1Field::Close]).is_empty()); assert_eq!(d.wf_window_sizes(), (24, 12, 12)); // full_window equals window_of over the showcase stream (byte-unchanged source) let s: Vec> = vec![Box::new(VecSource::new(showcase_prices()))]; assert_eq!(d.full_window(&env), window_of(&s).unwrap()); } + /// A multi-column blueprint over synthetic data (a single close walk) + /// refuses honestly through the builders' exit-free Err contract, naming + /// the beyond-close columns and the --real remedy — never a panic from a + /// source-count mismatch. High/low-consuming, closed (mc requires it), + /// with the mandatory `bias` output. + const OHLC_REFUSAL_BLUEPRINT: &str = r#"{ + "format_version": 1, + "blueprint": { + "name": "hl_range", + "nodes": [ {"primitive":{"type":"Sub"}} ], + "edges": [], + "input_roles": [ + {"name":"high","targets":[{"node":0,"slot":0}],"source":"F64"}, + {"name":"low","targets":[{"node":0,"slot":1}],"source":"F64"} + ], + "output": [{"node":0,"field":0,"name":"bias"}] + } +}"#; + + #[test] + fn synthetic_data_refuses_a_multi_column_blueprint() { + let env = project::Env::std(); + let err = blueprint_mc_family(OHLC_REFUSAL_BLUEPRINT, 2, &DataSource::Synthetic, &env) + .expect_err("a high/low blueprint cannot run over the synthetic close walk"); + assert_eq!( + err, + "strategy \"hl_range\" consumes columns beyond close (high, low) — synthetic \ + data generates a close series only; run with --real " + ); + let err = blueprint_sweep_family( + OHLC_REFUSAL_BLUEPRINT, + &[("x".to_string(), vec![Scalar::i64(1)])], + &DataSource::Synthetic, + &env, + ) + .expect_err("the synthetic sweep path refuses the same shape"); + assert!(err.contains("consumes columns beyond close"), "got: {err}"); + } + #[test] fn wf_real_roller_sizes_are_90_30_30_days_in_ns() { // Independent expected value: a day reconstructed from its time units @@ -4134,7 +4202,7 @@ mod tests { reload(), &[], &space, - data.run_sources(&env), + data.run_sources(&env, &binding.columns()), window, 0, pip, diff --git a/crates/aura-cli/tests/cli_run.rs b/crates/aura-cli/tests/cli_run.rs index 3105881..7ede2cd 100644 --- a/crates/aura-cli/tests/cli_run.rs +++ b/crates/aura-cli/tests/cli_run.rs @@ -1427,6 +1427,162 @@ fn reproduce_real_walkforward_family_does_not_panic() { ); } +// Fixture blueprints for the verb-level multi-column binding coverage below +// (#231 task 4 quality follow-up): a CLOSED high/low blueprint (`run` +// requires zero free knobs) and an OPEN high/low blueprint with two +// SMA-length axes (`walkforward` requires >= 1 axis). Written to a fresh +// temp file per test rather than a tracked `examples/` fixture — these exist +// only to drive the refusal / real-open contract, not as authoring samples. +const HL_RANGE_CLOSED_BLUEPRINT: &str = r#"{ + "format_version": 1, + "blueprint": { + "name": "hl_range", + "nodes": [ {"primitive":{"type":"Sub"}} ], + "edges": [], + "input_roles": [ + {"name":"high","targets":[{"node":0,"slot":0}],"source":"F64"}, + {"name":"low","targets":[{"node":0,"slot":1}],"source":"F64"} + ], + "output": [{"node":0,"field":0,"name":"bias"}] + } +}"#; +const HL_SIGNAL_OPEN_BLUEPRINT: &str = r#"{"format_version":1,"blueprint":{"name":"hl_signal","nodes":[{"primitive":{"type":"SMA","name":"fast"}},{"primitive":{"type":"SMA","name":"slow"}},{"primitive":{"type":"Sub"}},{"primitive":{"type":"Bias","name":"bias","bound":[{"pos":0,"name":"scale","kind":"F64","value":{"F64":0.5}}]}}],"edges":[{"from":0,"to":2,"slot":0,"from_field":0},{"from":1,"to":2,"slot":1,"from_field":0},{"from":2,"to":3,"slot":0,"from_field":0}],"input_roles":[{"name":"high","targets":[{"node":0,"slot":0}],"source":"F64"},{"name":"low","targets":[{"node":1,"slot":0}],"source":"F64"}],"output":[{"node":3,"field":0,"name":"bias"}]}}"#; + +/// Property (#231 task 4 quality follow-up): `aura run` over a multi-column +/// blueprint with NO `--real` (the synthetic default) refuses through the +/// verb-level exit-1 + stderr contract — `binding::synthetic_refusal`'s exact +/// prose, not merely its `Result` shape (the unit-mod test one layer down +/// covers the `Err`-returning builders only). Data-free: the guard fires +/// directly after binding resolution, before any archive access, so no +/// gating is needed. +#[test] +fn run_synthetic_refuses_a_multi_column_blueprint_before_data_access() { + let dir = temp_cwd("run_synthetic_multicolumn_refusal"); + let fixture = dir.join("hl_range.json"); + std::fs::write(&fixture, HL_RANGE_CLOSED_BLUEPRINT).expect("write fixture"); + let out = Command::new(BIN).args(["run", fixture.to_str().unwrap()]).output().expect("spawn aura run"); + assert_eq!(out.status.code(), Some(1), "status: {:?}", out.status); + let stderr = String::from_utf8_lossy(&out.stderr); + assert_eq!( + stderr.trim_end(), + "aura: strategy \"hl_range\" consumes columns beyond close (high, low) — synthetic \ + data generates a close series only; run with --real " + ); + let _ = std::fs::remove_dir_all(&dir); +} + +/// Property (#231 task 4 quality follow-up): the sibling exit-process arm at +/// `blueprint_walkforward_family` — `aura walkforward` over a multi-column +/// blueprint with no `--real` — refuses identically. Data-free (the guard +/// fires directly after binding resolution, before any window/archive +/// access); run inside a temp cwd purely by the project's verb-test +/// convention (no store write is expected on this refusal path either). +#[test] +fn walkforward_synthetic_refuses_a_multi_column_blueprint_before_data_access() { + let dir = temp_cwd("walkforward_synthetic_multicolumn_refusal"); + let fixture = dir.join("hl_signal.json"); + std::fs::write(&fixture, HL_SIGNAL_OPEN_BLUEPRINT).expect("write fixture"); + let out = Command::new(BIN) + .args([ + "walkforward", fixture.to_str().unwrap(), + "--axis", "hl_signal.fast.length=2,4", + "--axis", "hl_signal.slow.length=8", + ]) + .current_dir(&dir) + .output() + .expect("spawn aura walkforward"); + assert_eq!(out.status.code(), Some(1), "status: {:?}", out.status); + let stderr = String::from_utf8_lossy(&out.stderr); + assert_eq!( + stderr.trim_end(), + "aura: strategy \"hl_signal\" consumes columns beyond close (high, low) — synthetic \ + data generates a close series only; run with --real " + ); + assert!(!dir.join("runs").exists(), "a refused synthetic family must not start a run store"); + let _ = std::fs::remove_dir_all(&dir); +} + +/// Property (#231 task 4 quality follow-up): the positive counterpart of the +/// two refusals above — a multi-column blueprint over `--real GER40` opens +/// BOTH the high and low columns end to end (`open_real_source` -> +/// `resolve_run_data` -> `run_signal_r`) and produces a real `RunReport`, +/// never the single-close weld the pre-task-4 code hardcoded. Gated on the +/// local GER40 archive (skip-on-no-data convention). +#[test] +fn run_real_multi_column_blueprint_opens_high_and_low_columns() { + if !local_data_present() { + eprintln!("skip: no local data at {}", data_server::DEFAULT_DATA_PATH); + return; + } + let dir = temp_cwd("run_real_multicolumn_open"); + let fixture = dir.join("hl_range.json"); + std::fs::write(&fixture, HL_RANGE_CLOSED_BLUEPRINT).expect("write fixture"); + let out = Command::new(BIN) + .args([ + "run", fixture.to_str().unwrap(), + "--real", "GER40", + "--from", GER40_SEPT2024_FROM_MS, + "--to", GER40_SEPT2024_TO_MS, + ]) + .output() + .expect("spawn aura run"); + assert_eq!( + out.status.code(), Some(0), + "a high/low blueprint over real GER40 data must run cleanly, got status={:?} stderr={}", + out.status, String::from_utf8_lossy(&out.stderr) + ); + let stdout = String::from_utf8_lossy(&out.stdout); + assert!( + stdout.trim_start().starts_with("{\"manifest\":"), + "exit 0 must carry a JSON report, got: {stdout}" + ); + let _ = std::fs::remove_dir_all(&dir); +} + +/// Property (#231 task 4, "every real-data path" — `run --real` above covers +/// only `open_real_source`/`run_signal_r`; `sweep --real` walks a DISTINCT +/// real-data path, `blueprint_sweep_family` -> `blueprint_sweep_over` -> +/// `DataSource::windowed_sources`, which independently needed `binding. +/// columns()` threaded through in task 4): a multi-column (high/low) OPEN +/// blueprint swept over `--real GER40` opens both columns per member and +/// produces one member line per grid point, not a source-count-mismatch panic +/// or a fallback to a single close source. Gated on the local GER40 archive. +#[test] +fn sweep_real_multi_column_blueprint_opens_high_and_low_columns() { + if !local_data_present() { + eprintln!("skip: no local data at {}", data_server::DEFAULT_DATA_PATH); + return; + } + let (dir, _g) = fresh_project(); + let fixture_dir = temp_cwd("sweep_real_multicolumn_open"); + let fixture = fixture_dir.join("hl_signal.json"); + std::fs::write(&fixture, HL_SIGNAL_OPEN_BLUEPRINT).expect("write fixture"); + let out = std::process::Command::new(BIN) + .args([ + "sweep", fixture.to_str().unwrap(), + "--real", "GER40", + "--from", GER40_SEPT2024_FROM_MS, + "--to", GER40_SEPT2024_TO_MS, + "--axis", "hl_signal.fast.length=2,4", + "--axis", "hl_signal.slow.length=8", + ]) + .current_dir(dir) + .output() + .unwrap(); + assert!(out.status.success(), "stderr: {}", String::from_utf8_lossy(&out.stderr)); + let stdout = String::from_utf8_lossy(&out.stdout).into_owned(); + let lines: Vec<&str> = stdout.lines().collect(); + assert_eq!(lines.len(), 2, "one member line per grid point: {stdout}"); + for line in &lines { + let v: serde_json::Value = serde_json::from_str(line).expect("member line parses as JSON"); + assert!( + v["report"]["manifest"].is_object(), + "each member carries a real report, not a source-mismatch fault: {line}" + ); + } + let _ = std::fs::remove_dir_all(&fixture_dir); +} + /// Property (#210 c0110 fieldtest, "the wrapped-name refusal mangles the axis /// the user typed"): the dissolved real-data blueprint sweep types `--axis` /// against the WRAPPED probe namespace (`sma_signal.fast.length`, what diff --git a/crates/aura-engine/tests/ohlc_from_blueprint_data.rs b/crates/aura-engine/tests/ohlc_from_blueprint_data.rs new file mode 100644 index 0000000..ab68929 --- /dev/null +++ b/crates/aura-engine/tests/ohlc_from_blueprint_data.rs @@ -0,0 +1,137 @@ +//! Acceptance pin for the harness-input-binding spec (#231): a MULTI-COLUMN +//! blueprint (open/high/low/close input roles) composes from blueprint data +//! and runs over four inline `VecSource` columns to hand-computed values — +//! pinning the role-declaration/merge-order contract (role i is fed by source +//! i; same-timestamp ties break by source index, i.e. canonical column order — +//! C4) without any archive. Mirrors `rsi_from_blueprint_data.rs` (the same +//! public seam: a JSON document + `blueprint_from_json` + an injected +//! vocabulary, C24). Pure engine machinery — no CLI, no binding module. + +use std::sync::mpsc; + +use aura_core::{Firing, Scalar, ScalarKind, Timestamp}; +use aura_engine::{blueprint_from_json, BlueprintNode, Composite, Edge, Role, Target, VecSource}; +use aura_std::{std_vocabulary, Recorder}; + +// value = (high - low) + (close - open): every column is load-bearing, so any +// mis-pairing of roles to sources (or a wrong merge tie-break) changes the +// recorded rows. Node indices: 0 Sub(range = high - low) · +// 1 Sub(body = close - open) · 2 Add(range + body). +const OHLC_BLUEPRINT_JSON: &str = r#"{ + "format_version": 1, + "blueprint": { + "name": "ohlc_probe", + "nodes": [ + {"primitive":{"type":"Sub"}}, + {"primitive":{"type":"Sub"}}, + {"primitive":{"type":"Add"}} + ], + "edges": [ + {"from":0,"to":2,"slot":0,"from_field":0}, + {"from":1,"to":2,"slot":1,"from_field":0} + ], + "input_roles": [ + {"name":"open","targets":[{"node":1,"slot":1}]}, + {"name":"high","targets":[{"node":0,"slot":0}]}, + {"name":"low","targets":[{"node":0,"slot":1}]}, + {"name":"close","targets":[{"node":1,"slot":0}]} + ], + "output": [{"node":2,"field":0,"name":"value"}] + } +}"#; + +// Nest the loaded signal under a Rust root that declares the four source roles +// in canonical column order (open, high, low, close) and records the single +// output (the rsi_from_blueprint_data idiom — sinks are outside the std +// vocabulary, so the recorder is added in Rust). Root role slot i targets the +// nested composite's role i (input_roles order). +fn run_recording( + signal: Composite, + columns: [Vec<(Timestamp, Scalar)>; 4], +) -> Vec<(Timestamp, Vec)> { + let (tx, rx) = mpsc::channel(); + let roles: Vec = ["open", "high", "low", "close"] + .iter() + .enumerate() + .map(|(slot, name)| Role { + name: (*name).to_string(), + targets: vec![Target { node: 0, slot }], + source: Some(ScalarKind::F64), + }) + .collect(); + let root = Composite::new( + "h", + vec![ + BlueprintNode::Composite(signal), + Recorder::builder(vec![ScalarKind::F64], Firing::Any, tx).into(), + ], + vec![Edge { from: 0, to: 1, slot: 0, from_field: 0 }], + roles, + vec![], + ); + let mut h = root.bootstrap_with_params(vec![]).expect("bootstraps (no open params)"); + let sources: Vec> = columns + .into_iter() + .map(|c| Box::new(VecSource::new(c)) as Box) + .collect(); + h.run(sources); + rx.try_iter().collect() +} + +fn col(values: &[f64]) -> Vec<(Timestamp, Scalar)> { + values + .iter() + .enumerate() + .map(|(i, &v)| (Timestamp(i as i64 + 1), Scalar::f64(v))) + .collect() +} + +/// Hand-computed contract. The four sources at one bar timestamp arrive in +/// FOUR consecutive engine cycles (the k-way merge breaks the timestamp tie +/// by source index = canonical column order: open, high, low, close); every +/// `Firing::Any` node re-evaluates per fresh input over the newest value of +/// each leg (sample-and-hold on the stale leg). Bars (o,h,l,c): +/// (1,10,5,3), (2,20,8,6), (4,40,16,12). +/// +/// Bar 1: open=1 -> body fires, close leg cold -> None; high=10 -> range +/// fires, low leg cold -> None; low=5 -> range = 10-5 = 5, add fires, body +/// leg cold -> None; close=3 -> body = 3-1 = 2, add = 5+2 = 7 -> (t1, 7). +/// Bar 2: open=2 -> body = 3-2 = 1 (close held at 3), add = 5+1 = 6; +/// high=20 -> range = 20-5 = 15 (low held), add = 15+1 = 16; +/// low=8 -> range = 20-8 = 12, add = 12+1 = 13; +/// close=6 -> body = 6-2 = 4, add = 12+4 = 16. +/// Bar 3: open=4 -> body = 6-4 = 2, add = 12+2 = 14; high=40 -> range = +/// 40-8 = 32, add = 32+2 = 34; low=16 -> range = 40-16 = 24, add = 26; +/// close=12 -> body = 12-4 = 8, add = 24+8 = 32. +#[test] +fn ohlc_blueprint_runs_over_four_columns_in_canonical_order() { + let signal = blueprint_from_json(OHLC_BLUEPRINT_JSON, &|t| std_vocabulary(t)) + .expect("loads through the std vocabulary"); + let trace = run_recording( + signal, + [ + col(&[1.0, 2.0, 4.0]), // open + col(&[10.0, 20.0, 40.0]), // high + col(&[5.0, 8.0, 16.0]), // low + col(&[3.0, 6.0, 12.0]), // close + ], + ); + let want: Vec<(Timestamp, Vec)> = [ + (1, 7.0), + (2, 6.0), + (2, 16.0), + (2, 13.0), + (2, 16.0), + (3, 14.0), + (3, 34.0), + (3, 26.0), + (3, 32.0), + ] + .iter() + .map(|&(t, v)| (Timestamp(t), vec![Scalar::f64(v)])) + .collect(); + assert_eq!( + trace, want, + "role i must be fed by source i in canonical column order", + ); +}