use std::collections::HashMap;
use std::path::PathBuf;
use std::sync::{Arc, RwLock};
use indexmap::IndexMap;
use crate::activity::{Activity, ActivityConfig};
use crate::adapter::DriverAdapter;
use crate::opseq::{OpSequence, SequencerType};
use crate::synthesis::OpBuilder;
use nmbrs_metrics::cadence_reporter::CadenceReporter;
use nmbrs_metrics::component::{self, Component, ComponentState};
use nmbrs_metrics::labels::Labels;
use nmbrs_workload::model::{ScenarioNode, WorkloadPhase};
use polydat::kernel::{ScopeCoord, format_scope_coordinate_path};
pub(crate) fn resolve_stop_scope(
at: Option<nmbrs_workload::model::ScopeLevel>,
each: &[nmbrs_workload::model::ScopeLevel],
) -> crate::stop_conditions::StopScope {
use crate::stop_conditions::StopScope as S;
use nmbrs_workload::model::ScopeLevel as L;
fn rank(l: L) -> u8 {
match l {
L::SelfScope => 0,
L::Op => 1,
L::Phase => 2,
L::Scenario => 3,
L::Workload => 4,
}
}
let level = at.or_else(|| each.iter().copied().min_by_key(|l| rank(*l)));
match level {
Some(L::Scenario) => S::Scenario,
Some(L::Workload) => S::Workload,
_ => S::Phase,
}
}
#[derive(Clone)]
pub struct ExecCtx {
pub phases: HashMap<String, WorkloadPhase>,
pub phase_param_overrides: Arc<Vec<crate::phase_params::PhaseParamOverride>>,
pub workload_readouts: nmbrs_workload::model::ReadoutsBindings,
pub cli_readout_override: Option<String>,
pub workload_params: HashMap<String, String>,
pub wrappers_override: Option<Vec<String>>,
pub wrap_default_order: Option<Vec<String>>,
pub workload_scope: Arc<crate::scope_kernel::ScopeKernel>,
pub polydat_lib_paths: Vec<PathBuf>,
pub workload_dir: Option<PathBuf>,
pub strict: bool,
pub driver: String,
pub merged_params: HashMap<String, String>,
pub dry_run: Option<&'static str>,
pub phase_filter: Option<Arc<crate::phase_filter::PhasePattern>>,
pub refine_plan: Option<Arc<crate::refine_plan::RefinePlan>>,
pub diag: crate::runner::DiagnosticConfig,
pub pre_map_only: bool,
pub seq_type: SequencerType,
pub concurrency: usize,
pub rate: Option<f64>,
pub error_spec: String,
pub tries: Option<u32>,
pub error_rate_max: Option<f64>,
pub error_policy: Arc<crate::error_policy::ErrorPolicy>,
pub session_id: String,
pub exec_id: u64,
pub workload_name: String,
pub label_stack: Vec<(String, String)>,
pub session_component: Arc<RwLock<Component>>,
pub cadence_reporter: Arc<CadenceReporter>,
pub stop_handle: Arc<nmbrs_metrics::scheduler::StopHandle>,
pub observer: Arc<dyn crate::observer::RunObserver>,
pub scope_tree: Arc<crate::scope_tree::ScopeTree>,
pub schedule_spec: Arc<crate::scheduler::ScheduleSpec>,
pub current_parent_kernel: Option<Arc<crate::scope_kernel::ScopeKernel>>,
pub workload_source: Option<Arc<WorkloadSource>>,
pub checkpoint_writer: Option<Arc<crate::checkpoint::CheckpointWriter>>,
pub resume_plan: Arc<crate::checkpoint::ResumePlan>,
pub sqlite_reporter:
Arc<std::sync::Mutex<Option<nmbrs_metrics::reporters::sqlite::SqliteReporter>>>,
pub resource_pool: Arc<crate::resource_pool::ResourcePool>,
pub scene_tree_parent_id: crate::scene_tree::SceneNodeId,
pub scene_tree_path: Vec<crate::checkpoint::PathSegment>,
pub current_scope_idx: crate::scope_tree::ScopeNodeIdx,
pub workload_shell: Arc<crate::workload_shell::WorkloadShell>,
pub workload_stop_when: Vec<nmbrs_workload::model::StopConditionSpec>,
pub daemon_stop: Option<Arc<std::sync::atomic::AtomicBool>>,
pub optimize_objective: Option<String>,
pub optimize_objective_value: Option<f64>,
pub optimize_servo: Option<crate::optimize::servo::ServoSpec>,
}
pub struct WorkloadSource {
pub path: String,
pub text: String,
}
impl WorkloadSource {
pub fn locate(&self, needle: &str) -> Option<(usize, usize)> {
let idx = self.text.find(needle)?;
let prefix = &self.text[..idx];
let line = prefix.bytes().filter(|b| *b == b'\n').count() + 1;
let col = prefix.rfind('\n').map(|nl| idx - nl).unwrap_or(idx + 1);
Some((line, col))
}
}
pub(crate) fn enrich_with_yaml_location(ctx: &ExecCtx, needle: &str, err: String) -> String {
let Some(src) = ctx.workload_source.as_ref() else {
return err;
};
if err.starts_with(&format!("{}:", src.path)) {
return err;
}
let Some((line, col)) = src.locate(needle) else {
return err;
};
format!("{}:{line}:{col}: {err}", src.path)
}
fn enrich_outcome(
ctx: &ExecCtx,
needle: &str,
mut o: crate::phase_outcome::Outcome,
) -> crate::phase_outcome::Outcome {
if o.is_failure() {
let m = enrich_with_yaml_location(ctx, needle, o.reason.take().unwrap_or_default());
o = o.with_reason(m);
}
o
}
impl ExecCtx {
pub fn labels(&self) -> Labels {
let mut labels = Labels::of("session", &self.session_id)
.with("exec_id", self.exec_id.to_string())
.with("workload", &self.workload_name);
for (k, v) in &self.label_stack {
labels = labels.with(k, v);
}
labels
}
pub fn incremental_labels(&self) -> Labels {
let mut labels = Labels::empty();
for (k, v) in &self.label_stack {
labels = labels.with(k, v);
}
labels
}
pub fn push_label(&mut self, key: &str, value: &str) {
self.label_stack.push((key.to_string(), value.to_string()));
}
pub fn pop_label(&mut self) {
self.label_stack.pop();
}
}
pub fn execute_tree<'a>(
ctx: &'a mut ExecCtx,
nodes: &'a [ScenarioNode],
) -> std::pin::Pin<Box<dyn std::future::Future<Output = Result<(), String>> + Send + 'a>> {
Box::pin(async move {
let o = execute_tree_at(ctx, nodes, 0).await;
if o.is_failure() {
Err(o
.reason
.clone()
.unwrap_or_else(|| "scenario: a unit failed".to_string()))
} else {
Ok(())
}
})
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum ShellAction {
Continue,
Stop,
}
#[derive(Debug, Clone, Copy)]
struct ShellHandler;
impl ShellHandler {
fn scenario_default() -> Self {
ShellHandler
}
fn decide(&self, child: &crate::phase_outcome::Outcome) -> ShellAction {
if child.is_failure() {
ShellAction::Stop
} else {
ShellAction::Continue
}
}
}
#[allow(dead_code)]
pub(crate) trait ExecShell: Send + Sync {
fn run<'a>(
&'a self,
ctx: &'a mut ExecCtx,
) -> std::pin::Pin<
Box<dyn std::future::Future<Output = crate::phase_outcome::Outcome> + Send + 'a>,
>;
fn poll_stop(&self, ctx: &ExecCtx) -> Option<crate::session_signals::StopCause> {
use crate::session_signals::{self as sig, StopCause};
if sig::fault_stop_requested() {
Some(StopCause::Fault)
} else if sig::stop_requested() || ctx.workload_shell.should_stop() {
Some(StopCause::Interrupt)
} else {
None
}
}
fn shell_kind(&self) -> ShellKind;
}
#[allow(dead_code)] #[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum ShellKind {
Session,
Scenario,
Phase,
Stanza,
Op,
}
#[allow(dead_code)] trait CompositeShell: ExecShell {
fn handler(&self) -> &ShellHandler;
fn dispatch<'a>(
&'a self,
ctx: &'a mut ExecCtx,
nodes: &'a [ScenarioNode],
depth: usize,
) -> std::pin::Pin<
Box<dyn std::future::Future<Output = crate::phase_outcome::Outcome> + Send + 'a>,
>;
}
struct ScenarioShell<'n> {
handler: ShellHandler,
nodes: &'n [ScenarioNode],
depth: usize,
}
impl<'n> ScenarioShell<'n> {
fn scenario(nodes: &'n [ScenarioNode], depth: usize) -> Self {
Self {
handler: ShellHandler::scenario_default(),
nodes,
depth,
}
}
}
impl<'n> ExecShell for ScenarioShell<'n> {
fn run<'a>(
&'a self,
ctx: &'a mut ExecCtx,
) -> std::pin::Pin<
Box<dyn std::future::Future<Output = crate::phase_outcome::Outcome> + Send + 'a>,
> {
Box::pin(async move { self.dispatch(ctx, self.nodes, self.depth).await })
}
fn shell_kind(&self) -> ShellKind {
ShellKind::Scenario
}
}
impl<'n> CompositeShell for ScenarioShell<'n> {
fn handler(&self) -> &ShellHandler {
&self.handler
}
fn dispatch<'a>(
&'a self,
ctx: &'a mut ExecCtx,
nodes: &'a [ScenarioNode],
depth: usize,
) -> std::pin::Pin<
Box<dyn std::future::Future<Output = crate::phase_outcome::Outcome> + Send + 'a>,
> {
Box::pin(run_scenario_body(&self.handler, ctx, nodes, depth))
}
}
fn join_outcome(
res: Result<crate::phase_outcome::Outcome, tokio::task::JoinError>,
) -> (crate::phase_outcome::Outcome, Option<String>) {
use crate::phase_outcome::Outcome;
match res {
Err(join_err) => {
let msg = format!("concurrent task panicked: {join_err}");
(Outcome::failed().with_reason(msg.clone()), Some(msg))
}
Ok(outcome) => {
let reason = outcome.reason.clone();
(outcome, reason)
}
}
}
fn fold_child(
res: Result<crate::phase_outcome::Outcome, tokio::task::JoinError>,
handler: &ShellHandler,
first_failure: &mut Option<String>,
any_failed_reason: &mut Option<String>,
) {
let (child, reason) = join_outcome(res);
if child.is_failure() && any_failed_reason.is_none() {
*any_failed_reason = reason.clone();
}
if matches!(handler.decide(&child), ShellAction::Stop) && first_failure.is_none() {
*first_failure = reason;
}
}
fn fold_aggregate(
first_failure: Option<String>,
any_failed_reason: Option<String>,
should_stop: bool,
) -> crate::phase_outcome::Outcome {
use crate::phase_outcome::{Disposition, Outcome, Validity};
let cut_short = first_failure.is_some() || should_stop;
let disposition = if cut_short {
Disposition::Interrupted
} else {
Disposition::Completed
};
let validity = if first_failure.is_some() || any_failed_reason.is_some() {
Validity::Failed
} else {
Validity::Succeeded
};
let mut outcome = Outcome::new(disposition, validity);
if let Some(reason) = first_failure.or(any_failed_reason) {
outcome = outcome.with_reason(reason);
}
outcome
}
struct PhaseShell<'p> {
name: &'p str,
node_id: crate::scene_tree::SceneNodeId,
}
impl<'p> PhaseShell<'p> {
fn new(name: &'p str, node_id: crate::scene_tree::SceneNodeId) -> Self {
Self { name, node_id }
}
}
impl<'p> ExecShell for PhaseShell<'p> {
fn run<'a>(
&'a self,
ctx: &'a mut ExecCtx,
) -> std::pin::Pin<
Box<dyn std::future::Future<Output = crate::phase_outcome::Outcome> + Send + 'a>,
> {
Box::pin(async move { run_phase(ctx, self.name, self.node_id).await })
}
fn shell_kind(&self) -> ShellKind {
ShellKind::Phase
}
}
async fn run_phase_layered(
ctx: &mut ExecCtx,
name: &str,
node_id: crate::scene_tree::SceneNodeId,
) -> crate::phase_outcome::Outcome {
let leaf = PhaseShell::new(name, node_id);
match crate::wrappers::interval::for_phase(&ctx.phases, &ctx.workload_params, name) {
Some(spec) => {
crate::wrappers::interval::IntervalShell::new(&leaf, spec, name)
.run(ctx)
.await
}
None => leaf.run(ctx).await,
}
}
#[derive(Debug)]
#[allow(dead_code)] struct OpBodyPayload(Box<dyn crate::adapter::ResultBody>);
impl crate::phase_outcome::Payload for OpBodyPayload {}
#[allow(dead_code)] struct OpShell;
#[allow(dead_code)]
impl OpShell {
fn project(
res: Result<crate::adapter::OpResult, crate::adapter::ExecutionError>,
) -> crate::phase_outcome::Outcome {
use crate::phase_outcome::Outcome;
use std::sync::Arc;
match res {
Ok(r) if r.skipped => Outcome::skipped(),
Ok(r) => match r.body {
Some(body) => Outcome::completed().with_payload(Arc::new(OpBodyPayload(body))),
None => Outcome::completed(),
},
Err(e) => Outcome::completed_failed().with_reason(e.to_string()),
}
}
}
#[cfg(test)]
mod srd92_leaf_shell_tests {
use super::*;
use crate::phase_outcome::{Disposition, Validity};
#[test]
fn op_shell_projects_quadrants_and_payload() {
use crate::adapter::{AdapterError, ExecutionError, OpResult, TextBody};
let o = OpShell::project(Ok(OpResult::default()));
assert_eq!(o.disposition, Disposition::Completed);
assert_eq!(o.validity, Validity::Succeeded);
assert!(o.payload.is_none());
assert_eq!(
OpShell::project(Ok(OpResult::skipped())).disposition,
Disposition::Skipped
);
let err = ExecutionError::Op(AdapterError {
error_name: "test".into(),
message: "boom".into(),
retryable: false,
});
let o = OpShell::project(Err(err));
assert_eq!(o.disposition, Disposition::Completed);
assert_eq!(o.validity, Validity::Failed);
assert!(o.reason.as_deref().unwrap_or_default().contains("boom"));
let r = OpResult {
body: Some(Box::new(TextBody("hi".into()))),
skipped: false,
};
let o = OpShell::project(Ok(r));
assert_eq!(o.disposition, Disposition::Completed);
assert!(o.payload.is_some());
}
#[test]
fn fold_aggregate_two_latch_quadrants() {
let o = fold_aggregate(Some("boom".into()), Some("boom".into()), false);
assert_eq!(
(o.disposition, o.validity),
(Disposition::Interrupted, Validity::Failed)
);
assert_eq!(o.reason.as_deref(), Some("boom"));
let o = fold_aggregate(None, Some("soft".into()), false);
assert_eq!(
(o.disposition, o.validity),
(Disposition::Completed, Validity::Failed)
);
assert_eq!(o.reason.as_deref(), Some("soft"));
let o = fold_aggregate(None, None, true);
assert_eq!(
(o.disposition, o.validity),
(Disposition::Interrupted, Validity::Succeeded)
);
let o = fold_aggregate(None, None, false);
assert_eq!(
(o.disposition, o.validity),
(Disposition::Completed, Validity::Succeeded)
);
}
}
fn execute_tree_at<'a>(
ctx: &'a mut ExecCtx,
nodes: &'a [ScenarioNode],
depth: usize,
) -> std::pin::Pin<Box<dyn std::future::Future<Output = crate::phase_outcome::Outcome> + Send + 'a>>
{
Box::pin(async move {
let shell = ScenarioShell::scenario(nodes, depth);
shell.dispatch(ctx, nodes, depth).await
})
}
async fn run_scenario_body(
handler: &ShellHandler,
ctx: &mut ExecCtx,
nodes: &[ScenarioNode],
depth: usize,
) -> crate::phase_outcome::Outcome {
use crate::scheduler::ConcurrencyLimit;
let limit = ctx.schedule_spec.limit_at(depth);
let preview_useful = !matches!(limit, ConcurrencyLimit::Bounded(1));
if preview_useful {
let scheduled_phases: Vec<(usize, String)> = nodes
.iter()
.filter_map(|node| match node {
ScenarioNode::Phase(name) => Some(name.clone()),
_ => None,
})
.filter_map(|name| {
crate::scene_tree::current().and_then(|t| {
t.dfs_phases()
.find(|n| n.name == name)
.and_then(|n| n.seq)
.map(|seq| (seq, name.clone()))
})
})
.collect();
if !scheduled_phases.is_empty() {
let limit_disp = match limit {
ConcurrencyLimit::Bounded(n) => format!("limit={n}"),
ConcurrencyLimit::Unlimited => "limit=*".to_string(),
};
let total = crate::scene_tree::current()
.map(|t| t.total_phases())
.unwrap_or(scheduled_phases.len());
let listing: Vec<String> = scheduled_phases
.iter()
.map(|(seq, name)| format!("[{seq}/{total}] {name}"))
.collect();
crate::diag!(
crate::observer::LogLevel::Info,
"concurrent dispatch ({limit_disp}): {}",
listing.join(", ")
);
}
}
let sem: Option<Arc<tokio::sync::Semaphore>> = match limit {
ConcurrencyLimit::Bounded(n) => Some(Arc::new(tokio::sync::Semaphore::new(n as usize))),
ConcurrencyLimit::Unlimited => None,
};
let parent_scope_idx = ctx.current_scope_idx;
let child_scope_indices: Vec<crate::scope_tree::ScopeNodeIdx> =
ctx.scope_tree.nodes[parent_scope_idx].children.clone();
if child_scope_indices.len() != nodes.len() {
return crate::phase_outcome::Outcome::failed().with_reason(format!(
"scope-tree/scenario-tree drift: scope node {parent_scope_idx} has {} \
child scopes but the walker is dispatching {} sibling scenario nodes",
child_scope_indices.len(),
nodes.len(),
));
}
if ctx.workload_shell.should_stop() {
return crate::phase_outcome::Outcome::interrupted();
}
let daemon_flags: Vec<bool> = nodes
.iter()
.map(|n| match n {
ScenarioNode::Phase(name) => ctx.phases.get(name).map(|p| p.daemon).unwrap_or(false),
_ => false,
})
.collect();
let any_daemon = daemon_flags.iter().any(|&d| d);
let daemon_stop = Arc::new(std::sync::atomic::AtomicBool::new(false));
let mut daemon_set = tokio::task::JoinSet::new();
for (i, node) in nodes.iter().enumerate() {
if !daemon_flags[i] {
continue;
}
if let ScenarioNode::Phase(dname) = node
&& ctx
.phases
.get(dname.as_str())
.map(|p| p.for_each.is_none())
.unwrap_or(true)
{
let op_names: Vec<String> = ctx
.phases
.get(dname.as_str())
.map(|p| p.ops.iter().map(|op| op.name.clone()).collect())
.unwrap_or_default();
let phase_labels = canonical_phase_label(
&ctx.current_parent_kernel
.as_ref()
.map(|k| {
k.scope_coordinates()
.iter()
.rev()
.cloned()
.collect::<Vec<_>>()
})
.unwrap_or_default(),
);
let mut phase_path = ctx.scene_tree_path.clone();
phase_path.push(crate::checkpoint::PathSegment::Phase(dname.clone()));
let _ = push_phase_scene_node(
ctx.scene_tree_parent_id,
phase_path,
dname.clone(),
phase_labels,
op_names,
);
}
let node_scope_idx = child_scope_indices[i];
let node = node.clone();
let mut task_ctx = ctx.clone();
task_ctx.daemon_stop = Some(daemon_stop.clone());
daemon_set.spawn(crate::execution_context::propagate(async move {
execute_node(&mut task_ctx, &node, node_scope_idx, depth).await
}));
}
let mut set = tokio::task::JoinSet::new();
let mut first_failure: Option<String> = None;
let mut any_failed_reason: Option<String> = None;
use crate::child_source::{Child, ChildSource, CountedSource, Drive, select_drive};
let mut foreground = CountedSource::new(nodes.len());
debug_assert_eq!(
select_drive(foreground.realizability()),
Drive::BoundedSpawn
);
while let Some(Child::Node(i)) = foreground.poll_next() {
if daemon_flags[i] {
continue;
} let node = &nodes[i];
let node_scope_idx = child_scope_indices[i];
let permit = match sem.as_ref() {
Some(s) => match s.clone().acquire_owned().await {
Ok(p) => Some(p),
Err(e) => {
return crate::phase_outcome::Outcome::failed().with_reason(e.to_string());
}
},
None => None,
};
while let Some(res) = set.try_join_next() {
fold_child(res, handler, &mut first_failure, &mut any_failed_reason);
}
if first_failure.is_some() {
drop(permit);
break;
}
if ctx.workload_shell.should_stop() {
crate::diag!(
crate::observer::LogLevel::Debug,
"workload shell stopped — halting dispatch of remaining \
siblings at depth {depth}"
);
drop(permit);
break;
}
let node = node.clone();
let mut task_ctx = ctx.clone();
set.spawn(crate::execution_context::propagate(async move {
let _permit = permit;
execute_node(&mut task_ctx, &node, node_scope_idx, depth).await
}));
}
while let Some(res) = set.join_next().await {
fold_child(res, handler, &mut first_failure, &mut any_failed_reason);
}
if any_daemon {
daemon_stop.store(true, std::sync::atomic::Ordering::Relaxed);
while let Some(res) = daemon_set.join_next().await {
fold_child(res, handler, &mut first_failure, &mut any_failed_reason);
}
}
fold_aggregate(
first_failure,
any_failed_reason,
ctx.workload_shell.should_stop(),
)
}
fn push_phase_scene_node(
parent_id: crate::scene_tree::SceneNodeId,
yaml_path: Vec<crate::checkpoint::PathSegment>,
name: String,
labels: String,
op_names: Vec<String>,
) -> crate::scene_tree::SceneNodeId {
let mut id: crate::scene_tree::SceneNodeId = 0;
crate::scene_tree::with_global_mut(|t| {
id = t.push(parent_id, crate::scene_tree::NodeKind::Phase, name, labels);
t.set_phase_op_names(id, op_names);
t.set_yaml_path(id, yaml_path);
});
id
}
fn push_scope_scene_node(
parent_id: crate::scene_tree::SceneNodeId,
yaml_path: Vec<crate::checkpoint::PathSegment>,
header: String,
own_names: Vec<String>,
) -> crate::scene_tree::SceneNodeId {
let mut id: crate::scene_tree::SceneNodeId = 0;
crate::scene_tree::with_global_mut(|t| {
id = t.push(
parent_id,
crate::scene_tree::NodeKind::Scope,
header,
String::new(),
);
if !own_names.is_empty() {
t.set_own_names(id, own_names);
}
t.set_yaml_path(id, yaml_path);
});
id
}
fn format_iter_label(bindings: &[(String, polydat::ast::Value)]) -> String {
bindings
.iter()
.map(|(k, v)| format!("{k}={}", v.to_display_string()))
.collect::<Vec<_>>()
.join(", ")
}
fn canonical_phase_label(parent_coords: &[ScopeCoord]) -> String {
let leaf_first: Vec<_> = parent_coords.iter().rev().cloned().collect();
format_scope_coordinate_path(&leaf_first)
}
fn do_loop_own_names(
ctx: &ExecCtx,
condition: &str,
counter: Option<&str>,
invert: bool,
) -> Vec<String> {
let idx = ctx
.scope_tree
.iter_dfs()
.find_map(|(idx, n)| match &n.kind {
crate::scope_tree::ScopeKind::DoWhile {
condition: c,
counter: ct,
} if !invert && c == condition && ct.as_deref() == counter => Some(idx),
crate::scope_tree::ScopeKind::DoUntil {
condition: c,
counter: ct,
} if invert && c == condition && ct.as_deref() == counter => Some(idx),
_ => None,
});
idx.and_then(|i| ctx.scope_tree.nodes[i].cached_kernel.get().cloned())
.map(|k| {
k.program()
.own_output_names()
.into_iter()
.map(String::from)
.collect::<Vec<String>>()
})
.unwrap_or_default()
}
fn effective_parent_kernel(
ctx: &ExecCtx,
scope_idx: usize,
) -> Option<std::sync::Arc<crate::scope_kernel::ScopeKernel>> {
ctx.current_parent_kernel
.clone()
.or_else(|| ctx.scope_tree.nearest_installed_ancestor_kernel(scope_idx))
}
fn subtree_has_active_phase(
node: &ScenarioNode,
pattern: &crate::phase_filter::PhasePattern,
) -> bool {
match node {
ScenarioNode::Phase(name) => pattern.is_match(name),
ScenarioNode::Comprehension { children, .. }
| ScenarioNode::DoWhile { children, .. }
| ScenarioNode::DoUntil { children, .. }
| ScenarioNode::IncludedScenario { children, .. } => children
.iter()
.any(|c| subtree_has_active_phase(c, pattern)),
_ => false,
}
}
fn execute_node<'a>(
ctx: &'a mut ExecCtx,
node: &'a ScenarioNode,
node_scope_idx: crate::scope_tree::ScopeNodeIdx,
depth: usize,
) -> std::pin::Pin<Box<dyn std::future::Future<Output = crate::phase_outcome::Outcome> + Send + 'a>>
{
Box::pin(async move {
use crate::checkpoint::PathSegment;
if let Some(pat) = ctx.phase_filter.clone() {
let is_scope = !matches!(node, ScenarioNode::Phase(_));
if is_scope && !subtree_has_active_phase(node, &pat) {
crate::diag!(
crate::observer::LogLevel::Debug,
"phases=<filter>: eliding scope (no descendant phase matches)"
);
return crate::phase_outcome::Outcome::skipped();
}
}
match node {
ScenarioNode::Phase(name) => {
let phase_fe = ctx
.phases
.get(name.as_str())
.and_then(|p| p.for_each.clone());
let op_names: Vec<String> = ctx
.phases
.get(name.as_str())
.map(|p| p.ops.iter().map(|op| op.name.clone()).collect())
.unwrap_or_default();
if let Some(spec) = phase_fe {
let scope_idx = match ctx.scope_tree.phase_node_by_name(name) {
Some(v) => v,
None => {
return crate::phase_outcome::Outcome::failed().with_reason(format!(
"phase '{name}' for_each '{spec}': no matching scope-tree entry."
));
}
};
let canonical =
match ctx.scope_tree.nodes[scope_idx].cached_kernel.get().cloned() {
Some(v) => v,
None => return crate::phase_outcome::Outcome::failed().with_reason(
format!(
"phase '{name}' for_each '{spec}': scope at index {scope_idx} \
has no installed phase-for_each kernel."
),
),
};
let parent = match effective_parent_kernel(ctx, scope_idx) {
Some(v) => v,
None => {
return crate::phase_outcome::Outcome::failed().with_reason(format!(
"phase '{name}' for_each '{spec}': no installed ancestor kernel."
));
}
};
let comprehension =
match polydat::iteration::comprehension::spec::parse_inline(&spec) {
Ok(v) => v,
Err(e) => {
return crate::phase_outcome::Outcome::failed()
.with_reason(format!("phase '{name}' for_each '{spec}': {e}"));
}
};
let iter_vars: Vec<String> = comprehension
.coordinate_specs()
.into_iter()
.map(|(v, _)| v)
.collect();
let var_label = iter_vars.join(", ");
let needle = spec.clone();
let parent_coords = ctx
.current_parent_kernel
.as_ref()
.map(|k| {
k.scope_coordinates()
.iter()
.rev()
.cloned()
.collect::<Vec<_>>()
})
.unwrap_or_default();
let optimize_block = ctx
.phases
.get(name.as_str())
.and_then(|p| p.optimize.clone());
let continuous_axes = optimize_block
.as_ref()
.and_then(|_| continuous_axis_intervals(&comprehension));
let steps = if continuous_axes.is_some() {
Vec::new()
} else {
match runtime_iterate(
ctx,
&canonical,
&parent,
&parent_coords,
&comprehension,
) {
Ok(v) => v,
Err(e) => {
return crate::phase_outcome::Outcome::failed()
.with_reason(enrich_with_yaml_location(ctx, &needle, e));
}
}
};
let mut scope_path = ctx.scene_tree_path.clone();
scope_path.push(PathSegment::ForEach {
var: var_label.clone(),
});
let header = if let Some(b) = &optimize_block {
let axes = if continuous_axes.is_some() {
var_label.clone()
} else {
steps
.first()
.map(|s| {
s.bindings
.iter()
.map(|(k, _)| k.as_str())
.collect::<Vec<_>>()
.join(", ")
})
.unwrap_or_default()
};
format!(
"search · {} · maximize {} · {{{axes}}} · ≤{} evals",
b.method, b.objective, b.max_evals
)
} else {
format!(
"phase.for_each {var_label} in [{}]",
steps
.iter()
.filter_map(|s| s
.bindings
.first()
.map(|(_, v)| v.to_display_string()))
.collect::<Vec<_>>()
.join(", ")
)
};
let scope_id = push_scope_scene_node(
ctx.scene_tree_parent_id,
scope_path.clone(),
header,
Vec::new(),
);
let saved_parent = ctx.scene_tree_parent_id;
let saved_path = std::mem::replace(&mut ctx.scene_tree_path, scope_path);
ctx.scene_tree_parent_id = scope_id;
let phase_path_for_iters: Vec<PathSegment> = {
let mut p = ctx.scene_tree_path.clone();
p.push(PathSegment::Phase(name.clone()));
p
};
let res = if let Some(b) = optimize_block {
let (space, coord_eval) = if let Some(intervals) = continuous_axes {
let names = iter_vars.clone();
(
search_space_continuous(&names, &intervals),
CoordEval::Synthesized {
axis_names: names,
canonical: canonical.clone(),
parent: parent.clone(),
parent_coords: parent_coords.clone(),
},
)
} else {
let space = search_space_from_steps(&steps);
let index = index_steps(&steps);
(space, CoordEval::Enumerated { steps, index })
};
dispatch_optimization(
ctx,
space,
coord_eval,
b,
name,
depth + 1,
Some((name.clone(), op_names, phase_path_for_iters)),
)
.await
} else {
let phase_continue_if = ctx
.phases
.get(name.as_str())
.and_then(|p| p.continue_if.clone());
let coord_sample =
steps.first().map(|s| s.bindings.as_slice()).unwrap_or(&[]);
let gate = match resolve_continue_if(
phase_continue_if,
&parent,
coord_sample,
ctx.strict,
) {
Ok(g) => g,
Err(e) => {
return crate::phase_outcome::Outcome::failed().with_reason(e);
}
};
dispatch_comprehension(
ctx,
steps,
TerminalAction::Phase(name),
depth + 1,
false,
"for_each",
Some((name.clone(), op_names, phase_path_for_iters)),
gate,
)
.await
};
ctx.scene_tree_parent_id = saved_parent;
ctx.scene_tree_path = saved_path;
return enrich_outcome(ctx, &needle, res); } else {
let phase_labels = canonical_phase_label(
&ctx.current_parent_kernel
.as_ref()
.map(|k| {
k.scope_coordinates()
.iter()
.rev()
.cloned()
.collect::<Vec<_>>()
})
.unwrap_or_default(),
);
let mut phase_path = ctx.scene_tree_path.clone();
phase_path.push(PathSegment::Phase(name.clone()));
let phase_labels_for_gate = phase_labels.clone();
let phase_node_id = push_phase_scene_node(
ctx.scene_tree_parent_id,
phase_path,
name.clone(),
phase_labels,
op_names,
);
let phase_labels = phase_labels_for_gate;
let is_prereq_class = ctx
.phases
.get(name)
.and_then(|p| p.checkpoint.as_ref())
.map(|c| c.idempotent)
.unwrap_or(false);
let prereq_exempt = ctx.phase_filter.is_some() && is_prereq_class;
let pattern_active = ctx
.phase_filter
.as_ref()
.map(|pat| pat.is_match(name) || prereq_exempt)
.unwrap_or(true);
if prereq_exempt
&& !ctx
.phase_filter
.as_ref()
.map(|pat| pat.is_match(name))
.unwrap_or(true)
{
crate::diag!(
crate::observer::LogLevel::Info,
"phases=<filter>: phase '{name}' kept as an \
idempotent prerequisite (checkpoint: idempotent) \
— resume/refine provenance decides skip-or-run"
);
}
let explicitly_selected = ctx
.phase_filter
.as_ref()
.map(|pat| pat.is_match(name))
.unwrap_or(false);
let refine_missing_skip = ctx
.refine_plan
.as_ref()
.filter(|p| p.scope == crate::refine_plan::RefineScope::Missing)
.map(|p| p.is_completed(name, &phase_labels))
.unwrap_or(false)
&& !explicitly_selected
&& !is_prereq_class;
if !pattern_active {
crate::diag!(
crate::observer::LogLevel::Debug,
"phases=<filter>: skipping phase '{name}' (does not match)"
);
if !ctx.pre_map_only {
crate::scene_tree::with_global_mut(|t| {
t.set_phase_running_at(phase_node_id, 0);
t.set_phase_completed_at(phase_node_id, 0.0);
});
ctx.observer.phase_starting(
phase_node_id,
name,
&phase_labels,
0,
0,
0,
);
ctx.observer
.phase_completed(phase_node_id, name, &phase_labels, 0.0);
crate::phase_end_triggers::fire_phase_completed(
name,
&phase_labels,
0.0,
);
}
} else if ctx.diag.depth < crate::runner::ExecDepth::Dispenser {
if !ctx.pre_map_only {
crate::scene_tree::with_global_mut(|t| {
t.set_phase_running_at(phase_node_id, 0);
t.set_phase_completed_at(phase_node_id, 0.0);
});
ctx.observer.phase_starting(
phase_node_id,
name,
&phase_labels,
0,
0,
0,
);
ctx.observer
.phase_completed(phase_node_id, name, &phase_labels, 0.0);
crate::phase_end_triggers::fire_phase_completed(
name,
&phase_labels,
0.0,
);
}
} else if refine_missing_skip {
crate::diag!(
crate::observer::LogLevel::Info,
"refine: skipping phase '{name}' [{phase_labels}] \
(prior completed outcome)"
);
crate::scene_tree::with_global_mut(|t| {
t.set_phase_running_at(phase_node_id, 0);
t.set_phase_completed_at(phase_node_id, 0.0);
});
ctx.observer
.phase_starting(phase_node_id, name, &phase_labels, 0, 0, 0);
ctx.observer
.phase_completed(phase_node_id, name, &phase_labels, 0.0);
crate::phase_end_triggers::fire_phase_completed(name, &phase_labels, 0.0);
} else if ctx.diag.depth >= crate::runner::ExecDepth::Dispenser {
let __o = run_phase_layered(ctx, name, phase_node_id).await;
if __o.is_failure() {
return __o; }
}
}
}
ScenarioNode::Comprehension {
comprehension,
children,
continue_if,
anchor: _,
} => {
let label = crate::scope_tree::ScopeKind::Comprehension {
comprehension: comprehension.clone(),
}
.label();
let scope_idx = node_scope_idx;
debug_assert!(
matches!(
&ctx.scope_tree.nodes[scope_idx].kind,
crate::scope_tree::ScopeKind::Comprehension { comprehension: c }
if c == comprehension
),
"{label}: positional scope node {scope_idx} is not the matching comprehension",
);
let canonical = match ctx.scope_tree.nodes[scope_idx].cached_kernel.get().cloned() {
Some(v) => v,
None => {
return crate::phase_outcome::Outcome::failed().with_reason(format!(
"{label}: scope at index {scope_idx} has no installed kernel.",
));
}
};
let parent = match effective_parent_kernel(ctx, scope_idx) {
Some(v) => v,
None => {
return crate::phase_outcome::Outcome::failed()
.with_reason(format!("{label}: no installed ancestor kernel.",));
}
};
let own_names: Vec<String> = canonical
.program()
.own_output_names()
.into_iter()
.map(String::from)
.collect();
let parent_coords = ctx
.current_parent_kernel
.as_ref()
.map(|k| {
k.scope_coordinates()
.iter()
.rev()
.cloned()
.collect::<Vec<_>>()
})
.unwrap_or_default();
let coord_names = comprehension.coordinate_names();
let needle = coord_names.first().cloned().unwrap_or_default();
let steps = match runtime_iterate(
ctx,
&canonical,
&parent,
&parent_coords,
comprehension,
) {
Ok(v) => v,
Err(e) => {
return crate::phase_outcome::Outcome::failed()
.with_reason(enrich_with_yaml_location(ctx, &needle, e));
}
};
let mut scope_path = ctx.scene_tree_path.clone();
let (header, kind) = if coord_names.len() == 1 {
let var = &coord_names[0];
scope_path.push(PathSegment::ForEach { var: var.clone() });
(format!("each {var}"), "for_each")
} else {
scope_path.push(PathSegment::ForCombinations {
vars: coord_names.clone(),
});
let summary = coord_names.join(", ");
(format!("each {summary}"), "for_combinations")
};
let scope_id = push_scope_scene_node(
ctx.scene_tree_parent_id,
scope_path.clone(),
header,
own_names.clone(),
);
let coord_sample = steps.first().map(|s| s.bindings.as_slice()).unwrap_or(&[]);
let continue_if_gate = match resolve_continue_if(
continue_if.clone(),
&parent,
coord_sample,
ctx.strict,
) {
Ok(g) => g,
Err(e) => {
return crate::phase_outcome::Outcome::failed()
.with_reason(enrich_with_yaml_location(ctx, &needle, e));
}
};
let saved_parent = ctx.scene_tree_parent_id;
let saved_path = std::mem::replace(&mut ctx.scene_tree_path, scope_path);
let saved_scope_idx = ctx.current_scope_idx;
ctx.scene_tree_parent_id = scope_id;
ctx.current_scope_idx = scope_idx;
let res = dispatch_comprehension(
ctx,
steps,
TerminalAction::Children(children),
depth + 1,
false,
kind,
None,
continue_if_gate,
)
.await;
ctx.scene_tree_parent_id = saved_parent;
ctx.scene_tree_path = saved_path;
ctx.current_scope_idx = saved_scope_idx;
return enrich_outcome(ctx, &needle, res); }
ScenarioNode::IncludedScenario { name, children } => {
crate::diag!(
crate::observer::LogLevel::Debug,
"include scenario '{name}' ({} children)",
children.len()
);
let mut scope_path = ctx.scene_tree_path.clone();
scope_path.push(PathSegment::ScenarioInclude(name.clone()));
let scope_id = push_scope_scene_node(
ctx.scene_tree_parent_id,
scope_path.clone(),
format!("scenario '{name}'"),
Vec::new(),
);
let saved_parent = ctx.scene_tree_parent_id;
let saved_path = std::mem::replace(&mut ctx.scene_tree_path, scope_path);
let saved_scope_idx = ctx.current_scope_idx;
ctx.scene_tree_parent_id = scope_id;
ctx.current_scope_idx = node_scope_idx;
let res = execute_tree_at(ctx, children, depth + 1).await;
ctx.scene_tree_parent_id = saved_parent;
ctx.scene_tree_path = saved_path;
ctx.current_scope_idx = saved_scope_idx;
if res.is_failure() {
return res;
} }
ScenarioNode::DoWhile {
condition,
counter,
children,
} => {
crate::diag!(
crate::observer::LogLevel::Debug,
"=== do_while: {condition} ==="
);
let mut scope_path = ctx.scene_tree_path.clone();
scope_path.push(PathSegment::DoWhile {
counter: counter.clone(),
});
let own_names = do_loop_own_names(ctx, condition, counter.as_deref(), false);
let scope_id = push_scope_scene_node(
ctx.scene_tree_parent_id,
scope_path.clone(),
format!("do_while {condition}"),
own_names,
);
let saved_parent = ctx.scene_tree_parent_id;
let saved_path = std::mem::replace(&mut ctx.scene_tree_path, scope_path);
let saved_scope_idx = ctx.current_scope_idx;
ctx.scene_tree_parent_id = scope_id;
ctx.current_scope_idx = node_scope_idx; fire_scope_lifecycle(
ctx,
crate::lifecycle::EventType::ScopeStart,
&format!("do_while {condition}"),
depth,
);
let r = run_do_loop(
ctx,
condition,
counter.as_deref(),
false,
children,
depth + 1,
)
.await;
fire_scope_lifecycle(
ctx,
crate::lifecycle::EventType::ScopeEnd,
&format!("do_while {condition}"),
depth,
);
ctx.scene_tree_parent_id = saved_parent;
ctx.scene_tree_path = saved_path;
ctx.current_scope_idx = saved_scope_idx;
if let Err(e) = r {
return crate::phase_outcome::Outcome::failed().with_reason(e);
}
}
ScenarioNode::DoUntil {
condition,
counter,
children,
} => {
crate::diag!(
crate::observer::LogLevel::Debug,
"=== do_until: {condition} ==="
);
let mut scope_path = ctx.scene_tree_path.clone();
scope_path.push(PathSegment::DoUntil {
counter: counter.clone(),
});
let own_names = do_loop_own_names(ctx, condition, counter.as_deref(), true);
let scope_id = push_scope_scene_node(
ctx.scene_tree_parent_id,
scope_path.clone(),
format!("do_until {condition}"),
own_names,
);
let saved_parent = ctx.scene_tree_parent_id;
let saved_path = std::mem::replace(&mut ctx.scene_tree_path, scope_path);
let saved_scope_idx = ctx.current_scope_idx;
ctx.scene_tree_parent_id = scope_id;
ctx.current_scope_idx = node_scope_idx; fire_scope_lifecycle(
ctx,
crate::lifecycle::EventType::ScopeStart,
&format!("do_until {condition}"),
depth,
);
let r = run_do_loop(
ctx,
condition,
counter.as_deref(),
true,
children,
depth + 1,
)
.await;
fire_scope_lifecycle(
ctx,
crate::lifecycle::EventType::ScopeEnd,
&format!("do_until {condition}"),
depth,
);
ctx.scene_tree_parent_id = saved_parent;
ctx.scene_tree_path = saved_path;
ctx.current_scope_idx = saved_scope_idx;
if let Err(e) = r {
return crate::phase_outcome::Outcome::failed().with_reason(e);
}
}
ScenarioNode::Bindings { source, children } => {
let scope_idx = node_scope_idx;
debug_assert!(
matches!(
&ctx.scope_tree.nodes[scope_idx].kind,
crate::scope_tree::ScopeKind::Bindings { source: s } if s == source
),
"bindings: positional scope node {scope_idx} is not the matching bindings",
);
let installed = match ctx.scope_tree.nodes[scope_idx].cached_kernel.get().cloned() {
Some(v) => v,
None => {
return crate::phase_outcome::Outcome::failed().with_reason(format!(
"bindings scope at index {scope_idx} has no installed \
kernel — install-spec walker missed this node",
));
}
};
let one_line = source
.lines()
.map(str::trim)
.find(|l| !l.is_empty())
.unwrap_or("");
crate::diag!(
crate::observer::LogLevel::Debug,
"bindings: {one_line} ({} children)",
children.len()
);
let chained = match ctx.current_parent_kernel.as_ref() {
Some(parent) => match installed.bind_under(parent.kernel(), &[]) {
Ok(v) => std::sync::Arc::new(v),
Err(e) => {
return crate::phase_outcome::Outcome::failed().with_reason(format!(
"bindings scope at index {scope_idx}: \
chain to current parent kernel: {e}",
));
}
},
None => installed,
};
let label = crate::scope_tree::bindings_label(source);
let mut scope_path = ctx.scene_tree_path.clone();
scope_path.push(PathSegment::ScenarioInclude(label.clone()));
let scope_id = push_scope_scene_node(
ctx.scene_tree_parent_id,
scope_path.clone(),
label,
Vec::new(),
);
let prior_parent = ctx.current_parent_kernel.take();
ctx.current_parent_kernel = Some(chained);
let saved_scene_parent = ctx.scene_tree_parent_id;
let saved_scene_path = std::mem::replace(&mut ctx.scene_tree_path, scope_path);
let saved_scope_idx = ctx.current_scope_idx;
ctx.scene_tree_parent_id = scope_id;
ctx.current_scope_idx = scope_idx;
let res = execute_tree_at(ctx, children, depth + 1).await;
ctx.scene_tree_parent_id = saved_scene_parent;
ctx.scene_tree_path = saved_scene_path;
ctx.current_scope_idx = saved_scope_idx;
ctx.current_parent_kernel = prior_parent;
if res.is_failure() {
return res;
} }
}
crate::phase_outcome::Outcome::completed()
})
}
#[derive(Clone, Debug)]
pub struct IterationStep {
pub bindings: Vec<(String, polydat::ast::Value)>,
pub bound_kernel: std::sync::Arc<crate::scope_kernel::ScopeKernel>,
pub coord_path: Vec<polydat::kernel::ScopeCoord>,
}
fn runtime_iterate(
ctx: &ExecCtx,
canonical: &std::sync::Arc<crate::scope_kernel::ScopeKernel>,
parent: &std::sync::Arc<crate::scope_kernel::ScopeKernel>,
parent_coords: &[ScopeCoord],
comprehension: &polydat::iteration::comprehension::Comprehension,
) -> Result<Vec<IterationStep>, String> {
use polydat::iteration::comprehension::{NoneRead, evaluate_for_iteration_with_none_reads};
use polydat::kernel::ScopeCoord;
let (evaluated, none_reads) =
evaluate_for_iteration_with_none_reads(comprehension, parent.as_ref())
.map_err(|e| e.to_string())?;
let mut reported = std::collections::HashSet::new();
for (i, clause) in evaluated.clauses.iter().enumerate() {
if let Some(NoneRead::Unbound(name)) = none_reads
.clause(i)
.iter()
.find(|read| matches!(read, NoneRead::Unbound(_)))
{
let spec = clause.source.as_deref().unwrap_or("?");
let msg = format!(
"unresolved placeholder '{{{name}}}' in for_each clause '{var} in {spec}': \
no workload param, outer iter-var, or inherited binding is named '{name}'",
var = clause.var
);
if ctx.strict {
return Err(format!("strict: {msg}"));
}
if !ctx.pre_map_only {
crate::diag!(
crate::observer::LogLevel::Warn,
"warning: {msg}; the clause yields no values"
);
}
reported.insert(i);
}
}
for (i, clause) in evaluated.clauses.iter().enumerate() {
if clause.evaluations == 0 || clause.values > 0 || reported.contains(&i) {
continue;
}
let label = match &clause.source {
Some(spec) => format!("for_each clause '{var} in {spec}'", var = clause.var),
None => format!("for_each clause '{var}'", var = clause.var),
};
let msg = format!("{label}: produced no values");
if ctx.strict {
return Err(format!("strict: {msg}"));
}
if !ctx.pre_map_only {
crate::diag!(crate::observer::LogLevel::Warn, "warning: {msg}");
}
}
let tuples = evaluated.tuples;
let mut steps = Vec::with_capacity(tuples.len());
for tuple in tuples {
let bound_kernel = std::sync::Arc::new(canonical.bind_under(parent.kernel(), &tuple)?);
let mut coord_path = parent_coords.to_vec();
coord_path.push(ScopeCoord::from(tuple.iter().cloned()));
steps.push(IterationStep {
bindings: tuple,
bound_kernel,
coord_path,
});
}
Ok(steps)
}
#[derive(Clone)]
pub enum TerminalAction<'a> {
Children(&'a [ScenarioNode]),
Phase(&'a str),
}
#[derive(Clone)]
enum OwnedTerminal {
Children(std::sync::Arc<Vec<ScenarioNode>>),
Phase(String),
}
impl OwnedTerminal {
fn borrow(&self) -> TerminalAction<'_> {
match self {
OwnedTerminal::Children(arc) => TerminalAction::Children(arc.as_slice()),
OwnedTerminal::Phase(name) => TerminalAction::Phase(name.as_str()),
}
}
}
struct ContinueIfGate {
spec: nmbrs_workload::model::ContinueIfSpec,
gate_canonical: std::sync::Arc<crate::scope_kernel::ScopeKernel>,
parent: std::sync::Arc<crate::scope_kernel::ScopeKernel>,
}
fn resolve_continue_if(
spec: Option<nmbrs_workload::model::ContinueIfSpec>,
parent: &std::sync::Arc<crate::scope_kernel::ScopeKernel>,
coord_sample: &[(String, polydat::ast::Value)],
strict: bool,
) -> Result<Option<ContinueIfGate>, String> {
match spec {
None => Ok(None),
Some(spec) => {
let gate_canonical =
crate::stop_conditions::compile_continue_if(&spec.when, coord_sample, strict)?;
Ok(Some(ContinueIfGate {
spec,
gate_canonical,
parent: parent.clone(),
}))
}
}
}
fn dispatch_comprehension<'a>(
ctx: &'a mut ExecCtx,
steps: Vec<IterationStep>,
terminal: TerminalAction<'a>,
depth: usize,
sequential_only: bool,
kind: &'static str,
phase_terminal_meta: Option<(String, Vec<String>, Vec<crate::checkpoint::PathSegment>)>,
continue_if: Option<ContinueIfGate>,
) -> std::pin::Pin<Box<dyn std::future::Future<Output = crate::phase_outcome::Outcome> + Send + 'a>>
{
use crate::scheduler::ConcurrencyLimit;
Box::pin(async move {
if steps.is_empty() {
return crate::phase_outcome::Outcome::completed();
}
let limit = if sequential_only {
ConcurrencyLimit::Bounded(1)
} else {
ctx.schedule_spec.limit_at(depth)
};
let sem: Option<std::sync::Arc<tokio::sync::Semaphore>> = match limit {
ConcurrencyLimit::Bounded(n) => {
Some(std::sync::Arc::new(tokio::sync::Semaphore::new(n as usize)))
}
ConcurrencyLimit::Unlimited => None,
};
let owned_terminal = match &terminal {
TerminalAction::Children(c) => OwnedTerminal::Children(std::sync::Arc::new(c.to_vec())),
TerminalAction::Phase(name) => OwnedTerminal::Phase(name.to_string()),
};
let inner_own_names: Vec<String> = if phase_terminal_meta.is_none() {
crate::scene_tree::current()
.and_then(|t| {
t.nodes
.get(ctx.scene_tree_parent_id)
.map(|n| n.own_names.clone())
})
.unwrap_or_default()
} else {
Vec::new()
};
let mut set = tokio::task::JoinSet::new();
use crate::child_source::{Child, ChildSource, CountedSource, Drive, select_drive};
let mut csrc = CountedSource::new(steps.len());
debug_assert_eq!(select_drive(csrc.realizability()), Drive::BoundedSpawn);
let mut steps_iter = steps.into_iter();
let mut continue_if_halt: Option<String> = None;
while let Some(Child::Node(_)) = csrc.poll_next() {
let step = steps_iter
.next()
.expect("CountedSource length matches steps");
if let Some(gate) = continue_if.as_ref() {
if ctx.diag.depth >= crate::runner::ExecDepth::Op && !ctx.pre_map_only {
match crate::stop_conditions::eval_continue_if(
&gate.gate_canonical,
&gate.parent,
&step.bindings,
) {
Ok(true) => {} Ok(false) => {
let coord = format_iter_label(&step.bindings);
let reason =
format!("continue_if: {} — halted at {coord}", gate.spec.when);
crate::diag!(
crate::observer::LogLevel::Info,
"sweep halted (continue_if): `{}` false at {coord}",
gate.spec.when
);
crate::session_signals::request_graceful_stop();
if gate
.spec
.each
.iter()
.any(|l| matches!(l, nmbrs_workload::model::ScopeLevel::Workload))
{
ctx.workload_shell.request_stop(
crate::phase_outcome::Outcome::interrupted()
.with_reason(reason.clone()),
reason.clone(),
);
}
continue_if_halt = Some(reason);
break;
}
Err(e) => {
return crate::phase_outcome::Outcome::failed().with_reason(e);
}
}
}
}
let permit = match sem.as_ref() {
Some(s) => match s.clone().acquire_owned().await {
Ok(p) => Some(p),
Err(e) => {
return crate::phase_outcome::Outcome::failed().with_reason(e.to_string());
}
},
None => None,
};
if ctx.workload_shell.should_stop() {
drop(permit);
break;
}
let per_iter_scene_id = match &phase_terminal_meta {
Some((phase_name, op_names, phase_yaml_path)) => {
let iter_scope = push_scope_scene_node(
ctx.scene_tree_parent_id,
ctx.scene_tree_path.clone(),
format_iter_label(&step.bindings),
Vec::new(),
);
let labels = canonical_phase_label(&step.coord_path);
push_phase_scene_node(
iter_scope,
phase_yaml_path.clone(),
phase_name.clone(),
labels,
op_names.clone(),
)
}
None => push_scope_scene_node(
ctx.scene_tree_parent_id,
ctx.scene_tree_path.clone(),
format_iter_label(&step.bindings),
inner_own_names.clone(),
),
};
let mut task_ctx = ctx.clone();
task_ctx.scene_tree_parent_id = per_iter_scene_id;
let owned_terminal = owned_terminal.clone();
set.spawn(crate::execution_context::propagate(async move {
let _permit = permit;
let terminal = owned_terminal.borrow();
run_one_iteration(&mut task_ctx, &step, &terminal, depth, kind).await
}));
}
let mut first_err: Option<String> = None;
while let Some(res) = set.join_next().await {
match res {
Err(join_err) => {
if first_err.is_none() {
first_err = Some(format!(
"concurrent comprehension iteration panicked: {join_err}"
));
}
}
Ok(o) => {
if o.is_failure() && first_err.is_none() {
first_err = Some(
o.reason
.unwrap_or_else(|| "comprehension iteration failed".to_string()),
);
}
}
}
}
if let Some(reason) = continue_if_halt {
if first_err.is_none() {
return crate::phase_outcome::Outcome::interrupted().with_reason(reason);
}
}
fold_aggregate(None, first_err, ctx.workload_shell.should_stop())
})
}
fn dispatch_optimization<'a>(
ctx: &'a mut ExecCtx,
space: crate::optimize::SearchSpace,
coord_eval: CoordEval,
block: nmbrs_workload::model::OptimizeBlock,
phase_name: &'a str,
depth: usize,
phase_meta: Option<(String, Vec<String>, Vec<crate::checkpoint::PathSegment>)>,
) -> std::pin::Pin<Box<dyn std::future::Future<Output = crate::phase_outcome::Outcome> + Send + 'a>>
{
Box::pin(async move {
let __opt_result: Result<(), String> = async move {
if space.dims() == 0 {
return Ok(());
}
if ctx.diag.depth < crate::runner::ExecDepth::Op {
if let Some(step) = coord_eval.representative(&space) {
return run_one_eval(ctx, &step, phase_name, depth, &phase_meta).await;
}
return Ok(());
}
let phase_concurrency = ctx.phases.get(phase_name).and_then(|p| p.concurrency.clone());
let phase_has_rate = ctx.phases.get(phase_name).is_some_and(|p| p.rate.is_some());
let control_axes = classify_control_axes(
&space,
&block.servo,
phase_concurrency.as_deref(),
phase_has_rate,
)?;
if !control_axes.is_empty() {
if control_axes.len() == space.dims() {
return run_control_search(
ctx, space, coord_eval, control_axes, block, phase_name, depth, phase_meta,
)
.await;
}
return run_hybrid_search(
ctx, space, coord_eval, control_axes, block, phase_name, depth, phase_meta,
)
.await;
}
let mut params = crate::optimize::OptimizerParams::new();
for (k, v) in &block.params {
params = params.with(k.clone(), *v);
}
let optimizer = crate::optimize::by_name(&block.method, ¶ms).ok_or_else(|| {
format!("phase '{phase_name}': unknown optimizer method '{}'", block.method)
})?;
let budget = crate::optimize::Budget::seeded(block.max_evals, block.seed);
let lex: Box<dyn crate::optimize::PullSource> =
Box::new(crate::optimize::LexSource::new(&space));
let mut src = optimizer.coordinate_source(&space, &budget, lex);
ctx.optimize_objective = Some(crate::scope::objective_wire(&block.objective).to_string());
let mut best_value = f64::NEG_INFINITY;
let mut best_coord: Option<crate::optimize::Coord> = None;
let mut evals = 0usize;
let mut batch = source_next(&mut src, &[]);
'outer: while let Some(coords) = batch.take() {
let mut evaluated: Vec<(crate::optimize::Coord, f64)> = Vec::new();
for coord in coords {
if evals >= block.max_evals {
break 'outer;
}
let Some(step) = coord_eval.materialize(&coord) else {
crate::diag!(crate::observer::LogLevel::Warn,
"optimizer '{method}' on '{phase_name}': coordinate [{key}] is not in the \
enumerable grid — skipping",
method = block.method, key = coord_key(&coord));
continue;
};
run_one_eval(ctx, &step, phase_name, depth, &phase_meta).await?;
evals += 1;
let value = ctx.optimize_objective_value.ok_or_else(|| {
format!(
"phase '{phase_name}': optimize objective '{}' produced no value — it must \
be a numeric wire fully qualified on the phase node",
block.objective
)
})?;
if value > best_value {
best_value = value;
best_coord = Some(coord.clone());
}
evaluated.push((coord, value));
}
batch = source_next(&mut src, &evaluated);
}
ctx.optimize_objective = None;
let best_disp = best_coord
.as_ref()
.map(|c| c.iter().map(|v| v.to_string()).collect::<Vec<_>>().join(", "))
.unwrap_or_else(|| "<none>".into());
crate::diag!(crate::observer::LogLevel::Info,
"optimizer '{method}' on '{phase_name}': best [{best_disp}] → {objective}={best_value} \
after {evals} evals",
method = block.method, objective = block.objective);
Ok(())
}.await;
match __opt_result {
Ok(()) => crate::phase_outcome::Outcome::completed(),
Err(reason) => crate::phase_outcome::Outcome::failed().with_reason(reason),
}
})
}
fn classify_control_axes(
space: &crate::optimize::SearchSpace,
servo: &[String],
phase_concurrency: Option<&str>,
phase_has_rate: bool,
) -> Result<Vec<crate::optimize::servo::ControlAxis>, String> {
use crate::optimize::servo::ControlAxis;
const KNOWN_CONTROLS: &[&str] = &["concurrency", "rate"];
let mut controls = Vec::new();
for var in servo {
let Some(i) = space.axes.iter().position(|ax| &ax.name == var) else {
return Err(format!(
"optimize `servo: {var}` is not a search axis — name a var that appears in the \
phase's `for_each`"
));
};
let control = if KNOWN_CONTROLS.contains(&var.as_str()) {
if var == "rate" && !phase_has_rate {
return Err(
"optimize `servo: rate` but the phase declares no `rate:` field, so there is no \
rate control to servo — add a `rate:` to the phase (its value is the warmup the \
servo retargets from)"
.to_string(),
);
}
var.clone()
} else if phase_concurrency.is_some_and(|c| c.contains(&format!("{{{var}}}"))) {
"concurrency".to_string()
} else {
return Err(format!(
"optimize `servo: {var}` but '{var}' is neither a live control nor wired to one — \
name a control directly (`servo: concurrency`) or wire the var \
(`concurrency: \"{{{var}}}\"`); or drop it from `servo:` to step through its \
values by re-running the phase"
));
};
controls.push(ControlAxis {
axis_idx: i,
control,
});
}
Ok(controls)
}
fn require_windowed_objective(
ctx: &ExecCtx,
phase_name: &str,
objective: &str,
) -> Result<(), String> {
let phase_kernel = ctx
.scope_tree
.phase_node_by_name(phase_name)
.and_then(|idx| ctx.scope_tree.nodes[idx].cached_kernel.get().cloned());
if let Some(pk) = &phase_kernel
&& !crate::optimize::settle::program_reads_live_metrics(pk.program())
{
return Err(format!(
"phase '{phase_name}': Control-class optimize needs a windowed objective \
('{objective}') to settle per setting, but the phase reads no live metric. \
Use metric_window(...) or metricsql_scalar(rate(...[W]))."
));
}
Ok(())
}
#[allow(clippy::too_many_arguments)] async fn run_servo_cell(
ctx: &mut ExecCtx,
step: &IterationStep,
space: crate::optimize::SearchSpace,
controls: Vec<crate::optimize::servo::ControlAxis>,
block: &nmbrs_workload::model::OptimizeBlock,
phase_name: &str,
depth: usize,
phase_meta: &Option<(String, Vec<String>, Vec<crate::checkpoint::PathSegment>)>,
) -> Result<crate::optimize::servo::ServoOutcome, String> {
let result = std::sync::Arc::new(arc_swap::ArcSwap::from_pointee(
crate::optimize::servo::ServoOutcome::default(),
));
let spec = crate::optimize::servo::ServoSpec {
method: block.method.clone(),
params: block.params.iter().map(|(k, v)| (k.clone(), *v)).collect(),
objective: crate::scope::objective_wire(&block.objective).to_string(),
max_evals: block.max_evals,
seed: block.seed,
space,
controls,
result: result.clone(),
};
ctx.optimize_objective = None;
ctx.optimize_servo = Some(spec);
run_one_eval(ctx, step, phase_name, depth, phase_meta).await?;
ctx.optimize_servo = None;
Ok((**result.load()).clone())
}
#[allow(clippy::too_many_arguments)] fn run_control_search<'a>(
ctx: &'a mut ExecCtx,
mut space: crate::optimize::SearchSpace,
coord_eval: CoordEval,
control_axes: Vec<crate::optimize::servo::ControlAxis>,
block: nmbrs_workload::model::OptimizeBlock,
phase_name: &'a str,
depth: usize,
phase_meta: Option<(String, Vec<String>, Vec<crate::checkpoint::PathSegment>)>,
) -> std::pin::Pin<Box<dyn std::future::Future<Output = Result<(), String>> + Send + 'a>> {
Box::pin(async move {
require_windowed_objective(ctx, phase_name, &block.objective)?;
for ca in &control_axes {
space.axes[ca.axis_idx].changeover = crate::optimize::Changeover::Control;
}
let center = coord_eval.representative(&space).ok_or_else(|| {
format!("phase '{phase_name}': control optimize has no representative coordinate")
})?;
let outcome = run_servo_cell(
ctx,
¢er,
space,
control_axes,
&block,
phase_name,
depth,
&phase_meta,
)
.await?;
let (best_disp, best_value) = match &outcome.best {
Some(b) => (
b.coord
.iter()
.map(|v| v.to_string())
.collect::<Vec<_>>()
.join(", "),
b.value,
),
None => ("<none>".to_string(), f64::NEG_INFINITY),
};
crate::diag!(
crate::observer::LogLevel::Info,
"optimizer '{method}' on '{phase_name}' (control): best [{best_disp}] → \
{objective}={best_value} after {evals} evals",
method = block.method,
objective = block.objective,
evals = outcome.evals
);
Ok(())
})
}
#[allow(clippy::too_many_arguments)] fn run_hybrid_search<'a>(
ctx: &'a mut ExecCtx,
space: crate::optimize::SearchSpace,
coord_eval: CoordEval,
control_axes: Vec<crate::optimize::servo::ControlAxis>,
block: nmbrs_workload::model::OptimizeBlock,
phase_name: &'a str,
depth: usize,
phase_meta: Option<(String, Vec<String>, Vec<crate::checkpoint::PathSegment>)>,
) -> std::pin::Pin<Box<dyn std::future::Future<Output = Result<(), String>> + Send + 'a>> {
Box::pin(async move {
require_windowed_objective(ctx, phase_name, &block.objective)?;
let control_idx: std::collections::HashSet<usize> =
control_axes.iter().map(|c| c.axis_idx).collect();
let coord_indices: Vec<usize> = (0..space.axes.len())
.filter(|i| !control_idx.contains(i))
.collect();
let k_axes: Vec<crate::optimize::Axis> = control_axes
.iter()
.map(|ca| {
let mut ax = space.axes[ca.axis_idx].clone();
ax.changeover = crate::optimize::Changeover::Control;
ax
})
.collect();
let k_space = crate::optimize::SearchSpace::new(k_axes);
let k_controls: Vec<crate::optimize::servo::ControlAxis> = control_axes
.iter()
.enumerate()
.map(|(k_pos, ca)| crate::optimize::servo::ControlAxis {
axis_idx: k_pos,
control: ca.control.clone(),
})
.collect();
let CoordEval::Enumerated { steps, .. } = &coord_eval else {
return Err(format!(
"phase '{phase_name}': hybrid coordinate+control optimize requires enumerated \
coordinate axes (a continuous coordinate axis alongside a control axis is not \
yet supported)"
));
};
let mut seen = std::collections::HashSet::new();
let mut cells: Vec<IterationStep> = Vec::new();
for step in steps {
let key: String = coord_indices
.iter()
.map(|&i| step.bindings[i].1.to_display_string())
.collect::<Vec<_>>()
.join("\u{1f}");
if seen.insert(key) {
cells.push(step.clone());
}
}
let mut best_value = f64::NEG_INFINITY;
let mut best_disp = "<none>".to_string();
let mut total_evals = 0usize;
for cell in &cells {
let outcome = run_servo_cell(
ctx,
cell,
k_space.clone(),
k_controls.clone(),
&block,
phase_name,
depth,
&phase_meta,
)
.await?;
total_evals += outcome.evals;
if let Some(b) = &outcome.best
&& b.value > best_value
{
best_value = b.value;
let c_disp = coord_indices
.iter()
.map(|&i| {
format!(
"{}={}",
cell.bindings[i].0,
cell.bindings[i].1.to_display_string()
)
})
.collect::<Vec<_>>()
.join(", ");
let k_disp = b
.coord
.iter()
.map(|v| v.to_string())
.collect::<Vec<_>>()
.join(", ");
best_disp = format!("{c_disp}; {k_disp}");
}
}
crate::diag!(
crate::observer::LogLevel::Info,
"optimizer '{method}' on '{phase_name}' (hybrid {ncells} coordinate cells × control): \
best [{best_disp}] → {objective}={best_value} after {total_evals} evals",
method = block.method,
objective = block.objective,
ncells = cells.len()
);
Ok(())
})
}
async fn run_one_eval(
ctx: &mut ExecCtx,
step: &IterationStep,
phase_name: &str,
depth: usize,
phase_meta: &Option<(String, Vec<String>, Vec<crate::checkpoint::PathSegment>)>,
) -> Result<(), String> {
let scene_id = match phase_meta {
Some((pn, op_names, ypath)) => {
let labels = canonical_phase_label(&step.coord_path);
push_phase_scene_node(
ctx.scene_tree_parent_id,
ypath.clone(),
pn.clone(),
labels,
op_names.clone(),
)
}
None => ctx.scene_tree_parent_id,
};
let saved = ctx.scene_tree_parent_id;
ctx.scene_tree_parent_id = scene_id;
let terminal = TerminalAction::Phase(phase_name);
let res = run_one_iteration(ctx, step, &terminal, depth, "optimize").await;
ctx.scene_tree_parent_id = saved;
if res.is_failure() {
Err(res.reason.clone().unwrap_or_default())
} else {
Ok(())
}
}
pub(crate) fn source_next(
src: &mut Box<dyn crate::optimize::CoordinateSource>,
evaluated: &[(crate::optimize::Coord, f64)],
) -> Option<Vec<crate::optimize::Coord>> {
if let Some(f) = src.as_feedback() {
f.step(evaluated)
} else if let Some(p) = src.as_pull() {
p.pull()
} else {
None
}
}
fn coord_key(coord: &crate::optimize::Coord) -> String {
coord
.iter()
.map(|v| v.to_string())
.collect::<Vec<_>>()
.join("\u{1f}")
}
fn index_steps(steps: &[IterationStep]) -> std::collections::HashMap<String, usize> {
let mut m = std::collections::HashMap::new();
for (i, step) in steps.iter().enumerate() {
let key = step
.bindings
.iter()
.map(|(_, v)| polydat_value_to_axis(v).to_string())
.collect::<Vec<_>>()
.join("\u{1f}");
m.entry(key).or_insert(i);
}
m
}
fn polydat_value_to_axis(v: &polydat::ast::Value) -> crate::optimize::AxisValue {
use polydat::ast::Value;
match v {
Value::F64(f) => crate::optimize::AxisValue::Num(*f),
Value::U64(u) => crate::optimize::AxisValue::Num(*u as f64),
Value::Bool(b) => crate::optimize::AxisValue::Bool(*b),
Value::Str(s) => crate::optimize::AxisValue::Label(s.to_string()),
other => crate::optimize::AxisValue::Label(other.to_display_string()),
}
}
fn search_space_from_steps(steps: &[IterationStep]) -> crate::optimize::SearchSpace {
use crate::optimize::{Axis, AxisKind, AxisValue, Changeover, SearchSpace};
let n_axes = steps.first().map(|s| s.bindings.len()).unwrap_or(0);
let mut axes = Vec::with_capacity(n_axes);
for a in 0..n_axes {
let name = steps[0].bindings[a].0.clone();
let mut seen: Vec<String> = Vec::new();
let mut values: Vec<AxisValue> = Vec::new();
let mut all_numeric = true;
for step in steps {
let av = polydat_value_to_axis(&step.bindings[a].1);
let key = av.to_string();
if !seen.contains(&key) {
seen.push(key);
if !matches!(av, AxisValue::Num(_)) {
all_numeric = false;
}
values.push(av);
}
}
let kind = if all_numeric {
AxisKind::Discrete { detents: values }
} else {
AxisKind::Categorical { options: values }
};
axes.push(Axis {
name,
kind,
changeover: Changeover::Coordinate,
});
}
SearchSpace::new(axes)
}
fn axis_value_to_polydat(v: &crate::optimize::AxisValue) -> polydat::ast::Value {
use crate::optimize::AxisValue;
match v {
AxisValue::Num(f) => polydat::ast::Value::F64(*f),
AxisValue::Bool(b) => polydat::ast::Value::Bool(*b),
AxisValue::Label(s) => polydat::ast::Value::Str(s.as_str().into()),
}
}
fn continuous_axis_intervals(
comp: &polydat::iteration::comprehension::Comprehension,
) -> Option<Vec<polydat::iteration::comprehension::Interval>> {
use polydat::iteration::comprehension::IndexFn;
match comp.metadata().index_addressable {
Some(IndexFn::Continuous { intervals, .. }) => Some(intervals),
_ => None,
}
}
fn search_space_continuous(
axis_names: &[String],
intervals: &[polydat::iteration::comprehension::Interval],
) -> crate::optimize::SearchSpace {
use crate::optimize::{Axis, AxisKind, Changeover, SearchSpace};
SearchSpace::new(
axis_names
.iter()
.zip(intervals)
.map(|(name, iv)| Axis {
name: name.clone(),
kind: AxisKind::Continuous {
lo: iv.lo,
hi: iv.hi,
min_step: 0.0,
},
changeover: Changeover::Coordinate,
})
.collect(),
)
}
enum CoordEval {
Enumerated {
steps: Vec<IterationStep>,
index: std::collections::HashMap<String, usize>,
},
Synthesized {
axis_names: Vec<String>,
canonical: std::sync::Arc<crate::scope_kernel::ScopeKernel>,
parent: std::sync::Arc<crate::scope_kernel::ScopeKernel>,
parent_coords: Vec<polydat::kernel::ScopeCoord>,
},
}
impl CoordEval {
fn materialize(&self, coord: &crate::optimize::Coord) -> Option<IterationStep> {
match self {
CoordEval::Enumerated { steps, index } => {
index.get(&coord_key(coord)).map(|&i| steps[i].clone())
}
CoordEval::Synthesized {
axis_names,
canonical,
parent,
parent_coords,
} => {
let tuple: Vec<(String, polydat::ast::Value)> = axis_names
.iter()
.zip(coord)
.map(|(n, av)| (n.clone(), axis_value_to_polydat(av)))
.collect();
let bound_kernel = match canonical.bind_under(parent.kernel(), &tuple) {
Ok(k) => std::sync::Arc::new(k),
Err(e) => {
crate::diag!(
crate::observer::LogLevel::Warn,
"coordinate {tuple:?} could not be bound: {e}"
);
return None;
}
};
let mut coord_path = parent_coords.clone();
coord_path.push(polydat::kernel::ScopeCoord::from(tuple.iter().cloned()));
Some(IterationStep {
bindings: tuple,
bound_kernel,
coord_path,
})
}
}
}
fn representative(&self, space: &crate::optimize::SearchSpace) -> Option<IterationStep> {
match self {
CoordEval::Enumerated { steps, .. } => steps.first().cloned(),
CoordEval::Synthesized { .. } => self.materialize(&space.center()),
}
}
}
fn read_objective_at_completion(
parent: &std::sync::Arc<crate::scope_kernel::ScopeKernel>,
phase_kernel: &std::sync::Arc<crate::scope_kernel::ScopeKernel>,
objective: &str,
) -> Option<f64> {
use polydat::ast::Value;
let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
let mut k = phase_kernel.bind_under(parent.kernel(), &[]).ok()?;
Some(k.pull(objective))
}));
match result {
Ok(Some(Value::F64(f))) => Some(f),
Ok(Some(Value::U64(u))) => Some(u as f64),
Ok(Some(Value::Bool(b))) => Some(if b { 1.0 } else { 0.0 }),
_ => None,
}
}
fn fire_scope_lifecycle(
ctx: &ExecCtx,
event: crate::lifecycle::EventType,
spec: &str,
depth: usize,
) {
let depth_indent = " ".repeat(depth.saturating_sub(1));
let display_labels: String = {
let parent_coords: Vec<_> = ctx
.current_parent_kernel
.as_ref()
.map(|k| k.scope_coordinates().iter().rev().cloned().collect())
.unwrap_or_default();
polydat::kernel::format_scope_coordinate_path(&parent_coords)
};
let scope_ctx = crate::readout_context::LifecycleContext {
event,
subject_name: spec.to_string(),
subject_labels: display_labels,
depth_indent,
use_color: crate::observer::use_color(),
stick_reattached: String::new(),
};
crate::readout_context::fire_lifecycle(
event,
&ctx.workload_readouts,
None,
&scope_ctx,
Some(&ctx.sqlite_reporter),
);
}
async fn run_do_loop(
ctx: &mut ExecCtx,
condition: &str,
counter: Option<&str>,
invert: bool,
children: &[ScenarioNode],
depth: usize,
) -> Result<(), String> {
let scope_idx = ctx
.scope_tree
.iter_dfs()
.find_map(|(idx, node)| match &node.kind {
crate::scope_tree::ScopeKind::DoWhile {
condition: c,
counter: ct,
} => {
if c == condition && ct.as_deref() == counter && !invert {
Some(idx)
} else {
None
}
}
crate::scope_tree::ScopeKind::DoUntil {
condition: c,
counter: ct,
} => {
if c == condition && ct.as_deref() == counter && invert {
Some(idx)
} else {
None
}
}
_ => None,
})
.ok_or_else(|| {
format!(
"do-loop '{condition}': no matching scope-tree entry — \
scenario/scope-tree drift bug."
)
})?;
let canonical = ctx.scope_tree.nodes[scope_idx]
.cached_kernel
.get()
.cloned()
.ok_or_else(|| {
format!("do-loop '{condition}': scope at index {scope_idx} has no installed kernel.")
})?;
let parent = effective_parent_kernel(ctx, scope_idx)
.ok_or_else(|| format!("do-loop '{condition}': no installed ancestor kernel."))?;
let mut loop_kernel = canonical
.bind_under(parent.kernel(), &[])
.map_err(|e| format!("do-loop '{condition}': {e}"))?;
let mut counter_value: u64 = 0;
loop {
if let Some(c) = counter
&& let Some(idx) = loop_kernel.program().find_input(c)
{
crate::wires::write_input(
loop_kernel.kernel_mut(),
idx,
c,
polydat::ast::Value::U64(counter_value),
)
.map_err(|e| format!("do-loop '{condition}': counter '{c}': {e}"))?;
loop_kernel
.kernel_mut()
.init()
.map_err(|e| format!("do-loop '{condition}': {e}"))?;
}
let interpolated = polydat::kernel::interp::interpolate_via_kernel(condition, &loop_kernel)
.map_err(|e| format!("do-loop condition '{condition}': {e}"))?;
let cond_value = polydat::dsl::compile::eval_const_expr(&interpolated)
.map_err(|e| format!("do-loop condition '{condition}': {e}"))?;
let cond_true = match cond_value {
polydat::ast::Value::Bool(b) => b,
polydat::ast::Value::U64(n) => n != 0,
polydat::ast::Value::F64(n) => n != 0.0,
other => {
return Err(format!(
"do-loop condition '{condition}': expected bool/u64/f64, got {other:?}",
));
}
};
let should_continue = if invert { !cond_true } else { cond_true };
if !should_continue {
break;
}
let prior_parent = ctx.current_parent_kernel.take();
if let Some(c) = counter {
ctx.push_label(c, &counter_value.to_string());
}
let arc_loop = std::sync::Arc::new(std::mem::replace(
&mut loop_kernel,
canonical.fork(),
));
ctx.current_parent_kernel = Some(arc_loop.clone());
let kind: &'static str = if invert { "do_until" } else { "do_while" };
let mut iter_coords = std::collections::BTreeMap::new();
if let Some(c) = counter {
iter_coords.insert(c.to_string(), serde_json::Value::from(counter_value));
}
let path: Vec<std::collections::BTreeMap<String, serde_json::Value>> = arc_loop
.scope_coordinates()
.iter()
.rev()
.filter(|c| !c.is_empty())
.map(coord_to_btree)
.collect();
if let Some(writer) = ctx.checkpoint_writer.as_ref() {
writer.emit_scope_enter(kind, iter_coords.clone(), path.clone());
}
let res = execute_tree_at(ctx, children, depth).await;
if let Some(writer) = ctx.checkpoint_writer.as_ref() {
let outcome = if !res.is_failure() {
"completed"
} else {
"interrupted"
};
writer.emit_scope_exit(kind, iter_coords, path, outcome);
}
ctx.current_parent_kernel = prior_parent;
if counter.is_some() {
ctx.pop_label();
}
loop_kernel = std::sync::Arc::try_unwrap(arc_loop).map_err(|_| {
format!(
"do-loop '{condition}' iteration {counter_value}: persistent kernel \
still referenced after children completed — concurrency bug."
)
})?;
if res.is_failure() {
return Err(res.reason.clone().unwrap_or_default());
}
counter_value = counter_value.saturating_add(1);
}
Ok(())
}
async fn run_one_iteration(
ctx: &mut ExecCtx,
step: &IterationStep,
terminal: &TerminalAction<'_>,
depth: usize,
kind: &'static str,
) -> crate::phase_outcome::Outcome {
let prior_parent = ctx.current_parent_kernel.take();
ctx.current_parent_kernel = Some(step.bound_kernel.clone());
for (var, value) in &step.bindings {
ctx.push_label(var, &value.to_display_string());
}
let (enter_coords, enter_path) = scope_event_coords(&step.coord_path);
if let Some(writer) = ctx.checkpoint_writer.as_ref() {
writer.emit_scope_enter(kind, enter_coords.clone(), enter_path.clone());
}
let iter_label = step
.bindings
.iter()
.map(|(k, v)| format!("{k}={}", v.to_display_string()))
.collect::<Vec<_>>()
.join(", ");
let display_labels: String = {
let parent_coords: Vec<_> = ctx
.current_parent_kernel
.as_ref()
.map(|k| k.scope_coordinates().iter().rev().cloned().collect())
.unwrap_or_default();
polydat::kernel::format_scope_coordinate_path(&parent_coords)
};
let depth_indent = " ".repeat(depth.saturating_sub(1));
let each_ctx = crate::readout_context::LifecycleContext {
event: crate::lifecycle::EventType::EachStart,
subject_name: iter_label.clone(),
subject_labels: display_labels.clone(),
depth_indent: depth_indent.clone(),
use_color: crate::observer::use_color(),
stick_reattached: String::new(),
};
crate::readout_context::fire_lifecycle(
crate::lifecycle::EventType::EachStart,
&ctx.workload_readouts,
None,
&each_ctx,
Some(&ctx.sqlite_reporter),
);
let res = match terminal {
TerminalAction::Children(children) => execute_tree_at(ctx, children, depth).await,
TerminalAction::Phase(name) => {
if ctx.diag.depth >= crate::runner::ExecDepth::Op {
let phase_node_id = ctx.scene_tree_parent_id;
run_phase_layered(ctx, name, phase_node_id).await
} else {
crate::phase_outcome::Outcome::skipped()
}
}
};
let each_end_ctx = crate::readout_context::LifecycleContext {
event: crate::lifecycle::EventType::EachEnd,
subject_name: iter_label,
subject_labels: display_labels,
depth_indent,
use_color: crate::observer::use_color(),
stick_reattached: String::new(),
};
crate::readout_context::fire_lifecycle(
crate::lifecycle::EventType::EachEnd,
&ctx.workload_readouts,
None,
&each_end_ctx,
Some(&ctx.sqlite_reporter),
);
if let Some(writer) = ctx.checkpoint_writer.as_ref() {
let outcome = if !res.is_failure() {
"completed"
} else {
"interrupted"
};
writer.emit_scope_exit(kind, enter_coords, enter_path, outcome);
}
for _ in &step.bindings {
ctx.pop_label();
}
ctx.current_parent_kernel = prior_parent;
res
}
fn scope_event_coords(
coord_path: &[polydat::kernel::ScopeCoord],
) -> (
std::collections::BTreeMap<String, serde_json::Value>,
Vec<std::collections::BTreeMap<String, serde_json::Value>>,
) {
let coords = coord_path.last().map(coord_to_btree).unwrap_or_default();
let path: Vec<_> = if coord_path.len() <= 1 {
Vec::new()
} else {
coord_path[..coord_path.len() - 1]
.iter()
.rev()
.filter(|c| !c.is_empty())
.map(coord_to_btree)
.collect()
};
(coords, path)
}
fn coord_to_btree(
coord: &polydat::kernel::ScopeCoord,
) -> std::collections::BTreeMap<String, serde_json::Value> {
coord
.vars
.iter()
.map(|(k, v)| (k.clone(), v.to_json_value()))
.collect()
}
static EMPTY_BINDINGS: std::sync::LazyLock<HashMap<String, String>> =
std::sync::LazyLock::new(HashMap::new);
async fn run_phase(
ctx: &mut ExecCtx,
phase_name: &str,
scene_node_id: crate::scene_tree::SceneNodeId,
) -> crate::phase_outcome::Outcome {
let phase_start_ms = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0);
let outcome = crate::execution_context::with_current_phase(
scene_node_id,
phase_start_ms,
run_phase_inner(ctx, phase_name, scene_node_id),
)
.await;
if outcome.is_failure() {
let recorded =
crate::scene_tree::with_global(|t| t.phase_outcome_present_at(scene_node_id))
.unwrap_or(false);
if !recorded {
record_early_phase_failure(ctx, scene_node_id, phase_name, &outcome);
}
}
outcome
}
fn record_early_phase_failure(
ctx: &mut ExecCtx,
scene_node_id: crate::scene_tree::SceneNodeId,
phase_name: &str,
outcome: &crate::phase_outcome::Outcome,
) {
let reason = outcome
.reason
.clone()
.unwrap_or_else(|| "phase configuration failed".to_string());
crate::diag!(
crate::observer::LogLevel::Error,
"phase '{phase_name}': {reason} — failing phase (config)"
);
ctx.observer
.phase_failed(scene_node_id, phase_name, "", &reason);
let now_nanos = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_nanos() as u64)
.unwrap_or(0);
let phase_outcome = crate::phase_outcome::PhaseOutcome::failed(
crate::phase_outcome::PhaseIdentity::new(phase_name, ""),
0.0,
vec![crate::phase_outcome::PhaseErrorDetail {
class: "config".into(),
message: reason.clone(),
op_name: None,
cycle: None,
op_template: None,
op_resolved: None,
at_nanos: now_nanos,
retryable: false,
}],
);
if let Ok(mut guard) = ctx.sqlite_reporter.lock()
&& let Some(reporter) = guard.as_mut()
{
let row = phase_outcome.to_sqlite_row(&ctx.session_id, ctx.exec_id, now_nanos as i64);
reporter.write_phase_outcome(&row);
}
crate::scene_tree::with_global_mut(|t| {
t.set_phase_failed_at(scene_node_id, &reason);
t.set_phase_outcome_at(scene_node_id, phase_outcome);
});
if let Some(writer) = ctx.checkpoint_writer.as_ref() {
let identity = phase_identity_for(phase_name, "");
writer.phase_failed(&identity, &reason);
let _ = writer.flush();
}
if let Some((trip_outcome, trip_reason)) = ctx.workload_shell.record_phase(true, 0, 0) {
let cause = if trip_outcome.is_failure() {
crate::session_signals::StopCause::Fault
} else {
crate::session_signals::StopCause::Interrupt
};
crate::session_signals::request_shell_stop(cause);
crate::diag!(
crate::observer::LogLevel::Warn,
"scenario stop-on-error ({trip_reason}) after phase '{phase_name}' — halting remaining walk"
);
}
}
async fn run_phase_inner(
ctx: &mut ExecCtx,
phase_name: &str,
scene_node_id: crate::scene_tree::SceneNodeId,
) -> crate::phase_outcome::Outcome {
let phase_start = std::time::Instant::now();
let phase_start_nanos: i64 = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_nanos() as i64)
.unwrap_or(0);
let phase = match ctx.phases.get(phase_name) {
Some(p) => p.clone(),
None => {
return crate::phase_outcome::Outcome::failed()
.with_reason(format!("phase '{phase_name}' not found"));
}
};
let has_bindings = phase.ops.iter().any(|op| !op.bindings.is_empty());
let iter_var_values: IndexMap<String, String> = {
let mut out = IndexMap::new();
if let Some(parent) = ctx.current_parent_kernel.as_ref() {
for name in parent.program().output_names() {
if let Some(v) = parent.lookup(name) {
out.insert(name.to_string(), v.to_display_string());
}
}
for name in parent.program().input_names() {
if out.contains_key(&name) {
continue;
}
if let Some(v) = parent.lookup(&name) {
out.insert(name.clone(), v.to_display_string());
}
}
}
out
};
let is_iter = !iter_var_values.is_empty();
crate::diag!(
crate::observer::LogLevel::Debug,
"=== phase: {phase_name} ==="
);
if is_iter && let Some(parent) = ctx.current_parent_kernel.as_ref() {
let prog = parent.program();
for (var, val) in &iter_var_values {
if !val.is_empty() && !prog.is_inherited(var) {
crate::diag!(crate::observer::LogLevel::Debug, " {var}={val}");
}
}
}
let early_phase_labels = match ctx.current_parent_kernel.as_ref() {
Some(parent) => format_scope_coordinate_path(parent.scope_coordinates()),
None => String::new(),
};
let early_identity = phase_identity_for(phase_name, &early_phase_labels);
match ctx.resume_plan.action_for(&early_identity) {
crate::checkpoint::ResumeAction::Skip if ctx.refine_plan.is_some() => {
crate::diag!(
crate::observer::LogLevel::Info,
"phase '{phase_name}' [checkpoint-resume skip deferred \
to the refine hash gate]"
);
}
crate::checkpoint::ResumeAction::Skip => {
crate::diag!(
crate::observer::LogLevel::Info,
"phase '{phase_name}' [skipped — checkpoint resume]"
);
crate::scene_tree::with_global_mut(|t| {
t.set_phase_running_at(scene_node_id, 0);
t.set_phase_completed_at(scene_node_id, 0.0);
});
ctx.observer
.phase_starting(scene_node_id, phase_name, &early_phase_labels, 0, 0, 0);
ctx.observer
.phase_completed(scene_node_id, phase_name, &early_phase_labels, 0.0);
crate::phase_end_triggers::fire_phase_completed(phase_name, &early_phase_labels, 0.0);
return crate::phase_outcome::Outcome::skipped();
}
crate::checkpoint::ResumeAction::IdentityMismatch { reason } => {
crate::diag!(
crate::observer::LogLevel::Warn,
"phase '{phase_name}' [resume: {reason} — re-running]"
);
}
crate::checkpoint::ResumeAction::CursorResume { .. } => {
crate::diag!(
crate::observer::LogLevel::Info,
"phase '{phase_name}' [resume: cursor-state available — \
Tier 2 restore not yet wired, re-running from scratch]"
);
}
crate::checkpoint::ResumeAction::ReRun => {}
}
if crate::observer::skipped_phase_display() == crate::observer::SkippedPhaseDisplay::Prune
&& let Some(phase_def) = ctx.phases.get(phase_name)
&& !phase_def.ops.is_empty()
&& phase_def.ops.iter().all(|op| op.condition.is_some())
&& let Some(parent) = ctx.current_parent_kernel.as_ref()
{
let phase_nodes: Vec<crate::scope_tree::ScopeNodeIdx> = ctx.scope_tree.nodes
[ctx.current_scope_idx]
.children
.iter()
.copied()
.filter(|c| {
matches!(&ctx.scope_tree.nodes[*c].kind,
crate::scope_tree::ScopeKind::Phase { name } if name == phase_name)
})
.collect();
let mut all_false = phase_nodes.len() == 1;
let mut gates_checked = 0usize;
let op_children: Vec<crate::scope_tree::ScopeNodeIdx> = phase_nodes
.first()
.map(|p| ctx.scope_tree.nodes[*p].children.clone())
.unwrap_or_default();
for child in &op_children {
let node = &ctx.scope_tree.nodes[*child];
let crate::scope_tree::ScopeKind::OpTemplate { name } = &node.kind else {
continue;
};
let Some(op) = phase_def.ops.iter().find(|o| o.name == *name) else {
continue;
};
let Some(cond) = op.condition.as_ref() else {
continue;
};
let cond_name = cond.trim().trim_start_matches('{').trim_end_matches('}');
let Some(installed) = node.cached_kernel.get() else {
all_false = false; break;
};
let Ok(mut gate_kernel) = installed.bind_under(parent.kernel(), &[]) else {
all_false = false;
break;
};
let truthy = match gate_kernel.pull(cond_name) {
polydat::ast::Value::None => false,
polydat::ast::Value::U64(v) => v != 0,
polydat::ast::Value::F64(v) => v != 0.0,
polydat::ast::Value::Bool(v) => v,
polydat::ast::Value::Str(s) => !s.is_empty(),
_ => true,
};
gates_checked += 1;
if truthy {
all_false = false;
break;
}
}
if all_false && gates_checked != phase_def.ops.len() {
all_false = false;
}
if all_false {
crate::diag!(
crate::observer::LogLevel::Debug,
"phase '{phase_name}' [pruned — every op gate false at entry \
(skipped_phases=prune)]"
);
crate::scene_tree::with_global_mut(|t| {
t.set_phase_completed_at(scene_node_id, 0.0);
});
return crate::phase_outcome::Outcome::skipped();
}
}
if phase.ops.is_empty() {
return crate::phase_outcome::Outcome::failed().with_reason(format!(
"phase '{phase_name}' reached dispatch with no ops — \
a load-time validation gap; every phase needs inline \
ops, a matching `tags:` selector, or a bound \
implementation",
));
}
let (
iter_op_builder,
iter_ops,
runtime_cursor_extents,
runtime_cursor_min_ms,
runtime_cursor_min_passes,
runtime_cursor_min_count,
runtime_cursor_delta,
runtime_cursor_partition,
activation_scope,
) = if is_iter || has_bindings {
let mut ops = phase.ops.clone();
let parent_kernel = match ctx.current_parent_kernel.as_ref() {
Some(k) => k,
None => {
return crate::phase_outcome::Outcome::failed().with_reason(format!(
"phase '{phase_name}': no current_parent_kernel — \
single-resolution-path requires the populated parent kernel",
));
}
};
let validation_kernel = ctx
.scope_tree
.phase_node_by_name(phase_name)
.and_then(|idx| ctx.scope_tree.nodes[idx].cached_kernel.get())
.map(|k| k.as_ref())
.unwrap_or(parent_kernel);
let ancestors = ctx
.scope_tree
.phase_node_by_name(phase_name)
.map(|idx| ctx.scope_tree.ancestor_kernels(idx))
.unwrap_or_default();
let enclosing: Vec<&polydat::kernel::PolydatProgram> =
std::iter::once(parent_kernel.program().as_ref())
.chain(ancestors.iter().map(|k| k.program().as_ref()))
.collect();
if let Err(e) = crate::scope::validate_placeholders_via_kernel(
&ops,
validation_kernel.kernel(),
&enclosing,
)
.map_err(|e| format!("phase '{phase_name}': {e}"))
{
return crate::phase_outcome::Outcome::failed().with_reason(e);
}
crate::scope::rewrite_inline_exprs(&mut ops);
let parent_kernel = match ctx.current_parent_kernel.as_ref() {
Some(k) => k,
None => {
return crate::phase_outcome::Outcome::failed().with_reason(format!(
"phase '{phase_name}': no current_parent_kernel — \
workload root install missed at session start (internal bug)."
));
}
};
let classifier_kernel: &crate::scope_kernel::ScopeKernel = ctx
.scope_tree
.phase_node_by_name(phase_name)
.and_then(|idx| ctx.scope_tree.nodes[idx].cached_kernel.get())
.map(|k| k.as_ref())
.unwrap_or(parent_kernel);
let effective_manifest = crate::runner::extract_manifest(classifier_kernel.program());
let scope = match crate::scope::build_scope(
&ops,
&EMPTY_BINDINGS,
&effective_manifest,
&EMPTY_BINDINGS,
&ctx.phases,
phase.cycles.as_deref(),
&[], Some(classifier_kernel),
) {
Ok(v) => v,
Err(e) => return crate::phase_outcome::Outcome::failed().with_reason(e),
};
let polydat_context = if iter_var_values.is_empty() {
format!("phase '{phase_name}'")
} else {
let vars: Vec<String> = iter_var_values
.iter()
.map(|(k, v)| format!("{k}={v}"))
.collect();
format!("phase '{phase_name}' ({})", vars.join(", "))
};
if let Err(e) = scope
.validate()
.map_err(|e| format!("{polydat_context}: {e}"))
{
return crate::phase_outcome::Outcome::failed().with_reason(e);
}
let cursor_limit: Option<u64> = ctx.merged_params.get("limit").and_then(|s| s.parse().ok());
let phase_idx = ctx.scope_tree.phase_node_by_name(phase_name);
let phase_pragmas = phase_idx
.map(|idx| ctx.scope_tree.nodes[idx].pragmas.clone())
.unwrap_or_default();
let compile_phase = || {
crate::bindings::compile_from_scope(
&scope,
ctx.workload_dir.as_deref(),
ctx.polydat_lib_paths.clone(),
ctx.strict,
&polydat_context,
cursor_limit,
&phase_pragmas,
)
.map_err(|e| format!("{polydat_context}: {e}"))
};
let phase_scope: std::sync::Arc<crate::scope_kernel::ScopeKernel> = match phase_idx {
Some(idx) => {
let node = &ctx.scope_tree.nodes[idx];
match node.cached_kernel.get() {
Some(canonical) => canonical.clone(),
None => {
let compiled = match compile_phase() {
Ok(v) => std::sync::Arc::new(v),
Err(e) => {
return crate::phase_outcome::Outcome::failed().with_reason(e);
}
};
let _ = node.cached_kernel.set(compiled.clone());
compiled
}
}
}
None => match compile_phase() {
Ok(v) => std::sync::Arc::new(v),
Err(e) => return crate::phase_outcome::Outcome::failed().with_reason(e),
},
};
let mut kernel = match phase_scope.bind_under(parent_kernel.kernel(), &[]) {
Ok(v) => v,
Err(e) => {
return crate::phase_outcome::Outcome::failed()
.with_reason(format!("{polydat_context}: {e}"));
}
};
let const_outputs: Vec<String> = kernel
.program()
.const_outputs()
.iter()
.map(|s| s.to_string())
.collect();
for init_name in &const_outputs {
let pull_result =
std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| kernel.pull(init_name)));
match pull_result {
Ok(v) if !matches!(v, polydat::ast::Value::None) => {}
Ok(_) => {
return crate::phase_outcome::Outcome::failed().with_reason(format!(
"{polydat_context}: init binding '{init_name}' violates the init contract: \
scope-init eval returned Value::None (per SRD 11 §\"Init Binding Contract\" \
Plan B). The eval function signaled failure or returned no value."
));
}
Err(payload) => {
let msg = panic_message(&payload);
return crate::phase_outcome::Outcome::failed().with_reason(format!(
"{polydat_context}: init binding '{init_name}' violates the init contract: \
scope-init eval panicked: {msg} (per SRD 11 §\"Init Binding Contract\" \
Plan B)."
));
}
}
}
let const_outputs: Vec<String> = kernel
.program()
.const_outputs()
.iter()
.map(|s| s.to_string())
.collect();
for final_name in &const_outputs {
if kernel.lookup(final_name).is_some() {
continue;
}
let pull_result =
std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| kernel.pull(final_name)));
match pull_result {
Ok(v) if !matches!(v, polydat::ast::Value::None) => {}
Ok(_) => {
return crate::phase_outcome::Outcome::failed().with_reason(format!(
"{polydat_context}: final binding '{final_name}' could not be \
materialised at scope activation: eval returned Value::None. \
If the RHS depends on a wire that's only available per cycle, \
use a non-modifier cycle binding (`{final_name} := …`) instead \
of `final`."
));
}
Err(payload) => {
let msg = panic_message(&payload);
return crate::phase_outcome::Outcome::failed().with_reason(format!(
"{polydat_context}: final binding '{final_name}' eval panicked at \
scope activation: {msg}"
));
}
}
}
if !ctx.phase_param_overrides.is_empty() {
let chosen = match crate::phase_params::resolve_for_phase(
&ctx.phase_param_overrides,
phase_name,
)
.map_err(|e| format!("{polydat_context}: {e}"))
{
Ok(v) => v,
Err(e) => return crate::phase_outcome::Outcome::failed().with_reason(e),
};
let overridden = !chosen.is_empty();
for (ov, _dialect) in chosen {
use crate::wires::HostWriteError;
use polydat::kernel::WriteError;
let written = match kernel.input_index(&ov.param) {
Some(idx) => crate::wires::write_input(
kernel.kernel_mut(),
idx,
&ov.param,
polydat::ast::Value::Str(ov.value.clone().into()),
),
None => Err(HostWriteError::Write(WriteError::UnknownWire {
key: ov.param.clone(),
known: Vec::new(),
})),
};
match written {
Ok(()) => {
crate::diag!(
crate::observer::LogLevel::Info,
"phase '{phase_name}': param `{}` overridden to `{}` \
(CLI `{}.{}=`)",
ov.param,
ov.value,
ov.pattern.source(),
ov.param
);
}
Err(HostWriteError::Write(WriteError::UnknownWire { .. })) => {
if ov.pattern.is_exact_literal() {
crate::diag!(
crate::observer::LogLevel::Warn,
"phase '{phase_name}': override `{}.{}=` names a \
param this phase does not consume — no wire \
named `{}` in its scope",
ov.pattern.source(),
ov.param,
ov.param
);
}
}
Err(e) => {
return crate::phase_outcome::Outcome::failed().with_reason(format!(
"{polydat_context}: phase-scoped override `{}.{}=`: {e}",
ov.pattern.source(),
ov.param,
));
}
}
}
if overridden && let Err(e) = kernel.kernel_mut().init() {
return crate::phase_outcome::Outcome::failed()
.with_reason(format!("{polydat_context}: phase-scoped overrides: {e}"));
}
}
let mut runtime_extents: HashMap<String, u64> = HashMap::new();
let mut runtime_min_ms: HashMap<String, u64> = HashMap::new();
let mut runtime_min_passes: HashMap<String, u64> = HashMap::new();
let mut runtime_min_count: HashMap<String, u64> = HashMap::new();
let mut runtime_delta: HashMap<String, u64> = HashMap::new();
let mut runtime_partition: HashMap<String, (u64, u64)> = HashMap::new();
type CursorSpec = (
String,
Option<(String, String)>,
Option<u64>,
polydat::iteration::source::CursorKind,
Option<String>,
);
let cursor_specs: Vec<CursorSpec> = kernel
.program()
.cursor_schemas()
.iter()
.map(|s| {
(
s.name.clone(),
s.extent_outputs.clone(),
s.extent_limit,
s.cursor_kind.clone(),
s.partition_output.clone(),
)
})
.collect();
for (name, outputs, limit, cursor_kind, partition_output) in cursor_specs {
if let Some((start_out, end_out)) = outputs {
let start = kernel.pull(&start_out).as_u64();
let end = kernel.pull(&end_out).as_u64();
let extent = end.saturating_sub(start);
let final_extent = limit.map(|l| extent.min(l)).unwrap_or(extent);
runtime_extents.insert(name.clone(), final_extent);
}
if let Some(out) = &partition_output {
let value = kernel.pull(out);
let cursor_extent = runtime_extents.get(&name).copied().unwrap_or_else(|| {
kernel
.program()
.cursor_schemas()
.iter()
.find(|s| s.name == name)
.and_then(|s| s.extent)
.unwrap_or(0)
});
let mut write_slot = |suffix: &str, v: polydat::ast::Value| {
let slot = format!("{name}__cursor{suffix}");
if let Some(idx) = kernel.program().find_input(&slot)
&& let Err(e) =
crate::wires::write_input(kernel.kernel_mut(), idx, &slot, v)
{
crate::diag!(
crate::observer::LogLevel::Warn,
"cursor '{name}': slot '{slot}' refused its value: {e}"
);
}
};
use polydat::ast::Value as PValue;
let open_extent =
!matches!(cursor_kind, polydat::iteration::source::CursorKind::Range,);
match resolve_over(&value, cursor_extent, open_extent) {
Ok(Some(partition)) => {
runtime_partition
.insert(name.clone(), (partition.start_ord, partition.end_ord));
write_slot("__idx", PValue::U64(partition.idx));
write_slot("__partition_count", PValue::U64(partition.count.max(1)));
write_slot("__start_pct", PValue::F64(partition.start_pct));
write_slot("__end_pct", PValue::F64(partition.end_pct));
write_slot("__start_ordinal", PValue::U64(partition.start_ord));
write_slot("__end_ordinal", PValue::U64(partition.end_ord));
if partition.count > 1 {
crate::diag!(
crate::observer::LogLevel::Info,
"phase '{phase_name}': partition {}/{} [{}..{})",
partition.idx + 1,
partition.count,
partition.start_ord,
partition.end_ord
);
}
write_slot("", PValue::from_partition(partition));
}
Ok(None) => {
write_slot("__end_ordinal", PValue::U64(cursor_extent));
}
Err(e) => {
return crate::phase_outcome::Outcome::failed().with_reason(format!(
"cursor '{name}': `over` clause failed to resolve: {e}"
));
}
}
if let Err(e) = kernel.kernel_mut().init() {
return crate::phase_outcome::Outcome::failed()
.with_reason(format!("cursor '{name}': {e}"));
}
}
use polydat::iteration::source::CursorKind::*;
match &cursor_kind {
Range => {}
ExtendingTimed {
min_ms_output,
delta_output,
} => {
runtime_min_ms.insert(name.clone(), kernel.pull(min_ms_output).as_u64());
if let Some(d) = delta_output {
runtime_delta.insert(name.clone(), kernel.pull(d).as_u64());
}
}
ExtendingPasses {
min_passes_output,
delta_output,
} => {
runtime_min_passes
.insert(name.clone(), kernel.pull(min_passes_output).as_u64());
if let Some(d) = delta_output {
runtime_delta.insert(name.clone(), kernel.pull(d).as_u64());
}
}
ExtendingCount {
min_count_output,
delta_output,
} => {
runtime_min_count.insert(name.clone(), kernel.pull(min_count_output).as_u64());
if let Some(d) = delta_output {
runtime_delta.insert(name.clone(), kernel.pull(d).as_u64());
}
}
ExtendingElapsedAndPasses {
min_ms_output,
min_passes_output,
delta_output,
}
| ExtendingElapsedOrPasses {
min_ms_output,
min_passes_output,
delta_output,
} => {
runtime_min_ms.insert(name.clone(), kernel.pull(min_ms_output).as_u64());
runtime_min_passes
.insert(name.clone(), kernel.pull(min_passes_output).as_u64());
if let Some(d) = delta_output {
runtime_delta.insert(name.clone(), kernel.pull(d).as_u64());
}
}
}
}
let activation_scope = Arc::new(kernel.fork());
let op_builder = {
let mut b = OpBuilder::new(kernel);
if let Some(phase_idx) = ctx.scope_tree.phase_node_by_name(phase_name) {
let map = ctx.scope_tree.op_template_programs_for_phase(phase_idx);
if !map.is_empty() {
b = b.with_op_template_programs(map);
}
b = b.with_op_template_modules(
ctx.scope_tree.op_template_modules_for_phase(phase_idx),
);
}
Arc::new(b)
};
(
op_builder,
ops,
runtime_extents,
runtime_min_ms,
runtime_min_passes,
runtime_min_count,
runtime_delta,
runtime_partition,
activation_scope,
)
} else {
let parent = ctx
.current_parent_kernel
.as_ref()
.expect("workload-kernel fallback requires an installed parent kernel");
let workload_subscope = match ctx.workload_scope.bind_under(parent.kernel(), &[]) {
Ok(v) => v,
Err(e) => {
return crate::phase_outcome::Outcome::failed()
.with_reason(format!("phase '{phase_name}': workload scope: {e}"));
}
};
let activation_scope = Arc::new(workload_subscope.fork());
let mut b = OpBuilder::new(workload_subscope);
if let Some(phase_idx) = ctx.scope_tree.phase_node_by_name(phase_name) {
let map = ctx.scope_tree.op_template_programs_for_phase(phase_idx);
if !map.is_empty() {
b = b.with_op_template_programs(map);
}
}
(
Arc::new(b),
phase.ops.clone(),
HashMap::new(),
HashMap::new(),
HashMap::new(),
HashMap::new(),
HashMap::new(),
HashMap::new(),
activation_scope,
)
};
let op_sequence = OpSequence::from_ops(iter_ops, ctx.seq_type);
if op_sequence.stanza_length() == 0 {
crate::diag!(
crate::observer::LogLevel::Warn,
"warning: phase '{phase_name}' has no ops, skipping"
);
return crate::phase_outcome::Outcome::skipped();
}
let stanza_len = op_sequence.stanza_length() as u64;
let spec = phase.cycles.as_deref().unwrap_or("");
let phase_cycles = if spec == "==auto" {
crate::diag!(
crate::observer::LogLevel::Debug,
" cycles: auto ({stanza_len} ops = {stanza_len} cycles)"
);
stanza_len
} else if spec == "===auto" || spec.is_empty() {
stanza_len
} else if let Some(rest) = spec.strip_prefix("==ops:") {
let mut expanded = rest.to_string();
for (v, val) in &iter_var_values {
expanded = expanded.replace(&format!("{{{v}}}"), val);
}
expanded = crate::runner::expand_workload_params(&expanded, &ctx.workload_params);
crate::runner::parse_count(&expanded)
.or_else(|| {
if expanded.starts_with('{') && expanded.ends_with('}') {
let inner = &expanded[1..expanded.len() - 1];
polydat::dsl::compile::eval_const_expr(inner)
.ok()
.map(|v| v.as_u64())
} else {
None
}
})
.unwrap_or(stanza_len)
} else {
let mut expanded = spec.to_string();
for (v, val) in &iter_var_values {
expanded = expanded.replace(&format!("{{{v}}}"), val);
}
expanded = crate::runner::expand_workload_params(&expanded, &ctx.workload_params);
match resolve_stanza_count(&expanded, &activation_scope) {
Ok(stanzas) => stanzas * stanza_len,
Err(e) => {
return crate::phase_outcome::Outcome::failed()
.with_reason(format!("phase '{phase_name}': cycles: {e}"));
}
}
};
if ctx.diag.show_wiring {
let note = if is_iter {
let pairs: Vec<String> = iter_var_values
.iter()
.map(|(k, v)| format!("{k}={v}"))
.collect();
format!(" ({})", pairs.join(", "))
} else {
String::new()
};
crate::describe::print_wiring_analysis(phase_name, ¬e, &iter_op_builder.program());
}
let phase_concurrency = match phase.concurrency.as_ref() {
Some(s) => {
let mut exp = crate::runner::expand_workload_params(s, &ctx.workload_params);
for (v, val) in &iter_var_values {
exp = exp.replace(&format!("{{{v}}}"), val);
}
let parsed = exp.parse::<usize>().ok().or_else(|| {
let bare = exp.trim();
iter_var_values
.get(bare)
.cloned()
.or_else(|| {
ctx.current_parent_kernel
.as_ref()
.and_then(|k| k.lookup(bare))
.map(|v| match v {
polydat::ast::Value::U64(n) => n.to_string(),
polydat::ast::Value::F64(f) => (f as u64).to_string(),
polydat::ast::Value::Str(s) => s.to_string(),
other => other.to_display_string(),
})
})
.or_else(|| ctx.workload_params.get(bare).cloned())
.and_then(|resolved| resolved.trim().parse::<usize>().ok())
});
match parsed {
Some(v) => v,
None => {
return crate::phase_outcome::Outcome::failed().with_reason(format!(
"concurrency '{exp}' is neither an integer nor a resolvable \
name (checked iteration variables, the parent scope, and \
workload params)"
));
}
}
}
None => ctx.concurrency,
};
let phase_labels = match ctx.current_parent_kernel.as_ref() {
Some(parent) => format_scope_coordinate_path(parent.scope_coordinates()),
None => String::new(),
};
let stanza_len = op_sequence.stanza_length();
let activity_name = {
let leaf_label = ctx
.current_parent_kernel
.as_ref()
.and_then(|k| k.scope_coordinates().first())
.filter(|c| !c.is_empty())
.map(|c| {
c.vars
.iter()
.map(|(k, v)| format!("{k}={}", v.to_display_string()))
.collect::<Vec<_>>()
.join(", ")
})
.unwrap_or_default();
if leaf_label.is_empty() {
phase_name.to_string()
} else {
format!("{phase_name} ({leaf_label})")
}
};
let source_factory: Option<Arc<dyn polydat::iteration::source::DataSourceFactory>> = {
let program = iter_op_builder.program();
let schemas = program.cursor_schemas();
if let Some(schema) = schemas.first() {
let extent = runtime_cursor_extents
.get(&schema.name)
.copied()
.or(schema.extent)
.unwrap_or(phase_cycles);
let partitioned = runtime_cursor_partition.contains_key(&schema.name);
let (range_start, range_end) = runtime_cursor_partition
.get(&schema.name)
.copied()
.unwrap_or((0, extent));
let effective_extent = range_end.saturating_sub(range_start);
use polydat::iteration::source as src;
let chunk = extent.min(effective_extent);
let delta = runtime_cursor_delta
.get(&schema.name)
.copied()
.unwrap_or(chunk);
let bound = |f: src::ExtendingRangeSourceFactory| {
if partitioned { f.bounded(range_end) } else { f }
};
match &schema.cursor_kind {
src::CursorKind::Range => Some(Arc::new(src::RangeSourceFactory::named(
&schema.name,
range_start,
range_end,
))
as Arc<dyn src::DataSourceFactory>),
src::CursorKind::ExtendingTimed { .. } => {
let min_ms = runtime_cursor_min_ms
.get(&schema.name)
.copied()
.unwrap_or(0);
let policy: Arc<dyn src::ExtensionPolicy> =
Arc::new(src::UntilElapsedPolicy { min_ms, delta });
Some(Arc::new(bound(src::ExtendingRangeSourceFactory::new(
&schema.name,
range_start,
chunk,
policy,
))) as Arc<dyn src::DataSourceFactory>)
}
src::CursorKind::ExtendingPasses { .. } => {
let min_passes = runtime_cursor_min_passes
.get(&schema.name)
.copied()
.unwrap_or(0);
let policy: Arc<dyn src::ExtensionPolicy> =
Arc::new(src::UntilPassesPolicy { min_passes, delta });
Some(Arc::new(bound(src::ExtendingRangeSourceFactory::new(
&schema.name,
range_start,
chunk,
policy,
))) as Arc<dyn src::DataSourceFactory>)
}
src::CursorKind::ExtendingCount { .. } => {
let min_count = runtime_cursor_min_count
.get(&schema.name)
.copied()
.unwrap_or(0);
let policy: Arc<dyn src::ExtensionPolicy> =
Arc::new(src::UntilCountPolicy { min_count, delta });
Some(Arc::new(bound(src::ExtendingRangeSourceFactory::new(
&schema.name,
range_start,
chunk,
policy,
))) as Arc<dyn src::DataSourceFactory>)
}
src::CursorKind::ExtendingElapsedAndPasses { .. } => {
let min_ms = runtime_cursor_min_ms
.get(&schema.name)
.copied()
.unwrap_or(0);
let min_passes = runtime_cursor_min_passes
.get(&schema.name)
.copied()
.unwrap_or(0);
let elapsed: Arc<dyn src::ExtensionPolicy> =
Arc::new(src::UntilElapsedPolicy { min_ms, delta });
let passes: Arc<dyn src::ExtensionPolicy> =
Arc::new(src::UntilPassesPolicy { min_passes, delta });
let policy: Arc<dyn src::ExtensionPolicy> = Arc::new(src::AndPolicy {
policies: vec![elapsed, passes],
});
Some(Arc::new(bound(src::ExtendingRangeSourceFactory::new(
&schema.name,
range_start,
chunk,
policy,
))) as Arc<dyn src::DataSourceFactory>)
}
src::CursorKind::ExtendingElapsedOrPasses { .. } => {
let min_ms = runtime_cursor_min_ms
.get(&schema.name)
.copied()
.unwrap_or(0);
let min_passes = runtime_cursor_min_passes
.get(&schema.name)
.copied()
.unwrap_or(0);
let elapsed: Arc<dyn src::ExtensionPolicy> =
Arc::new(src::UntilElapsedPolicy { min_ms, delta });
let passes: Arc<dyn src::ExtensionPolicy> =
Arc::new(src::UntilPassesPolicy { min_passes, delta });
let policy: Arc<dyn src::ExtensionPolicy> = Arc::new(src::OrPolicy {
policies: vec![elapsed, passes],
});
Some(Arc::new(bound(src::ExtendingRangeSourceFactory::new(
&schema.name,
range_start,
chunk,
policy,
))) as Arc<dyn src::DataSourceFactory>)
}
}
} else {
None
}
};
let progress_extent = source_factory
.as_ref()
.and_then(|f| f.global_extent())
.unwrap_or(phase_cycles);
let progress_cursor_name = source_factory
.as_ref()
.map(|f| f.schema().name.clone())
.unwrap_or_else(|| "cycles".into());
let progress_fibers = phase_concurrency;
let progress_source_factory = source_factory.clone();
fn live_extent_of(
factory: &Option<Arc<dyn polydat::iteration::source::DataSourceFactory>>,
fallback: u64,
) -> u64 {
factory
.as_ref()
.and_then(|f| f.global_extent())
.unwrap_or(fallback)
}
fn live_rows_of(
factory: &Option<Arc<dyn polydat::iteration::source::DataSourceFactory>>,
) -> (u64, u64) {
match factory {
Some(f) => (f.global_consumed(), f.global_extent().unwrap_or(0)),
None => (0, 0),
}
}
let stanza_len_u64 = stanza_len as u64;
if stanza_len_u64 > 0 && progress_extent > 0 && !progress_extent.is_multiple_of(stanza_len_u64)
{
let remainder = progress_extent % stanza_len_u64;
let full_stanzas = progress_extent / stanza_len_u64;
crate::diag!(
crate::observer::LogLevel::Warn,
"phase '{phase_name}': cursor extent ({progress_extent}) is not an \
even multiple of stanza length ({stanza_len_u64}) — boundary stanza \
will be {remainder}/{stanza_len_u64} of a full stride after \
{full_stanzas} clean stanza(s). If the op-sequence assumes \
complete stanzas (e.g. for relevancy or aggregation evaluation), \
align by sizing the cursor to a multiple of {stanza_len_u64}, or \
by adjusting stanza_concurrency / ops."
);
}
crate::scene_tree::with_global_mut(|t| {
t.set_phase_running_at(scene_node_id, stanza_len);
});
ctx.observer.phase_starting(
scene_node_id,
phase_name,
&phase_labels,
stanza_len,
progress_extent,
phase_concurrency,
);
{
let display_labels: String = {
let parent_coords: Vec<_> = ctx
.current_parent_kernel
.as_ref()
.map(|k| k.scope_coordinates().iter().rev().cloned().collect())
.unwrap_or_default();
polydat::kernel::format_scope_coordinate_path(&parent_coords)
};
let depth_indent = crate::scene_tree::running_phase_indent();
let phase_ctx = crate::readout_context::LifecycleContext {
event: crate::lifecycle::EventType::PhaseStart,
subject_name: phase_name.to_string(),
subject_labels: display_labels.clone(),
depth_indent,
use_color: crate::observer::use_color(),
stick_reattached: String::new(),
};
crate::readout_context::fire_lifecycle(
crate::lifecycle::EventType::PhaseStart,
&ctx.workload_readouts,
None,
&phase_ctx,
Some(&ctx.sqlite_reporter),
);
}
let phase_config_text = crate::checkpoint::phase_config_canonical_text(&phase);
let phase_hash_bytes = crate::checkpoint::compose_phase_hash(
crate::checkpoint::ancestor_chain_hash(&ctx.scope_tree, phase_name),
crate::checkpoint::config_text_hash(&phase_config_text),
);
let phase_hash_hex: String = phase_hash_bytes
.iter()
.map(|b| format!("{b:02x}"))
.collect();
let params_consumed_json: Option<String> = {
let phase_idx = ctx.scope_tree.phase_node_by_name(phase_name);
let ancestors_below = phase_idx
.map(|idx| ctx.scope_tree.ancestor_kernels_split(idx).0)
.unwrap_or_default();
let op_templates: Vec<std::sync::Arc<polydat::kernel::PolydatProgram>> = phase_idx
.map(|idx| {
ctx.scope_tree
.op_template_programs_for_phase(idx)
.into_values()
.collect()
})
.unwrap_or_default();
let consumed = crate::checkpoint::params_scope::consumed_params(
&iter_op_builder.program(),
&op_templates,
&ancestors_below,
&phase_config_text,
&ctx.workload_params,
);
serde_json::to_string(&consumed).ok()
};
if let Some(writer) = ctx.checkpoint_writer.clone() {
let identity = phase_identity_for(phase_name, &phase_labels);
writer.update_phase_hash(&identity, phase_hash_bytes, params_consumed_json.clone());
if ctx.resume_plan.is_resume {
ctx.push_label("phase", phase_name);
let labels_for_purge = ctx.labels();
ctx.pop_label();
let mut guard = ctx
.sqlite_reporter
.lock()
.unwrap_or_else(|e| e.into_inner());
if let Some(reporter) = guard.as_mut() {
let n = reporter.purge_samples_with_labels(&labels_for_purge);
if n > 0 {
crate::diag!(
crate::observer::LogLevel::Info,
"resume: purged {n} prior sample rows for phase '{phase_name}'"
);
}
}
}
writer.phase_started(&identity);
if let Err(e) = writer.flush() {
crate::diag!(
crate::observer::LogLevel::Warn,
"checkpoint flush at phase '{phase_name}' start: {e}"
);
}
}
let explicitly_selected = ctx
.phase_filter
.as_ref()
.map(|pat| pat.is_match(phase_name))
.unwrap_or(false);
let is_prereq_class = phase
.checkpoint
.as_ref()
.map(|c| c.idempotent)
.unwrap_or(false);
if let Some(plan) = ctx.refine_plan.as_ref()
&& !explicitly_selected
&& (plan.scope == crate::refine_plan::RefineScope::Changed
|| (plan.scope == crate::refine_plan::RefineScope::Missing && is_prereq_class))
{
match plan.unchanged_verdict(
phase_name,
&phase_labels,
&phase_hash_hex,
&ctx.workload_params,
) {
Ok(()) => {
crate::diag!(
crate::observer::LogLevel::Info,
"refine: skipping phase '{phase_name}' [{phase_labels}] \
(prior completed outcome, hash unchanged)"
);
crate::scene_tree::with_global_mut(|t| {
t.set_phase_running_at(scene_node_id, 0);
t.set_phase_completed_at(scene_node_id, 0.0);
});
ctx.observer
.phase_starting(scene_node_id, phase_name, &phase_labels, 0, 0, 0);
ctx.observer
.phase_completed(scene_node_id, phase_name, &phase_labels, 0.0);
crate::phase_end_triggers::fire_phase_completed(phase_name, &phase_labels, 0.0);
return crate::phase_outcome::Outcome::skipped();
}
Err(crate::refine_plan::SkipBlocker::NoPrior) => {}
Err(blocker) => {
crate::diag!(
crate::observer::LogLevel::Info,
"refine: re-running phase '{phase_name}' \
[{phase_labels}] — {blocker}"
);
}
}
}
let timeout_guard = match phase.timeout.as_deref() {
None => None,
Some(raw) => {
let mut t = crate::runner::expand_workload_params(raw, &ctx.workload_params);
for (v, val) in &iter_var_values {
t = t.replace(&format!("{{{v}}}"), val);
}
match crate::timeval::parse_time_ms(&t) {
Ok(ms) => {
let guard = crate::stop_conditions::StopConditionDecl::timeout_guard(ms);
crate::diag!(
crate::observer::LogLevel::Info,
"phase '{phase_name}': timeout={t} → stop_when: {} \
(synthesized)",
guard.when
);
Some(guard)
}
Err(e) => {
return crate::phase_outcome::Outcome::failed().with_reason(format!(
"phase '{phase_name}': `timeout: \"{raw}\"` \
(resolved: \"{t}\"): {e}"
));
}
}
}
};
let phase_rate: Option<f64> = match phase.rate.as_deref() {
None => None,
Some(raw) => {
let mut r = crate::runner::expand_workload_params(raw, &ctx.workload_params);
for (v, val) in &iter_var_values {
r = r.replace(&format!("{{{v}}}"), val);
}
match r.trim().parse::<f64>() {
Ok(f) => Some(f),
Err(_) => {
return crate::phase_outcome::Outcome::failed().with_reason(format!(
"phase '{phase_name}': `rate: \"{raw}\"` (resolved: \
\"{r}\") is not a number of ops/sec"
));
}
}
}
};
let config = ActivityConfig {
name: activity_name,
cycles: phase_cycles,
concurrency: phase_concurrency,
rate: phase_rate.or(ctx.rate),
sequencer: ctx.seq_type,
error_spec: phase
.errors
.clone()
.unwrap_or_else(|| ctx.error_spec.clone()),
error_rate_max: phase.error_rate_max.or(ctx.error_rate_max),
throttle: phase.throttle.as_ref().and_then(|t| t.to_spec()),
stop_when: phase
.error_rate_max
.or(ctx.error_rate_max)
.map(|max| {
let guard = crate::stop_conditions::StopConditionDecl::error_rate_guard(max);
crate::diag!(
crate::observer::LogLevel::Info,
"phase '{phase_name}': error_rate_max={max} → stop_when: {} \
(synthesized)",
guard.when
);
guard
})
.into_iter()
.chain(timeout_guard)
.chain(
phase
.stop_when
.iter()
.filter(|c| {
c.each.iter().any(|l| {
matches!(
l,
nmbrs_workload::model::ScopeLevel::SelfScope
| nmbrs_workload::model::ScopeLevel::Phase
)
})
})
.chain(ctx.workload_stop_when.iter().filter(|c| {
c.each
.iter()
.any(|l| matches!(l, nmbrs_workload::model::ScopeLevel::Phase))
}))
.map(|c| {
let mut when =
crate::runner::expand_workload_params(&c.when, &ctx.workload_params);
for (v, val) in &iter_var_values {
when = when.replace(&format!("{{{v}}}"), val);
}
crate::stop_conditions::StopConditionDecl {
when,
effect: crate::stop_conditions::StopConditionDecl::effect_from_str(
c.effect.as_deref(),
crate::phase_outcome::Outcome::failed(),
),
reason: None,
target: resolve_stop_scope(c.at, &c.each),
cancel_ops:
crate::stop_conditions::StopConditionDecl::action_cancels_ops(
c.effect.as_deref(),
),
}
}),
)
.collect(),
tries: phase.tries.or(ctx.tries),
tries_backoff: phase.tries_backoff.clone(),
stanza_concurrency: 1,
source_factory,
suppress_status_line: ctx.observer.live_suppress_flag().unwrap_or_else(|| {
Arc::new(std::sync::atomic::AtomicBool::new(
ctx.observer.suppresses_stderr(),
))
}),
status_metrics: phase.status_metrics.clone(),
phase_labels: {
let parent_coords: Vec<_> = ctx
.current_parent_kernel
.as_ref()
.map(|k| k.scope_coordinates().iter().rev().cloned().collect())
.unwrap_or_default();
polydat::kernel::format_scope_coordinate_path(&parent_coords)
},
phase_seq: crate::scene_tree::current().and_then(|t| {
t.nodes
.get(scene_node_id)
.and_then(|n| n.seq)
.map(|s| (s, t.total_phases()))
}),
readouts: ctx.workload_readouts.clone(),
cli_readout_override: ctx.cli_readout_override.clone(),
snapshot_writer: Some(ctx.sqlite_reporter.clone()),
dry_run_mode: ctx.dry_run.map(String::from),
stop_after_dispenser_init: ctx.diag.depth < crate::runner::ExecDepth::Cycle,
};
let phase_driver_owned = phase.adapter.clone().unwrap_or_else(|| ctx.driver.clone());
let phase_driver = phase_driver_owned.as_str();
let mut adapter_names = std::collections::HashSet::new();
adapter_names.insert(phase_driver.to_string());
for t in op_sequence.templates() {
if let Some(a) = t.params.get("adapter").and_then(|v| v.as_str())
&& a != phase_driver
{
adapter_names.insert(a.to_string());
}
}
let mut adapters: HashMap<String, Arc<dyn DriverAdapter>> = HashMap::new();
let mut attach_guards: Vec<crate::resource_pool::AttachGuard> = Vec::new();
let phase_seq_label = format!("{:?}", ctx.label_stack);
for aname in &adapter_names {
let aname_owned = aname.clone();
let merged_params = ctx.merged_params.clone();
let dry_run = ctx.dry_run;
let aname_for_factory = aname_owned.clone();
let shared_reg = if dry_run.is_none() {
let driver_name = crate::adapter::resolve_driver_name(
&aname_owned,
&resolve_selector_param(&aname_owned),
&merged_params,
);
driver_name.and_then(|d| crate::adapter::find_shared_driver(&aname_owned, d))
} else {
None
};
let (adapter, guard) = if let Some(reg) = shared_reg {
let key = match (reg.resource_key)(&merged_params) {
Ok(v) => v,
Err(e) => return crate::phase_outcome::Outcome::failed().with_reason(e),
};
match crate::resource_pool::attach_shared_adapter(
&ctx.resource_pool,
&aname_owned,
phase_name,
key,
move || async move {
crate::runner::create_adapter(&aname_for_factory, &merged_params).await
},
)
.await
{
Ok(v) => v,
Err(e) => return crate::phase_outcome::Outcome::failed().with_reason(e),
}
} else {
match crate::resource_pool::attach_legacy_adapter(
&ctx.resource_pool,
&aname_owned,
phase_name,
&[
("__phase", phase_name),
("__phase_seq", phase_seq_label.as_str()),
],
move || async move {
crate::runner::create_adapter(&aname_for_factory, &merged_params).await
},
)
.await
{
Ok(v) => v,
Err(e) => return crate::phase_outcome::Outcome::failed().with_reason(e),
}
};
adapters.insert(aname_owned.clone(), adapter);
attach_guards.push(guard);
}
ctx.push_label("phase", phase_name);
let labels = ctx.labels();
let phase_own_labels = ctx.incremental_labels();
ctx.pop_label();
let phase_component = Arc::new(RwLock::new(Component::new(
phase_own_labels,
HashMap::new(),
)));
component::attach(&ctx.session_component, &phase_component);
let sigdigs = nmbrs_metrics::instruments::histogram::resolve_hdr_sigdigs(
&phase_component.read().unwrap_or_else(|e| e.into_inner()),
);
let phase_error_policy =
ctx.error_policy
.resolve_child(Some(crate::error_policy::PolicyConfig::new(
config.error_spec.clone(),
config.error_rate_max,
)));
let phase_kernel = Some(activation_scope.clone());
let metric_detail = crate::activity::metric_detail_from_params(&ctx.merged_params);
let mut activity = Activity::with_params_and_sigdigs(
config,
&labels,
op_sequence,
ctx.workload_params.clone(),
sigdigs,
phase_error_policy,
phase_kernel,
&metric_detail,
);
activity.walk_stop = Some(ctx.workload_shell.walk_stop_flag());
activity.daemon_stop = ctx.daemon_stop.clone();
activity.set_wrappers_override(ctx.wrappers_override.clone());
activity.set_wrap_default_order(ctx.wrap_default_order.clone());
activity.attach_component(phase_component.clone());
{
let mut pc = phase_component.write().unwrap_or_else(|e| e.into_inner());
pc.set_state(ComponentState::Running);
}
crate::activity::declare_adapter_controls(&adapters, &phase_component);
if ctx.diag.depth < crate::runner::ExecDepth::Dispenser {
ctx.observer
.phase_completed(scene_node_id, phase_name, &phase_labels, 0.0);
crate::phase_end_triggers::fire_phase_completed(phase_name, &phase_labels, 0.0);
crate::scene_tree::with_global_mut(|t| {
t.set_phase_completed_at(scene_node_id, 0.0);
});
return crate::phase_outcome::Outcome::skipped();
}
let validation_frame = activity.validation_frame.clone();
let observer_for_progress = ctx.observer.clone();
let progress_metrics = activity.shared_metrics();
let progress_start = std::time::Instant::now();
let progress_running = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(true));
let progress_flag = progress_running.clone();
{
let (seq, depth_indent) = crate::readout_context::resolve_phase_coord_by_id(scene_node_id);
let update_bodies = {
let phase_status_default = {
let readout = crate::readouts::Registry::lookup("phase_status")
.expect("phase_status registered");
crate::readouts::BakedBody::from_single(readout, crate::readouts::Lod::Labeled)
};
match crate::readouts::build_event_binder_with_cli(
&activity.config.readouts,
crate::lifecycle::EventType::Update,
phase_status_default,
activity.config.cli_readout_override.as_deref(),
) {
Ok(mut binder) => binder.take_bodies(crate::lifecycle::EventType::Update),
Err(e) => {
crate::diag!(
crate::observer::LogLevel::Warn,
"readouts: failed to bind on_update — {e}"
);
Vec::new()
}
}
};
observer_for_progress.phase_render_attach(crate::observer::PhaseRenderHandle {
exec_id: crate::execution_context::current_exec_id(),
name: phase_name.to_string(),
labels: phase_labels.clone(),
activity_name: activity.config.name.clone(),
metrics: progress_metrics.clone(),
bodies: std::sync::Arc::new(update_bodies),
memo: activity.memo.clone(),
gutter: activity.gutter.clone(),
status_metrics: activity.config.status_metrics.clone().into(),
concurrency: activity.config.concurrency,
seq,
depth_indent,
});
}
if observer_for_progress.suppresses_stderr() {
let (rows_consumed, rows_total) = live_rows_of(&progress_source_factory);
observer_for_progress.phase_progress(&crate::observer::PhaseProgressUpdate {
exec_id: crate::execution_context::current_exec_id(),
name: phase_name.to_string(),
labels: phase_labels.clone(),
cursor_name: progress_cursor_name.clone(),
cursor_extent: live_extent_of(&progress_source_factory, progress_extent),
daemon: phase.daemon,
rows_consumed,
rows_total,
fibers: progress_fibers,
ops_started: 0,
ops_finished: 0,
ops_ok: 0,
skips: 0,
errors: 0,
retries: 0,
ops_per_sec: 0.0,
adapter_counters: Vec::new(),
rows_per_batch: 0.0,
relevancy: Vec::new(),
});
}
let _progress_thread = if observer_for_progress.suppresses_stderr() {
let obs = observer_for_progress.clone();
let cursor_name_for_thread = progress_cursor_name.clone();
let fibers_for_thread = progress_fibers;
let name_for_thread = phase_name.to_string();
let labels_for_thread = phase_labels.clone();
let progress_metrics = progress_metrics.clone();
let factory_for_thread = progress_source_factory.clone();
let daemon_for_thread = phase.daemon;
Some(std::thread::spawn(move || {
let progress_cursor_name = cursor_name_for_thread;
let progress_fibers = fibers_for_thread;
let phase_name = name_for_thread;
let phase_labels = labels_for_thread;
while progress_flag.load(std::sync::atomic::Ordering::Relaxed) {
std::thread::sleep(std::time::Duration::from_millis(500));
if !progress_flag.load(std::sync::atomic::Ordering::Relaxed) {
break;
}
let started = progress_metrics
.ops_started
.load(std::sync::atomic::Ordering::Relaxed);
let finished = progress_metrics
.ops_finished
.load(std::sync::atomic::Ordering::Relaxed);
let successes = progress_metrics.result_success.count();
let errors = progress_metrics.errors_total.get();
let elapsed = progress_start.elapsed().as_secs_f64();
let ops_per_sec = if elapsed > 0.0 {
finished as f64 / elapsed
} else {
0.0
};
let adapter_counters: Vec<(String, u64, f64)> = progress_metrics
.collect_status_counters()
.into_iter()
.map(|(name, total)| {
let rate = if elapsed > 0.0 {
total as f64 / elapsed
} else {
0.0
};
(name, total, rate)
})
.collect();
let stanzas = progress_metrics.stanzas_total.get();
let find_counter = |want: &str| {
adapter_counters
.iter()
.find(|(n, _, _)| n == want)
.map(|(_, t, _)| *t)
};
let rows_per_batch = crate::readout_context::rows_per_batch(
find_counter("rows_inserted"),
find_counter("_batch_writes"),
stanzas,
)
.unwrap_or(0.0);
let relevancy = progress_metrics.collect_relevancy_live();
let (rows_consumed, rows_total) = live_rows_of(&factory_for_thread);
obs.phase_progress(&crate::observer::PhaseProgressUpdate {
exec_id: crate::execution_context::current_exec_id(),
name: phase_name.clone(),
labels: phase_labels.clone(),
cursor_name: progress_cursor_name.clone(),
cursor_extent: live_extent_of(&factory_for_thread, progress_extent),
daemon: daemon_for_thread,
rows_consumed,
rows_total,
fibers: progress_fibers,
ops_started: started,
ops_finished: finished,
ops_ok: successes,
skips: progress_metrics.skips_total.get(),
errors,
retries: errors.saturating_sub(finished.saturating_sub(successes)),
ops_per_sec,
adapter_counters,
rows_per_batch,
relevancy,
});
}
}))
} else {
None
};
if let Some(poll_spec) = phase.poll.as_ref() {
let resolve_u64 = |raw: Option<&str>, field: &str, default: u64| -> Result<u64, String> {
let Some(raw) = raw else { return Ok(default) };
let mut r = crate::runner::expand_workload_params(raw, &ctx.workload_params);
for (v, val) in &iter_var_values {
r = r.replace(&format!("{{{v}}}"), val);
}
r.trim().parse::<u64>().map_err(|_| {
format!(
"phase '{phase_name}': poll `{field}: \"{raw}\"` (resolved: \
\"{r}\") is not a whole number of milliseconds"
)
})
};
let resolved = (|| -> Result<(u64, u64, u64), String> {
Ok((
resolve_u64(poll_spec.interval_ms.as_deref(), "interval_ms", 1000)?,
resolve_u64(poll_spec.timeout_ms.as_deref(), "timeout_ms", 300_000)?,
resolve_u64(
poll_spec.max_error_retries.as_deref(),
"max_error_retries",
0,
)?,
))
})();
let (interval_ms, timeout_ms, max_error_retries) = match resolved {
Ok(v) => v,
Err(e) => return crate::phase_outcome::Outcome::failed().with_reason(e),
};
let interval = std::time::Duration::from_millis(interval_ms);
let timeout = std::time::Duration::from_millis(timeout_ms);
let phase_kernel = match ctx
.scope_tree
.phase_node_by_name(phase_name)
.and_then(|idx| ctx.scope_tree.nodes[idx].cached_kernel.get().cloned())
{
Some(k) => k,
None => {
return crate::phase_outcome::Outcome::failed().with_reason(format!(
"phase '{phase_name}': SRD-75 phase-poll requires the phase \
scope kernel to be installed, but no cached kernel was found. \
This is a synthesis bug — every phase with `poll:` should \
land via the `Bindings` install spec.",
));
}
};
let started_at = std::time::Instant::now();
let on_timeout_policy = match poll_spec.on_timeout.as_deref() {
Some("abort") => crate::activity::PhasePollTimeoutPolicy::Abort,
_ => crate::activity::PhasePollTimeoutPolicy::Error,
};
activity.phase_poll = Some(crate::activity::PhasePollContext {
kernel: phase_kernel,
interval,
deadline: started_at + timeout,
started_at,
metric_name: poll_spec.metric_name.clone(),
max_error_retries: max_error_retries.min(u32::MAX as u64) as u32,
on_timeout: on_timeout_policy,
require: poll_spec.require.clone(),
require_grace: started_at
+ ctx.cadence_reporter.declared_cadences().smallest()
+ interval,
});
}
crate::diag!(
crate::observer::LogLevel::Debug,
"phase '{phase_name}': activity starting (concurrency={phase_concurrency})"
);
let stop_reason = activity.stop_reason.clone();
let activity_stop_outcome = activity.stop_outcome.clone();
let activity_phase_errors = activity.phase_errors.clone();
let settle = ctx.optimize_objective.clone().and_then(|obj| {
let parent = ctx.current_parent_kernel.clone()?;
let phase_kernel = ctx
.scope_tree
.phase_node_by_name(phase_name)
.and_then(|idx| ctx.scope_tree.nodes[idx].cached_kernel.get().cloned())?;
match crate::optimize::settle::start_settle(
&parent,
&phase_kernel,
&obj,
&ctx.cadence_reporter,
activity.stop_flag.clone(),
) {
Ok(handle) => Some(handle),
Err(
crate::optimize::settle::SettleSkip::NotWindowed
| crate::optimize::settle::SettleSkip::CadenceDisabled,
) => None,
Err(e @ crate::optimize::settle::SettleSkip::Failed(_)) => {
crate::diag!(
crate::observer::LogLevel::Warn,
"phase '{phase_name}': objective '{obj}' is read once at completion: {e}"
);
None
}
}
});
let mut servo_completed = false;
let stopped = if let Some(servo_spec) = ctx.optimize_servo.take() {
let parent = ctx.current_parent_kernel.clone();
let phase_kernel = ctx
.scope_tree
.phase_node_by_name(phase_name)
.and_then(|idx| ctx.scope_tree.nodes[idx].cached_kernel.get().cloned());
match (parent, phase_kernel) {
(Some(parent), Some(phase_kernel)) => {
let stop_flag = activity.stop_flag.clone();
let reporter = ctx.cadence_reporter.clone();
let pc = phase_component.clone();
let phase_done = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
let pd = phase_done.clone();
let act = async {
let s = crate::runner::run_activity_simple(
activity,
adapters,
phase_driver,
iter_op_builder,
)
.await;
pd.store(true, std::sync::atomic::Ordering::Relaxed);
s
};
let servoed = crate::optimize::servo::servo(
servo_spec,
stop_flag,
reporter,
parent,
phase_kernel,
pc,
phase_done,
);
let (stopped, servo_res) = tokio::join!(act, servoed);
match servo_res {
Ok(()) => servo_completed = true,
Err(e) => crate::diag!(
crate::observer::LogLevel::Warn,
"phase '{phase_name}': optimizer servoing error: {e}"
),
}
stopped
}
_ => {
crate::runner::run_activity_simple(
activity,
adapters,
phase_driver,
iter_op_builder,
)
.await
}
}
} else {
crate::runner::run_activity_simple(activity, adapters, phase_driver, iter_op_builder).await
};
crate::diag!(
crate::observer::LogLevel::Debug,
"phase '{phase_name}': activity returned (stopped={stopped})"
);
let settle_succeeded = settle.as_ref().is_some_and(|h| {
matches!(
&**h.outcome.load(), Some(o) if o.validity == crate::phase_outcome::Validity::Succeeded
)
});
if let Some(h) = &settle {
ctx.cadence_reporter.unsubscribe(h.subscriber);
}
progress_running.store(false, std::sync::atomic::Ordering::Relaxed);
if ctx.diag.depth < crate::runner::ExecDepth::Cycle {
ctx.observer
.phase_completed(scene_node_id, phase_name, &phase_labels, 0.0);
crate::phase_end_triggers::fire_phase_completed(phase_name, &phase_labels, 0.0);
crate::scene_tree::with_global_mut(|t| {
t.set_phase_completed_at(scene_node_id, 0.0);
});
return crate::phase_outcome::Outcome::skipped();
}
if ctx.observer.suppresses_stderr() {
let started_total = progress_metrics
.ops_started
.load(std::sync::atomic::Ordering::Relaxed);
let finished_total = progress_metrics
.ops_finished
.load(std::sync::atomic::Ordering::Relaxed);
let successes = progress_metrics.result_success.count();
let errors = progress_metrics.errors_total.get();
let elapsed = progress_start.elapsed().as_secs_f64();
let ops_per_sec = if elapsed > 0.0 {
finished_total as f64 / elapsed
} else {
0.0
};
let adapter_counters: Vec<(String, u64, f64)> = progress_metrics
.collect_status_counters()
.into_iter()
.map(|(name, total)| {
let rate = if elapsed > 0.0 {
total as f64 / elapsed
} else {
0.0
};
(name, total, rate)
})
.collect();
let stanzas = progress_metrics.stanzas_total.get();
let find_counter = |want: &str| {
adapter_counters
.iter()
.find(|(n, _, _)| n == want)
.map(|(_, t, _)| *t)
};
let rows_per_batch = crate::readout_context::rows_per_batch(
find_counter("rows_inserted"),
find_counter("_batch_writes"),
stanzas,
)
.unwrap_or(0.0);
let relevancy = progress_metrics.collect_relevancy_live();
let (rows_consumed, rows_total) = live_rows_of(&progress_source_factory);
ctx.observer
.phase_progress(&crate::observer::PhaseProgressUpdate {
exec_id: crate::execution_context::current_exec_id(),
name: phase_name.to_string(),
labels: phase_labels.clone(),
cursor_name: progress_cursor_name.clone(),
cursor_extent: live_extent_of(&progress_source_factory, progress_extent),
daemon: phase.daemon,
rows_consumed,
rows_total,
fibers: progress_fibers,
ops_started: started_total,
ops_finished: finished_total,
ops_ok: successes,
skips: progress_metrics.skips_total.get(),
errors,
retries: errors.saturating_sub(finished_total.saturating_sub(successes)),
ops_per_sec,
adapter_counters,
rows_per_batch,
relevancy,
});
}
let graceful_stop = stopped
&& activity_stop_outcome
.lock()
.ok()
.and_then(|g| g.clone())
.is_some_and(|o| !o.is_failure());
if (!stopped || graceful_stop) && !phase.metrics.is_empty() {
let phase_start_epoch_ms = (phase_start_nanos / 1_000_000) as u64;
if let Err(e) = emit_phase_metrics(
&ctx.scope_tree,
phase_name,
&phase,
phase_start_epoch_ms,
&phase_component,
) {
crate::diag!(
crate::observer::LogLevel::Warn,
"phase '{phase_name}': phase-metric emission: {e}"
);
}
}
ctx.optimize_objective_value = None;
if let Some(h) = &settle {
ctx.optimize_objective_value = Some(h.register.load().value);
} else if !stopped
&& let (Some(obj), Some(parent), Some(phase_kernel)) = (
ctx.optimize_objective.clone(),
ctx.current_parent_kernel.clone(),
ctx.scope_tree
.phase_node_by_name(phase_name)
.and_then(|idx| ctx.scope_tree.nodes[idx].cached_kernel.get().cloned()),
)
{
ctx.optimize_objective_value = read_objective_at_completion(&parent, &phase_kernel, &obj);
}
{
let mut final_delta = phase_component
.read()
.unwrap_or_else(|e| e.into_inner())
.capture_delta_auto(std::time::Duration::from_secs(1));
final_delta.mark_partial();
ctx.cadence_reporter.ingest(&labels, final_delta.clone());
ctx.stop_handle.report_frame(&final_delta);
if let Some(mut vframe) = validation_frame
.lock()
.unwrap_or_else(|e| e.into_inner())
.take()
{
vframe.mark_partial();
if crate::observer::trace_enabled() {
crate::observer::trace(
&labels,
&format!(
"event=validation_frame.ingest family_count={}",
vframe.len()
),
);
}
ctx.cadence_reporter.ingest(&labels, vframe.clone());
ctx.stop_handle.report_frame(&vframe);
}
ctx.cadence_reporter.close_path(&labels);
}
{
let mut pc = phase_component.write().unwrap_or_else(|e| e.into_inner());
pc.set_state(ComponentState::Stopped);
}
let phase_duration = phase_start.elapsed().as_secs_f64();
let phase_op_count = progress_metrics.cycles_completed();
let phase_error_count = progress_metrics.result_failure.count();
if stopped && !settle_succeeded && !servo_completed && !graceful_stop {
let reason = stop_reason
.lock()
.ok()
.and_then(|g| g.clone())
.unwrap_or_else(|| "stopped by error handler".to_string());
let detail_msg = format!("stopped by error handler: {reason}");
let depth_indent = crate::scene_tree::running_phase_indent();
let color = crate::observer::use_color();
let bold = if color { "\x1b[1m" } else { "" };
let red = if color { "\x1b[31m" } else { "" };
let dim = if color { "\x1b[2m" } else { "" };
let reset = if color { "\x1b[0m" } else { "" };
let error_head_consumed: usize = depth_indent.chars().count()
+ "phase '".chars().count()
+ phase_name.chars().count()
+ "' ".chars().count();
let coords_part = crate::readouts::builtins::phase_outcome::format_coords_block(
&phase_labels,
color,
error_head_consumed,
&format!("{depth_indent} "),
true,
);
crate::diag!(
crate::observer::LogLevel::Error,
"{depth_indent}phase '{bold}{phase_name}{reset}'{coords_part} {red}{detail_msg}{reset} {dim}({phase_duration:.2}s){reset}"
);
ctx.observer
.phase_failed(scene_node_id, phase_name, &phase_labels, &detail_msg);
crate::phase_end_triggers::fire_phase_failed(phase_name, &phase_labels, &detail_msg);
let errors = activity_phase_errors
.lock()
.ok()
.map(|mut g| std::mem::take(&mut *g))
.unwrap_or_default();
let errors = if errors.is_empty() {
vec![crate::phase_outcome::PhaseErrorDetail {
class: "phase_failed".into(),
message: reason.clone(),
op_name: None,
cycle: None,
op_template: None,
op_resolved: None,
at_nanos: std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_nanos() as u64)
.unwrap_or(0),
retryable: false,
}]
} else {
errors
};
let outcome = crate::phase_outcome::PhaseOutcome::failed(
crate::phase_outcome::PhaseIdentity::new(phase_name, phase_labels.as_str()),
phase_duration,
errors,
)
.with_phase_hash(phase_hash_hex.clone())
.with_params_consumed(params_consumed_json.clone());
if let Ok(mut guard) = ctx.sqlite_reporter.lock()
&& let Some(reporter) = guard.as_mut()
{
let row = outcome.to_sqlite_row(&ctx.session_id, ctx.exec_id, phase_start_nanos);
reporter.write_phase_outcome(&row);
}
crate::scene_tree::with_global_mut(|t| {
t.set_phase_failed_at(scene_node_id, &detail_msg);
t.set_phase_outcome_at(scene_node_id, outcome);
});
if let Some(writer) = ctx.checkpoint_writer.as_ref() {
let identity = phase_identity_for(phase_name, &phase_labels);
writer.phase_failed(&identity, &detail_msg);
if let Err(e) = writer.flush() {
crate::diag!(
crate::observer::LogLevel::Warn,
"checkpoint flush after phase '{phase_name}' failed: {e}"
);
}
}
if let Some((outcome, reason)) =
ctx.workload_shell
.record_phase(true, phase_op_count, phase_error_count)
{
let cause = if outcome.is_failure() {
crate::session_signals::StopCause::Fault
} else {
crate::session_signals::StopCause::Interrupt
};
crate::session_signals::request_shell_stop(cause);
crate::diag!(
crate::observer::LogLevel::Warn,
"scenario stop-on-error ({reason}) after phase \
'{phase_name}' — halting remaining walk"
);
}
return crate::phase_outcome::Outcome::failed()
.with_reason(format!("phase '{phase_name}' {detail_msg}"));
}
ctx.observer
.phase_completed(scene_node_id, phase_name, &phase_labels, phase_duration);
crate::phase_end_triggers::fire_phase_completed(phase_name, &phase_labels, phase_duration);
let success_errors = activity_phase_errors
.lock()
.ok()
.map(|mut g| std::mem::take(&mut *g))
.unwrap_or_default();
let success_outcome = crate::phase_outcome::PhaseOutcome {
phase_id: crate::phase_outcome::PhaseIdentity::new(phase_name, phase_labels.as_str()),
disposition: if graceful_stop {
crate::phase_outcome::Disposition::Interrupted
} else {
crate::phase_outcome::Disposition::Completed
},
validity: crate::phase_outcome::Validity::Succeeded,
duration_secs: phase_duration,
errors: success_errors,
resume_cursor: None,
phase_hash: Some(phase_hash_hex.clone()),
params_consumed: params_consumed_json.clone(),
};
if let Ok(mut guard) = ctx.sqlite_reporter.lock()
&& let Some(reporter) = guard.as_mut()
{
let row = success_outcome.to_sqlite_row(&ctx.session_id, ctx.exec_id, phase_start_nanos);
reporter.write_phase_outcome(&row);
}
crate::scene_tree::with_global_mut(|t| {
t.set_phase_completed_at(scene_node_id, phase_duration);
t.set_phase_outcome_at(scene_node_id, success_outcome);
});
if let Some(writer) = ctx.checkpoint_writer.as_ref() {
let identity = phase_identity_for(phase_name, &phase_labels);
writer.phase_completed(&identity, phase_duration);
if let Err(e) = writer.flush() {
crate::diag!(
crate::observer::LogLevel::Warn,
"checkpoint flush after phase '{phase_name}' completed: {e}"
);
}
}
if let Some((outcome, reason)) =
ctx.workload_shell
.record_phase(false, phase_op_count, phase_error_count)
{
let actual = ctx.workload_shell.describe_state();
if outcome.is_failure() {
crate::diag!(
crate::observer::LogLevel::Error,
"workload stop condition tripped ({reason}) — actual: {actual} \
— after phase '{phase_name}' — failing session"
);
return crate::phase_outcome::Outcome::failed().with_reason(format!(
"workload stop condition tripped: {reason} — actual: {actual}"
));
}
crate::session_signals::request_graceful_stop();
crate::diag!(
crate::observer::LogLevel::Warn,
"workload stop condition tripped ({reason}) — actual: {actual} \
— after phase '{phase_name}' — halting remaining walk"
);
}
if graceful_stop {
let mut oc = crate::phase_outcome::Outcome::interrupted();
if let Some(r) = stop_reason.lock().ok().and_then(|g| g.clone()) {
oc = oc.with_reason(r);
}
return oc;
}
crate::phase_outcome::Outcome::completed()
}
fn emit_phase_metrics(
scope_tree: &crate::scope_tree::ScopeTree,
phase_name: &str,
phase: &nmbrs_workload::model::WorkloadPhase,
phase_start_epoch_ms: u64,
phase_component: &Arc<RwLock<nmbrs_metrics::component::Component>>,
) -> Result<(), String> {
use polydat::ast::Value;
fn to_f64(v: &Value) -> Option<f64> {
match v {
Value::F64(f) => Some(*f),
Value::U64(u) => Some(*u as f64),
Value::Bool(b) => Some(if *b { 1.0 } else { 0.0 }),
_ => None,
}
}
let phase_kernel = scope_tree
.phase_node_by_name(phase_name)
.and_then(|idx| scope_tree.nodes[idx].cached_kernel.get().cloned())
.ok_or_else(|| {
format!(
"phase '{phase_name}' has `metrics:` but no cached phase scope \
kernel was installed — synthesis bug (a phase with metrics \
classifies as PolydatMatter::Definitions and must install a \
kernel via the `Bindings` install spec)"
)
})?;
let metric_pull_error = |e: String| format!("phase '{phase_name}': metric-pull subscope: {e}");
let mut k = phase_kernel
.bind_under(phase_kernel.kernel(), &[])
.map_err(metric_pull_error)?;
polydat::kernel::propagate_inputs(phase_kernel.kernel(), k.kernel_mut())
.map_err(|e| metric_pull_error(e.to_string()))?;
if let Some(idx) = k.program().find_input("phase_start") {
crate::wires::write_input(
k.kernel_mut(),
idx,
"phase_start",
Value::U64(phase_start_epoch_ms),
)
.map_err(|e| metric_pull_error(e.to_string()))?;
k.kernel_mut()
.init()
.map_err(|e| metric_pull_error(e.to_string()))?;
}
let mut entries: Vec<_> = phase.metrics.iter().collect();
entries.sort_by(|a, b| a.0.cmp(b.0));
let mut recorded: Vec<(
String,
nmbrs_workload::model::MetricSpec,
f64,
Option<nmbrs_metrics::labels::Labels>,
)> = Vec::new();
'metrics: for (name, spec) in entries {
let binding = crate::scope::synthesize_metric_binding_name(name);
let value = k.pull(&binding);
let Some(raw) = to_f64(&value) else {
crate::diag!(
crate::observer::LogLevel::Warn,
"phase '{phase_name}' metric '{name}': value `{expr}` resolved to \
a non-numeric {disc:?}; skipping (metric values must be U64 / F64 / Bool)",
expr = spec.value,
disc = std::mem::discriminant(&value)
);
continue;
};
let sanitised = match &spec.format {
Some(f) => match nmbrs_workload::metric_format::parse_format_spec(f) {
Ok(fs) => fs.apply(raw),
Err(e) => {
crate::diag!(
crate::observer::LogLevel::Warn,
"phase '{phase_name}' metric '{name}' format '{f}': {e}; \
recording unformatted value"
);
raw
}
},
None => raw,
};
let coord = if spec.cell.is_empty() {
None
} else {
let mut c = nmbrs_metrics::labels::Labels::default();
for dim in spec.cell.keys() {
let wire = crate::scope::synthesize_cell_binding_name(name, dim);
let v = k.pull(&wire);
let Value::Str(text) = &v else {
crate::diag!(
crate::observer::LogLevel::Warn,
"phase '{phase_name}' metric '{name}' cell '{dim}': \
coordinate resolved to a non-string {disc:?}; skipping \
the metric (dimension values are label values — convert \
the expression explicitly)",
disc = std::mem::discriminant(&v)
);
continue 'metrics;
};
c = c.with(dim.clone(), text.to_string());
}
Some(c)
};
recorded.push((name.clone(), spec.clone(), sanitised, coord));
}
for (name, spec, value, coord) in recorded {
use nmbrs_metrics::component::InstrumentRef;
use nmbrs_workload::model::MetricKind;
let family = spec.family.clone().unwrap_or_else(|| name.clone());
let kind = spec.kind.unwrap_or_default();
let target = match &coord {
None => phase_component.clone(),
Some(c) => nmbrs_metrics::cells::resolve_under(phase_component, c),
};
let target_labels = {
let g = target.read().unwrap_or_else(|e| e.into_inner());
g.effective_labels().clone()
};
let instr_labels = target_labels.with("family", family.clone());
let instrument: InstrumentRef = match kind {
MetricKind::Gauge => {
let g = std::sync::Arc::new(nmbrs_metrics::instruments::gauge::ValueGauge::new(
instr_labels,
));
g.set(value);
InstrumentRef::Gauge(g)
}
MetricKind::Histogram => {
let h = std::sync::Arc::new(nmbrs_metrics::instruments::histogram::Histogram::new(
instr_labels,
));
h.record(value as u64);
InstrumentRef::Histogram(h)
}
MetricKind::Counter => {
let c = std::sync::Arc::new(nmbrs_metrics::instruments::counter::Counter::new(
instr_labels,
));
if value > 0.0 {
c.inc_by(value as u64);
}
InstrumentRef::Counter(c)
}
};
let mut g = target.write().unwrap_or_else(|e| e.into_inner());
if let Err(e) =
g.register_instrument_with_unit(family.clone(), spec.unit.clone(), instrument)
{
crate::diag!(
crate::observer::LogLevel::Warn,
"phase '{phase_name}' metric '{name}': instrument registration \
for family '{family}': {e}"
);
}
}
Ok(())
}
fn resolve_stanza_count(
spec: &str,
scope: &crate::scope_kernel::ScopeKernel,
) -> Result<u64, String> {
if let Some(n) = crate::runner::parse_count(spec) {
return Ok(n);
}
let Some(inner) = spec.strip_prefix('{').and_then(|s| s.strip_suffix('}')) else {
return Err(format!("`{spec}` is neither a count nor a `{{binding}}`"));
};
let inner = inner.trim();
if let Some(v) = scope.pull_value(inner) {
return crate::validation::value_to_u64_for_count(v.clone())
.ok_or_else(|| format!("`{{{inner}}}` is {v:?}, not a count"));
}
let v = polydat::dsl::compile::eval_const_expr(inner).map_err(|e| {
format!(
"`{{{inner}}}` is not a binding in this phase's scope, nor a constant expression: {e}"
)
})?;
crate::validation::value_to_u64_for_count(v.clone())
.ok_or_else(|| format!("`{{{inner}}}` is {v:?}, not a count"))
}
fn resolve_selector_param(adapter: &str) -> String {
format!("{adapter}driver")
}
fn phase_identity_for(phase_name: &str, phase_labels: &str) -> crate::checkpoint::PhaseIdentity {
let yaml_path = crate::scene_tree::current()
.and_then(|t| {
t.find_phase(phase_name, phase_labels, None)
.and_then(|id| t.nodes.get(id).map(|n| n.yaml_path.clone()))
})
.unwrap_or_default();
crate::checkpoint::PhaseIdentity {
yaml_path,
coords: phase_labels.to_string(),
phase_hash: None,
}
}
fn resolve_over(
value: &polydat::ast::Value,
extent: u64,
open_extent: bool,
) -> Result<Option<polydat::iteration::cursor_partition::Partition>, String> {
use polydat::ast::Value;
use polydat::iteration::cursor_partition::{Partition, parse, resolve};
let single_of = |parts: Vec<Partition>| -> Result<Partition, String> {
match parts.len() {
0 => Err("partition list is empty — spec produced no partitions".to_string()),
1 => Ok(parts.into_iter().next().unwrap()),
n => Err(format!(
"spec resolves to {n} partitions, but this cursor consumes it \
directly (no enclosing `for:` iteration). Iterate the list with \
`for: \"p in <param>.partitions\"` and declare the cursor \
`over p`, or supply a single-partition spec"
)),
}
};
let reproject = |p: &Partition| -> Partition {
if open_extent || p.base_extent == extent || extent == 0 {
return *p;
}
let start_ord = ((p.start_pct / 100.0) * extent as f64).round() as u64;
let end_ord = ((p.end_pct / 100.0) * extent as f64).round() as u64;
Partition {
idx: p.idx,
count: p.count,
start_ord,
end_ord,
start_pct: p.start_pct,
end_pct: p.end_pct,
base_extent: extent,
}
};
let reject_spec_for_open = || -> String {
"an open-extent cursor (`until_*`) has no extent to resolve a \
partition spec against — its declared size is just the per-pass \
base chunk. Resolve the spec against an explicit reference extent \
first (`for: \"p in partitions(<spec>, <extent>)\"`) and declare \
the cursor `over p`"
.to_string()
};
match value {
Value::None => Ok(None),
Value::Str(s) => {
if open_extent {
return Err(reject_spec_for_open());
}
let spec = parse(s.as_ref())?;
let parts = resolve(&spec, 0, extent)?;
Ok(Some(single_of(parts)?))
}
Value::Ext(b) => {
if let Some(p) = value.as_partition() {
Ok(Some(reproject(p)))
} else if let Some(spec) = value.as_partition_spec() {
if open_extent {
return Err(reject_spec_for_open());
}
let parts = resolve(spec, 0, extent)?;
Ok(Some(single_of(parts)?))
} else if let Some(list) = value.as_partition_list() {
let single = single_of(list.as_slice().to_vec())?;
Ok(Some(reproject(&single)))
} else {
Err(format!(
"`over` expression produced an Ext value of unexpected type `{}` — \
expected Partition, PartitionSpec, or PartitionList",
b.type_name(),
))
}
}
other => Err(format!(
"`over` expression produced unsupported value type — expected Str or partition-typed Ext, got {other:?}"
)),
}
}
fn panic_message(payload: &Box<dyn std::any::Any + Send>) -> String {
if let Some(s) = payload.downcast_ref::<&str>() {
(*s).to_string()
} else if let Some(s) = payload.downcast_ref::<String>() {
s.clone()
} else {
"<non-string panic payload>".to_string()
}
}