use std::collections::HashMap;
use std::sync::Arc;
use crate::activity::Activity;
use crate::adapter::{
find_adapter_registration, registered_adapter_params, registered_driver_names,
};
use crate::bindings::build_workload_root_kernel;
use crate::opseq::SequencerType;
use crate::synthesis::OpBuilder;
use nmbrs_metrics::labels::Labels;
use nmbrs_metrics::scheduler::Reporter;
use nmbrs_workload::tags::TagFilter;
static KNOWN_PARAMS: std::sync::OnceLock<Vec<&'static str>> = std::sync::OnceLock::new();
pub fn install_known_params(keys: Vec<&'static str>) {
let _ = KNOWN_PARAMS.set(keys);
}
fn known_params() -> Option<&'static [&'static str]> {
KNOWN_PARAMS.get().map(|v| v.as_slice())
}
pub(crate) fn is_cli_param(name: &str) -> bool {
known_params().map(|p| p.contains(&name)).unwrap_or(true)
}
static KNOWN_BARE_FLAGS: std::sync::OnceLock<Vec<&'static str>> = std::sync::OnceLock::new();
static KNOWN_VALUE_FLAGS: std::sync::OnceLock<Vec<&'static str>> = std::sync::OnceLock::new();
pub fn install_known_flags(bare: Vec<&'static str>, value: Vec<&'static str>) {
let _ = KNOWN_BARE_FLAGS.set(bare);
let _ = KNOWN_VALUE_FLAGS.set(value);
}
pub(crate) fn is_recognized_bare_flag(arg: &str) -> bool {
arg == "--refine"
|| KNOWN_BARE_FLAGS
.get()
.map(|v| v.iter().any(|f| *f == arg))
.unwrap_or_else(|| RECOGNIZED_BARE_FLAGS.contains(&arg))
}
pub(crate) fn known_value_flags() -> &'static [&'static str] {
KNOWN_VALUE_FLAGS
.get()
.map(|v| v.as_slice())
.unwrap_or(SESSION_DIR_FLAGS)
}
fn cli_flag_value(args: &[String], flag: &str) -> Option<String> {
let eq_prefix = format!("{flag}=");
let mut iter = args.iter();
while let Some(arg) = iter.next() {
if let Some(rest) = arg.strip_prefix(&eq_prefix) {
return Some(rest.to_string());
}
if arg == flag {
return iter.next().cloned();
}
}
None
}
pub const DEFERRED_STDOUT_FILE: &str = ".report_stdout.md";
pub fn report_config_from_summary(
config: &nmbrs_workload::model::SummaryConfig,
exec_id_filter: Option<u64>,
) -> nmbrs_metrics::reporters::sqlite::ReportConfig {
nmbrs_metrics::reporters::sqlite::ReportConfig {
columns: config.columns.clone(),
row_filters: config.row_filters.clone(),
aggregates: config
.aggregates
.iter()
.map(|a| nmbrs_metrics::reporters::sqlite::ReportAggregate {
function: a.function.to_string(),
column_pattern: a.column_pattern.clone(),
label_key: a.label_key.clone(),
label_pattern: a.label_pattern.clone(),
group_by: a.group_by.clone(),
})
.collect(),
show_details: config.show_details,
exec_id_filter,
}
}
pub fn resolve_workload_file_public(name: &str) -> Option<String> {
resolve_workload_file(name)
}
pub fn scenarios_in_workload_file(path: &str) -> Vec<String> {
let Ok(src) = std::fs::read_to_string(path) else {
return Vec::new();
};
let Ok(doc) = serde_yaml::from_str::<serde_yaml::Value>(&src) else {
return Vec::new();
};
let Some(scenarios) = doc.get("scenarios") else {
return Vec::new();
};
let Some(map) = scenarios.as_mapping() else {
return Vec::new();
};
map.keys()
.filter_map(|k| k.as_str().map(String::from))
.collect()
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord)]
pub enum ExecDepth {
Phase,
Dispenser,
Op,
Cycle,
Full,
}
#[derive(Clone)]
pub struct DiagnosticConfig {
pub depth: ExecDepth,
pub show_wiring: bool,
pub show_labels: bool,
pub list_controls: bool,
}
impl DiagnosticConfig {
pub fn normal() -> Self {
Self {
depth: ExecDepth::Full,
show_wiring: false,
show_labels: false,
list_controls: false,
}
}
pub fn parse(spec: &str) -> Self {
let mut config = Self::normal();
let mut depth_set = false;
for flag in spec.split(',') {
match flag.trim() {
"phase" => {
config.depth = ExecDepth::Phase;
depth_set = true;
}
"dispenser" => {
config.depth = ExecDepth::Dispenser;
depth_set = true;
}
"op" => {
config.depth = ExecDepth::Op;
depth_set = true;
}
"cycle" => {
config.depth = ExecDepth::Cycle;
depth_set = true;
}
"full" => {
config.depth = ExecDepth::Full;
depth_set = true;
}
"wiring" => {
config.show_wiring = true;
if !depth_set {
config.depth = ExecDepth::Op;
depth_set = true;
}
}
"labels" => config.show_labels = true,
"controls" => {
config.list_controls = true;
config.depth = ExecDepth::Phase;
depth_set = true;
}
"fields" | "silent" | "json" => {}
"kernels" => {}
_ => crate::diag!(
crate::observer::LogLevel::Warn,
"warning: unknown dryrun flag '{flag}'"
),
}
}
if !depth_set {
config.depth = ExecDepth::Phase;
}
config
}
}
pub fn render_controls_tree(
root: &std::sync::Arc<std::sync::RwLock<nmbrs_metrics::component::Component>>,
out: &mut dyn std::io::Write,
) -> std::io::Result<()> {
use nmbrs_metrics::component::find;
use nmbrs_metrics::selector::Selector;
writeln!(out, "Declared dynamic controls (SRD 23):")?;
let all = find(root, &Selector::new());
let mut entries: Vec<(String, String, String, String, String, String)> = Vec::new();
for comp in all {
let guard = match comp.read() {
Ok(g) => g,
Err(_) => continue,
};
let path = guard
.effective_labels()
.iter()
.map(|(k, v)| format!("{k}={v}"))
.collect::<Vec<_>>()
.join(",");
for ctl in guard.controls().list() {
let scope = match ctl.branch_scope() {
nmbrs_metrics::controls::BranchScope::Local => "local",
nmbrs_metrics::controls::BranchScope::Subtree => "subtree",
};
let final_marker = match ctl.final_scope() {
Some(s) => format!("final@{s}"),
None => "-".to_string(),
};
entries.push((
if path.is_empty() {
"<root>".into()
} else {
path.clone()
},
ctl.name().to_string(),
ctl.value_type_name().to_string(),
ctl.value_string(),
format!(
"scope={scope}, {final_marker}, appliers={}",
ctl.applier_count()
),
if ctl.accepts_f64_writes() {
"f64-writable".into()
} else {
"no-f64".into()
},
));
}
}
if entries.is_empty() {
writeln!(out, " (no controls declared)")?;
return Ok(());
}
entries.sort();
for (path, name, ty, value, meta, write) in entries {
writeln!(
out,
" {path}\n {name}: {value} [{ty}] {meta} {write}",
)?;
}
Ok(())
}
pub fn render_scope_elision_summary(
tree: &crate::scope_tree::ScopeTree,
out: &mut dyn std::io::Write,
) -> std::io::Result<()> {
let summary = crate::scope_elision::elision_summary(tree);
let name_width = summary
.iter()
.map(|(_, _, _, name, _)| name.len())
.max()
.unwrap_or(0)
.max(48);
writeln!(out, "scope elision summary")?;
writeln!(out, "------------------------")?;
for (idx, _depth, materialised, logical_name, _kind) in &summary {
match materialised {
Some(true) => {
writeln!(
out,
"{:<width$} materialised=true",
logical_name,
width = name_width
)?;
}
Some(false) => {
let elides_to = tree
.nearest_materialised(*idx)
.map(|p| tree.nodes[p].logical_name.clone())
.unwrap_or_else(|| "<unknown>".to_string());
writeln!(
out,
"{:<width$} materialised=false elides-to={}",
logical_name,
elides_to,
width = name_width
)?;
}
None => {
writeln!(
out,
"{:<width$} materialised=unknown",
logical_name,
width = name_width
)?;
}
}
}
Ok(())
}
pub async fn run(args: &[String]) -> Result<(), String> {
let stripped: &[String] = match args.first().map(|s| s.as_str()) {
Some("run") => &args[1..],
_ => args,
};
let cli_params = parse_params(stripped);
let min_level = cli_params
.get("loglevel")
.or_else(|| cli_params.get("loglevel-display"))
.or_else(|| cli_params.get("loglevel_display"))
.and_then(|s| parse_log_level(s))
.unwrap_or(crate::observer::LogLevel::Info);
let retain_level = cli_params
.get("loglevel-retain")
.or_else(|| cli_params.get("loglevel_retain"))
.and_then(|s| parse_log_level(s))
.unwrap_or(crate::observer::LogLevel::Debug);
crate::observer::set_retain_level(retain_level);
crate::observer::set_display_level(min_level);
run_with_observer(
args,
Arc::new(crate::observer::StderrObserver::with_min_level(min_level)),
)
.await
}
pub fn parse_log_level(s: &str) -> Option<crate::observer::LogLevel> {
use crate::observer::LogLevel;
match s.trim().to_ascii_lowercase().as_str() {
"trace" | "trc" => Some(LogLevel::Trace),
"debug" | "dbg" => Some(LogLevel::Debug),
"info" | "inf" => Some(LogLevel::Info),
"warn" | "wrn" | "warning" => Some(LogLevel::Warn),
"error" | "err" => Some(LogLevel::Error),
_ => None,
}
}
pub async fn run_with_observer(
args: &[String],
observer: Arc<dyn crate::observer::RunObserver>,
) -> Result<(), String> {
let args: &[String] = match args.first().map(|s| s.as_str()) {
Some("run") => &args[1..],
Some(cmd) if !cmd.contains('=') && !cmd.ends_with(".yaml") && !cmd.ends_with(".yml") => {
return Err(format!(
"unknown command '{cmd}'. Use 'run' or pass a workload file."
));
}
_ => args,
};
detect_conflicting_duplicate_params(args)?;
run_impl(args, observer).await
}
fn build_session_metrics(
session: &crate::session::Session,
sqlite_reporter: &std::sync::Arc<
std::sync::Mutex<Option<nmbrs_metrics::reporters::sqlite::SqliteReporter>>,
>,
observer: &Arc<dyn crate::observer::RunObserver>,
merged_params: &HashMap<String, String>,
openmetrics_url: &Option<String>,
args: &[String],
params: &HashMap<String, String>,
) -> Result<
(
std::sync::Arc<nmbrs_metrics::cadence_reporter::CadenceReporter>,
nmbrs_metrics::cadence::CadenceTree,
std::sync::Arc<nmbrs_metrics::metrics_query::MetricsQuery>,
std::sync::Arc<nmbrs_metrics::scheduler::StopHandle>,
),
String,
> {
let (base_interval, cadences) = resolve_cadence_config(merged_params, observer)?;
let cadence_tree = nmbrs_metrics::cadence::CadenceTree::plan_validated(
cadences,
nmbrs_metrics::cadence::DEFAULT_MAX_FAN_IN,
base_interval,
)
.map_err(|e| format!("cadence tree: {e}"))?;
let cadence_reporter = Arc::new(nmbrs_metrics::cadence_reporter::CadenceReporter::new(
cadence_tree.clone(),
));
let metrics_query = Arc::new(nmbrs_metrics::metrics_query::MetricsQuery::new(
cadence_reporter.clone(),
session.component.clone(),
));
session.set_metrics_query(metrics_query.clone());
nmbrs_metrics::polydat_nodes::set_global_query(metrics_query.clone());
let mem_access = std::sync::Arc::new(nmbrs_metrics::queryapi::MetricsQueryAccess::new(
metrics_query.clone(),
));
let composed: std::sync::Arc<dyn nmbrs_metrics::queryapi::MetricAccess> = {
let db = session.output_dir.join("metrics.db");
match nmbrs_metrics::queryapi::sqlite::SqliteDataSource::open(&db) {
Ok(cold) => {
let cold = cold.with_execution_selection(
nmbrs_metrics::queryapi::sqlite::ExecutionSelection::All,
);
let mem_for_horizon = mem_access.clone();
std::sync::Arc::new(nmbrs_metrics::queryapi::HybridStore::new(vec![
nmbrs_metrics::queryapi::Tier::new(
mem_access.clone(),
std::sync::Arc::new(move || mem_for_horizon.earliest_ms()),
),
nmbrs_metrics::queryapi::Tier::unbounded(std::sync::Arc::new(cold)),
]))
}
Err(e) => {
crate::diag!(
crate::observer::LogLevel::Debug,
"metrics: hybrid sqlite tail unavailable ({e}); live reads are in-memory only"
);
mem_access.clone()
}
}
};
nmbrs_metrics::queryapi::install_live_access(std::sync::Arc::new(
nmbrs_metrics::queryapi::ExecScopedAccess::new(composed),
));
nmbrs_metrics::queryapi::install_read_exec_id_hook(|| {
crate::execution_context::try_current().map(|c| c.exec_id)
});
observer.on_metrics_query(metrics_query.clone());
let session_for_capture = session.component.clone();
let mut sched_builder = nmbrs_metrics::scheduler::SchedulerBuilder::new()
.base_interval(base_interval)
.with_cadence_reporter(cadence_reporter.clone())
.with_cadence_tree(cadence_tree.clone());
let sqlite_cadence = cadence_tree.align_to_declared(std::time::Duration::from_secs(30));
if let (Some(cadence), Ok(guard)) = (sqlite_cadence, sqlite_reporter.lock())
&& guard.is_some()
{
drop(guard);
let sqlite_for_sub = sqlite_reporter.clone();
match cadence_reporter.subscribe(
cadence,
Box::new(MutexReporter(sqlite_for_sub)),
nmbrs_metrics::cadence_reporter::SubscriptionOpts::default(),
) {
Ok(_) => {
crate::diag!(
crate::observer::LogLevel::Info,
"metrics: SQLite writes every {:?} (WAL mode)",
cadence
);
}
Err(e) => {
crate::diag!(
crate::observer::LogLevel::Warn,
"metrics: SQLite subscription failed: {e}"
);
}
}
}
observer.session_dir_ready(&session.output_dir);
if let Some(setting) = params.get("sysmon") {
let selection = crate::sysmon::parse_selection(setting)?;
let interval = match params.get("sysmon-interval") {
None => std::time::Duration::from_secs(5),
Some(v) => match v.trim_end_matches('s').parse::<f64>() {
Ok(secs) if secs > 0.0 => std::time::Duration::from_secs_f64(secs),
_ => {
return Err(format!(
"sysmon-interval: expected seconds (e.g. `sysmon-interval=5`), got '{v}'"
));
}
},
};
let membw_peak_bytes_per_s = match params.get("sysmon-membw-gbps") {
None => None,
Some(v) => match v.parse::<f64>() {
Ok(gbps) if gbps > 0.0 => Some(gbps * 1e9),
_ => {
return Err(format!(
"sysmon-membw-gbps: expected the host's peak memory bandwidth \
in GB/s, got '{v}'"
));
}
},
};
let mut config = crate::sysmon::SysmonConfig {
cats: crate::sysmon::Categories::ALL,
interval,
membw_peak_bytes_per_s,
};
match selection {
crate::sysmon::Selection::Any => {
let (cats, skipped) = crate::sysmon::resolve_any(&config);
config.cats = cats;
for reason in skipped {
crate::diag!(
crate::observer::LogLevel::Warn,
"sysmon: skipping a subsystem: {reason}"
);
}
}
crate::sysmon::Selection::Cats(cats) => {
config.cats = cats;
crate::sysmon::check_rambw_requirements(&config)?;
}
}
crate::sysmon::spawn(config, session.component.clone(), observer.clone())?;
crate::diag!(
crate::observer::LogLevel::Info,
"sysmon: sampling {setting} every {interval:?}"
);
}
let metrics_log_setting: Option<String> = args
.iter()
.find_map(|a| {
a.strip_prefix("--metrics-log")
.map(|rest| rest.strip_prefix('=').unwrap_or("true").to_string())
})
.or_else(|| params.get("metrics-log").cloned())
.or_else(|| std::env::var("NMBRS_METRICS_LOG").ok());
let metrics_log_path = match metrics_log_setting.as_deref() {
None => None,
Some("0") | Some("false") | Some("no") | Some("off") | Some("") => None,
Some("1") | Some("true") | Some("yes") | Some("on") => {
Some(session.output_dir.join("metrics.jsonl"))
}
Some(explicit) => {
let from_flag = args.iter().any(|a| a.starts_with("--metrics-log"));
if from_flag {
Some(std::path::PathBuf::from(explicit))
} else {
Some(crate::session::confine_to_dir(&session.output_dir, explicit).map_err(
|e| format!("metrics-log: {e} (an explicit path outside the session directory is only accepted from the --metrics-log= flag)"),
)?)
}
}
};
if let Some(log_path) = metrics_log_path {
match nmbrs_metrics::reporters::metrics_log::MetricsLogReporter::new(&log_path) {
Ok(reporter) => {
if let Some(cadence) =
cadence_tree.align_to_declared(std::time::Duration::from_secs(30))
{
match cadence_reporter.subscribe(
cadence,
Box::new(reporter),
nmbrs_metrics::cadence_reporter::SubscriptionOpts::default(),
) {
Ok(_) => {
crate::diag!(
crate::observer::LogLevel::Info,
"metrics: JSONL log every {:?} -> {} (session db unaffected)",
cadence,
log_path.display()
);
}
Err(e) => {
crate::diag!(
crate::observer::LogLevel::Warn,
"metrics: metrics log subscribe failed: {e}"
);
}
}
}
}
Err(e) => {
crate::diag!(
crate::observer::LogLevel::Warn,
"metrics: metrics log disabled: {e}"
);
}
}
}
let per_instance_enabled = args.iter().any(|a| a == "--per-instance-metrics")
|| params
.get("per-instance-metrics")
.map(|s| matches!(s.as_str(), "1" | "true" | "yes" | "on"))
.unwrap_or(false)
|| std::env::var("NMBRS_PER_INSTANCE_METRICS")
.ok()
.map(|s| matches!(s.as_str(), "1" | "true" | "yes" | "on"))
.unwrap_or(false);
if per_instance_enabled {
let per_instance_dir = session.output_dir.join("metrics");
match nmbrs_metrics::reporters::per_instance::PerInstanceReporter::new(&per_instance_dir) {
Ok(reporter) => {
if let Some(cadence) =
cadence_tree.align_to_declared(std::time::Duration::from_secs(30))
{
match cadence_reporter.subscribe(
cadence,
Box::new(reporter),
nmbrs_metrics::cadence_reporter::SubscriptionOpts::default(),
) {
Ok(_) => {
crate::diag!(
crate::observer::LogLevel::Info,
"metrics: per-instance JSONL writes every {:?} into {}",
cadence,
per_instance_dir.display()
);
}
Err(e) => {
crate::diag!(
crate::observer::LogLevel::Warn,
"metrics: per-instance subscription failed: {e}"
);
}
}
}
}
Err(e) => {
crate::diag!(
crate::observer::LogLevel::Warn,
"metrics: per-instance reporter disabled ({}): {e}",
per_instance_dir.display()
);
}
}
}
if let Some(url) = openmetrics_url.as_ref()
&& let Some(cadence) = cadence_tree.align_to_declared(std::time::Duration::from_secs(10))
{
let jobname = merged_params
.get("jobname")
.cloned()
.unwrap_or_else(|| "default".to_string());
let instance = merged_params
.get("instance")
.cloned()
.unwrap_or_else(|| "default".to_string());
let mut vm =
match nmbrs_metrics::reporters::victoriametrics::VictoriaMetricsReporter::from_spec(url)
{
Ok(r) => r,
Err(_) => {
nmbrs_metrics::reporters::victoriametrics::VictoriaMetricsReporter::new(url)
}
};
vm = vm.with_jobname(jobname).with_instance(instance);
if let Some(token_path) = merged_params.get("prompush_apikeyfile") {
match vm.with_bearer_token_file(token_path) {
Ok(r) => vm = r,
Err(e) => {
crate::diag!(
crate::observer::LogLevel::Warn,
"prompush_apikeyfile '{token_path}': {e}"
);
vm = nmbrs_metrics::reporters::victoriametrics
::VictoriaMetricsReporter::from_spec(url)
.unwrap_or_else(|_| nmbrs_metrics::reporters::victoriametrics
::VictoriaMetricsReporter::new(url))
.with_jobname(
merged_params.get("jobname").cloned()
.unwrap_or_else(|| "default".to_string()),
)
.with_instance(
merged_params.get("instance").cloned()
.unwrap_or_else(|| "default".to_string()),
);
}
}
}
let _ = cadence_reporter.subscribe(
cadence,
Box::new(vm),
nmbrs_metrics::cadence_reporter::SubscriptionOpts::default(),
);
}
for (interval, reporter) in observer.reporters() {
sched_builder = sched_builder.add_reporter(interval, BoxedReporter(reporter));
}
let scheduler = sched_builder.build(Box::new(move || {
nmbrs_metrics::component::capture_tree(&session_for_capture, base_interval)
}));
let stop_handle = Arc::new(scheduler.start());
crate::session_signals::install_signal_handler();
let dump_root = session.component.clone();
crate::session_signals::set_diag_dump_hook(Box::new(move || {
fn walk(
node: &std::sync::Arc<std::sync::RwLock<nmbrs_metrics::component::Component>>,
out: &mut Vec<String>,
) {
let g = node.read().unwrap_or_else(|e| e.into_inner());
if g.state() == nmbrs_metrics::component::ComponentState::Running {
out.push(format!(
"{} ({} instrument(s))",
g.effective_labels(),
g.instruments().len(),
));
}
for child in g.children() {
walk(child, out);
}
}
let mut lines = Vec::new();
walk(&dump_root, &mut lines);
crate::diag!(
crate::observer::LogLevel::Info,
"session: SIGQUIT inventory — {} running component(s)",
lines.len()
);
for line in lines {
crate::diag!(crate::observer::LogLevel::Info, " {line}");
}
}));
Ok((cadence_reporter, cadence_tree, metrics_query, stop_handle))
}
struct SessionHost {
session: crate::session::Session,
sqlite_reporter:
std::sync::Arc<std::sync::Mutex<Option<nmbrs_metrics::reporters::sqlite::SqliteReporter>>>,
cadence_reporter: std::sync::Arc<nmbrs_metrics::cadence_reporter::CadenceReporter>,
#[allow(dead_code)]
cadence_tree: nmbrs_metrics::cadence::CadenceTree,
#[allow(dead_code)]
metrics_query: std::sync::Arc<nmbrs_metrics::metrics_query::MetricsQuery>,
stop_handle: std::sync::Arc<nmbrs_metrics::scheduler::StopHandle>,
refine_plan: Option<Arc<crate::refine_plan::RefinePlan>>,
resume_target: Option<std::path::PathBuf>,
refine_requested: bool,
refine_scope: Option<String>,
stick_reattached: Option<String>,
profiler: Option<crate::profiler::ProfileGuard>,
sqlite_guard: nmbrs_metrics::reporters::sqlite::SqliteShutdownGuard,
checkpoint_writer: std::sync::Arc<crate::checkpoint::CheckpointWriter>,
saved_doc: Option<crate::checkpoint::Checkpoint>,
}
impl SessionHost {
fn setup(
args: &[String],
observer: Arc<dyn crate::observer::RunObserver>,
) -> Result<SessionHost, String> {
crate::observer::set_global_observer(observer.clone());
nmbrs_errorhandler::handlers::set_log_fn(|msg| {
let indent = crate::scene_tree::running_phase_indent();
crate::observer::log(crate::observer::LogLevel::Debug, &format!("{indent}{msg}"));
});
nmbrs_metrics::diag::set_warn_fn(|msg| {
let indent = crate::scene_tree::running_phase_indent();
crate::observer::log(crate::observer::LogLevel::Warn, &format!("{indent}{msg}"));
});
nmbrs_metrics::diag::set_info_fn(|msg| {
let indent = crate::scene_tree::running_phase_indent();
crate::observer::log(crate::observer::LogLevel::Info, &format!("{indent}{msg}"));
});
let args = normalize_args(args);
let params = parse_params(&args);
let eff_params = effective_params(&args);
let openmetrics_url: Option<String> = cli_flag_value(&args[..], "--report-openmetrics-to")
.or_else(|| {
args.iter()
.find_map(|a| a.strip_prefix("report-openmetrics-to="))
.map(|s| s.to_string())
});
let stick_reattach: Option<std::path::PathBuf> = {
let cli_stick: Option<bool> =
params.get("stick_session").map(|v| v == "true" || v == "1");
let stick_on =
cli_stick.unwrap_or_else(|| peek_stick_session(¶ms, &args).unwrap_or(false));
let reattach_expressed = args
.iter()
.any(|a| a == "--refine" || a == "--resume-latest")
|| params.contains_key("resume")
|| params.contains_key("resume_latest");
if !stick_on || reattach_expressed || crate::session::args_request_dryrun(&args) {
None
} else {
let spec = crate::session::resolve_session_dir(&args);
let operator_selected = !spec.is_empty()
|| spec.force_new
|| spec.reuse != crate::session::SessionReuse::Error;
if operator_selected {
None
} else {
let latest = crate::session::default_sessions_root().join("latest");
std::fs::read_link(&latest)
.ok()
.map(|t| {
if t.is_absolute() {
t
} else {
crate::session::default_sessions_root().join(t)
}
})
.filter(|d| d.join("checkpoint.jsonl").is_file())
}
}
};
let stick_reattached: Option<String> = stick_reattach
.as_ref()
.and_then(|d| d.file_name())
.and_then(|s| s.to_str())
.map(String::from);
if let Some(id) = stick_reattached.as_deref() {
crate::diag!(
crate::observer::LogLevel::Info,
"stick_session: re-attaching to {id} — pass `--session new` to start fresh"
);
}
let resume_target: Option<std::path::PathBuf> = {
let explicit = params.get("resume").filter(|s| !s.is_empty()).map(|s| {
let p = std::path::PathBuf::from(s);
if p.is_file() {
p
} else if p.is_dir() {
p.join("checkpoint.jsonl")
} else {
crate::session::default_sessions_root()
.join(s)
.join("checkpoint.jsonl")
}
});
let resume_latest = params
.get("resume_latest")
.map(|s| s != "false" && s != "0")
.unwrap_or(false)
|| args.iter().any(|a| a == "--resume-latest")
|| stick_reattach.is_some();
if resume_latest {
let latest = crate::session::default_sessions_root().join("latest");
let resolved = std::fs::read_link(&latest)
.ok()
.map(|target| {
if target.is_absolute() {
target
} else {
crate::session::default_sessions_root().join(target)
}
})
.map(|d| d.join("checkpoint.jsonl"));
explicit.or(resolved)
} else {
explicit
}
};
let refine_requested = args.iter().any(|a| a == "--refine") || stick_reattach.is_some();
let scenario_for_session = params
.get("scenario")
.map(|s| s.as_str())
.unwrap_or("default");
let refine_scope: Option<&str> = params
.get("scope")
.map(|s| s.as_str())
.filter(|_| refine_requested);
let refine_plan: Option<Arc<crate::refine_plan::RefinePlan>> = if refine_requested {
resume_target
.as_ref()
.and_then(|p| p.parent().map(|d| d.to_path_buf()))
.and_then(|prior_dir| {
if !prior_dir.exists() {
crate::diag!(
crate::observer::LogLevel::Warn,
"refine: prior session dir not found ({}); \
running every phase as if this were a fresh `nmbrs run`",
prior_dir.display()
);
return None;
}
let mut plan =
crate::refine_plan::RefinePlan::load_from_session_dir(&prior_dir);
if plan.is_none() {
crate::diag!(
crate::observer::LogLevel::Warn,
"refine: no readable phase_outcomes in {}; \
running every phase as if this were a fresh `nmbrs run`",
prior_dir.display()
);
}
if let Some(p) = plan.as_mut() {
p.scope = match refine_scope {
Some("all") => {
crate::diag!(
crate::observer::LogLevel::Info,
"refine: scope=all — every phase will run \
under exec_id={}",
p.next_exec_id
);
crate::refine_plan::RefineScope::All
}
Some("changed") => {
crate::diag!(
crate::observer::LogLevel::Info,
"refine: scope=changed — comparing each \
phase's program hash against the prior \
outcome; unchanged phases skip, changed \
phases re-run under exec_id={}",
p.next_exec_id
);
crate::refine_plan::RefineScope::Changed
}
_ => crate::refine_plan::RefineScope::Missing,
};
}
plan.map(Arc::new)
})
} else {
None
};
let session = match (refine_plan.as_ref(), resume_target.as_ref()) {
(Some(plan), Some(p)) if p.exists() => {
let prior_dir = p
.parent()
.map(|d| d.to_path_buf())
.unwrap_or_else(crate::session::latest_session_dir);
crate::diag!(
crate::observer::LogLevel::Info,
"refine: attached to session {}; \
prior outcomes={}, completed phases to skip={}, \
next exec_id={}",
prior_dir.display(),
plan.prior_outcomes_seen,
plan.completed.len(),
plan.next_exec_id
);
crate::session::Session::reattach(prior_dir, scenario_for_session)
}
(_, Some(p)) if p.exists() => {
let prior_dir = p
.parent()
.map(|d| d.to_path_buf())
.unwrap_or_else(crate::session::latest_session_dir);
crate::session::Session::reattach(prior_dir, scenario_for_session)
}
_ => crate::session::Session::new_with_args(scenario_for_session, &args),
};
let session_log_path = session.output_dir.join("session.log");
if let Err(e) = crate::observer::set_log_file(&session_log_path) {
crate::diag!(
crate::observer::LogLevel::Warn,
"warning: failed to open session log {}: {e}",
session_log_path.display()
);
}
let trace_specs = collect_repeated_flag(&args, "trace");
match crate::trace_router::init(&trace_specs, &session.output_dir) {
Ok(0) => {} Ok(n) => crate::diag!(
crate::observer::LogLevel::Info,
"trace router: {n} route(s) configured"
),
Err(e) => crate::diag!(
crate::observer::LogLevel::Warn,
"trace router init failed: {e}"
),
}
crate::diag!(
crate::observer::LogLevel::Info,
"session: {} ({})",
session.id,
session.output_dir.display()
);
polydat::set_panic_reporting_downstream(true);
polydat::audit::set_log_fn(|level, msg| {
use polydat::audit::LogLevel as AuditLevel;
let mapped = match level {
AuditLevel::Trace | AuditLevel::Debug => crate::observer::LogLevel::Debug,
AuditLevel::Info => crate::observer::LogLevel::Info,
AuditLevel::Warn => crate::observer::LogLevel::Warn,
AuditLevel::Error => crate::observer::LogLevel::Error,
};
crate::observer::log(mapped, &format!("[lib] {msg}"));
});
let sqlite_path = session.metrics_path();
let sqlite_reporter = nmbrs_metrics::reporters::sqlite::SqliteReporter::new(&sqlite_path)
.map(|mut r| {
r.set_metadata("session", &session.id);
crate::diag!(
crate::observer::LogLevel::Info,
"metrics: {}",
sqlite_path.display()
);
r
})
.map_err(|e| {
crate::diag!(
crate::observer::LogLevel::Warn,
"warning: SQLite metrics disabled: {e}"
)
})
.ok();
let sqlite_reporter = std::sync::Arc::new(std::sync::Mutex::new(sqlite_reporter));
let _sqlite_shutdown_guard =
nmbrs_metrics::reporters::sqlite::SqliteShutdownGuard::new(sqlite_reporter.clone());
{
let reporter = sqlite_reporter.clone();
tokio::spawn(async move {
let mut interval = tokio::time::interval(std::time::Duration::from_secs(60));
interval.tick().await;
loop {
interval.tick().await;
if let Ok(g) = reporter.lock()
&& let Some(r) = g.as_ref()
{
r.passive_checkpoint();
}
}
});
}
let (cadence_reporter, cadence_tree, metrics_query, stop_handle) = build_session_metrics(
&session,
&sqlite_reporter,
&observer,
&eff_params,
&openmetrics_url,
&args,
&eff_params,
)?;
crate::session_signals::install_signal_handler();
let _profiler =
crate::profiler::ProfileGuard::maybe_start(¶ms, Some(&session.output_dir));
let checkpoint_path = session.output_dir.join("checkpoint.jsonl");
let saved_doc = match resume_target.as_ref() {
Some(p) => match crate::checkpoint::storage::read(p) {
Ok(Some(doc)) => Some(doc),
Ok(None) => {
crate::diag!(
crate::observer::LogLevel::Warn,
"resume: no checkpoint found at {} — fresh session",
p.display()
);
None
}
Err(e) => return Err(format!("resume: {e}")),
},
None => None,
};
let invocation = saved_doc.as_ref().map(|d| d.invocation + 1).unwrap_or(1);
let started_at = saved_doc
.as_ref()
.map(|d| d.started_at.clone())
.unwrap_or_else(crate::checkpoint::storage::now_rfc3339);
let checkpoint_writer =
std::sync::Arc::new(if crate::session::args_request_dryrun(&args) {
crate::checkpoint::CheckpointWriter::disabled(checkpoint_path.clone())
} else {
match saved_doc.as_ref() {
Some(_doc) => crate::checkpoint::CheckpointWriter::from_existing(
checkpoint_path.clone(),
saved_doc.clone().unwrap(),
crate::checkpoint::storage::now_rfc3339(),
invocation,
),
None => crate::checkpoint::CheckpointWriter::new(
checkpoint_path.clone(),
session.id.clone(),
started_at,
invocation,
),
}
});
Ok(SessionHost {
session,
sqlite_reporter,
cadence_reporter,
cadence_tree,
metrics_query,
stop_handle,
refine_plan,
resume_target,
refine_requested,
refine_scope: refine_scope.map(|s| s.to_string()),
stick_reattached,
profiler: _profiler,
sqlite_guard: _sqlite_shutdown_guard,
checkpoint_writer,
saved_doc,
})
}
async fn shutdown(self) {
if let Some(mut profiler) = self.profiler {
profiler.finish();
}
self.stop_handle.stop().await;
let _teardown_t = std::time::Instant::now();
self.cadence_reporter.shutdown().await;
crate::diag!(
crate::observer::LogLevel::Debug,
"shutdown: cadence reporter flush+join {:?}",
_teardown_t.elapsed()
);
nmbrs_metrics::queryapi::uninstall_live_access();
crate::diag!(
crate::observer::LogLevel::Info,
"shutting down — consolidating metrics.db WAL"
);
self.sqlite_guard.consume();
crate::diag!(crate::observer::LogLevel::Info, "shutdown complete");
crate::session_signals::mark_shutdown_complete();
}
}
async fn run_execution(
host: &SessionHost,
args: &[String],
observer: Arc<dyn crate::observer::RunObserver>,
) -> Result<(), String> {
let session = &host.session;
let session_id = host.session.id.clone();
let sqlite_reporter = host.sqlite_reporter.clone();
let cadence_reporter = host.cadence_reporter.clone();
let stop_handle = host.stop_handle.clone();
let resume_target = host.resume_target.clone();
let refine_plan = host.refine_plan.clone();
let refine_requested = host.refine_requested;
let refine_scope = host.refine_scope.as_deref();
let stick_reattached = host.stick_reattached.clone();
let mut diag = DiagnosticConfig::normal();
let args = normalize_args(args);
let mut params = parse_params(&args);
if let Some(driver_name) = params.get("driver").cloned()
&& let Some((manifest, library_ref)) = resolve_driver_manifest(&driver_name)?
{
let workload_given = params.contains_key("workload")
|| args
.iter()
.any(|a| a.ends_with(".yaml") || a.ends_with(".yml"));
apply_driver_manifest(
&mut params,
&driver_name,
manifest,
library_ref,
workload_given,
)?;
}
let params = params;
let mut workload_file: Option<String> = None;
let mut workload_is_bundled = false;
let mut workload_source_text: Option<String> = None;
let mut workload = if let Some(op_str) = params.get("op") {
if params.contains_key("workload") {
crate::diag!(
crate::observer::LogLevel::Warn,
"warning: op= overrides workload="
);
}
nmbrs_workload::inline::synthesize_inline_workload(op_str)
.map_err(|e| format!("inline workload: {e}"))?
} else {
let workload_raw = params
.get("workload")
.cloned()
.or_else(|| {
args.iter()
.find(|a| a.ends_with(".yaml") || a.ends_with(".yml"))
.cloned()
})
.ok_or("no workload specified. Use workload=file.yaml or op=\"...\"")?;
match resolve_workload(&workload_raw)? {
ResolvedWorkload::Path(workload_path) => {
workload_file = Some(workload_path.clone());
let yaml_source = std::fs::read_to_string(&workload_path)
.map_err(|e| format!("read workload '{workload_path}': {e}"))?;
let workload = nmbrs_workload::parse::parse_workload_from_path(
std::path::Path::new(&workload_path),
¶ms,
)
.map_err(|e| format!("parse workload: {e}"))?;
workload_source_text = Some(yaml_source);
workload
}
ResolvedWorkload::Bundled(bundled) => {
workload_file = Some(bundled.name.to_string());
workload_is_bundled = true;
let (merged, res_warnings) =
nmbrs_workload::extends::load_and_merge_bundled(bundled)
.map_err(|e| format!("bundled workload `{}`: {e}", bundled.name))?;
let mut workload = nmbrs_workload::parse::parse_workload(&merged, ¶ms)
.map_err(|e| format!("parse bundled workload `{}`: {e}", bundled.name))?;
workload.resolution_warnings.extend(res_warnings);
workload_source_text = Some(bundled.source.to_string());
workload
}
}
};
let impl_param = params.get("impl").cloned();
if let Some(target) = workload.implements.clone() {
if impl_param.is_some() {
return Err(format!(
"workload '{}' is an implementation (declares \
`implements:`), so `impl=` does not apply — invoke \
either the implementation directly or the blueprint \
with impl=",
workload_file.as_deref().unwrap_or("<inline>")
));
}
let base_dir = (!workload_is_bundled)
.then(|| {
workload_file
.as_deref()
.map(std::path::Path::new)
.and_then(|p| p.parent().map(|d| d.to_path_buf()))
})
.flatten();
let bundled_origin = workload_is_bundled
.then(|| workload_file.as_deref())
.flatten();
let (mut blueprint, blueprint_id) =
load_secondary_workload(&target, ¶ms, base_dir.as_deref(), bundled_origin)?;
crate::diag!(
crate::observer::LogLevel::Info,
"implements: binding '{}' into blueprint '{blueprint_id}'",
workload_file.as_deref().unwrap_or("<inline>")
);
nmbrs_workload::implements::bind_implementation(&mut blueprint, workload)
.map_err(|e| format!("implements binding: {e}"))?;
workload = blueprint;
} else if let Some(impl_ref) = impl_param {
let base_dir = (!workload_is_bundled)
.then(|| {
workload_file
.as_deref()
.map(std::path::Path::new)
.and_then(|p| p.parent().map(|d| d.to_path_buf()))
})
.flatten();
let bundled_origin = workload_is_bundled
.then(|| workload_file.as_deref())
.flatten();
let (implementation, impl_id) =
load_secondary_workload(&impl_ref, ¶ms, base_dir.as_deref(), bundled_origin)?;
let declared = implementation.implements.clone().ok_or_else(|| {
format!(
"impl='{impl_ref}' resolves to '{impl_id}', which declares no \
`implements:` — an implementation module must name its \
blueprint"
)
})?;
let impl_dir = std::path::Path::new(&impl_id)
.parent()
.map(|d| d.to_path_buf());
let declared_id =
workload_ref_identity(&declared, impl_dir.as_deref(), Some(impl_id.as_str()))?;
let invoked_id = workload_file
.clone()
.map(|f| canonical_identity(&f))
.unwrap_or_default();
if declared_id != invoked_id {
return Err(format!(
"impl='{impl_id}' declares implements='{declared}' \
(resolves to '{declared_id}'), which is not the invoked \
workload '{invoked_id}'"
));
}
crate::diag!(
crate::observer::LogLevel::Info,
"implements: binding '{impl_id}' into blueprint '{invoked_id}'"
);
nmbrs_workload::implements::bind_implementation(&mut workload, implementation)
.map_err(|e| format!("implements binding: {e}"))?;
}
{
let unbound = nmbrs_workload::implements::unbound_abstract_slots(&workload);
if !unbound.is_empty() {
return Err(format!(
"abstract op slot(s) [{}] unbound — pass impl=<workload> \
or invoke an implementing workload (SRD-108)",
unbound.join(", ")
));
}
}
workload.params = overlay_cli_params(std::mem::take(&mut workload.params), ¶ms);
workload.synthesize_default_phase();
let merged_params = workload.params.clone();
let driver = merged_params
.get("adapter")
.or_else(|| merged_params.get("driver"))
.cloned()
.unwrap_or_else(|| "stdout".into());
let explicit_cycles: Option<u64> = merged_params.get("cycles").and_then(|s| parse_count(s));
let concurrency: usize = match merged_params.get("concurrency") {
Some(s) => s
.parse()
.map_err(|_| format!("concurrency value '{s}' is not a valid integer"))?,
None => 1,
};
let rate: Option<f64> = match merged_params.get("rate") {
Some(s) => Some(
s.parse()
.map_err(|_| format!("rate value '{s}' is not a valid number"))?,
),
None => None,
};
let tries: Option<u32> = match merged_params.get("tries") {
Some(s) => Some(
s.parse()
.map_err(|_| format!("tries value '{s}' is not a valid integer"))?,
),
None => None,
};
let tag_filter = merged_params.get("tags").cloned();
let seq_type = merged_params
.get("seq")
.map(|s| SequencerType::parse(s).unwrap_or(SequencerType::Bucket))
.unwrap_or(SequencerType::Bucket);
let mut error_spec = merged_params
.get("errors")
.cloned()
.unwrap_or_else(|| ".*:warn,stop".to_string());
let error_rate_max: Option<f64> = match merged_params.get("error_rate_max") {
Some(s) => match s.trim().parse::<f64>() {
Ok(v) if v >= 0.0 => Some(v),
_ => {
eprintln!(
"error: error_rate_max must be a non-negative number \
(e.g. 0.1 = 10%); got '{s}'"
);
std::process::exit(2);
}
},
None => None,
};
let force_retry_failed = params
.get("force_retry_failed")
.map(|s| s != "false" && s != "0")
.unwrap_or(false)
|| args.iter().any(|a| a == "--force-retry-failed");
if force_retry_failed {
error_spec = format!(".*:retry,warn;{error_spec}");
crate::diag!(
crate::observer::LogLevel::Info,
"--force-retry-failed: errors cascade prefixed with '.*:retry,warn'"
);
}
if let Some(cli_params) = known_params() {
let adapter_params = registered_adapter_params();
let all_known: Vec<&str> = cli_params
.iter()
.copied()
.chain(adapter_params.iter().copied())
.chain(workload.declared_params.iter().map(|s| s.as_str()))
.collect();
for key in params.keys() {
if !all_known.contains(&key.as_str()) {
let suggestion = closest_match(key, &all_known);
if let Some(closest) = suggestion {
return Err(format!(
"unrecognized parameter '{key}='. Did you mean '{closest}='?"
));
} else {
return Err(format!("unrecognized parameter '{key}='"));
}
}
}
}
{
let mut invalid_polydat_braces: Vec<PolydatBraceFinding> =
collect_polydat_brace_refs(&workload);
if !invalid_polydat_braces.is_empty() {
invalid_polydat_braces.sort();
invalid_polydat_braces.dedup();
let file_path = workload_file.as_deref().unwrap_or("<inline>");
let lines: Vec<String> = invalid_polydat_braces
.iter()
.map(|f| {
let yaml_line = workload_source_text
.as_deref()
.and_then(|src| find_yaml_line_for_brace(src, &f.placeholder));
let prefix = match yaml_line {
Some(n) => format!("{file_path}:{n}"),
None => file_path.to_string(),
};
format!(
" {prefix}: in {} — `{{{}}}`. Use bare `{}`.",
f.location, f.placeholder, f.placeholder
)
})
.collect();
return Err(format!(
"`{{...}}` braces in Polydat expression context (invalid syntax).\n\
Polydat accepts bare identifiers; braces are only for YAML string\n\
interpolation (op `prepared:`/`raw:`, `cycles:`, etc.).\n{}",
lines.join("\n"),
));
}
let referenced = collect_param_references(&workload);
let adapter_params: std::collections::HashSet<&'static str> =
registered_adapter_params().into_iter().collect();
let iter_var_names = collect_iter_var_names(&workload);
let wire_names = collect_polydat_binding_names(&workload);
for name in &workload.declared_params {
if is_cli_param(name) || adapter_params.contains(name.as_str()) {
continue;
}
if !referenced.contains(name) {
return Err(format!(
"workload declares param '{name}' but it is never referenced as '{{{}}}' \
in any op, phase, or binding. Remove it or use it.",
name
));
}
}
let declared_set: std::collections::HashSet<&str> = workload
.declared_params
.iter()
.map(|s| s.as_str())
.collect();
let mut undeclared: Vec<&str> = referenced
.placeholders
.iter()
.map(|s| s.as_str())
.filter(|name| !declared_set.contains(*name))
.filter(|name| !is_cli_param(name))
.filter(|name| !adapter_params.contains(name))
.filter(|name| !iter_var_names.contains(*name))
.filter(|name| !wire_names.contains(*name))
.collect();
if !undeclared.is_empty() {
undeclared.sort();
return Err(format!(
"workload references undeclared placeholder{plural} {names} — \
add to the `params:` block, or check for a typo. Recognised \
sources for `{{name}}` placeholders: workload `params:`, \
runner/adapter built-ins, scenario-tree iter-vars from \
`for_each:`/`for_combinations:`, and wire names from Polydat \
`bindings:`.",
plural = if undeclared.len() == 1 { "" } else { "s" },
names = undeclared
.iter()
.map(|n| format!("`{{{n}}}`"))
.collect::<Vec<_>>()
.join(", "),
));
}
}
let declared: std::collections::HashSet<&String> = workload.declared_params.iter().collect();
let workload_params: HashMap<String, String> = workload
.params
.iter()
.filter(|(k, _)| declared.contains(*k))
.map(|(k, v)| (k.clone(), v.clone()))
.collect();
drop(declared);
let mut phases = workload.phases;
for phase in phases.values_mut() {
crate::scope::rewrite_inline_exprs(&mut phase.ops);
}
let phase_order = workload.phase_order;
let scenarios = workload.scenarios;
let workload_readouts = workload.readouts.clone();
let workload_wrappers_override: Option<Vec<String>> = workload
.wrappers
.as_ref()
.filter(|c| !c.order.is_empty())
.map(|c| c.order.clone());
let cli_readout_override = crate::session::resolve_flag(&args[..], "--readout");
let cli_wrap_order: Option<Vec<String>> =
crate::session::resolve_flag(&args[..], "--wrap-order")
.map(|s| {
s.split(',')
.map(|t| t.trim().to_string())
.filter(|t| !t.is_empty())
.collect()
})
.filter(|v: &Vec<String>| !v.is_empty());
let cli_wrap_default_order: Option<Vec<String>> =
crate::session::resolve_flag(&args[..], "--wrap-default-order")
.map(|s| {
s.split(',')
.map(|t| t.trim().to_string())
.filter(|t| !t.is_empty())
.collect()
})
.filter(|v: &Vec<String>| !v.is_empty());
let workload_wrappers_override: Option<Vec<String>> =
workload_wrappers_override.or(cli_wrap_order);
let workload_report = workload.report.clone();
let workload_summaries: HashMap<String, nmbrs_workload::model::SummaryConfig> = workload_report
.items()
.filter(|i| matches!(i.kind, nmbrs_workload::report::Kind::Table))
.map(|i| {
(
i.name.clone(),
nmbrs_workload::model::SummaryConfig::parse(&i.body),
)
})
.collect();
let mut ops = workload.ops;
if let Some(ref filter) = tag_filter {
ops =
TagFilter::filter_ops(&ops, filter).map_err(|e| format!("invalid tag filter: {e}"))?;
}
let mut phase_ops_for_compile: Vec<nmbrs_workload::model::ParsedOp> = Vec::new();
let mut phase_raw_ops: HashMap<String, Vec<nmbrs_workload::model::ParsedOp>> = HashMap::new();
let mut phases_needing_own_kernel: std::collections::HashSet<String> =
std::collections::HashSet::new();
for (name, phase) in &phases {
let has_own_bindings = phase.ops.iter().any(|op| !op.bindings.is_empty());
if phase.for_each.is_some() || has_own_bindings {
phase_raw_ops.insert(name.clone(), phase.ops.clone());
phases_needing_own_kernel.insert(name.clone());
} else {
phase_ops_for_compile.extend(phase.ops.iter().cloned());
}
}
if ops.is_empty() && phases.is_empty() {
return Err("no ops selected (tag filter may have excluded all ops)".into());
}
if phases.is_empty() {
crate::diag!(
crate::observer::LogLevel::Info,
"{} ops, {} cycles, concurrency={}, adapter={}",
ops.len(),
explicit_cycles
.map(|c| c.to_string())
.unwrap_or("auto".into()),
concurrency,
driver
);
} else {
crate::diag!(
crate::observer::LogLevel::Info,
"{} phases, {} top-level ops, adapter={}",
phases.len(),
ops.len(),
driver
);
}
let polydat_lib_paths: Vec<std::path::PathBuf> = args
.iter()
.filter_map(|a| a.strip_prefix("--polydat-lib="))
.map(std::path::PathBuf::from)
.collect();
let strict = args.iter().any(|a| a == "--strict")
|| matches!(
params.get("strict").map(String::as_str),
Some("true") | Some("1")
);
let kernel_opt: polydat::kernel::KernelOptLevel = {
let raw =
cli_flag_value(&args[..], "--kernel-opt").or_else(|| params.get("kernel_opt").cloned());
match raw {
None => polydat::kernel::KernelOptLevel::default(),
Some(s) => polydat::kernel::KernelOptLevel::parse(s.trim()).map_err(|bad| {
format!("unknown --kernel-opt value '{bad}' — use 'release' or 'diagnostic'")
})?,
}
};
if let Some(spec) = params.get("dryrun") {
diag = DiagnosticConfig::parse(spec);
}
if let Some(spec) = params.get("skipped_phases") {
match crate::observer::SkippedPhaseDisplay::parse(spec) {
Some(mode) => crate::observer::set_skipped_phase_display(mode),
None => {
return Err(format!(
"unknown skipped_phases value '{spec}' — use 'elide', 'mark', or 'prune'"
));
}
}
}
if let Some(spec) = params.get("completed_phases") {
match crate::observer::CompletedPhaseDisplay::parse(spec) {
Some(mode) => crate::observer::set_completed_phase_display(mode),
None => {
return Err(format!(
"unknown completed_phases value '{spec}' — use 'full' or 'headers'"
));
}
}
}
if !workload.scenario_parse_errors.is_empty() {
return Err(format!(
"scenario parse error{plural} — workload is malformed:\n - {msgs}",
plural = if workload.scenario_parse_errors.len() == 1 {
""
} else {
"s"
},
msgs = workload.scenario_parse_errors.join("\n - "),
));
}
if !workload.report_warnings.is_empty() {
if strict {
return Err(format!(
"report-block warnings (strict mode promotes to errors):\n - {}",
workload.report_warnings.join("\n - "),
));
}
for w in &workload.report_warnings {
crate::observer::log_tagged(
crate::observer::LogLevel::Warn,
crate::observer::EventTag::in_flight(crate::observer::EventCategory::Report),
&format!("report: {w}"),
);
}
}
if !workload.resolution_warnings.is_empty() {
if strict {
return Err(format!(
"reference-resolution warnings (strict mode promotes to errors):\n - {}",
workload.resolution_warnings.join("\n - "),
));
}
for w in &workload.resolution_warnings {
crate::observer::log_tagged(
crate::observer::LogLevel::Warn,
crate::observer::EventTag::in_flight(crate::observer::EventCategory::Resolution),
&format!("resolve: {w}"),
);
}
}
let dry_run: Option<&str> = if diag.depth == ExecDepth::Cycle {
Some("cycle")
} else {
params.get("dryrun").and_then(|s| match s.as_str() {
"fields" => Some("fields"),
"silent" => Some("silent"),
"op" => Some("op"),
_ => None,
})
};
if dry_run.is_some() && diag.depth < ExecDepth::Cycle {
diag.depth = ExecDepth::Cycle;
}
let scenario_for_session = params
.get("scenario")
.map(|s| s.as_str())
.unwrap_or("default");
let on_removed_policy: &str = params
.get("on_removed")
.map(|s| s.as_str())
.unwrap_or("error");
if let Some(plan) = refine_plan.as_ref() {
let current_names: std::collections::HashSet<&str> =
phases.keys().map(|s| s.as_str()).collect();
let mut removed: Vec<&str> = plan
.seen_identities
.iter()
.map(|(name, _)| name.as_str())
.filter(|n| !current_names.contains(n))
.collect();
removed.sort();
removed.dedup();
if !removed.is_empty() {
match on_removed_policy {
"error" => {
return Err(format!(
"refine: workload removes {n} phase{plural} that have \
prior outcomes in this session:\n - {names}\n\
Pass `on_removed=keep` to retain the prior outcomes \
(no work, no error); `on_removed=drop` is reserved \
(not yet implemented). Default `error` refuses to \
proceed so accidental axis-trim doesn't drop history \
silently.",
n = removed.len(),
plural = if removed.len() == 1 { "" } else { "s" },
names = removed.join("\n - "),
));
}
"keep" => {
crate::diag!(
crate::observer::LogLevel::Info,
"refine: on_removed=keep — retaining prior outcomes \
for {n} removed phase(s): {names}",
n = removed.len(),
names = removed.join(", ")
);
}
"drop" => {
crate::diag!(
crate::observer::LogLevel::Warn,
"refine: on_removed=drop is reserved — not yet \
implemented. Treating as `keep` (retaining prior \
outcomes) for now: {names}",
names = removed.join(", ")
);
}
other => {
return Err(format!(
"refine: unknown on_removed= policy '{other}'; \
expected `error` (default), `keep`, or `drop`"
));
}
}
}
}
let (exec_verb, exec_id_seed): (&'static str, u64) =
match crate::execution_context::try_current() {
Some(ctx) => ("run", ctx.exec_id),
None => match (refine_plan.as_ref(), resume_target.as_ref()) {
(Some(plan), Some(p)) if p.exists() => ("refine", plan.next_exec_id),
(_, Some(p)) if p.exists() => ("resume", 1),
_ => ("run", 1),
},
};
let execution = crate::session::Execution::start(
&session,
workload_file.as_deref().unwrap_or("inline"),
scenario_for_session,
exec_verb,
exec_id_seed,
);
let exec_id = execution.exec_id;
{
let mut cli_keys: Vec<&String> = params.keys().collect();
cli_keys.sort();
let cli_text: String = cli_keys
.iter()
.filter_map(|k| params.get(*k).map(|v| format!("{k}={v}")))
.collect::<Vec<_>>()
.join("\n");
let scope_for_row: Option<&str> = if refine_requested {
Some(refine_scope.unwrap_or("missing"))
} else {
None
};
if let Ok(mut guard) = sqlite_reporter.lock()
&& let Some(r) = guard.as_mut()
{
let exec_id = execution.exec_id;
let sid = session.id.clone();
r.set_execution_metadata(&sid, exec_id, "workload", &execution.workload);
r.set_execution_metadata(&sid, exec_id, "scenario", &execution.scenario);
r.set_execution_metadata(
&sid,
exec_id,
"start_time",
&format!(
"{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_secs()
),
);
for (k, v) in &merged_params {
r.set_execution_metadata(&sid, exec_id, &format!("param.{k}"), v);
}
if let Some(yaml) = workload_source_text.as_deref() {
r.set_execution_metadata(&sid, exec_id, "workload_yaml", yaml);
}
r.set_execution_metadata(&sid, exec_id, "cli_params", &cli_text);
r.insert_execution_start(
&session.id,
execution.exec_id,
execution.verb,
scope_for_row,
execution.started_at_nanos,
workload_source_text.as_deref().unwrap_or(""),
&cli_text,
);
}
}
{
let session_ctx = crate::readout_context::LifecycleContext {
event: crate::lifecycle::EventType::SessionStart,
subject_name: session.id.clone(),
subject_labels: String::new(),
depth_indent: String::new(),
use_color: crate::observer::use_color(),
stick_reattached: stick_reattached.clone().unwrap_or_default(),
};
let default_body = if stick_reattached.is_some() {
crate::readouts::parse::bake("session_notice")
.map(|(body, _)| body)
.ok()
} else {
None
};
crate::readout_context::fire_lifecycle(
crate::lifecycle::EventType::SessionStart,
&workload_readouts,
default_body,
&session_ctx,
Some(&sqlite_reporter),
);
}
let _num_top_level_ops = ops.len();
let mut all_ops_for_compile: Vec<nmbrs_workload::model::ParsedOp> = ops;
all_ops_for_compile.extend(phase_ops_for_compile);
crate::scope::rewrite_inline_exprs(&mut all_ops_for_compile);
let workload_level_polydat: Option<String> = match &workload.bindings {
nmbrs_workload::model::BindingsDef::PolydatSource(s) if !s.trim().is_empty() => {
Some(s.clone())
}
nmbrs_workload::model::BindingsDef::Map(m) if !m.is_empty() => Some(
crate::bindings::legacy_chain_map_to_polydat_lines(m)
.map_err(|e| format!("workload-level bindings: {e}"))?,
),
_ => None,
};
let workload_dir: Option<&std::path::Path> = if workload_is_bundled {
Some(std::path::Path::new("."))
} else {
workload_file
.as_ref()
.and_then(|p| std::path::Path::new(p).parent())
.or_else(|| Some(std::path::Path::new(".")))
};
let mut config_refs: Vec<String> = params
.values()
.filter(|v| v.starts_with('{') && v.ends_with('}'))
.map(|v| {
let mut inner = v[1..v.len() - 1].to_string();
for (key, value) in &workload_params {
let placeholder = format!("{{{key}}}");
if inner.contains(&placeholder) {
inner = inner.replace(&placeholder, value);
}
}
inner
})
.collect();
for (name, phase) in &phases {
if phase.for_each.is_some() {
continue; }
if let Some(ref c) = phase.cycles
&& c.starts_with('{')
&& c.ends_with('}')
{
let mut inner = c[1..c.len() - 1].to_string();
for (key, value) in &workload_params {
let placeholder = format!("{{{key}}}");
if inner.contains(&placeholder) {
inner = inner.replace(&placeholder, value);
}
}
config_refs.push(inner);
}
let _ = name; }
let cursor_limit: Option<u64> = merged_params.get("limit").and_then(|s| s.parse().ok());
let params_kernel = crate::params::build_workload_params_kernel(&workload_params)?;
let workload_canonical_kernel: std::sync::Arc<crate::scope_kernel::ScopeKernel> =
std::sync::Arc::new(
build_workload_root_kernel(
¶ms_kernel,
&all_ops_for_compile,
workload_dir,
polydat_lib_paths.clone(),
strict,
&config_refs,
"outer workload bindings",
cursor_limit,
&workload_params,
workload_level_polydat.as_deref(),
)
.map_err(|e| format!("outer workload bindings: {e}"))?,
);
let kernel = workload_canonical_kernel.clone();
fn collect_grouped_phases(
nodes: &[nmbrs_workload::model::ScenarioNode],
in_group: bool,
out: &mut std::collections::HashSet<String>,
) {
for node in nodes {
match node {
nmbrs_workload::model::ScenarioNode::Phase(name) => {
if in_group {
out.insert(name.clone());
}
}
nmbrs_workload::model::ScenarioNode::Comprehension { children, .. }
| nmbrs_workload::model::ScenarioNode::DoWhile { children, .. }
| nmbrs_workload::model::ScenarioNode::DoUntil { children, .. } => {
collect_grouped_phases(children, true, out);
}
nmbrs_workload::model::ScenarioNode::IncludedScenario { children, .. } => {
collect_grouped_phases(children, in_group, out);
}
nmbrs_workload::model::ScenarioNode::Bindings { children, .. } => {
collect_grouped_phases(children, in_group, out);
}
}
}
}
let mut grouped_phases = std::collections::HashSet::new();
for nodes in scenarios.values() {
collect_grouped_phases(nodes, false, &mut grouped_phases);
}
let mut resolved_phase_cycles: HashMap<String, Option<u64>> = HashMap::new();
for (name, phase) in &phases {
if phase.for_each.is_some() || grouped_phases.contains(name) {
continue;
}
let resolved = phase.cycles.as_ref().and_then(|s| {
let expanded = expand_workload_params(s, &workload_params);
resolve_polydat_config(&expanded, &kernel)
});
resolved_phase_cycles.insert(name.clone(), resolved);
}
for op in &mut all_ops_for_compile {
op.params.remove("adapter");
op.params.remove("driver");
}
for ops in phase_raw_ops.values_mut() {
for op in ops.iter_mut() {
op.params.remove("adapter");
op.params.remove("driver");
}
}
let builder = Arc::new(OpBuilder::new(kernel));
{
let scenario_name = params
.get("scenario")
.map(|s| s.as_str())
.unwrap_or("default");
let scenario_nodes = resolve_scenario(&scenarios, &phase_order, scenario_name)?;
let scope_tree = {
let mut t = crate::scope_tree::ScopeTree::build(scenario_name, &scenario_nodes);
t.populate_pragmas(&phases);
let wp_names: std::collections::HashSet<String> =
workload_params.keys().cloned().collect();
t.validate_iter_var_uniqueness(&wp_names)?;
for w in crate::workload_lint::lint_workload(
&workload.stop_when,
phases.iter().map(|(k, v)| (k.as_str(), v)),
)? {
crate::diag!(crate::observer::LogLevel::Warn, "{w}");
}
t.extend_with_op_templates(&phases);
let classify_inputs = crate::scope_elision::ClassifyInputs {
bindings: &workload.bindings,
params: &workload.params,
phases: &phases,
};
crate::scope_elision::classify_and_mark(&mut t, &classify_inputs);
std::sync::Arc::new(t)
};
if params.get("dryrun").map(|s| s.as_str()) == Some("kernels") {
print_kernel_dump_header();
crate::scope_tree::set_kernel_install_visitor(Some(Box::new(|node, _idx, kernel| {
print_kernel_for_scope(node, kernel);
})));
}
scope_tree.install_kernel(scope_tree.root, std::sync::Arc::new(params_kernel));
scope_tree.install_kernel(scope_tree.workload_root_idx(), workload_canonical_kernel);
let workload_dir_owned: Option<std::path::PathBuf> = workload_dir.map(|p| p.to_path_buf());
#[allow(clippy::large_enum_variant)]
enum InstallSpec {
ForComprehension {
idx: crate::scope_tree::ScopeNodeIdx,
iter_vars: Vec<String>,
spec_exprs: Vec<String>,
phase_bindings: nmbrs_workload::model::BindingsDef,
},
DoLoop {
idx: crate::scope_tree::ScopeNodeIdx,
counter: Option<String>,
condition: String,
},
OpTemplate {
idx: crate::scope_tree::ScopeNodeIdx,
op: nmbrs_workload::model::ParsedOp,
},
Bindings {
idx: crate::scope_tree::ScopeNodeIdx,
bindings: nmbrs_workload::model::BindingsDef,
},
}
let install_specs: Vec<InstallSpec> = scope_tree
.iter_dfs()
.filter_map(|(idx, node)| match &node.kind {
crate::scope_tree::ScopeKind::Comprehension { comprehension } => {
let pairs = comprehension.coordinate_specs();
let vars: Vec<String> = pairs.iter().map(|(v, _)| v.clone()).collect();
let specs: Vec<String> = pairs.iter().map(|(_, e)| e.clone()).collect();
Some(InstallSpec::ForComprehension {
idx,
iter_vars: vars,
spec_exprs: specs,
phase_bindings: nmbrs_workload::model::BindingsDef::default(),
})
}
crate::scope_tree::ScopeKind::DoWhile { condition, counter } => {
Some(InstallSpec::DoLoop {
idx,
counter: counter.clone(),
condition: condition.clone(),
})
}
crate::scope_tree::ScopeKind::DoUntil { condition, counter } => {
Some(InstallSpec::DoLoop {
idx,
counter: counter.clone(),
condition: condition.clone(),
})
}
crate::scope_tree::ScopeKind::Phase { name } => {
let phase = phases.get(name.as_str())?;
if let Some(spec) = phase.for_each.as_ref() {
if !phase.metrics.is_empty() {
crate::diag!(
crate::observer::LogLevel::Error,
"phase '{name}': phase-level `metrics:` + phase-level \
`for_each:` is not supported yet. Move the for_each to \
scenario-tree level (so each iteration is its own phase \
activation, each with its own metrics), or drop one. \
Phase will be skipped.",
);
return None;
}
if phase.poll.is_some() {
crate::diag!(
crate::observer::LogLevel::Error,
"phase '{name}': `poll:` + phase-level `for_each:` is not \
supported in the initial ship of SRD-75. Move the for_each \
to scenario-tree level (so each iter is its own phase \
activation), or drop one of the two. Phase will be skipped.",
);
return None;
}
let comp = match polydat::iteration::comprehension::spec::parse_inline(spec)
{
Ok(c) => c,
Err(e) => {
crate::diag!(
crate::observer::LogLevel::Error,
"phase '{name}' for_each '{spec}': {e}"
);
return None;
}
};
let pairs = comp.coordinate_specs();
let iter_vars: Vec<String> = pairs.iter().map(|(v, _)| v.clone()).collect();
let spec_exprs: Vec<String> =
pairs.iter().map(|(_, e)| e.clone()).collect();
let phase_bindings =
match crate::scope::synthesize_phase_scope_bindings(phase) {
Ok(b) => b,
Err(e) => {
crate::diag!(
crate::observer::LogLevel::Error,
"phase '{name}': phase-scope synthesis: {e}"
);
return None;
}
};
Some(InstallSpec::ForComprehension {
idx,
iter_vars,
spec_exprs,
phase_bindings,
})
} else {
let synth = match crate::scope::synthesize_phase_scope_bindings(phase) {
Ok(b) => b,
Err(e) => {
crate::diag!(
crate::observer::LogLevel::Error,
"phase '{name}': SRD-75 phase-poll synthesis: {e}",
);
return None;
}
};
if !synth.is_empty() {
Some(InstallSpec::Bindings {
idx,
bindings: synth,
})
} else {
None
}
}
}
crate::scope_tree::ScopeKind::OpTemplate { name } => {
if node.materialised != Some(true) {
return None;
}
let owning_phase: Option<&str> = {
let mut cursor = scope_tree.nodes[idx].parent;
let mut found: Option<&str> = None;
while let Some(p) = cursor {
if let crate::scope_tree::ScopeKind::Phase { name: pname } =
&scope_tree.nodes[p].kind
{
found = Some(pname.as_str());
break;
}
cursor = scope_tree.nodes[p].parent;
}
found
};
owning_phase
.and_then(|pname| phases.get(pname))
.and_then(|phase| phase.ops.iter().find(|op| op.name == *name))
.cloned()
.map(|op| InstallSpec::OpTemplate { idx, op })
}
crate::scope_tree::ScopeKind::Bindings { source } => {
Some(InstallSpec::Bindings {
idx,
bindings: nmbrs_workload::model::BindingsDef::PolydatSource(source.clone()),
})
}
_ => None,
})
.collect();
for install_spec in install_specs {
let idx = match &install_spec {
InstallSpec::ForComprehension { idx, .. } => *idx,
InstallSpec::DoLoop { idx, .. } => *idx,
InstallSpec::OpTemplate { idx, .. } => *idx,
InstallSpec::Bindings { idx, .. } => *idx,
};
let parent_kernel = {
let mut cursor = scope_tree.nodes[idx].parent;
let mut found: Option<std::sync::Arc<crate::scope_kernel::ScopeKernel>> = None;
while let Some(p) = cursor {
if let Some(k) = scope_tree.nodes[p].cached_kernel.get() {
found = Some(k.clone());
break;
}
cursor = scope_tree.nodes[p].parent;
}
found.expect("workload root always has an installed kernel")
};
let parent_manifest = extract_manifest(parent_kernel.program());
let context = format!("scope idx {idx} ({})", scope_tree.nodes[idx].kind.label(),);
let result = match install_spec {
InstallSpec::ForComprehension {
iter_vars,
spec_exprs,
phase_bindings,
..
} => {
let bindings: Vec<(String, String)> = iter_vars
.iter()
.cloned()
.zip(spec_exprs.iter().cloned())
.collect();
let phase_bindings_source = match phase_bindings {
nmbrs_workload::model::BindingsDef::PolydatSource(s)
if !s.trim().is_empty() =>
{
Some(s)
}
nmbrs_workload::model::BindingsDef::Map(m) if !m.is_empty() => {
let mut out = String::new();
for (name, expr) in &m {
out.push_str(&format!("{name} := {expr}\n"));
}
Some(out)
}
_ => None,
};
crate::scope_synth::build_for_each_scope_kernel(
&bindings,
&parent_manifest,
&parent_kernel,
&workload_params,
polydat_lib_paths.clone(),
workload_dir_owned.as_deref(),
strict,
&context,
phase_bindings_source.as_deref(),
)
}
InstallSpec::DoLoop {
counter, condition, ..
} => crate::scope::build_do_loop_scope_kernel(
counter.as_deref(),
&condition,
&parent_manifest,
&parent_kernel,
&workload_params,
polydat_lib_paths.clone(),
workload_dir_owned.as_deref(),
strict,
&context,
),
InstallSpec::OpTemplate { op, .. } => {
crate::scope::build_op_template_scope_kernel(
&op,
&parent_manifest,
&parent_kernel,
&workload_params,
polydat_lib_paths.clone(),
workload_dir_owned.as_deref(),
strict,
kernel_opt,
&context,
)
}
InstallSpec::Bindings { bindings, .. } => {
crate::scope::build_phase_scope_kernel(
&bindings,
&parent_manifest,
&parent_kernel,
&workload_params,
polydat_lib_paths.clone(),
workload_dir_owned.as_deref(),
strict,
&context,
)
}
};
match result {
Ok(kernel) => {
let _ = scope_tree.install_kernel(idx, std::sync::Arc::new(kernel));
}
Err(e) => {
return Err(format!("scope kernel synthesis failed: {e}"));
}
}
}
crate::diag!(
crate::observer::LogLevel::Info,
"scenario '{scenario_name}':\n{}",
format_scenario_tree(&scenario_nodes, &phases)
);
let schedule_spec = std::sync::Arc::new(match params.get("schedule") {
Some(s) => crate::scheduler::ScheduleSpec::parse(s)
.map_err(|e| format!("schedule= param: {e}"))?,
None => crate::scheduler::ScheduleSpec::default_serial(),
});
let dry_run_static: Option<&'static str> = match dry_run {
Some("silent") => Some("silent"),
Some("fields") => Some("fields"),
Some("cycle") => Some("cycle"),
Some("op") => Some("op"),
_ => None,
};
let phase_filter: Option<Arc<crate::phase_filter::PhasePattern>> = match params
.get("phases")
.map(|s| s.as_str())
.filter(|s| !s.is_empty())
{
None => None,
Some(src) => {
let pat = crate::phase_filter::PhasePattern::parse(src)
.map_err(|e| format!("phases= param: {e}"))?;
crate::diag!(
crate::observer::LogLevel::Info,
"phases=<filter>: pattern '{src}' ({}{})",
if pat.negated() { "negated " } else { "" },
pat.dialect().as_str()
);
Some(Arc::new(pat))
}
};
let resource_pool = Arc::new(crate::resource_pool::ResourcePool::new());
crate::resource_pool::install_accessor(&resource_pool);
let initial_scene_tree_path = vec![crate::checkpoint::PathSegment::Scenario(
scenario_name.to_string(),
)];
let root_error_policy = crate::error_policy::ErrorPolicy::root(
crate::error_policy::PolicyConfig::new(error_spec.clone(), error_rate_max),
);
let workload_shell = {
use nmbrs_workload::model::ScopeLevel;
let mut conditions: Vec<crate::stop_conditions::StopConditionDecl> =
vec![crate::stop_conditions::StopConditionDecl {
when: "children_failed > 0".to_string(),
effect: crate::phase_outcome::Outcome::failed(),
reason: None,
target: crate::stop_conditions::StopScope::Workload,
cancel_ops: false,
}];
conditions.extend(
workload
.stop_when
.iter()
.filter(|c| {
c.each
.iter()
.any(|l| matches!(l, ScopeLevel::SelfScope | ScopeLevel::Workload))
})
.map(|c| crate::stop_conditions::StopConditionDecl {
when: expand_workload_params(&c.when, &workload_params),
effect: crate::stop_conditions::StopConditionDecl::effect_from_str(
c.effect.as_deref(),
crate::phase_outcome::Outcome::interrupted(),
),
reason: None,
target: crate::executor::resolve_stop_scope(c.at, &c.each),
cancel_ops: crate::stop_conditions::StopConditionDecl::action_cancels_ops(
c.effect.as_deref(),
),
}),
);
let set = match scope_tree.nodes[scope_tree.workload_root_idx()]
.cached_kernel
.get()
{
Some(root_kernel) => crate::stop_conditions::StopConditionSet::build_for_phase(
root_kernel,
&conditions,
)
.unwrap_or_else(|e| {
crate::diag!(
crate::observer::LogLevel::Error,
"workload stop-condition compile failed: {e}"
);
crate::stop_conditions::StopConditionSet::empty()
}),
_ => crate::stop_conditions::StopConditionSet::empty(),
};
std::sync::Arc::new(crate::workload_shell::WorkloadShell::new(set))
};
let phase_param_overrides =
std::sync::Arc::new(crate::phase_params::parse_overrides(&args)?);
crate::phase_params::validate_against_phases(
&phase_param_overrides,
phases.keys().map(|s| s.as_str()),
)?;
let mut exec_ctx = crate::executor::ExecCtx {
phases: phases.clone(),
optimize_objective: None,
optimize_objective_value: None,
optimize_servo: None,
phase_param_overrides,
workload_shell,
workload_stop_when: workload.stop_when.clone(),
daemon_stop: None,
workload_readouts: workload_readouts.clone(),
cli_readout_override: cli_readout_override.clone(),
workload_params: workload_params.clone(),
wrappers_override: workload_wrappers_override.clone(),
wrap_default_order: cli_wrap_default_order.clone(),
workload_scope: builder.source_kernel().clone(),
polydat_lib_paths: polydat_lib_paths.clone(),
workload_dir: workload_dir.map(|p| p.to_path_buf()),
strict,
driver: driver.clone(),
merged_params: merged_params.clone(),
dry_run: dry_run_static,
phase_filter: phase_filter.clone(),
refine_plan: refine_plan.clone(),
diag: {
let mut d = diag.clone();
d.depth = ExecDepth::Phase;
d
},
pre_map_only: true,
seq_type,
concurrency,
rate,
error_spec: error_spec.clone(),
tries,
error_rate_max,
error_policy: root_error_policy,
session_id: session_id.clone(),
exec_id,
workload_name: execution.workload.clone(),
label_stack: Vec::new(),
session_component: execution.component.clone(),
cadence_reporter: cadence_reporter.clone(),
stop_handle: stop_handle.clone(),
observer: observer.clone(),
scope_tree: scope_tree.clone(),
schedule_spec: schedule_spec.clone(),
current_parent_kernel: scope_tree.nodes[scope_tree.workload_root_idx()]
.cached_kernel
.get()
.cloned(),
workload_source: workload_file.as_ref().and_then(|path| {
workload_source_text.as_ref().map(|text| {
std::sync::Arc::new(crate::executor::WorkloadSource {
path: path.clone(),
text: text.clone(),
})
})
}),
checkpoint_writer: None,
resume_plan: std::sync::Arc::new(crate::checkpoint::ResumePlan::fresh()),
sqlite_reporter: sqlite_reporter.clone(),
resource_pool: resource_pool.clone(),
scene_tree_parent_id: 0,
scene_tree_path: initial_scene_tree_path.clone(),
current_scope_idx: 0,
};
exec_ctx.current_scope_idx = exec_ctx.scope_tree.scenario_root_idx();
crate::scene_tree::install_global(crate::scene_tree::SceneTree::new());
let pre_map_result = crate::executor::execute_tree(&mut exec_ctx, &scenario_nodes).await;
let pre_mapped_tree = match pre_map_result {
Ok(()) => {
let tree = crate::scene_tree::current();
if let Some(ref t) = tree {
observer.scenario_pre_mapped(t);
}
tree
}
Err(e) if strict => return Err(e),
Err(e) => {
crate::diag!(
crate::observer::LogLevel::Warn,
"pre-map walker failed (scope hierarchy will be flat in summaries / TUI): {e}"
);
None
}
};
if params.get("dryrun").map(|s| s.as_str()) == Some("kernels") {
print_kernel_dump_legend();
crate::scope_tree::set_kernel_install_visitor(None);
return Ok(());
}
let checkpoint_writer = host.checkpoint_writer.clone();
let saved_doc = host.saved_doc.clone();
let invocation = saved_doc.as_ref().map(|d| d.invocation + 1).unwrap_or(1);
let parent_for_keep_check = if let Some(p) = session.output_dir.parent() {
p.to_path_buf()
} else {
crate::session::default_sessions_root()
};
let session_keep = crate::session::resolve_session_dir(&args).session_keep;
struct EndOfRunNoticeGuard {
writer: std::sync::Arc<crate::checkpoint::CheckpointWriter>,
parent: std::path::PathBuf,
keep_cap: usize,
}
impl Drop for EndOfRunNoticeGuard {
fn drop(&mut self) {
if let Some(hint) = self.writer.resume_hint() {
for line in hint.lines() {
crate::diag!(crate::observer::LogLevel::Info, "{line}");
}
}
let n = crate::session::forecast_keep_purge(&self.parent, self.keep_cap);
if n > 0 {
crate::diag!(
crate::observer::LogLevel::Info,
"the next new nmbrs session will auto-purge {n} prior session \
director{plural} under {} due to --session-keep={cap}. \
To disable: --session-keep=0 (or NMBRS_SESSION_KEEP=0). \
To raise the cap: --session-keep=<bigger>.",
self.parent.display(),
plural = if n == 1 { "y" } else { "ies" },
cap = self.keep_cap,
);
}
}
}
let _eor_notice_guard = EndOfRunNoticeGuard {
writer: checkpoint_writer.clone(),
parent: parent_for_keep_check,
keep_cap: session_keep,
};
let resume_plan =
if let (Some(saved), Some(tree)) = (saved_doc.as_ref(), pre_mapped_tree.as_ref()) {
let candidates =
crate::checkpoint::scene_tree_resume_candidates(tree, &scope_tree, &phases);
std::sync::Arc::new(crate::checkpoint::ResumePlan::from_checkpoint(
saved,
&candidates,
&workload_params,
))
} else {
std::sync::Arc::new(crate::checkpoint::ResumePlan::fresh())
};
if let Some(tree) = pre_mapped_tree.as_ref() {
crate::checkpoint::declare_scene_tree_phases(&checkpoint_writer, tree, &phases);
}
if resume_plan.is_resume {
let skip = resume_plan.skip_count();
let mismatch = resume_plan.mismatch_count();
let cursor = resume_plan.cursor_resume_count();
crate::diag!(
crate::observer::LogLevel::Info,
"resume: invocation #{invocation} — \
{skip} skip, {mismatch} mismatched, {cursor} cursor-resume"
);
}
if let Some(tree) = pre_mapped_tree.as_ref() {
crate::resource_pool::pre_map_pending_uses(
&resource_pool,
tree,
&phases,
&driver,
&merged_params,
)?;
}
exec_ctx.checkpoint_writer = Some(checkpoint_writer.clone());
exec_ctx.resume_plan = resume_plan.clone();
exec_ctx.diag = diag.clone();
exec_ctx.scene_tree_parent_id = 0;
exec_ctx.scene_tree_path = initial_scene_tree_path.clone();
exec_ctx.current_scope_idx = exec_ctx.scope_tree.scenario_root_idx();
exec_ctx.pre_map_only = false;
let scheduler = crate::scheduler::build(&schedule_spec);
let scheduler_result = scheduler.run(&mut exec_ctx, &scenario_nodes).await;
exec_ctx.resource_pool.shutdown().await;
if scheduler_result.is_err() {
close_execution_row(&sqlite_reporter, &session_id, exec_id);
}
scheduler_result?;
cadence_reporter.close_path(&Labels::of("session", &session.id));
if diag.depth == ExecDepth::Dispenser {
let mut out = std::io::stdout();
if let Err(e) = render_scope_elision_summary(&scope_tree, &mut out) {
crate::diag!(
crate::observer::LogLevel::Warn,
"warning: rendering scope-elision summary: {e}"
);
}
}
checkpoint_writer.mark_run_reached_end();
}
cadence_reporter.close_path(&Labels::of("session", &session.id));
cadence_reporter
.quiesce(std::time::Duration::from_secs(30))
.await;
{
let session_ctx = crate::readout_context::LifecycleContext {
event: crate::lifecycle::EventType::SessionEnd,
subject_name: session.id.clone(),
subject_labels: String::new(),
depth_indent: String::new(),
use_color: crate::observer::use_color(),
stick_reattached: String::new(),
};
crate::readout_context::fire_lifecycle(
crate::lifecycle::EventType::SessionEnd,
&workload_readouts,
None,
&session_ctx,
Some(&sqlite_reporter),
);
}
observer.run_finished();
if dry_run.is_some() {
crate::diag!(crate::observer::LogLevel::Info, "dry-run complete.");
} else {
crate::diag!(crate::observer::LogLevel::Info, "done.");
}
let active_summaries: HashMap<String, nmbrs_workload::model::SummaryConfig> =
if let Some(cli_summary) = merged_params.get("summary") {
let mut m = HashMap::new();
m.insert(
"default".into(),
nmbrs_workload::model::SummaryConfig::parse(cli_summary),
);
m
} else {
workload_summaries.clone()
};
let summary_destinations: HashMap<String, Vec<nmbrs_workload::report::Destination>> = {
use nmbrs_workload::report::Destination as D;
if merged_params.contains_key("summary") {
let mut m = HashMap::new();
m.insert("default".to_string(), vec![D::SessionDir, D::Stdout]);
m
} else {
workload_report
.items()
.filter(|i| matches!(i.kind, nmbrs_workload::report::Kind::Table))
.map(|i| {
(
i.name.clone(),
i.style
.destinations
.clone()
.unwrap_or_else(|| vec![D::SessionDir]),
)
})
.collect()
}
};
if let Ok(mut guard) = sqlite_reporter.lock()
&& let Some(ref mut reporter) = *guard
{
let end_time = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0);
let sid = session.id.clone();
let exec_id = execution.exec_id;
reporter.set_execution_metadata(&sid, exec_id, "end_time", &end_time.to_string());
reporter.set_execution_metadata(&sid, exec_id, "phase_count", &phases.len().to_string());
reporter.set_execution_metadata(
&sid,
exec_id,
"scenario_count",
&scenarios.len().to_string(),
);
if let Some(wf) = workload_file.as_deref() {
reporter.set_execution_metadata(&sid, exec_id, "workload_file", wf);
}
reporter.set_execution_metadata(&sid, exec_id, "adapter", &driver);
}
if !active_summaries.is_empty() {
if let Ok(mut guard) = sqlite_reporter.lock()
&& let Some(ref mut reporter) = *guard
{
let sid = session.id.clone();
let exec_id = execution.exec_id;
for item in workload_report.items() {
let value = item.to_yaml_directive_string();
reporter.set_execution_metadata(
&sid,
exec_id,
&format!("report.{}", item.name),
&value,
);
}
let mut names: Vec<&String> = active_summaries.keys().collect();
names.sort();
for name in names {
let cfg = &active_summaries[name];
let (basename, format) =
nmbrs_metrics::reporters::sqlite::derive_name_and_format(name);
let report_config = report_config_from_summary(cfg, Some(exec_id));
let rendered = reporter.format_summary_with_format(&report_config, &format);
if rendered.is_empty() {
continue;
}
let dests = summary_destinations
.get(name.as_str())
.cloned()
.unwrap_or_else(|| vec![nmbrs_workload::report::Destination::SessionDir]);
let to_session = dests.contains(&nmbrs_workload::report::Destination::SessionDir);
let to_stdout = dests.contains(&nmbrs_workload::report::Destination::Stdout);
let to_stderr = dests.contains(&nmbrs_workload::report::Destination::Stderr);
let filename = format!("{basename}_summary.{format}");
let summary_path = session.output_dir.join(&filename);
if to_session {
if let Err(e) = std::fs::write(&summary_path, &rendered) {
crate::diag!(
crate::observer::LogLevel::Warn,
"warning: failed to write summary to {}: {e}",
summary_path.display()
);
} else {
crate::diag!(
crate::observer::LogLevel::Info,
"summary: {}",
summary_path.display()
);
}
}
if to_stderr {
eprint!("{rendered}");
}
if to_stdout {
if observer.suppresses_stderr() {
let deferred = session.output_dir.join(DEFERRED_STDOUT_FILE);
if let Err(e) = std::fs::OpenOptions::new()
.create(true)
.append(true)
.open(&deferred)
.and_then(|mut f| {
std::io::Write::write_all(&mut f, rendered.as_bytes())
})
{
crate::diag!(
crate::observer::LogLevel::Warn,
"warning: failed to defer summary stdout to {}: {e}",
deferred.display()
);
}
} else {
print!("{rendered}");
}
}
}
}
}
refresh_latest_file_links(&session);
if diag.list_controls {
let mut out = std::io::stdout();
if let Err(e) = render_controls_tree(&session.component, &mut out) {
crate::diag!(
crate::observer::LogLevel::Warn,
"warning: rendering controls: {e}"
);
}
}
close_execution_row(&sqlite_reporter, &session_id, exec_id);
Ok(())
}
fn close_execution_row(
sqlite_reporter: &std::sync::Arc<
std::sync::Mutex<Option<nmbrs_metrics::reporters::sqlite::SqliteReporter>>,
>,
session_id: &str,
exec_id: u64,
) {
let disposition =
crate::scene_tree::with_global(|t| t.session_disposition().label()).unwrap_or("UNKNOWN");
let ended_at_nanos = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_nanos() as i64)
.unwrap_or(0);
if let Ok(mut g) = sqlite_reporter.lock()
&& let Some(r) = g.as_mut()
{
r.update_execution_end(session_id, exec_id, ended_at_nanos, disposition);
}
}
async fn run_impl(
args: &[String],
observer: Arc<dyn crate::observer::RunObserver>,
) -> Result<(), String> {
let host = SessionHost::setup(args, observer.clone())?;
let result = run_execution(&host, args, observer).await;
host.shutdown().await;
result
}
pub struct ExecutionSpec {
pub args: Vec<String>,
pub observer: Arc<dyn crate::observer::RunObserver>,
pub channel: Option<Arc<dyn crate::output_channel::OutputChannel>>,
}
pub async fn run_executions(
session_args: &[String],
session_observer: Arc<dyn crate::observer::RunObserver>,
specs: Vec<ExecutionSpec>,
max_concurrent: usize,
) -> Result<Vec<Result<(), String>>, String> {
let host = std::sync::Arc::new(SessionHost::setup(session_args, session_observer)?);
let sem = std::sync::Arc::new(tokio::sync::Semaphore::new(max_concurrent.max(1)));
let futs = specs.into_iter().map(|spec| {
let host = host.clone();
let sem = sem.clone();
async move {
let _permit = sem.acquire().await.expect("semaphore not closed");
let ctx = match spec.channel.clone() {
Some(ch) => crate::execution_context::ExecutionContext::with_observer_and_channel(
spec.observer.clone(),
ch,
),
None => {
crate::execution_context::ExecutionContext::with_observer(spec.observer.clone())
}
};
crate::execution_context::scope(ctx, run_execution(&host, &spec.args, spec.observer))
.await
}
});
let results = futures::future::join_all(futs).await;
match std::sync::Arc::try_unwrap(host) {
Ok(h) => h.shutdown().await,
Err(_) => crate::diag!(
crate::observer::LogLevel::Warn,
"run_executions: session host still referenced at teardown; \
scheduler/WAL will close on drop"
),
}
Ok(results)
}
fn refresh_latest_file_links(session: &crate::session::Session) {
let logs_dir = std::path::Path::new("logs");
if !crate::session::target_is_under(logs_dir, &session.output_dir) {
return;
}
for file in ["metrics.db", "summary.md", "session.log"] {
let target = session.output_dir.join(file);
if !target.exists() {
continue;
}
let link = logs_dir.join(file);
let _ = std::fs::remove_file(&link);
let rel_target = std::path::Path::new("latest").join(file);
if let Err(e) = crate::session::symlink_any(&rel_target, &link) {
crate::diag!(
crate::observer::LogLevel::Warn,
"warning: failed to link {} → {}: {e}",
link.display(),
rel_target.display()
);
}
}
}
pub async fn create_adapter(
driver: &str,
params: &HashMap<String, String>,
) -> Result<Arc<dyn crate::adapter::DriverAdapter>, String> {
let reg = find_adapter_registration(driver).ok_or_else(|| {
let available = registered_driver_names();
format!(
"unknown adapter '{driver}' (available: {})",
available.join(", ")
)
})?;
(reg.create)(params.clone()).await
}
fn print_kernel_dump_header() {
let is_tty = std::io::IsTerminal::is_terminal(&std::io::stdout());
let (bold, reset) = if is_tty {
("\x1b[1m", "\x1b[0m")
} else {
("", "")
};
println!();
println!("{bold}Polydat Scope Kernels{reset}");
println!("{bold}═════════════════════{reset}");
println!();
}
fn print_kernel_for_scope(
node: &crate::scope_tree::ScopeNode,
kernel: &crate::scope_kernel::ScopeKernel,
) {
let is_tty = std::io::IsTerminal::is_terminal(&std::io::stdout());
let (bold, dim, reset, cyan, magenta, green) = if is_tty {
(
"\x1b[1m", "\x1b[2m", "\x1b[0m", "\x1b[36m", "\x1b[35m", "\x1b[32m",
)
} else {
("", "", "", "", "", "")
};
let logical = if node.logical_name.is_empty() {
"<unnamed scope>".to_string()
} else {
node.logical_name.clone()
};
let depth_indent = " ".repeat(node.depth);
println!(
"{depth_indent}{green}●{reset} {bold}{cyan}{logical}{reset} \
{dim}(depth={}, kind={:?}){reset}",
node.depth, node.kind
);
let source = kernel.program().source().trim_end();
if source.is_empty() {
println!("{depth_indent} {dim}(empty kernel — no own bindings){reset}");
} else {
for line in source.lines() {
println!("{depth_indent} {magenta}│{reset} {line}");
}
}
println!();
}
fn print_kernel_dump_legend() {
let is_tty = std::io::IsTerminal::is_terminal(&std::io::stdout());
let (dim, reset, green) = if is_tty {
("\x1b[2m", "\x1b[0m", "\x1b[32m")
} else {
("", "", "")
};
println!(
" {dim}Legend: {green}●{reset}{dim} kernel installed at this scope. \
Flattened scopes (those that inherit a parent's kernel) emit no entry.{reset}"
);
println!();
}
pub async fn run_activity_simple(
activity: Activity,
adapters: std::collections::HashMap<String, Arc<dyn crate::adapter::DriverAdapter>>,
default_adapter: &str,
op_builder: Arc<crate::synthesis::OpBuilder>,
) -> bool {
activity
.run_with_adapters(adapters, default_adapter, op_builder)
.await
}
struct MutexReporter(
std::sync::Arc<std::sync::Mutex<Option<nmbrs_metrics::reporters::sqlite::SqliteReporter>>>,
);
impl Reporter for MutexReporter {
fn report(&mut self, snapshot: &nmbrs_metrics::snapshot::MetricSet) {
if let Ok(mut guard) = self.0.lock()
&& let Some(ref mut r) = *guard
{
Reporter::report(r, snapshot);
}
}
fn flush(&mut self) {
if let Ok(mut guard) = self.0.lock()
&& let Some(ref mut r) = *guard
{
Reporter::flush(r);
}
}
}
struct BoxedReporter(Box<dyn Reporter>);
impl Reporter for BoxedReporter {
fn report(&mut self, snapshot: &nmbrs_metrics::snapshot::MetricSet) {
self.0.report(snapshot);
}
fn flush(&mut self) {
self.0.flush();
}
}
pub fn expand_workload_params(s: &str, params: &HashMap<String, String>) -> String {
let mut result = s.to_string();
for (key, value) in params {
let placeholder = format!("{{{key}}}");
if result.contains(&placeholder) {
result = result.replace(&placeholder, value);
}
}
result
}
#[derive(Default)]
struct ParamRefs {
placeholders: std::collections::HashSet<String>,
runtime_only_placeholders: std::collections::HashSet<String>,
expression_idents: std::collections::HashSet<String>,
templates: Vec<String>,
}
impl ParamRefs {
fn contains(&self, param: &str) -> bool {
if self.placeholders.contains(param) {
return true;
}
if self.runtime_only_placeholders.contains(param) {
return true;
}
if self.expression_idents.contains(param) {
return true;
}
self.templates
.iter()
.any(|tpl| template_matches(tpl, param))
}
}
fn template_matches(template: &str, param: &str) -> bool {
let t = template.as_bytes();
let p = param.as_bytes();
let mut ti = 0;
let mut pi = 0;
while ti < t.len() {
if t[ti] == b'{' {
let close = match template[ti + 1..].find('}') {
Some(n) => ti + 1 + n,
None => return false, };
let next_lit = close + 1;
if next_lit >= t.len() {
if pi >= p.len() {
return false;
}
return p[pi..]
.iter()
.all(|b| b.is_ascii_alphanumeric() || *b == b'_');
}
let stop = t[next_lit];
let mut consumed = 0;
while pi + consumed < p.len() && p[pi + consumed] != stop {
let b = p[pi + consumed];
if !(b.is_ascii_alphanumeric() || b == b'_') {
return false;
}
consumed += 1;
}
if consumed == 0 {
return false;
} pi += consumed;
ti = next_lit;
} else {
if pi >= p.len() || p[pi] != t[ti] {
return false;
}
ti += 1;
pi += 1;
}
}
pi == p.len()
}
fn scan_json_for_refs(v: &serde_json::Value, refs: &mut ParamRefs) {
match v {
serde_json::Value::String(s) => scan_param_refs(s, refs),
serde_json::Value::Array(a) => {
for item in a {
scan_json_for_refs(item, refs);
}
}
serde_json::Value::Object(m) => {
for item in m.values() {
scan_json_for_refs(item, refs);
}
}
_ => {} }
}
fn scan_param_refs(text: &str, refs: &mut ParamRefs) {
let bytes = text.as_bytes();
let mut i = 0;
while i < bytes.len() {
if bytes[i] != b'{' {
i += 1;
continue;
}
let body_start = i + 1;
let mut depth = 1;
let mut j = body_start;
while j < bytes.len() && depth > 0 {
match bytes[j] {
b'{' => depth += 1,
b'}' => depth -= 1,
_ => {}
}
if depth == 0 {
break;
}
j += 1;
}
if depth != 0 {
break;
}
let body = &text[body_start..j];
if body.contains('{') {
refs.templates.push(body.to_string());
scan_param_refs(body, refs);
} else if !body.is_empty()
&& body.bytes().all(|b| b.is_ascii_alphanumeric() || b == b'_')
&& !body.bytes().next().unwrap().is_ascii_digit()
{
refs.placeholders.insert(body.to_string());
} else {
scan_expression_idents(body, &mut refs.expression_idents);
}
i = j + 1;
}
}
fn scan_expression_idents(body: &str, out: &mut std::collections::HashSet<String>) {
let bytes = body.as_bytes();
let mut i = 0;
while i < bytes.len() {
let b = bytes[i];
if b == b'/' && i + 1 < bytes.len() && bytes[i + 1] == b'/' {
i += 2;
while i < bytes.len() && bytes[i] != b'\n' {
i += 1;
}
continue;
}
if b == b'"' || b == b'\'' {
let quote = b;
i += 1;
while i < bytes.len() {
if bytes[i] == b'\\' && i + 1 < bytes.len() {
i += 2;
continue;
}
if bytes[i] == quote {
i += 1;
break;
}
i += 1;
}
continue;
}
if b.is_ascii_alphabetic() || b == b'_' {
let start = i;
while i < bytes.len() && (bytes[i].is_ascii_alphanumeric() || bytes[i] == b'_') {
i += 1;
}
let ident = &body[start..i];
if ident != "true" && ident != "false" {
out.insert(ident.to_string());
}
continue;
}
i += 1;
}
}
fn collect_param_references(workload: &nmbrs_workload::model::Workload) -> ParamRefs {
let mut refs = ParamRefs::default();
fn scan_op(op: &nmbrs_workload::model::ParsedOp, refs: &mut ParamRefs) {
for value in op.op.values() {
if let serde_json::Value::String(s) = value {
scan_param_refs(s, refs);
}
}
if let Some(s) = &op.condition {
scan_param_refs(s, refs);
scan_expression_idents(s, &mut refs.expression_idents);
}
if let Some(spec) = &op.delay {
for name in spec.names() {
scan_param_refs(name, refs);
scan_expression_idents(name, &mut refs.expression_idents);
}
}
for (k, v) in op.params.iter() {
if k == "gutter" {
continue;
}
scan_json_for_refs(v, refs);
}
if let Some(rel) = op.params.get("relevancy").and_then(|v| v.as_object()) {
for key in ["actual", "expected", "k", "r"] {
if let Some(s) = rel.get(key).and_then(|v| v.as_str()) {
scan_expression_idents(s, &mut refs.expression_idents);
}
}
}
match &op.bindings {
nmbrs_workload::model::BindingsDef::PolydatSource(s) => {
scan_param_refs(s, refs);
scan_expression_idents(s, &mut refs.expression_idents);
}
nmbrs_workload::model::BindingsDef::Map(m) => {
for v in m.values() {
scan_param_refs(v, refs);
}
}
}
}
for op in &workload.ops {
scan_op(op, &mut refs);
}
match &workload.bindings {
nmbrs_workload::model::BindingsDef::PolydatSource(s) => {
scan_param_refs(s, &mut refs);
scan_expression_idents(s, &mut refs.expression_idents);
}
nmbrs_workload::model::BindingsDef::Map(m) => {
for v in m.values() {
scan_param_refs(v, &mut refs);
}
}
}
for c in &workload.stop_when {
scan_param_refs(&c.when, &mut refs);
}
for phase in workload.phases.values() {
if let Some(s) = &phase.cycles {
scan_param_refs(s, &mut refs);
}
if let Some(s) = &phase.timeout {
scan_param_refs(s, &mut refs);
}
if let Some(s) = &phase.interval {
scan_param_refs(s, &mut refs);
}
if let Some(s) = &phase.concurrency {
scan_param_refs(s, &mut refs);
}
if let Some(s) = &phase.for_each {
scan_param_refs(s, &mut refs);
}
if let Some(s) = &phase.rate {
scan_param_refs(s, &mut refs);
}
if let Some(poll) = &phase.poll {
scan_expression_idents(&poll.until, &mut refs.expression_idents);
for s in [&poll.interval_ms, &poll.timeout_ms, &poll.max_error_retries]
.into_iter()
.flatten()
{
scan_param_refs(s, &mut refs);
}
}
match &phase.bindings {
nmbrs_workload::model::BindingsDef::PolydatSource(s) => {
scan_param_refs(s, &mut refs);
scan_expression_idents(s, &mut refs.expression_idents);
}
nmbrs_workload::model::BindingsDef::Map(m) => {
for v in m.values() {
scan_param_refs(v, &mut refs);
}
}
}
for c in &phase.stop_when {
scan_param_refs(&c.when, &mut refs);
}
if let Some(ci) = &phase.continue_if {
scan_expression_idents(&ci.when, &mut refs.expression_idents);
}
for op in &phase.ops {
scan_op(op, &mut refs);
}
}
fn scan_scenario_nodes(nodes: &[nmbrs_workload::model::ScenarioNode], refs: &mut ParamRefs) {
for node in nodes {
match node {
nmbrs_workload::model::ScenarioNode::Phase(_) => {}
nmbrs_workload::model::ScenarioNode::Comprehension {
comprehension,
children,
..
} => {
for name in comprehension.referenced_source_names() {
if name.contains('{') {
refs.templates.push(name);
} else {
refs.runtime_only_placeholders.insert(name);
}
}
scan_scenario_nodes(children, refs);
}
nmbrs_workload::model::ScenarioNode::DoWhile {
condition,
children,
..
}
| nmbrs_workload::model::ScenarioNode::DoUntil {
condition,
children,
..
} => {
scan_param_refs(condition, refs);
scan_scenario_nodes(children, refs);
}
nmbrs_workload::model::ScenarioNode::IncludedScenario { children, .. } => {
scan_scenario_nodes(children, refs);
}
nmbrs_workload::model::ScenarioNode::Bindings { source, children } => {
let mut deferred = ParamRefs::default();
scan_param_refs(source, &mut deferred);
refs.runtime_only_placeholders.extend(deferred.placeholders);
refs.expression_idents.extend(deferred.expression_idents);
refs.templates.extend(deferred.templates);
scan_scenario_nodes(children, refs);
}
}
}
}
for nodes in workload.scenarios.values() {
scan_scenario_nodes(nodes, &mut refs);
}
refs
}
fn collect_iter_var_names(
workload: &nmbrs_workload::model::Workload,
) -> std::collections::HashSet<String> {
let mut out = std::collections::HashSet::new();
for nodes in workload.scenarios.values() {
for node in nodes {
collect_iter_vars_recursive(node, &mut out);
}
}
for phase in workload.phases.values() {
if let Some(text) = phase.for_each.as_deref()
&& let Ok(comp) =
polydat::iteration::comprehension::spec::parse_comprehension_algebra(text)
{
for name in comp.coordinate_names() {
out.insert(name.to_string());
}
}
}
out
}
fn collect_iter_vars_recursive(
node: &nmbrs_workload::model::ScenarioNode,
out: &mut std::collections::HashSet<String>,
) {
use nmbrs_workload::model::ScenarioNode::*;
match node {
Phase(_) => {}
Comprehension {
comprehension,
children,
..
} => {
for name in comprehension.coordinate_names() {
out.insert(name.to_string());
}
for child in children {
collect_iter_vars_recursive(child, out);
}
}
DoWhile {
children, counter, ..
}
| DoUntil {
children, counter, ..
} => {
if let Some(c) = counter {
out.insert(c.clone());
}
for child in children {
collect_iter_vars_recursive(child, out);
}
}
IncludedScenario { children, .. } => {
for child in children {
collect_iter_vars_recursive(child, out);
}
}
Bindings { children, .. } => {
for child in children {
collect_iter_vars_recursive(child, out);
}
}
}
}
fn collect_polydat_binding_names(
workload: &nmbrs_workload::model::Workload,
) -> std::collections::HashSet<String> {
let mut out = std::collections::HashSet::new();
use nmbrs_workload::model::BindingsDef;
let mut any_input_decl = false;
let mut scan_bindings =
|bindings: &BindingsDef, sink: &mut std::collections::HashSet<String>| {
match bindings {
BindingsDef::PolydatSource(s) => {
scan_polydat_binding_lhs(sink, s);
if s.lines().any(|l| l.trim_start().starts_with("input ")) {
any_input_decl = true;
}
}
BindingsDef::Map(m) => {
for name in m.keys() {
sink.insert(name.clone());
}
}
}
};
scan_bindings(&workload.bindings, &mut out);
for op in &workload.ops {
scan_bindings(&op.bindings, &mut out);
}
for phase in workload.phases.values() {
scan_bindings(&phase.bindings, &mut out);
for op in &phase.ops {
scan_bindings(&op.bindings, &mut out);
}
}
if !any_input_decl {
out.insert("cycle".to_string());
}
out
}
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
struct PolydatBraceFinding {
location: String,
placeholder: String,
}
fn collect_polydat_brace_refs(
workload: &nmbrs_workload::model::Workload,
) -> Vec<PolydatBraceFinding> {
use nmbrs_workload::model::BindingsDef;
let mut out: Vec<PolydatBraceFinding> = Vec::new();
let mut push_refs = |loc: &str, source: &str| {
let parses = polydat::dsl::lexer::lex(source)
.ok()
.and_then(|toks| polydat::dsl::parser::parse(toks).ok())
.is_some();
if parses {
return;
}
for name in scan_polydat_braced_refs(source) {
out.push(PolydatBraceFinding {
location: loc.to_string(),
placeholder: name,
});
}
};
if let BindingsDef::PolydatSource(s) = &workload.bindings {
push_refs("workload `bindings:`", s);
}
for (phase_name, phase) in &workload.phases {
if let BindingsDef::PolydatSource(s) = &phase.bindings {
push_refs(&format!("phase '{phase_name}' bindings"), s);
}
for op in &phase.ops {
if let BindingsDef::PolydatSource(s) = &op.bindings {
push_refs(&format!("phase '{phase_name}' op-bindings"), s);
}
}
}
out
}
fn find_yaml_line_for_brace(yaml_source: &str, placeholder: &str) -> Option<usize> {
let needle = format!("{{{placeholder}}}");
yaml_source
.lines()
.enumerate()
.find(|(_, line)| line.contains(&needle))
.map(|(idx, _)| idx + 1)
}
fn scan_polydat_braced_refs(source: &str) -> Vec<String> {
let bytes = source.as_bytes();
let mut out: Vec<String> = Vec::new();
let mut i = 0;
while i < bytes.len() {
let b = bytes[i];
if b == b'#' {
while i < bytes.len() && bytes[i] != b'\n' {
i += 1;
}
continue;
}
if b == b'"' || b == b'\'' {
let quote = b;
i += 1;
while i < bytes.len() {
if bytes[i] == b'\\' && i + 1 < bytes.len() {
i += 2;
continue;
}
if bytes[i] == quote {
i += 1;
break;
}
i += 1;
}
continue;
}
if b == b'{' {
let start = i + 1;
let mut depth = 1;
let mut j = start;
while j < bytes.len() && depth > 0 {
match bytes[j] {
b'{' => depth += 1,
b'}' => depth -= 1,
_ => {}
}
if depth == 0 {
break;
}
j += 1;
}
if depth == 0 && j > start {
let body = &source[start..j];
let trimmed = body.trim();
if !trimmed.is_empty() {
out.push(trimmed.to_string());
}
i = j + 1;
continue;
}
break;
}
i += 1;
}
out
}
fn scan_polydat_binding_lhs(out: &mut std::collections::HashSet<String>, source: &str) {
for raw_line in source.lines() {
let line = raw_line.trim();
if line.is_empty() || line.starts_with('#') {
continue;
}
if let Some(rest) = line.strip_prefix("input ") {
scan_input_decl_names(out, rest.trim());
continue;
}
if let Some(rest) = line.strip_prefix("extern ") {
scan_input_decl_names(out, rest.trim());
continue;
}
let mut body = line;
loop {
let stripped = body
.strip_prefix("const ")
.or_else(|| body.strip_prefix("cursor "))
.or_else(|| body.strip_prefix("shared "))
.or_else(|| body.strip_prefix("volatile "));
match stripped {
Some(rest) => body = rest,
None => break,
}
}
if body.starts_with('(') {
if let Some(close) = body.find(')') {
let after = body[close + 1..].trim_start();
if after.starts_with(":=") || after.starts_with('=') {
for raw in body[1..close].split(',') {
let name = raw.trim();
if !name.is_empty() {
out.insert(name.to_string());
}
}
}
}
continue;
}
let bytes = body.as_bytes();
let mut i = 0;
while i < bytes.len() && (bytes[i].is_ascii_alphanumeric() || bytes[i] == b'_') {
i += 1;
}
if i == 0 {
continue;
}
let mut j = i;
while j < bytes.len() && bytes[j].is_ascii_whitespace() {
j += 1;
}
let is_binding = bytes.get(j) == Some(&b'=')
|| (bytes.get(j) == Some(&b':') && bytes.get(j + 1) == Some(&b'='));
if is_binding {
out.insert(body[..i].to_string());
continue;
}
if bytes.get(j) == Some(&b':') {
let rest = body[j + 1..].trim_start();
let te = rest
.find(|c: char| !(c.is_ascii_alphanumeric() || c == '_'))
.unwrap_or(rest.len());
if te > 0 {
let after_ty = rest[te..].trim_start();
if after_ty.starts_with(":=") || after_ty.starts_with('=') {
out.insert(body[..i].to_string());
}
}
}
}
}
fn scan_input_decl_names(out: &mut std::collections::HashSet<String>, body: &str) {
let body = body.trim();
if let Some(inner) = body.strip_prefix('(').and_then(|s| s.strip_suffix(')')) {
for part in inner.split(',') {
let name = part.trim().split(':').next().unwrap_or("").trim();
if !name.is_empty() {
out.insert(name.to_string());
}
}
return;
}
let name = body.split(':').next().unwrap_or("").trim();
if !name.is_empty() {
out.insert(name.to_string());
}
}
pub fn resolve_polydat_config(
value: &str,
kernel: &crate::scope_kernel::ScopeKernel,
) -> Option<u64> {
if value.starts_with('{') && value.ends_with('}') {
let inner = &value[1..value.len() - 1];
if let Some(v) = kernel.lookup(inner) {
return Some(value_to_u64(&v));
}
match polydat::dsl::compile::eval_const_expr(inner) {
Ok(v) => Some(value_to_u64(&v)),
Err(e) => {
crate::diag!(
crate::observer::LogLevel::Error,
"error: const expression failed: '{{{inner}}}'"
);
crate::diag!(crate::observer::LogLevel::Error, " {e}");
None
}
}
} else {
parse_count(value)
}
}
fn value_to_u64(v: &polydat::ast::Value) -> u64 {
match v {
polydat::ast::Value::U64(n) => *n,
polydat::ast::Value::F64(f) => *f as u64,
polydat::ast::Value::Bool(b) => {
if *b {
1
} else {
0
}
}
_ => 0,
}
}
fn resolve_scenario(
scenarios: &HashMap<String, Vec<nmbrs_workload::model::ScenarioNode>>,
phase_order: &[String],
name: &str,
) -> Result<Vec<nmbrs_workload::model::ScenarioNode>, String> {
if let Some(nodes) = scenarios.get(name) {
return Ok(nodes.clone());
}
if name == "default" && !phase_order.is_empty() {
return Ok(phase_order
.iter()
.map(|n| nmbrs_workload::model::ScenarioNode::Phase(n.clone()))
.collect());
}
Err(format!("scenario '{name}' not found"))
}
fn format_scenario_tree(
nodes: &[nmbrs_workload::model::ScenarioNode],
phases: &std::collections::HashMap<String, nmbrs_workload::model::WorkloadPhase>,
) -> String {
let mut out = String::new();
format_scenario_nodes(nodes, phases, 0, &mut out);
if out.ends_with('\n') {
out.pop();
}
out
}
fn format_for_combinations(pairs: &[(String, String)], indent_prefix: &str, color: bool) -> String {
let kw_open = if color { "\x1b[1;36m" } else { "" };
let kw_close = if color { "\x1b[0m" } else { "" };
let bracket_open = if color { "\x1b[2m" } else { "" };
let bracket_close = if color { "\x1b[0m" } else { "" };
let widths: Vec<usize> = pairs
.iter()
.map(|(v, s)| v.chars().count().max(s.chars().count()))
.collect();
let pad = |entry: &str, idx: usize, last: bool| -> String {
if last {
entry.to_string()
} else {
let with_comma = format!("{entry},");
let width_target = widths[idx] + 1; let visible = with_comma.chars().count();
if visible >= width_target {
with_comma
} else {
format!("{with_comma}{:<pad$}", "", pad = width_target - visible)
}
}
};
let last_idx = pairs.len().saturating_sub(1);
let vars_line: String = pairs
.iter()
.enumerate()
.map(|(i, (v, _))| pad(v, i, i == last_idx))
.collect::<Vec<_>>()
.join(" ");
let specs_line: String = pairs
.iter()
.enumerate()
.map(|(i, (_, s))| pad(s, i, i == last_idx))
.collect::<Vec<_>>()
.join(" ");
format!(
"{kw_open}for{kw_close} {bracket_open}[{bracket_close}{vars_line}{bracket_open}]{bracket_close}\n\
{indent_prefix} {kw_open}in{kw_close} {bracket_open}[{bracket_close}{specs_line}{bracket_open}]{bracket_close}"
)
}
fn format_scenario_nodes(
nodes: &[nmbrs_workload::model::ScenarioNode],
phases: &std::collections::HashMap<String, nmbrs_workload::model::WorkloadPhase>,
depth: usize,
out: &mut String,
) {
use nmbrs_workload::model::ScenarioNode::*;
let indent = " ".repeat(depth);
for node in nodes {
match node {
Phase(name) => {
let suffix = phases
.get(name)
.map(format_phase_config_suffix)
.unwrap_or_default();
if suffix.is_empty() {
out.push_str(&format!("{indent}{name}\n"));
} else {
out.push_str(&format!("{indent}{name:<32} {suffix}\n",));
}
}
Comprehension {
comprehension,
children,
..
} => {
use polydat::iteration::comprehension::Comprehension as Comp;
let mut body = comprehension;
while let Comp::Order { child, .. } | Comp::Filter { child, .. } = body {
body = child;
}
let header = match body {
Comp::Union { children } => {
let names = comprehension.coordinate_names().join(", ");
format!("for_each_union [{}] ({} sub-spaces)", names, children.len())
}
_ => {
let pairs = comprehension.coordinate_specs();
if pairs.len() == 1 {
let (var, spec) = &pairs[0];
format!("for_each {var} in {spec}")
} else {
format_for_combinations(&pairs, &indent, crate::observer::use_color())
}
}
};
out.push_str(&format!("{indent}{header}\n"));
format_scenario_nodes(children, phases, depth + 1, out);
}
DoWhile {
condition,
counter,
children,
} => {
let ctr = counter
.as_deref()
.map(|c| format!(" (counter={c})"))
.unwrap_or_default();
out.push_str(&format!("{indent}do_while '{condition}'{ctr}\n"));
format_scenario_nodes(children, phases, depth + 1, out);
}
DoUntil {
condition,
counter,
children,
} => {
let ctr = counter
.as_deref()
.map(|c| format!(" (counter={c})"))
.unwrap_or_default();
out.push_str(&format!("{indent}do_until '{condition}'{ctr}\n"));
format_scenario_nodes(children, phases, depth + 1, out);
}
IncludedScenario { name, children } => {
out.push_str(&format!("{indent}scenario '{name}'\n"));
format_scenario_nodes(children, phases, depth + 1, out);
}
Bindings { source, children } => {
let summary = source
.lines()
.map(str::trim)
.find(|l| !l.is_empty())
.unwrap_or("");
if source.lines().filter(|l| !l.trim().is_empty()).count() > 1 {
out.push_str(&format!("{indent}bindings: {summary} …\n"));
} else {
out.push_str(&format!("{indent}bindings: {summary}\n"));
}
format_scenario_nodes(children, phases, depth + 1, out);
}
}
}
}
fn format_phase_config_suffix(phase: &nmbrs_workload::model::WorkloadPhase) -> String {
let mut parts: Vec<String> = Vec::new();
if let Some(c) = phase.cycles.as_deref()
&& !c.is_empty()
{
parts.push(format!("cycles: {c}"));
}
if let Some(c) = phase.concurrency.as_deref()
&& !c.is_empty()
{
parts.push(format!("concurrency: {c}"));
}
if parts.is_empty() {
String::new()
} else {
format!("({})", parts.join(", "))
}
}
fn load_secondary_workload(
reference: &str,
params: &HashMap<String, String>,
base_dir: Option<&std::path::Path>,
bundled_origin: Option<&str>,
) -> Result<(nmbrs_workload::model::Workload, String), String> {
match resolve_secondary_ref(reference, base_dir, bundled_origin)? {
ResolvedWorkload::Path(path) => {
let workload = nmbrs_workload::parse::parse_workload_from_path(
std::path::Path::new(&path),
params,
)
.map_err(|e| format!("parse workload '{path}': {e}"))?;
Ok((workload, canonical_identity(&path)))
}
ResolvedWorkload::Bundled(bundled) => {
let (merged, res_warnings) =
nmbrs_workload::extends::load_and_merge_bundled(bundled)
.map_err(|e| format!("bundled workload `{}`: {e}", bundled.name))?;
let mut workload = nmbrs_workload::parse::parse_workload(&merged, params)
.map_err(|e| format!("parse bundled workload `{}`: {e}", bundled.name))?;
workload.resolution_warnings.extend(res_warnings);
Ok((workload, bundled.name.to_string()))
}
}
}
fn workload_ref_identity(
reference: &str,
base_dir: Option<&std::path::Path>,
bundled_origin: Option<&str>,
) -> Result<String, String> {
match resolve_secondary_ref(reference, base_dir, bundled_origin)? {
ResolvedWorkload::Path(path) => Ok(canonical_identity(&path)),
ResolvedWorkload::Bundled(bundled) => Ok(bundled.name.to_string()),
}
}
fn resolve_driver_manifest(
name: &str,
) -> Result<Option<(nmbrs_workload::drivers::DriverManifest, String)>, String> {
let local_path = std::path::Path::new("drivers")
.join(name)
.join("driver.yaml");
let catalog_name = format!("drivers/{name}/driver");
let bundled = nmbrs_workload::catalog::lookup(&catalog_name);
match (local_path.is_file(), bundled) {
(true, Some(_)) => {
crate::observer::log_tagged(
crate::observer::LogLevel::Warn,
crate::observer::EventTag::in_flight(crate::observer::EventCategory::Resolution),
&format!(
"resolve: driver '{name}' matches multiple resources — local \
manifest {} AND bundled driver `{catalog_name}` — using the \
local manifest (filesystem-first). Same-named resources in \
multiple places invite confusion: prefer a unique name.",
local_path.display()
),
);
let source = std::fs::read_to_string(&local_path)
.map_err(|e| format!("read driver manifest {}: {e}", local_path.display()))?;
let manifest = nmbrs_workload::drivers::parse_driver_manifest(
&source,
&local_path.display().to_string(),
)?;
verify_driver_identity(&manifest, name)?;
let dir = local_path.parent().expect("manifest path has a parent");
let library_ref = resolve_driver_library_local(dir, &manifest.library)?;
Ok(Some((manifest, library_ref)))
}
(true, None) => {
let source = std::fs::read_to_string(&local_path)
.map_err(|e| format!("read driver manifest {}: {e}", local_path.display()))?;
let manifest = nmbrs_workload::drivers::parse_driver_manifest(
&source,
&local_path.display().to_string(),
)?;
verify_driver_identity(&manifest, name)?;
let dir = local_path.parent().expect("manifest path has a parent");
let library_ref = resolve_driver_library_local(dir, &manifest.library)?;
Ok(Some((manifest, library_ref)))
}
(false, Some(entry)) => {
let manifest =
nmbrs_workload::drivers::parse_driver_manifest(entry.source, &catalog_name)?;
verify_driver_identity(&manifest, name)?;
let stem = manifest
.library
.strip_suffix(".yaml")
.or_else(|| manifest.library.strip_suffix(".yml"))
.unwrap_or(&manifest.library);
let library_ref = format!("drivers/{name}/{stem}");
Ok(Some((manifest, library_ref)))
}
(false, None) => Ok(None),
}
}
fn verify_driver_identity(
manifest: &nmbrs_workload::drivers::DriverManifest,
name: &str,
) -> Result<(), String> {
if manifest.driver != name {
return Err(format!(
"driver manifest for '{name}' declares `driver: {}` — the \
manifest must be named for its directory",
manifest.driver
));
}
Ok(())
}
fn resolve_driver_library_local(dir: &std::path::Path, library: &str) -> Result<String, String> {
let as_written = dir.join(library);
if as_written.is_file() {
return Ok(as_written.display().to_string());
}
let with_ext = dir.join(format!("{library}.yaml"));
if with_ext.is_file() {
return Ok(with_ext.display().to_string());
}
Err(format!(
"driver manifest {}: library '{library}' not found beside the manifest",
dir.display()
))
}
fn apply_driver_manifest(
params: &mut HashMap<String, String>,
driver_name: &str,
manifest: nmbrs_workload::drivers::DriverManifest,
library_ref: String,
workload_given: bool,
) -> Result<(), String> {
if params.contains_key("impl") {
return Err(format!(
"driver={driver_name} supplies the implementation library \
('{}') — impl= conflicts; pass one or the other",
manifest.library
));
}
if workload_given {
params.insert("impl".into(), library_ref.clone());
} else {
params.insert("workload".into(), library_ref.clone());
}
params
.entry("adapter".to_string())
.or_insert_with(|| manifest.adapter.clone());
for (k, v) in &manifest.default_params {
params.entry(k.clone()).or_insert_with(|| v.clone());
}
crate::diag!(
crate::observer::LogLevel::Info,
"driver: {driver_name} → adapter={}, library={library_ref}",
manifest.adapter
);
Ok(())
}
fn resolve_secondary_ref(
reference: &str,
base_dir: Option<&std::path::Path>,
bundled_origin: Option<&str>,
) -> Result<ResolvedWorkload, String> {
let pinned = reference.starts_with("./") || reference.starts_with("../");
let origin_file: Option<String> = base_dir
.map(|dir| dir.join(reference))
.filter(|c| c.is_file())
.map(|c| c.display().to_string());
if pinned && base_dir.is_some() && origin_file.is_some() {
return Ok(ResolvedWorkload::Path(origin_file.unwrap()));
}
let cwd_file: Option<String> = resolve_workload_file(reference).filter(|p| {
origin_file.as_deref().map(canonical_identity) != Some(canonical_identity(p))
});
let stem = reference
.strip_suffix(".yaml")
.or_else(|| reference.strip_suffix(".yml"))
.unwrap_or(reference);
let stem = stem.strip_prefix("./").unwrap_or(stem);
let ns_stem = bundled_origin
.and_then(|o| o.rsplit_once('/'))
.map(|(ns, _)| format!("{ns}/{stem}"));
let bundled: Option<&'static nmbrs_workload::catalog::BundledWorkload> =
nmbrs_workload::catalog::lookup(reference)
.or_else(|| ns_stem.as_deref().and_then(nmbrs_workload::catalog::lookup))
.or_else(|| nmbrs_workload::catalog::lookup(stem));
let mut names: Vec<String> = Vec::new();
if let Some(p) = &origin_file {
names.push(format!("file {p} (beside the referring document)"));
}
if let Some(p) = &cwd_file {
names.push(format!("local file {p}"));
}
if let Some(b) = bundled {
names.push(format!("bundled workload `{}`", b.name));
}
if names.len() > 1 {
crate::observer::log_tagged(
crate::observer::LogLevel::Warn,
crate::observer::EventTag::in_flight(crate::observer::EventCategory::Resolution),
&format!(
"resolve: reference '{reference}' matches multiple resources — {} — \
using the nearest ({}). Same-named resources in multiple places \
invite confusion: prefer a unique name, or pin the intent with a \
`./` path / full catalog name.",
names.join(" AND "),
names[0]
),
);
}
if let Some(p) = origin_file {
return Ok(ResolvedWorkload::Path(p));
}
if let Some(p) = cwd_file {
return Ok(ResolvedWorkload::Path(p));
}
if let Some(b) = bundled {
return Ok(ResolvedWorkload::Bundled(b));
}
Err(format!(
"workload not found: '{reference}'. Not a local file, and no bundled \
workload by that name — `nmbrs describe workloads` lists what this \
binary carries.{}",
nmbrs_workload::suggest::did_you_mean(&nmbrs_workload::suggest::suggest_workloads(
reference
)),
))
}
fn canonical_identity(reference: &str) -> String {
std::fs::canonicalize(reference)
.map(|p| p.display().to_string())
.unwrap_or_else(|_| reference.to_string())
}
fn peek_stick_session(params: &HashMap<String, String>, args: &[String]) -> Option<bool> {
if params.contains_key("op") {
return None; }
let workload_raw = params.get("workload").cloned().or_else(|| {
args.iter()
.find(|a| a.ends_with(".yaml") || a.ends_with(".yml"))
.cloned()
})?;
let merged = match resolve_workload(&workload_raw).ok()? {
ResolvedWorkload::Path(p) => {
nmbrs_workload::extends::load_and_merge(std::path::Path::new(&p))
.ok()?
.0
}
ResolvedWorkload::Bundled(b) => nmbrs_workload::extends::load_and_merge_bundled(b).ok()?.0,
};
let doc: serde_yaml::Value = serde_yaml::from_str(&merged).ok()?;
doc.get("stick_session")?.as_bool()
}
pub enum ResolvedWorkload {
Path(String),
Bundled(&'static nmbrs_workload::catalog::BundledWorkload),
}
pub fn resolve_workload(name: &str) -> Result<ResolvedWorkload, String> {
let local = resolve_workload_file(name);
let bundled = nmbrs_workload::catalog::lookup(name);
match (local, bundled) {
(Some(local_path), Some(b)) => {
if !name.starts_with("./") && !name.starts_with("../") {
crate::observer::log_tagged(
crate::observer::LogLevel::Warn,
crate::observer::EventTag::in_flight(
crate::observer::EventCategory::Resolution,
),
&format!(
"resolve: workload '{name}' matches multiple resources — \
local file {local_path} AND bundled workload `{}` — using \
the local file (filesystem-first). Same-named resources in \
multiple places invite confusion: prefer a unique name, or \
pin the intent with a `./` path.",
b.name
),
);
}
Ok(ResolvedWorkload::Path(local_path))
}
(Some(local_path), None) => Ok(ResolvedWorkload::Path(local_path)),
(None, Some(b)) => Ok(ResolvedWorkload::Bundled(b)),
(None, None) => Err(format!(
"workload not found: '{name}'. Not a local file, and no bundled \
workload by that name — `nmbrs describe workloads` lists what \
this binary carries.{}",
nmbrs_workload::suggest::did_you_mean(&nmbrs_workload::suggest::suggest_workloads(
name
),)
)),
}
}
fn resolve_workload_file(name: &str) -> Option<String> {
let p = std::path::Path::new(name);
if p.exists() {
return Some(name.to_string());
}
if name.ends_with(".yaml") || name.ends_with(".yml") {
let under = format!("workloads/{name}");
if std::path::Path::new(&under).exists() {
return Some(under);
}
return None;
}
for ext in [".yaml", ".yml"] {
let with_ext = format!("{name}{ext}");
if std::path::Path::new(&with_ext).exists() {
return Some(with_ext);
}
}
for ext in ["", ".yaml", ".yml"] {
let under = format!("workloads/{name}{ext}");
if std::path::Path::new(&under).exists() {
return Some(under);
}
}
None
}
pub fn normalize_args(args: &[String]) -> Vec<String> {
let value_flags = known_value_flags();
let mut result = Vec::new();
let mut workload_seen = false;
let mut scenario_set = false;
let mut iter = args.iter().peekable();
while let Some(arg) = iter.next() {
if value_flags.iter().any(|f| *f == arg) {
result.push(arg.clone());
if let Some(next) = iter.next() {
result.push(next.clone());
}
continue;
}
if !workload_seen
&& (arg.ends_with(".yaml") || arg.ends_with(".yml") || arg.contains("workload="))
{
workload_seen = true;
result.push(arg.clone());
} else if workload_seen && !scenario_set && !arg.contains('=') && !arg.starts_with('-') {
result.push(format!("scenario={arg}"));
scenario_set = true;
} else {
result.push(arg.clone());
}
}
result
}
const RECOGNIZED_BARE_FLAGS: &[&str] = &[
"--strict", "--resume-latest", "--force-retry-failed", "--refine", ];
pub(crate) fn elide_outer_quotes(s: &str) -> &str {
let bytes = s.as_bytes();
if bytes.len() < 2 {
return s;
}
let first = bytes[0];
let last = bytes[bytes.len() - 1];
if (first == b'\'' || first == b'"') && first == last {
&s[1..s.len() - 1]
} else {
s
}
}
const SESSION_DIR_FLAGS: &[&str] = &[
"--session",
"--session-name",
"--session-path",
"--session-reuse",
"--session-keep",
"--session-shelflife",
"--readout",
];
pub fn detect_conflicting_duplicate_params(args: &[String]) -> Result<(), String> {
let mut seen: HashMap<String, String> = HashMap::new();
let mut iter = args.iter().peekable();
while let Some(arg) = iter.next() {
if known_value_flags()
.iter()
.any(|p| arg == p || arg.starts_with(&format!("{p}=")))
{
if !arg.contains('=') {
let _consumed = iter.next();
}
continue;
}
let unquoted = elide_outer_quotes(arg.as_str());
let stripped = unquoted.trim_start_matches('-');
let Some(eq_pos) = stripped.find('=') else {
continue;
};
let key = stripped[..eq_pos].to_string();
if key.contains('.') && !key.contains('/') && !key.contains('\\') {
continue;
}
let value = elide_outer_quotes(&stripped[eq_pos + 1..]).to_string();
match seen.get(&key) {
Some(prev) if *prev != value => {
return Err(format!(
"parameter '{key}' specified more than once with conflicting \
values ('{prev}' and '{value}') — pass it exactly once"
));
}
Some(_) => {} None => {
seen.insert(key, value);
}
}
}
Ok(())
}
pub fn parse_params(args: &[String]) -> HashMap<String, String> {
let mut params = HashMap::new();
let mut iter = args.iter().peekable();
while let Some(arg) = iter.next() {
if known_value_flags()
.iter()
.any(|p| arg == p || arg.starts_with(&format!("{p}=")))
{
if !arg.contains('=') {
let _consumed = iter.next();
}
continue;
}
let unquoted = elide_outer_quotes(arg.as_str());
let stripped = unquoted.trim_start_matches('-');
if let Some(eq_pos) = stripped.find('=') {
let key = stripped[..eq_pos].to_string();
if key.contains('.') && !key.contains('/') && !key.contains('\\') {
continue;
}
let value = elide_outer_quotes(&stripped[eq_pos + 1..]).to_string();
params.insert(key, value);
} else if arg.ends_with(".yaml") || arg.ends_with(".yml") {
} else if is_recognized_bare_flag(arg.as_str()) || arg.starts_with("--polydat-lib=") {
} else {
crate::diag!(
crate::observer::LogLevel::Error,
"error: unrecognized argument '{arg}'. Expected key=value format."
);
std::process::exit(1);
}
}
params
}
fn overlay_cli_params(
mut base: HashMap<String, String>,
cli: &HashMap<String, String>,
) -> HashMap<String, String> {
for (k, v) in cli {
let coerced = match base.get(k) {
Some(existing) => nmbrs_workload::magnitude::coerce_param_override(existing, v),
None => v.clone(),
};
base.insert(k.clone(), coerced);
}
base
}
pub fn effective_params(args: &[String]) -> HashMap<String, String> {
let cli = parse_params(args);
let workload_ref = cli.get("workload").cloned().or_else(|| {
args.iter()
.find(|a| (a.ends_with(".yaml") || a.ends_with(".yml")) && !a.contains('='))
.cloned()
});
let base = workload_ref
.as_deref()
.and_then(nmbrs_workload::verify::declared_params)
.unwrap_or_default();
overlay_cli_params(base, &cli)
}
pub const SESSION_PARAMS: &[&str] = &["metrics_cadence"];
pub fn session_param_signature(reference: &str) -> Vec<(String, String)> {
let declared = nmbrs_workload::verify::declared_params(reference).unwrap_or_default();
let mut sig: Vec<(String, String)> = SESSION_PARAMS
.iter()
.filter_map(|k| declared.get(*k).map(|v| ((*k).to_string(), v.clone())))
.collect();
sig.sort();
sig
}
fn resolve_cadence_config(
params: &HashMap<String, String>,
observer: &Arc<dyn crate::observer::RunObserver>,
) -> Result<(std::time::Duration, nmbrs_metrics::cadence::Cadences), String> {
use std::time::Duration;
let Some(raw) = params.get("metrics_cadence") else {
let cadences = observer
.cadences()
.unwrap_or_else(nmbrs_metrics::cadence::Cadences::defaults);
return Ok((Duration::from_secs(1), cadences));
};
let floor = nmbrs_metrics::cadence::parse_duration(raw).map_err(|_| {
format!("metrics_cadence: invalid duration `{raw}` (use e.g. `100ms`, `200ms`, `1s`)")
})?;
if floor.is_zero() {
return Err("metrics_cadence: must be greater than zero".to_string());
}
let mut layers = vec![floor];
for secs in [1u64, 10, 30, 60, 300] {
let d = Duration::from_secs(secs);
if d > floor {
layers.push(d);
}
}
let cadences = nmbrs_metrics::cadence::Cadences::new(&layers)
.map_err(|e| format!("metrics_cadence `{raw}`: cannot build a cadence ladder: {e:?}"))?;
Ok((floor, cadences))
}
pub fn collect_repeated_flag(args: &[String], name: &str) -> Vec<String> {
let mut out = Vec::new();
let mut iter = args.iter().peekable();
let long_eq = format!("--{name}=");
let bare_eq = format!("{name}=");
while let Some(arg) = iter.next() {
let unquoted = elide_outer_quotes(arg.as_str());
if let Some(v) = unquoted.strip_prefix(&long_eq) {
out.push(elide_outer_quotes(v).to_string());
} else if let Some(v) = unquoted.strip_prefix(&bare_eq) {
out.push(elide_outer_quotes(v).to_string());
} else if unquoted == format!("--{name}")
&& let Some(v) = iter.next()
{
out.push(elide_outer_quotes(v.as_str()).to_string());
}
}
out
}
pub fn parse_count(s: &str) -> Option<u64> {
let s = s.trim().to_uppercase();
if let Some(n) = s.strip_suffix('K') {
n.trim().parse::<u64>().ok().map(|v| v * 1_000)
} else if let Some(n) = s.strip_suffix('M') {
n.trim().parse::<u64>().ok().map(|v| v * 1_000_000)
} else if let Some(n) = s.strip_suffix('B') {
n.trim().parse::<u64>().ok().map(|v| v * 1_000_000_000)
} else {
s.parse().ok()
}
}
fn closest_match<'a>(input: &str, candidates: &[&'a str]) -> Option<&'a str> {
let mut best: Option<(&str, usize)> = None;
for &candidate in candidates {
let d = levenshtein(input, candidate);
if best.is_none() || d < best.unwrap().1 {
best = Some((candidate, d));
}
}
best.filter(|(_, d)| *d <= (input.len() / 2).max(2))
.map(|(s, _)| s)
}
fn levenshtein(a: &str, b: &str) -> usize {
let a: Vec<char> = a.chars().collect();
let b: Vec<char> = b.chars().collect();
let (m, n) = (a.len(), b.len());
let mut prev = (0..=n).collect::<Vec<_>>();
let mut curr = vec![0; n + 1];
for i in 1..=m {
curr[0] = i;
for j in 1..=n {
let cost = if a[i - 1] == b[j - 1] { 0 } else { 1 };
curr[j] = (prev[j] + 1).min(curr[j - 1] + 1).min(prev[j - 1] + cost);
}
std::mem::swap(&mut prev, &mut curr);
}
prev[n]
}
pub use polydat::kernel::{ManifestEntry, extract_manifest};
#[cfg(test)]
mod tests {
use super::*;
fn pp(args: &[&str]) -> HashMap<String, String> {
parse_params(&args.iter().map(|s| s.to_string()).collect::<Vec<_>>())
}
#[test]
fn parse_params_bare_unchanged() {
let m = pp(&["cursor=0..53%"]);
assert_eq!(m.get("cursor").map(String::as_str), Some("0..53%"));
}
#[test]
fn parse_params_value_single_quoted_stripped() {
let m = pp(&["cursor='0..53%'"]);
assert_eq!(m.get("cursor").map(String::as_str), Some("0..53%"));
}
#[test]
fn parse_params_value_double_quoted_stripped() {
let m = pp(&["cursor=\"0..53%\""]);
assert_eq!(m.get("cursor").map(String::as_str), Some("0..53%"));
}
#[test]
fn parse_params_whole_arg_single_quoted_stripped() {
let m = pp(&["'cursor=0..53%'"]);
assert_eq!(m.get("cursor").map(String::as_str), Some("0..53%"));
}
#[test]
fn parse_params_whole_arg_double_quoted_stripped() {
let m = pp(&["\"cursor=0..53%\""]);
assert_eq!(m.get("cursor").map(String::as_str), Some("0..53%"));
}
#[test]
fn parse_params_bracket_value_with_quotes() {
let m = pp(&["cursor='[0..53%)'"]);
assert_eq!(m.get("cursor").map(String::as_str), Some("[0..53%)"));
}
#[test]
fn parse_params_mismatched_quotes_not_stripped() {
let m = pp(&["cursor='0..53%\""]);
assert_eq!(m.get("cursor").map(String::as_str), Some("'0..53%\""));
}
fn dup(args: &[&str]) -> Result<(), String> {
let owned: Vec<String> = args.iter().map(|s| s.to_string()).collect();
detect_conflicting_duplicate_params(&owned)
}
#[test]
fn duplicate_conflicting_scenario_is_rejected() {
let err = dup(&["workload=x.yaml", "scenario=reset", "scenario=idx_sweep"]).unwrap_err();
assert!(
err.contains("scenario") && err.contains("reset") && err.contains("idx_sweep"),
"expected a conflicting-duplicate error naming both values, got: {err}"
);
}
#[test]
fn duplicate_identical_value_is_allowed() {
assert!(dup(&["scenario=idx_sweep", "scenario=idx_sweep"]).is_ok());
}
#[test]
fn distinct_params_are_allowed() {
assert!(dup(&["workload=x.yaml", "scenario=reset", "cycles=10", "host=h"]).is_ok());
}
#[test]
fn conflicting_duplicate_any_param_is_rejected() {
assert!(dup(&["cycles=10", "cycles=20"]).is_err());
}
#[test]
fn duplicate_check_elides_quotes_before_comparing() {
assert!(dup(&["scenario=reset", "scenario='reset'"]).is_ok());
assert!(dup(&["scenario='reset'", "scenario=\"idx_sweep\""]).is_err());
}
#[test]
fn duplicate_check_skips_session_flags_and_dotted_overrides() {
assert!(dup(&["--session-path", "/a", "--session-path", "/b"]).is_ok());
assert!(dup(&["phase1.cycles=10", "phase2.cycles=20"]).is_ok());
}
#[test]
fn parse_params_equals_in_value_preserved() {
let m = pp(&["key='a=b'"]);
assert_eq!(m.get("key").map(String::as_str), Some("a=b"));
}
#[test]
fn parse_params_multiple_params_independent() {
let m = pp(&["dataset=example", "cursor='0..1%'", "concurrency=\"100\""]);
assert_eq!(m.get("dataset").map(String::as_str), Some("example"));
assert_eq!(m.get("cursor").map(String::as_str), Some("0..1%"));
assert_eq!(m.get("concurrency").map(String::as_str), Some("100"));
}
fn scan_to_set(src: &str) -> std::collections::HashSet<String> {
let mut out = std::collections::HashSet::new();
scan_polydat_binding_lhs(&mut out, src);
out
}
#[test]
fn polydat_brace_guard_allows_valid_if_block() {
let src = "extern segments: u64 = 0\nmean := if segments > 0 { 100 } else { 0 }\n";
let parses = polydat::dsl::lexer::lex(src)
.ok()
.and_then(|t| polydat::dsl::parser::parse(t).ok())
.is_some();
assert!(parses, "if-block source must parse: {src}");
}
#[test]
fn polydat_brace_guard_still_catches_stray_placeholder() {
let src = "const passes := multiples_at_least({min_query_cycles}, base)\n";
let parses = polydat::dsl::lexer::lex(src)
.ok()
.and_then(|t| polydat::dsl::parser::parse(t).ok())
.is_some();
assert!(!parses, "stray placeholder must fail to parse");
let refs = scan_polydat_braced_refs(src);
assert!(refs.iter().any(|r| r == "min_query_cycles"), "got {refs:?}");
}
#[test]
fn scan_polydat_braced_refs_flags_expression_position_braces() {
let refs = scan_polydat_braced_refs(
"const passes := multiples_at_least({min_query_cycles}, base)\n",
);
assert_eq!(refs, vec!["min_query_cycles".to_string()]);
}
#[test]
fn scan_polydat_braced_refs_ignores_braces_inside_double_quotes() {
let refs = scan_polydat_braced_refs(
"const prebuffered := dataset_prebuffer(\"{dataset}:{profile}\")\n",
);
assert!(
refs.is_empty(),
"must not flag `{{dataset}}` / `{{profile}}` inside string \
literal — string interpolation handles them, got {refs:?}"
);
}
#[test]
fn scan_polydat_braced_refs_ignores_braces_inside_single_quotes() {
let refs = scan_polydat_braced_refs("tag := assert_eq(actual, '{expected}')\n");
assert!(
refs.is_empty(),
"single-quoted strings get the same treatment: {refs:?}"
);
}
#[test]
fn scan_polydat_braced_refs_handles_escaped_quotes_in_strings() {
let refs = scan_polydat_braced_refs("x := concat(\"prefix \\\"{embedded}\\\" suffix\")\n");
assert!(refs.is_empty(), "escaped quotes inside strings: {refs:?}");
}
#[test]
fn scan_polydat_braced_refs_ignores_comments() {
let refs = scan_polydat_braced_refs(
"# this is a comment with {fake} placeholder\n\
const real := 1\n",
);
assert!(
refs.is_empty(),
"`{{fake}}` inside a comment must not be flagged: {refs:?}"
);
}
#[test]
fn scan_polydat_braced_refs_catches_multiple_invalid_braces() {
let refs = scan_polydat_braced_refs(
"a := foo({x}, {y})\n\
b := bar({z})\n",
);
assert_eq!(
refs,
vec!["x".to_string(), "y".to_string(), "z".to_string(),]
);
}
#[test]
fn scan_polydat_braced_refs_handles_mixed_string_and_expression_braces() {
let refs = scan_polydat_braced_refs("x := concat(\"foo {inside}\", {outside})\n");
assert_eq!(refs, vec!["outside".to_string()]);
}
#[test]
fn scan_polydat_binding_lhs_handles_typed_cell_form() {
let mut out = std::collections::HashSet::new();
scan_polydat_binding_lhs(
&mut out,
"shared sstables: u64 := 0\nshared measured: f64 := 1.0\nplain := 2\n",
);
assert!(out.contains("sstables"), "{out:?}");
assert!(out.contains("measured"), "{out:?}");
assert!(out.contains("plain"), "{out:?}");
}
#[test]
fn scan_polydat_binding_lhs_handles_tuple_destructure() {
let names = scan_to_set("(y, mo, d, h, mi, s, ms) := date_components(0)\n");
for expected in ["y", "mo", "d", "h", "mi", "s", "ms"] {
assert!(
names.contains(expected),
"tuple-LHS scanner missed `{expected}` — got {names:?}"
);
}
}
#[test]
fn scan_polydat_binding_lhs_finds_all_recognised_shapes() {
let names = scan_to_set(
"const prebuffered := dataset_prebuffer(\"foo\")\n\
cursor q = range(0, 100)\n\
query_vector := query_vector_at(prebuffered, q)\n\
shared query_passes := set_or_get(query_passes, 7)\n\
const tag := \"label_00\"\n",
);
for expected in ["prebuffered", "q", "query_vector", "query_passes", "tag"] {
assert!(
names.contains(expected),
"scanner missed `{expected}` — got {names:?}"
);
}
}
#[test]
fn scan_polydat_binding_lhs_skips_comments_and_blank_lines() {
let names = scan_to_set(
"# comment\n\
\n\
const real_binding := 1\n\
# another comment\n",
);
assert_eq!(names.len(), 1);
assert!(names.contains("real_binding"));
}
#[test]
fn scan_polydat_binding_lhs_ignores_non_binding_lines() {
let names = scan_to_set(
"foo(1, 2)\n\
bar.baz\n\
const real := 1\n",
);
assert_eq!(names.len(), 1);
assert!(names.contains("real"));
}
#[test]
fn scan_polydat_binding_lhs_picks_up_input_decl_bare() {
let names = scan_to_set("input cycle: u64\nx := hash(cycle)\n");
assert!(names.contains("cycle"));
assert!(names.contains("x"));
}
#[test]
fn scan_polydat_binding_lhs_picks_up_input_decl_untyped() {
let names = scan_to_set("input cycle\n");
assert!(names.contains("cycle"));
}
#[test]
fn scan_polydat_binding_lhs_picks_up_input_decl_tuple() {
let names = scan_to_set("input (cycle: u64, q: f64)\n");
assert!(names.contains("cycle"));
assert!(names.contains("q"));
}
#[test]
fn scan_polydat_binding_lhs_picks_up_extern_decl() {
let names = scan_to_set(
"extern active_compactions: u64 = 0\n\
extern completion_ratio: f64 = 0.0\n\
extern (a: u64, b: f64)\n",
);
assert!(names.contains("active_compactions"));
assert!(names.contains("completion_ratio"));
assert!(names.contains("a"));
assert!(names.contains("b"));
}
#[test]
fn parse_dryrun_controls_sets_list_flag() {
let cfg = DiagnosticConfig::parse("controls");
assert!(cfg.list_controls);
assert_eq!(cfg.depth, ExecDepth::Phase);
}
#[test]
fn parse_dryrun_controls_combines_with_other_flags() {
let cfg = DiagnosticConfig::parse("controls,labels");
assert!(cfg.list_controls);
assert!(cfg.show_labels);
}
#[test]
fn parse_dryrun_unknown_flag_does_not_set_controls() {
let cfg = DiagnosticConfig::parse("phase,bogus");
assert!(!cfg.list_controls);
}
#[test]
fn parse_dryrun_op_sets_op_depth() {
let cfg = DiagnosticConfig::parse("op");
assert_eq!(cfg.depth, ExecDepth::Op);
}
#[test]
fn parse_dryrun_phase_still_sets_phase_depth() {
let cfg = DiagnosticConfig::parse("phase");
assert_eq!(cfg.depth, ExecDepth::Phase);
}
#[test]
fn parse_dryrun_cycle_still_sets_cycle_depth() {
let cfg = DiagnosticConfig::parse("cycle");
assert_eq!(cfg.depth, ExecDepth::Cycle);
}
#[test]
fn parse_dryrun_op_combines_with_wiring_flag() {
let cfg = DiagnosticConfig::parse("op,wiring");
assert_eq!(cfg.depth, ExecDepth::Op);
assert!(cfg.show_wiring);
}
#[test]
fn parse_dryrun_wiring_alone_bumps_depth_to_op() {
let cfg = DiagnosticConfig::parse("wiring");
assert_eq!(cfg.depth, ExecDepth::Op);
assert!(cfg.show_wiring);
}
#[test]
fn parse_dryrun_wiring_does_not_override_explicit_depth() {
let cfg = DiagnosticConfig::parse("phase,wiring");
assert_eq!(cfg.depth, ExecDepth::Phase);
assert!(cfg.show_wiring);
}
#[test]
fn exec_depth_ordering_matches_srd_13d() {
assert!(ExecDepth::Phase < ExecDepth::Op);
assert!(ExecDepth::Op < ExecDepth::Cycle);
assert!(ExecDepth::Cycle < ExecDepth::Full);
assert!(ExecDepth::Phase < ExecDepth::Cycle);
assert!(ExecDepth::Op < ExecDepth::Full);
}
#[test]
fn exec_depth_phase_and_op_short_circuit_before_cycles() {
assert!(ExecDepth::Phase < ExecDepth::Cycle);
assert!(ExecDepth::Op < ExecDepth::Cycle);
assert!((ExecDepth::Cycle >= ExecDepth::Cycle));
assert!((ExecDepth::Full >= ExecDepth::Cycle));
}
#[test]
fn render_scope_elision_summary_shows_materialised_and_elides_to() {
use nmbrs_workload::model::{BindingsDef, ScenarioNode, WorkloadPhase};
use std::collections::HashMap;
let phase = WorkloadPhase {
key_metrics: Vec::new(),
dimensions: Default::default(),
cycles: None,
concurrency: None,
rate: None,
daemon: false,
adapter: None,
errors: None,
tries: None,
tries_backoff: None,
interval: None,
repeat: None,
error_rate_max: None,
timeout: None,
stop_when: Vec::new(),
throttle: None,
continue_if: None,
tags: None,
ops: vec![],
for_each: None,
loop_scope: None,
iter_scope: None,
checkpoint: None,
status_metrics: vec![],
metrics: Default::default(),
bindings: BindingsDef::default(),
poll: None,
optimize: None,
};
let mut phases = HashMap::new();
phases.insert("predict".to_string(), phase);
let mut tree = crate::scope_tree::ScopeTree::build(
"default",
&[ScenarioNode::Phase("predict".into())],
);
let inputs = crate::scope_elision::ClassifyInputs {
bindings: &BindingsDef::default(),
params: &HashMap::new(),
phases: &phases,
};
crate::scope_elision::classify_and_mark(&mut tree, &inputs);
let mut buf: Vec<u8> = Vec::new();
render_scope_elision_summary(&tree, &mut buf).unwrap();
let s = String::from_utf8(buf).unwrap();
assert!(s.contains("scope elision summary"), "missing header: {s}");
assert!(
s.contains("workload") && s.contains("materialised=true"),
"expected materialised=true line for workload root: {s}"
);
assert!(
s.contains("elides-to=workload"),
"expected elides-to=workload for empty phase: {s}"
);
assert!(
s.contains("workload.scenario.default"),
"expected scenario logical name: {s}"
);
assert!(
s.contains("workload.scenario.default.phase.predict"),
"expected phase logical name: {s}"
);
}
#[test]
fn render_controls_tree_empty_session_writes_placeholder() {
let root = nmbrs_metrics::component::Component::root(
nmbrs_metrics::labels::Labels::of("session", "t"),
std::collections::HashMap::new(),
);
let mut buf: Vec<u8> = Vec::new();
render_controls_tree(&root, &mut buf).unwrap();
let s = String::from_utf8(buf).unwrap();
assert!(s.contains("no controls declared"), "got: {s}");
}
#[test]
fn render_controls_tree_lists_session_root_controls() {
let root = nmbrs_metrics::component::Component::root(
nmbrs_metrics::labels::Labels::of("session", "t"),
std::collections::HashMap::new(),
);
root.read().unwrap().controls().declare(
nmbrs_metrics::controls::ControlBuilder::new("log_level", 1u32)
.reify_as_gauge(|v| Some(*v as f64))
.branch_scope(nmbrs_metrics::controls::BranchScope::Subtree)
.from_f64(|v| Ok(v as u32))
.final_at_scope("session_root")
.build(),
);
let mut buf: Vec<u8> = Vec::new();
render_controls_tree(&root, &mut buf).unwrap();
let s = String::from_utf8(buf).unwrap();
assert!(s.contains("log_level"), "missing name: {s}");
assert!(s.contains("scope=subtree"), "missing scope: {s}");
assert!(
s.contains("final@session_root"),
"missing final marker: {s}"
);
assert!(s.contains("f64-writable"), "missing write surface: {s}");
}
fn s(v: &[&str]) -> Vec<String> {
v.iter().map(|x| x.to_string()).collect()
}
#[test]
fn normalize_args_session_path_space_form_value_passes_through() {
let out = normalize_args(&s(&[
"wl.yaml",
"cycles=2",
"--session-path",
"target/test-tmp/foo/session",
]));
assert!(
!out.iter().any(|a| a.starts_with("scenario=")),
"scenario= auto-promotion fired on a flag value: {out:?}"
);
assert_eq!(
out,
s(&[
"wl.yaml",
"cycles=2",
"--session-path",
"target/test-tmp/foo/session",
])
);
}
#[test]
fn normalize_args_session_path_equals_form_unchanged() {
let out = normalize_args(&s(&[
"wl.yaml",
"--session-path=target/test-tmp/foo/session",
]));
assert_eq!(
out,
s(&["wl.yaml", "--session-path=target/test-tmp/foo/session",])
);
}
#[test]
fn normalize_args_real_scenario_positional_still_promotes() {
let out = normalize_args(&s(&["wl.yaml", "myscenario", "cycles=2"]));
assert_eq!(out, s(&["wl.yaml", "scenario=myscenario", "cycles=2",]));
}
#[test]
fn normalize_args_scenario_after_session_path_still_promotes() {
let out = normalize_args(&s(&["wl.yaml", "--session-path", "/tmp/x", "myscenario"]));
assert_eq!(
out,
s(&["wl.yaml", "--session-path", "/tmp/x", "scenario=myscenario",])
);
}
#[test]
fn normalize_args_readout_value_passes_through() {
let out = normalize_args(&s(&["wl.yaml", "--readout", "throughput ok_pct"]));
assert!(
!out.iter().any(|a| a.starts_with("scenario=")),
"readout body misread as scenario: {out:?}"
);
}
#[test]
fn format_for_combinations_aligns_columns() {
let pairs = vec![
("sm".to_string(), "{sm_values}".to_string()),
("mnc".to_string(), "{mnc_values}".to_string()),
(
"alf_label".to_string(),
"concat({alf_label_values})".to_string(),
),
];
let out = format_for_combinations(&pairs, "", false);
let lines: Vec<&str> = out.split('\n').collect();
assert_eq!(lines.len(), 2, "MUST produce exactly 2 lines: {out:?}");
assert!(
lines[0].starts_with("for ["),
"first line MUST start with `for [`: {:?}",
lines[0]
);
assert!(
lines[1].starts_with(" in ["),
"second line MUST start with ` in [`: {:?}",
lines[1]
);
let l0_bracket = lines[0].find('[').unwrap();
let l1_bracket = lines[1].find('[').unwrap();
assert_eq!(
l0_bracket, l1_bracket,
"`[` brackets MUST align: line0={l0_bracket}, line1={l1_bracket}"
);
let mnc_pos = lines[0].find("mnc").unwrap();
let mnc_values_pos = lines[1].find("{mnc_values}").unwrap();
assert_eq!(
mnc_pos, mnc_values_pos,
"column 2 MUST align: `mnc`@{mnc_pos} vs `{{mnc_values}}`@{mnc_values_pos}\n{out}"
);
let alf_pos = lines[0].find("alf_label").unwrap();
let alf_concat_pos = lines[1].find("concat(").unwrap();
assert_eq!(
alf_pos, alf_concat_pos,
"column 3 MUST align: `alf_label`@{alf_pos} vs `concat(...)`@{alf_concat_pos}\n{out}"
);
assert!(lines[0].ends_with(']'));
assert!(lines[1].ends_with(']'));
}
}