feat(cli): aura campaign run — the executor verb over the MemberRunner driver (0107 tasks 8-9)

campaign_run.rs: target resolution (file is register-then-run sugar;
bare 64-hex is the canonical id — #198 decision 1), project gate before
any store write, zero-fault referential gate, process fetch + v1
preflight, then aura_campaign::execute with the CLI MemberRunner over
the shipped loaded-blueprint convention: per-member blueprint reload,
reduce-mode wrap, unique suffix-join of raw campaign axis names onto
the wrapped param_space, windowed real data via the shipped ms->ns
source seam (absent archive/zero-bar windows -> NoData). Emission:
family_table / selection_report lines gated on the doc's emit list
(names debug_asserted against emit_vocabulary), the campaign_run
record line always, persist_taps deferred LOUDLY on stderr before
execution (F7 lesson), zero-survivor cells noted on stderr with exit 0.
exec_fault_prose/member_fault_prose keep the Debug-free house seam.

Seam tests: 8 campaign_run_* e2e tests incl. outside-project, bogus
target, unknown id, v1-boundary (mc process registers fine, run
refuses), persist_taps ordering, and the gated real-data
sweep->gate->walkforward e2e — which ran its full assert path on this
host (local GER40 2024-09 archive) rather than the data-less skip.
Also: exec.rs campaign-id guard tightened to lowercase hex, matching
is_content_id and the store's self-keyed form (review nit).

Gates: workspace 975/0, clippy -D warnings clean.

refs #198
This commit is contained in:
2026-07-03 21:33:29 +02:00
parent b15ad07208
commit aeb0366aed
7 changed files with 907 additions and 7 deletions
+493
View File
@@ -0,0 +1,493 @@
//! `aura campaign run` — the driver that turns a stored campaign document into
//! a realized run-set (#198). The execution *semantics* (preflight, cell loop,
//! stage sequencing, selection, registry writes) live in `aura-campaign`; this
//! module owns what is CLI-specific: target resolution (file sugar vs content
//! id), the project + referential gates, the [`MemberRunner`] implementation
//! over the shipped loaded-blueprint machinery (`wrap_r` reduce-mode member
//! runs via `run_blueprint_member` + `M1FieldSource` windowed real-data
//! binding), fault prose (`exec_fault_prose` lives HERE, beside its consumer,
//! not in `research_docs.rs` — it phrases `aura_campaign` types the doc-tier
//! module deliberately does not import), and stdout/stderr emission.
//!
//! Root items (`blueprint_axis_probe`, `run_blueprint_member`,
//! `family_member_line`) are reached via `crate::` — the `graph_construct`
//! submodule idiom: main.rs is the crate root, so its private items are
//! visible to child modules without promotion.
use std::path::{Path, PathBuf};
use std::sync::Arc;
use aura_campaign::{CellSpec, ExecFault, MemberFault, MemberRunner};
use aura_core::{Cell, ParamSpec, Scalar};
use aura_engine::{blueprint_from_json, FamilySelection, RunReport};
use aura_ingest::{instrument_geometry, unix_ms_to_epoch_ns, M1Field, M1FieldSource};
use aura_registry::CampaignRunRecord;
use aura_research::{
campaign_to_json, content_id_of, parse_campaign, parse_process, validate_campaign,
validate_process, DocRef,
};
use crate::project::Env;
use crate::research_docs::{
doc_error_prose, doc_fault_prose, fault_block, parse_valid_campaign, ref_fault_prose,
};
/// A bare store address: exactly 64 lowercase hex chars (the content-id key shape).
fn is_content_id(s: &str) -> bool {
s.len() == 64 && s.bytes().all(|b| matches!(b, b'0'..=b'9' | b'a'..=b'f'))
}
/// Phrase an [`ExecFault`] for stderr — Debug-leak-free, path-addressed
/// (stage index + block id), the `doc_fault_prose`/`ref_fault_prose` register.
fn exec_fault_prose(f: &ExecFault) -> String {
match f {
ExecFault::UnsupportedStage { stage, block } => format!(
"process stage {stage} ({block}) is not executable in v1; executable \
pipeline shape: std::sweep [std::gate]* [std::walk_forward]"
),
ExecFault::PipelineShape { detail } => {
format!("process pipeline is not executable: {detail}")
}
ExecFault::UnrankableMetric { stage, metric } => format!(
"process stage {stage}: metric \"{metric}\" is not rankable (winner \
selection needs one of the registry's rankable metrics)"
),
ExecFault::GateMetricNotPerMember { stage, metric } => format!(
"process stage {stage}: gate metric \"{metric}\" is not a per-member \
scalar (selection annotations cannot gate members)"
),
ExecFault::PlateauInWalkForward { stage } => format!(
"process stage {stage}: walk_forward cannot use a plateau select (a \
gated survivor subset has no parameter lattice)"
),
ExecFault::DeflatePlateauConflict { stage } => format!(
"process stage {stage}: sweep deflate: true composes only with select \"argmax\""
),
ExecFault::Window { stage, detail } => format!(
"process stage {stage}: walk_forward windows do not fit the campaign \
window: {detail}"
),
ExecFault::Member(MemberFault::NoData { instrument, window_ms }) => format!(
"no data for instrument {instrument} in window [{}, {}] (epoch-ms)",
window_ms.0, window_ms.1
),
ExecFault::Member(MemberFault::Bind(detail)) => {
format!("a member failed to bind: {detail}")
}
ExecFault::Member(MemberFault::Run(detail)) => {
format!("a member failed to run: {detail}")
}
ExecFault::Registry(e) => e.to_string(),
}
}
/// stdout wire form of one selection-bearing stage. Field order is the wire
/// contract (serde derives declaration order).
#[derive(serde::Serialize)]
struct SelectionReportLine<'a> {
selection_report: SelectionReportBody<'a>,
}
#[derive(serde::Serialize)]
struct SelectionReportBody<'a> {
family_id: &'a str,
stage: usize,
block: &'a str,
winner_ordinal: usize,
params: &'a [(String, Scalar)],
selection: &'a FamilySelection,
}
/// The always-on final line: the stored realization record under one key.
#[derive(serde::Serialize)]
struct CampaignRunLine<'a> {
campaign_run: &'a CampaignRunRecord,
}
/// Suffix-join each raw campaign axis onto exactly one wrapped param (wrapped
/// == raw, or wrapped ends with ".{raw}" — the established suffix-match
/// pattern), then require every wrapped slot to be covered. Pure and
/// independent of the env/data seam so the three fault arms are unit-testable
/// without a project, a store, or a loaded blueprint.
fn bind_axes(
space: &[ParamSpec],
strategy_id: &str,
params: &[(String, Scalar)],
) -> Result<Vec<Cell>, MemberFault> {
let mut slots: Vec<Option<Scalar>> = vec![None; space.len()];
for (raw, value) in params {
let hits: Vec<usize> = space
.iter()
.enumerate()
.filter(|(_, p)| p.name == *raw || p.name.ends_with(&format!(".{raw}")))
.map(|(i, _)| i)
.collect();
match hits.as_slice() {
[i] => slots[*i] = Some(*value),
[] => {
return Err(MemberFault::Bind(format!(
"axis \"{raw}\" matches no open param of the wrapped strategy {strategy_id}"
)));
}
_ => {
let names: Vec<&str> = hits.iter().map(|&i| space[i].name.as_str()).collect();
return Err(MemberFault::Bind(format!(
"axis \"{raw}\" is ambiguous in the wrapped param space of \
strategy {strategy_id}: matches {}",
names.join(", ")
)));
}
}
}
let mut point: Vec<Cell> = Vec::with_capacity(space.len());
for (spec, slot) in space.iter().zip(&slots) {
match slot {
Some(v) => point.push(v.cell()),
None => {
return Err(MemberFault::Bind(format!(
"open param \"{}\" of strategy {strategy_id} is bound by no campaign axis",
spec.name
)));
}
}
}
Ok(point)
}
/// The CLI's harness/data binding seam for `aura_campaign::execute`: members
/// run through the shipped loaded-blueprint machinery (`wrap_r` reduce-mode
/// via `crate::run_blueprint_member`) over windowed real M1 close bars
/// (`M1FieldSource::open_window` — the ms→ns crossing happens at exactly this
/// seam, via `unix_ms_to_epoch_ns`). All refusals are member faults for the
/// library to surface; never a process exit inside a sweep worker.
struct CliMemberRunner<'a> {
env: &'a Env,
server: Arc<data_server::DataServer>,
}
impl MemberRunner for CliMemberRunner<'_> {
fn run_member(
&self,
cell: &CellSpec,
params: &[(String, Scalar)],
window_ms: (i64, i64),
) -> Result<RunReport, MemberFault> {
// The wrapped axis namespace — the SAME probe the sweep verbs resolve
// against (single-sourced in main.rs; identical `false, true, None` wrap).
let space = crate::blueprint_axis_probe(&cell.blueprint_json, self.env).param_space();
let point = bind_axes(&space, &cell.strategy_id, params)?;
// Real windowed data — geometry BEFORE bar data (the shipped pre-data
// refusal order of `open_real_source`), both as member faults.
let geo = instrument_geometry(&self.server, &cell.instrument).ok_or_else(|| {
MemberFault::Run(format!(
"no recorded geometry for symbol '{}' at {} — refusing to run a \
real instrument with a guessed pip",
cell.instrument,
self.env.data_path()
))
})?;
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(
&self.server,
&cell.instrument,
Some(from),
Some(to),
M1Field::Close,
) {
Some(s) => s,
// No archived file overlaps the window at all.
None => {
return Err(MemberFault::NoData {
instrument: cell.instrument.clone(),
window_ms,
});
}
};
// 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() {
return Err(MemberFault::NoData {
instrument: cell.instrument.clone(),
window_ms,
});
}
// The shipped member recipe (reload — a Composite is !Clone — wrap,
// bind, run, summarize): `run_blueprint_member` verbatim. Seed 0
// (seed-free real-data runs), topology_hash = the strategy's content
// id, project provenance stamped inside the helper.
let signal = blueprint_from_json(&cell.blueprint_json, &|t| self.env.resolve(t))
.expect("stored blueprint passed the referential gate; reload is infallible");
let mut report = crate::run_blueprint_member(
signal,
&point,
&space,
vec![Box::new(source)],
(from, to),
0,
geo.pip_size,
&cell.strategy_id,
self.env,
);
report.manifest.instrument = Some(cell.instrument.clone());
Ok(report)
}
}
/// `aura campaign run <target>`: resolve, gate, execute, emit. Every `Err`
/// is a refusal `campaign_cmd` renders as `aura: {msg}` + exit 1.
pub(crate) fn run_campaign(target: &str, env: &Env) -> Result<(), String> {
// Project gate FIRST: nothing (not even the file-sugar registration)
// touches a store outside a project.
if env.provenance().is_none() {
let cwd = std::env::current_dir()
.map(|d| d.display().to_string())
.unwrap_or_default();
return Err(format!(
"campaign run needs a project: strategies resolve against the project \
store and vocabulary (no Aura.toml found up from {cwd})"
));
}
let registry = env.registry();
// Target resolution: a readable file is register-then-run sugar; a bare
// 64-hex token addresses the store directly; anything else refuses naming
// both readings.
let campaign_id = if Path::new(target).is_file() {
let doc = parse_valid_campaign(&PathBuf::from(target))?;
registry
.put_campaign(&campaign_to_json(&doc))
.map_err(|e| e.to_string())?
} else if is_content_id(target) {
target.to_string()
} else {
return Err(format!(
"'{target}' is neither a readable .json file nor a 64-hex content id"
));
};
// One uniform path from here: fetch the stored canonical bytes by id (so
// file addressing and id addressing produce the same realization by
// construction) and re-run the intrinsic tier on them.
let campaign_text = registry
.get_campaign(&campaign_id)
.map_err(|e| e.to_string())?
.ok_or_else(|| format!("no campaign {campaign_id} in the project store"))?;
let campaign = parse_campaign(&campaign_text)
.map_err(|e| doc_error_prose("stored campaign document", &e))?;
let faults = validate_campaign(&campaign);
if !faults.is_empty() {
return Err(fault_block(
"campaign document invalid:",
faults.iter().map(doc_fault_prose).collect(),
));
}
// Referential gate: zero faults or refuse (the campaign-validate seam).
let resolve = |t: &str| env.resolve(t);
let ref_faults = registry
.validate_campaign_refs(&campaign, &resolve)
.map_err(|e| e.to_string())?;
if !ref_faults.is_empty() {
return Err(fault_block(
"campaign references do not resolve:",
ref_faults.iter().map(ref_fault_prose).collect(),
));
}
// Process fetch + intrinsic validation (stored text, not a file path —
// parse_valid_process is file-based, so its constituents run here).
let DocRef::ContentId(process_id) = &campaign.process.r#ref else {
// validate_campaign already refuses identity process refs; defensive.
return Err("process.ref: a process is referenced by content id in this version".into());
};
let process_text = registry
.get_process(process_id)
.map_err(|e| e.to_string())?
.ok_or_else(|| format!("no process {process_id} in the project store"))?;
let process = parse_process(&process_text)
.map_err(|e| doc_error_prose("stored process document", &e))?;
let process_faults = validate_process(&process);
if !process_faults.is_empty() {
return Err(fault_block(
"process document invalid:",
process_faults.iter().map(doc_fault_prose).collect(),
));
}
// Strategies: canonical bytes from the store, index-aligned with the doc.
// The recorded id is the content id of those bytes (== the ref id for
// content refs; computed for identity refs).
let mut strategies: Vec<(String, String)> = Vec::with_capacity(campaign.strategies.len());
for entry in &campaign.strategies {
let canonical = match &entry.r#ref {
DocRef::ContentId(id) => registry
.get_blueprint(id)
.map_err(|e| e.to_string())?
.ok_or_else(|| format!("strategy {id} not found in the blueprint store"))?,
DocRef::IdentityId(id) => registry
.find_blueprint_by_identity(id, &resolve)
.map_err(|e| e.to_string())?
.ok_or_else(|| format!("identity id {id} matches no stored blueprint"))?,
};
strategies.push((content_id_of(&canonical), canonical));
}
// Loud deferral (#198 decision 6), once per run — pinned ORDER: this
// prints after the document gates and BEFORE any member executes, so a
// data refusal inside execute cannot swallow it.
if !campaign.presentation.persist_taps.is_empty() {
eprintln!(
"aura: persist_taps not yet honored ({} tap(s) ignored)",
campaign.presentation.persist_taps.len()
);
}
let runner = CliMemberRunner {
env,
server: Arc::new(data_server::DataServer::new(env.data_path())),
};
let outcome = aura_campaign::execute(
&campaign_id,
&campaign,
&process,
&strategies,
&runner,
&registry,
)
.map_err(|f| exec_fault_prose(&f))?;
// Zero-survivor stderr notes (exit stays 0 — a null result is a valid
// research result, #198 decision 8). Addressed by the fields the record
// already carries (strategy/instrument/window_ms), not by re-deriving a
// doc-order ordinal from the loop-nesting the executor happens to use
// today — that duplicated invariant would silently drift if the executor
// ever reorders its cell loop.
for cell in &outcome.record.cells {
for (stage_ix, st) in cell.stages.iter().enumerate() {
if matches!(&st.survivor_ordinals, Some(v) if v.is_empty()) {
eprintln!(
"aura: cell {}/{}/[{}, {}]: gate at stage {stage_ix} left no \
survivors; cell realization truncated",
cell.strategy, cell.instrument, cell.window_ms.0, cell.window_ms.1
);
}
}
}
// Emission: emit-gated family/selection lines per cell, then the
// always-on final record line. Matched by NAME against the two closed-
// vocabulary terms (self-evident at the call site, unlike positional
// indexing into the vocab slice); the debug_assert keeps the names
// honest against `aura_research::emit_vocabulary()` so a #190
// rename/extend of the closed set fails loudly here instead of silently
// misrouting emission.
debug_assert!(
aura_research::emit_vocabulary().contains(&"family_table")
&& aura_research::emit_vocabulary().contains(&"selection_report"),
"emit_vocabulary drifted from the names campaign_run matches by"
);
let emit_family = campaign.presentation.emit.iter().any(|e| e == "family_table");
let emit_selection = campaign.presentation.emit.iter().any(|e| e == "selection_report");
for cell_out in &outcome.cells {
if emit_family {
for fam in &cell_out.families {
for report in &fam.reports {
println!("{}", crate::family_member_line(&fam.family_id, report));
}
}
}
if emit_selection {
for sel in &cell_out.selections {
let line = SelectionReportLine {
selection_report: SelectionReportBody {
family_id: &sel.family_id,
stage: sel.stage,
block: sel.block,
winner_ordinal: sel.winner_ordinal,
params: &sel.params,
selection: &sel.selection,
},
};
println!(
"{}",
serde_json::to_string(&line).expect("selection report serializes")
);
}
}
}
println!(
"{}",
serde_json::to_string(&CampaignRunLine { campaign_run: &outcome.record })
.expect("campaign run record serializes")
);
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use aura_core::ScalarKind;
fn spec(name: &str) -> ParamSpec {
ParamSpec { name: name.to_string(), kind: ScalarKind::I64 }
}
#[test]
/// The only shape treated as a direct store address is a bare 64-char
/// lowercase-hex token; anything else (wrong length, uppercase, non-hex)
/// falls through to `run_campaign`'s file-path branch instead.
fn is_content_id_accepts_only_64_lowercase_hex() {
assert!(is_content_id(&"a".repeat(64)));
assert!(!is_content_id(&"A".repeat(64)));
assert!(!is_content_id(&"a".repeat(63)));
assert!(!is_content_id(&format!("{}g", "a".repeat(63))));
}
#[test]
/// bind_axes resolves a raw axis name against the ONE wrapped param whose
/// path ends with ".{raw}" (or equals it), producing a co-indexed point.
fn bind_axes_resolves_a_unique_suffix_match() {
let space = vec![spec("sma.length")];
let params = vec![("length".to_string(), Scalar::i64(14))];
let point = bind_axes(&space, "strat", &params).unwrap();
assert_eq!(point, vec![Scalar::i64(14).cell()]);
}
#[test]
/// A raw campaign axis that matches no open param of the wrapped
/// strategy is a named Bind fault naming the axis, never a silently
/// dropped binding.
fn bind_axes_refuses_an_unmatched_axis() {
let space = vec![spec("sma.length")];
let params = vec![("period".to_string(), Scalar::i64(14))];
let err = bind_axes(&space, "strat", &params).unwrap_err();
let MemberFault::Bind(m) = err else { panic!("expected Bind, got {err:?}") };
assert!(m.contains("period") && m.contains("matches no open param"));
}
#[test]
/// A raw axis matching MORE THAN ONE wrapped param is a named Bind fault
/// that lists every ambiguous candidate, never a silent first-match pick.
fn bind_axes_refuses_an_ambiguous_axis() {
let space = vec![spec("fast.length"), spec("slow.length")];
let params = vec![("length".to_string(), Scalar::i64(14))];
let err = bind_axes(&space, "strat", &params).unwrap_err();
let MemberFault::Bind(m) = err else { panic!("expected Bind, got {err:?}") };
assert!(m.contains("ambiguous") && m.contains("fast.length") && m.contains("slow.length"));
}
#[test]
/// An open wrapped param left unbound by any campaign axis is a named
/// Bind fault naming the param, not an implicit default.
fn bind_axes_refuses_an_uncovered_param() {
let space = vec![spec("sma.length"), spec("sma.threshold")];
let params = vec![("length".to_string(), Scalar::i64(14))];
let err = bind_axes(&space, "strat", &params).unwrap_err();
let MemberFault::Bind(m) = err else { panic!("expected Bind, got {err:?}") };
assert!(m.contains("threshold") && m.contains("bound by no campaign axis"));
}
}
+1
View File
@@ -14,6 +14,7 @@
mod render;
mod graph_construct;
mod project;
mod campaign_run;
mod research_docs;
mod scaffold;
use render::{ChartData, ChartMeta, ChartMode, ReduceKind, Series};
+9 -5
View File
@@ -67,7 +67,7 @@ fn guard_one_mode(cmd: &DocIntrospectCmd, family: &str) {
}
}
fn doc_error_prose(what: &str, e: &DocError) -> String {
pub(crate) fn doc_error_prose(what: &str, e: &DocError) -> String {
match e {
DocError::NotJson(msg) => format!("{what} is not JSON: {msg}"),
DocError::BadFormatVersion(v) => {
@@ -84,7 +84,7 @@ fn doc_error_prose(what: &str, e: &DocError) -> String {
}
}
fn doc_fault_prose(f: &DocFault) -> String {
pub(crate) fn doc_fault_prose(f: &DocFault) -> String {
match f {
DocFault::EmptyPipeline => "the pipeline is empty — a process needs at least one stage".into(),
DocFault::UnknownMetric { stage, metric } => {
@@ -121,7 +121,7 @@ fn read_doc(file: &PathBuf) -> Result<String, String> {
std::fs::read_to_string(file).map_err(|e| format!("cannot read {}: {e}", file.display()))
}
fn fault_block(header: &str, lines: Vec<String>) -> String {
pub(crate) fn fault_block(header: &str, lines: Vec<String>) -> String {
format!("{header}\n {}", lines.join("\n "))
}
@@ -244,9 +244,12 @@ pub enum CampaignSub {
Introspect(DocIntrospectCmd),
/// Register a valid campaign document into the store under the runs root.
Register { file: PathBuf },
/// Execute a stored campaign into a realized run-set (a .json file is
/// register-then-run sugar; the canonical address is the content id).
Run { target: String },
}
fn ref_fault_prose(f: &RefFault) -> String {
pub(crate) fn ref_fault_prose(f: &RefFault) -> String {
match f {
RefFault::ProcessNotFound(id) => format!("process {id} not found in the project store"),
RefFault::StrategyNotFound(id) => format!("strategy {id} not found in the blueprint store"),
@@ -272,6 +275,7 @@ pub fn campaign_cmd(cmd: CampaignCmd, env: &Env) {
introspect_campaign(i)
}
CampaignSub::Register { file } => register_campaign(file, env),
CampaignSub::Run { target } => crate::campaign_run::run_campaign(target, env),
};
if let Err(m) = result {
eprintln!("aura: {m}");
@@ -279,7 +283,7 @@ pub fn campaign_cmd(cmd: CampaignCmd, env: &Env) {
}
}
fn parse_valid_campaign(file: &PathBuf) -> Result<CampaignDoc, String> {
pub(crate) fn parse_valid_campaign(file: &PathBuf) -> Result<CampaignDoc, String> {
let text = read_doc(file)?;
let doc = parse_campaign(&text).map_err(|e| doc_error_prose("campaign document", &e))?;
let faults = validate_campaign(&doc);