use serde::{Deserialize, Serialize};
use std::sync::{Arc, Mutex, RwLock};
static GLOBAL_TREE: Mutex<Option<Arc<RwLock<SceneTree>>>> = Mutex::new(None);
pub fn install_global(tree: SceneTree) -> Arc<RwLock<SceneTree>> {
let arc = Arc::new(RwLock::new(tree));
if !crate::execution_context::install_scene_tree(arc.clone()) {
*GLOBAL_TREE.lock().unwrap_or_else(|e| e.into_inner()) = Some(arc.clone());
}
arc
}
fn active_handle() -> Option<Arc<RwLock<SceneTree>>> {
crate::execution_context::current_scene_tree().or_else(|| {
GLOBAL_TREE
.lock()
.unwrap_or_else(|e| e.into_inner())
.clone()
})
}
pub fn current() -> Option<SceneTree> {
active_handle().and_then(|t| t.read().ok().map(|g| g.clone()))
}
pub fn with_global_mut<F: FnOnce(&mut SceneTree)>(f: F) {
if let Some(arc) = active_handle()
&& let Ok(mut g) = arc.write()
{
f(&mut g);
}
}
pub fn with_global<R, F: FnOnce(&SceneTree) -> R>(f: F) -> Option<R> {
active_handle().and_then(|a| a.read().ok().map(|g| f(&g)))
}
pub type SceneNodeId = usize;
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub enum NodeKind {
Root,
Phase,
Scope,
}
#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
pub enum PhaseStatus {
Pending,
Running,
Completed,
Failed(String),
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct SceneNode {
pub id: SceneNodeId,
pub parent: Option<SceneNodeId>,
pub children: Vec<SceneNodeId>,
pub depth: usize,
pub kind: NodeKind,
pub name: String,
pub labels: String,
pub status: PhaseStatus,
pub op_count: usize,
pub duration_secs: Option<f64>,
#[serde(default)]
pub op_names: Vec<String>,
#[serde(default)]
pub own_names: Vec<String>,
#[serde(default)]
pub seq: Option<usize>,
#[serde(default)]
pub yaml_path: Vec<crate::checkpoint::PathSegment>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub outcome: Option<crate::phase_outcome::PhaseOutcome>,
#[serde(default = "default_active")]
pub active: bool,
}
fn default_active() -> bool {
true
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct SceneTree {
pub nodes: Vec<SceneNode>,
}
impl Default for SceneTree {
fn default() -> Self {
Self::new()
}
}
impl SceneTree {
pub fn new() -> Self {
let mut t = Self { nodes: Vec::new() };
t.nodes.push(SceneNode {
id: 0,
parent: None,
children: Vec::new(),
depth: 0,
kind: NodeKind::Root,
name: String::new(),
labels: String::new(),
status: PhaseStatus::Pending,
op_count: 0,
duration_secs: None,
op_names: Vec::new(),
own_names: Vec::new(),
seq: None,
yaml_path: Vec::new(),
outcome: None,
active: true,
});
t
}
pub fn root(&self) -> SceneNodeId {
0
}
pub fn push(
&mut self,
parent: SceneNodeId,
kind: NodeKind,
name: impl Into<String>,
labels: impl Into<String>,
) -> SceneNodeId {
let name: String = name.into();
let labels: String = labels.into();
if let Some(&existing) = self.nodes[parent]
.children
.iter()
.find(|&&c| self.nodes[c].kind == kind && self.nodes[c].name == name)
{
return existing;
}
let id = self.nodes.len();
let depth = self.nodes[parent].depth + 1;
let seq = match kind {
NodeKind::Phase => {
let count = self
.nodes
.iter()
.filter(|n| n.kind == NodeKind::Phase)
.count();
Some(count + 1)
}
_ => None,
};
self.nodes.push(SceneNode {
id,
parent: Some(parent),
children: Vec::new(),
depth,
kind,
name,
labels,
status: PhaseStatus::Pending,
op_count: 0,
duration_secs: None,
op_names: Vec::new(),
own_names: Vec::new(),
seq,
yaml_path: Vec::new(),
outcome: None,
active: true,
});
self.nodes[parent].children.push(id);
id
}
pub fn set_yaml_path(&mut self, id: SceneNodeId, path: Vec<crate::checkpoint::PathSegment>) {
if id < self.nodes.len() {
self.nodes[id].yaml_path = path;
}
}
pub fn total_phases(&self) -> usize {
self.nodes
.iter()
.filter(|n| n.kind == NodeKind::Phase)
.count()
}
pub fn apply_phase_filter(
&mut self,
pattern: Option<&crate::phase_filter::PhasePattern>,
) -> PhaseFilterStats {
let mut stats = PhaseFilterStats::default();
if pattern.is_none() {
stats.matched = self.total_phases();
return stats;
}
let pat = pattern.unwrap();
for n in self.nodes.iter_mut() {
if n.kind == NodeKind::Phase {
n.active = pat.is_match(&n.name);
stats.total += 1;
if n.active {
stats.matched += 1;
}
}
}
let n = self.nodes.len();
for i in (0..n).rev() {
if matches!(self.nodes[i].kind, NodeKind::Phase) {
continue;
}
let kids = self.nodes[i].children.clone();
let any_active = kids.iter().any(|c| self.nodes[*c].active);
self.nodes[i].active = any_active;
}
stats
}
pub fn is_phase_active(&self, id: SceneNodeId) -> bool {
self.nodes.get(id).map(|n| n.active).unwrap_or(false)
}
}
#[derive(Default, Debug, Clone, Copy)]
pub struct PhaseFilterStats {
pub matched: usize,
pub total: usize,
}
impl SceneTree {
pub fn set_phase_op_names(&mut self, id: SceneNodeId, names: Vec<String>) {
if id < self.nodes.len() {
self.nodes[id].op_names = names;
}
}
pub fn set_own_names(&mut self, id: SceneNodeId, names: Vec<String>) {
if id < self.nodes.len() {
self.nodes[id].own_names = names;
}
}
pub fn dfs(&self) -> DfsIter<'_> {
DfsIter {
tree: self,
stack: vec![0],
}
}
pub fn dfs_phases(&self) -> impl Iterator<Item = &SceneNode> {
self.dfs().filter(|n| n.kind == NodeKind::Phase)
}
pub fn find_phase(
&self,
name: &str,
_labels: &str,
want: Option<&PhaseStatus>,
) -> Option<SceneNodeId> {
self.dfs_phases()
.find(|n| n.name == name && want.is_none_or(|w| &n.status == w))
.map(|n| n.id)
}
pub fn set_phase_running_at(&mut self, id: SceneNodeId, op_count: usize) {
if let Some(n) = self.nodes.get_mut(id) {
n.status = PhaseStatus::Running;
n.op_count = op_count;
}
}
pub fn remove_node(&mut self, id: SceneNodeId) {
let Some(parent) = self.nodes.get(id).and_then(|n| n.parent) else {
return;
};
if let Some(p) = self.nodes.get_mut(parent) {
p.children.retain(|c| *c != id);
}
}
pub fn set_phase_completed_at(&mut self, id: SceneNodeId, duration_secs: f64) {
if let Some(n) = self.nodes.get_mut(id) {
n.status = PhaseStatus::Completed;
n.duration_secs = Some(duration_secs);
}
}
pub fn set_phase_failed_at(&mut self, id: SceneNodeId, error: &str) {
if let Some(n) = self.nodes.get_mut(id) {
n.status = PhaseStatus::Failed(error.to_string());
}
}
pub fn set_phase_running(&mut self, name: &str, labels: &str, op_count: usize) {
if let Some(id) = self.find_phase(name, labels, Some(&PhaseStatus::Pending)) {
self.set_phase_running_at(id, op_count);
}
}
pub fn set_phase_completed(&mut self, name: &str, labels: &str, duration_secs: f64) {
if let Some(id) = self.find_phase(name, labels, Some(&PhaseStatus::Running)) {
self.set_phase_completed_at(id, duration_secs);
}
}
pub fn set_phase_failed(&mut self, name: &str, labels: &str, error: &str) {
if let Some(id) = self.find_phase(name, labels, None) {
self.set_phase_failed_at(id, error);
}
}
pub fn set_phase_outcome(
&mut self,
name: &str,
labels: &str,
outcome: crate::phase_outcome::PhaseOutcome,
) {
let Some(id) = self.find_phase(name, labels, None) else {
return;
};
self.set_phase_outcome_at(id, outcome);
}
pub fn set_phase_outcome_at(
&mut self,
id: SceneNodeId,
outcome: crate::phase_outcome::PhaseOutcome,
) {
let Some(n) = self.nodes.get_mut(id) else {
return;
};
let lifecycle = match outcome.validity {
crate::phase_outcome::Validity::Succeeded => PhaseStatus::Completed,
crate::phase_outcome::Validity::Failed => {
let msg = outcome
.first_error_message()
.unwrap_or("unknown error")
.to_string();
PhaseStatus::Failed(msg)
}
};
n.status = lifecycle;
if outcome.duration_secs > 0.0 {
n.duration_secs = Some(outcome.duration_secs);
}
n.outcome = Some(outcome);
}
pub fn session_disposition(&self) -> crate::phase_outcome::SessionDisposition {
let any_failed = self
.nodes
.iter()
.filter(|n| matches!(n.kind, NodeKind::Phase))
.filter_map(|n| n.outcome.as_ref())
.any(|o| o.is_failure());
if any_failed {
crate::phase_outcome::SessionDisposition::Failure
} else {
crate::phase_outcome::SessionDisposition::Success
}
}
pub fn phase_outcome_present_at(&self, id: SceneNodeId) -> bool {
self.nodes
.get(id)
.map(|n| n.outcome.is_some())
.unwrap_or(false)
}
pub fn iter_phase_outcomes(&self) -> impl Iterator<Item = &crate::phase_outcome::PhaseOutcome> {
self.nodes
.iter()
.filter(|n| matches!(n.kind, NodeKind::Phase))
.filter_map(|n| n.outcome.as_ref())
}
pub fn aggregate_status(&self, id: SceneNodeId) -> PhaseStatus {
let n = &self.nodes[id];
if n.kind == NodeKind::Phase {
return n.status.clone();
}
let mut seen_phase = false;
let mut all_completed = true;
let mut any_running = false;
let mut first_failure: Option<String> = None;
for &child in &n.children {
let cs = self.aggregate_status(child);
match cs {
PhaseStatus::Failed(e) => {
if first_failure.is_none() {
first_failure = Some(e);
}
all_completed = false;
}
PhaseStatus::Running => {
any_running = true;
all_completed = false;
}
PhaseStatus::Pending => {
all_completed = false;
}
PhaseStatus::Completed => {}
}
if self.nodes[child].kind == NodeKind::Phase || self.descendants_contain_phase(child) {
seen_phase = true;
}
}
if let Some(e) = first_failure {
return PhaseStatus::Failed(e);
}
if any_running {
return PhaseStatus::Running;
}
if seen_phase && all_completed {
return PhaseStatus::Completed;
}
PhaseStatus::Pending
}
fn descendants_contain_phase(&self, id: SceneNodeId) -> bool {
let n = &self.nodes[id];
if n.kind == NodeKind::Phase {
return true;
}
n.children
.iter()
.any(|&c| self.descendants_contain_phase(c))
}
pub fn phase_count(&self) -> usize {
self.dfs_phases().count()
}
}
pub fn running_phase_indent() -> String {
let Some(tree) = current() else {
return String::new();
};
if let Some(id) = crate::execution_context::current_phase_node()
&& let Some(n) = tree.nodes.get(id)
{
return " ".repeat(n.depth.saturating_sub(1));
}
tree.dfs_phases()
.find(|n| matches!(n.status, PhaseStatus::Running))
.map(|n| " ".repeat(n.depth.saturating_sub(1)))
.unwrap_or_default()
}
pub struct DfsIter<'a> {
tree: &'a SceneTree,
stack: Vec<SceneNodeId>,
}
impl<'a> Iterator for DfsIter<'a> {
type Item = &'a SceneNode;
fn next(&mut self) -> Option<Self::Item> {
let id = self.stack.pop()?;
let node = &self.tree.nodes[id];
for &c in node.children.iter().rev() {
self.stack.push(c);
}
Some(node)
}
}
#[cfg(test)]
mod tests {
use super::*;
fn build_simple() -> SceneTree {
let mut t = SceneTree::new();
let s = t.push(t.root(), NodeKind::Scope, "for_each x=1", "");
let _ = t.push(s, NodeKind::Phase, "p", "x=1");
let _ = t.push(s, NodeKind::Phase, "q", "x=1");
let s2 = t.push(t.root(), NodeKind::Scope, "for_each x=2", "");
let _ = t.push(s2, NodeKind::Phase, "p", "x=2");
let _ = t.push(s2, NodeKind::Phase, "q", "x=2");
t
}
#[test]
fn dfs_yields_all_in_display_order() {
let t = build_simple();
let names: Vec<&str> = t.dfs().map(|n| n.name.as_str()).collect();
assert_eq!(
names,
vec!["", "for_each x=1", "p", "q", "for_each x=2", "p", "q"]
);
}
#[test]
fn dfs_phases_skips_root_and_scopes() {
let t = build_simple();
let names: Vec<&str> = t.dfs_phases().map(|n| n.name.as_str()).collect();
assert_eq!(names, vec!["p", "q", "p", "q"]);
}
#[test]
fn find_pending_then_running_progresses_through_iterations() {
let mut t = build_simple();
t.set_phase_running("p", "x=1", 3);
let n = t
.find_phase("p", "x=1", Some(&PhaseStatus::Running))
.unwrap();
assert_eq!(t.nodes[n].op_count, 3);
t.set_phase_completed("p", "x=1", 0.5);
t.set_phase_running("p", "x=2", 5);
let n2 = t
.find_phase("p", "x=2", Some(&PhaseStatus::Running))
.unwrap();
assert_ne!(n, n2);
assert_eq!(t.nodes[n2].op_count, 5);
}
#[test]
fn id_based_flips_attribute_to_correct_node_under_reordered_completion() {
let mut t = build_simple();
let p_x1 = t.dfs_phases().filter(|n| n.name == "p").next().unwrap().id;
let p_x2 = t.dfs_phases().filter(|n| n.name == "p").nth(1).unwrap().id;
assert_ne!(p_x1, p_x2);
t.set_phase_running_at(p_x1, 3);
t.set_phase_running_at(p_x2, 7);
t.set_phase_completed_at(p_x2, 10.0);
t.set_phase_completed_at(p_x1, 5.0);
assert_eq!(t.nodes[p_x1].op_count, 3);
assert_eq!(t.nodes[p_x2].op_count, 7);
assert_eq!(t.nodes[p_x1].duration_secs, Some(5.0));
assert_eq!(t.nodes[p_x2].duration_secs, Some(10.0));
assert_eq!(t.nodes[p_x1].status, PhaseStatus::Completed);
assert_eq!(t.nodes[p_x2].status, PhaseStatus::Completed);
}
#[test]
fn id_based_outcome_install_targets_the_dispatch_node() {
use crate::phase_outcome::{PhaseIdentity, PhaseOutcome};
let mut t = build_simple();
let p_x1 = t.dfs_phases().filter(|n| n.name == "p").next().unwrap().id;
let p_x2 = t.dfs_phases().filter(|n| n.name == "p").nth(1).unwrap().id;
t.set_phase_running_at(p_x1, 1);
t.set_phase_running_at(p_x2, 1);
t.set_phase_outcome_at(
p_x2,
PhaseOutcome::completed(PhaseIdentity::new("p", "x=2"), 9.0),
);
t.set_phase_outcome_at(
p_x1,
PhaseOutcome::completed(PhaseIdentity::new("p", "x=1"), 4.0),
);
assert_eq!(t.nodes[p_x1].duration_secs, Some(4.0));
assert_eq!(t.nodes[p_x2].duration_secs, Some(9.0));
assert!(t.nodes[p_x1].outcome.is_some());
assert!(t.nodes[p_x2].outcome.is_some());
}
#[test]
fn same_name_cells_under_one_parent_alias_to_one_node() {
let mut t = SceneTree::new();
let s = t.push(t.root(), NodeKind::Scope, "phase.for_each x", "");
let c1 = t.push(s, NodeKind::Phase, "p", "x=1");
let c2 = t.push(s, NodeKind::Phase, "p", "x=2");
assert_eq!(
c1, c2,
"push is idempotent by name — flat for_each / sweep cells collapse"
);
let s2 = t.push(t.root(), NodeKind::Scope, "phase.for_each y", "");
let c3 = t.push(s2, NodeKind::Phase, "p", "y=1");
assert_ne!(
c1, c3,
"same name under a DIFFERENT parent is a distinct node"
);
}
#[test]
fn aggregate_status_walks_descendants() {
let mut t = build_simple();
assert_eq!(t.aggregate_status(t.root()), PhaseStatus::Pending);
for (name, labels) in [("p", "x=1"), ("q", "x=1"), ("p", "x=2"), ("q", "x=2")] {
t.set_phase_running(name, labels, 1);
t.set_phase_completed(name, labels, 0.1);
}
assert_eq!(t.aggregate_status(t.root()), PhaseStatus::Completed);
}
#[test]
fn aggregate_propagates_failure() {
let mut t = build_simple();
t.set_phase_running("p", "x=1", 1);
t.set_phase_failed("p", "x=1", "boom");
let s = t.aggregate_status(t.root());
assert!(
matches!(s, PhaseStatus::Failed(ref e) if e == "boom"),
"got {s:?}"
);
}
#[test]
fn aggregate_running_when_any_running() {
let mut t = build_simple();
t.set_phase_running("p", "x=1", 1);
assert_eq!(t.aggregate_status(t.root()), PhaseStatus::Running);
}
#[test]
fn set_phase_outcome_installs_structured_and_mirrors_legacy() {
use crate::phase_outcome::{PhaseIdentity, PhaseOutcome};
let mut t = build_simple();
t.set_phase_running("p", "x=1", 1);
let outcome = PhaseOutcome::completed(PhaseIdentity::new("p", "x=1"), 2.5);
t.set_phase_outcome("p", "x=1", outcome.clone());
let phase_id = t.find_phase("p", "x=1", None).expect("phase found");
let n = &t.nodes[phase_id];
assert_eq!(n.outcome.as_ref(), Some(&outcome));
assert_eq!(n.status, PhaseStatus::Completed);
assert_eq!(n.duration_secs, Some(2.5));
let outcomes: Vec<_> = t.iter_phase_outcomes().collect();
assert_eq!(outcomes.len(), 1);
assert_eq!(outcomes[0], &outcome);
}
#[test]
fn failed_outcome_legacy_status_carries_first_error_message() {
use crate::phase_outcome::{PhaseErrorDetail, PhaseIdentity, PhaseOutcome};
let mut t = build_simple();
t.set_phase_running("p", "x=1", 1);
let outcome = PhaseOutcome::failed(
PhaseIdentity::new("p", "x=1"),
142.7,
vec![PhaseErrorDetail {
class: "poll_timeout".into(),
message: "deadline reached after 14400s".into(),
op_name: None,
cycle: None,
op_template: None,
op_resolved: None,
at_nanos: 1_000,
retryable: false,
}],
);
t.set_phase_outcome("p", "x=1", outcome);
let phase_id = t.find_phase("p", "x=1", None).expect("phase found");
match &t.nodes[phase_id].status {
PhaseStatus::Failed(msg) => assert_eq!(msg, "deadline reached after 14400s"),
other => panic!("expected Failed, got {other:?}"),
}
}
#[test]
fn session_disposition_failure_when_any_phase_failed() {
use crate::phase_outcome::{
PhaseErrorDetail, PhaseIdentity, PhaseOutcome, SessionDisposition,
};
let mut t = build_simple();
t.set_phase_outcome(
"p",
"x=1",
PhaseOutcome::completed(PhaseIdentity::new("p", "x=1"), 1.0),
);
assert_eq!(t.session_disposition(), SessionDisposition::Success);
t.set_phase_outcome(
"q",
"x=1",
PhaseOutcome::failed(
PhaseIdentity::new("q", "x=1"),
0.5,
vec![PhaseErrorDetail {
class: "BindError".into(),
message: "bad".into(),
op_name: None,
cycle: None,
op_template: None,
op_resolved: None,
at_nanos: 0,
retryable: false,
}],
),
);
assert_eq!(t.session_disposition(), SessionDisposition::Failure);
}
#[test]
fn session_disposition_success_when_no_phase_ran() {
use crate::phase_outcome::SessionDisposition;
let t = build_simple();
assert_eq!(t.session_disposition(), SessionDisposition::Success);
}
#[test]
fn session_disposition_skipped_and_cursor_suspended_are_success() {
use crate::phase_outcome::{PhaseIdentity, PhaseOutcome, SessionDisposition};
let mut t = build_simple();
t.set_phase_outcome(
"p",
"x=1",
PhaseOutcome::skipped(PhaseIdentity::new("p", "x=1")),
);
t.set_phase_outcome(
"q",
"x=1",
PhaseOutcome::interrupted(PhaseIdentity::new("q", "x=1"), 0.5, None),
);
assert_eq!(t.session_disposition(), SessionDisposition::Success);
}
}