fn state_inputs_exist_for_ready_set(
workspace_root: &Path,
artifact_root: &Path,
rhei: &rhei_core::ast::Rhei,
machine: &rhei_validator::StateMachine,
task: &rhei_core::ast::Task,
state_name: &str,
) -> bool {
let Some(state_def) = machine.states.get(state_name) else {
return false;
};
if state_def.inputs.is_empty() {
return true;
}
let settings = match load_merged_settings(workspace_root) {
Ok(settings) => settings,
Err(_) => return false,
};
let visit_count = Some(render_visit_count(
rhei.metadata.as_ref(),
&task.id,
state_name,
task.state.as_str(),
machine,
));
ensure_state_inputs_exist_for_transition(
artifact_root,
Some(task),
&task.id.to_string(),
state_name,
state_def,
visit_count,
machine,
&settings,
"",
)
.is_ok()
}
fn open_descendant_tasks<'a>(
task: &'a rhei_core::ast::Task,
machines: &rhei_validator::MachineSet,
) -> Vec<&'a rhei_core::ast::Task> {
fn recurse<'a>(
task: &'a rhei_core::ast::Task,
machines: &rhei_validator::MachineSet,
out: &mut Vec<&'a rhei_core::ast::Task>,
) {
for child in &task.children {
if !is_terminal_state(child.state.as_str(), machines.for_task(&child.id)) {
out.push(child);
}
recurse(child, machines, out);
}
}
let mut out = Vec::new();
recurse(task, machines, &mut out);
out
}
fn any_open_descendant(
task: &rhei_core::ast::Task,
machines: &rhei_validator::MachineSet,
) -> bool {
task.children.iter().any(|child| {
!is_terminal_state(child.state.as_str(), machines.for_task(&child.id))
|| any_open_descendant(child, machines)
})
}
fn descendants_are_terminal(
task: &rhei_core::ast::Task,
machines: &rhei_validator::MachineSet,
) -> bool {
!any_open_descendant(task, machines)
}
fn format_open_descendants(
open: &[&rhei_core::ast::Task],
machines: &rhei_validator::MachineSet,
) -> String {
let items: Vec<String> = open
.iter()
.take(3)
.map(|task| {
format!(
"Task {} ({})",
task.id,
normalized_state_name(task.state.as_str(), machines.for_task(&task.id))
)
})
.collect();
let suffix =
if open.len() > 3 { format!(" (+{} more)", open.len() - 3) } else { String::new() };
format!("{}{}", items.join(", "), suffix)
}
fn find_ready_tasks<'a>(
rhei: &'a rhei_core::ast::Rhei,
machines: &rhei_validator::MachineSet,
workspace_root: &Path,
task_roots: &std::collections::HashMap<String, std::path::PathBuf>,
spawned: &HashSet<String>,
) -> Vec<&'a rhei_core::ast::Task> {
use std::collections::HashMap;
let mut all_tasks = Vec::new();
collect_plan_tasks(&rhei.tasks, &mut all_tasks);
let index = task_index(&all_tasks);
let state_map: HashMap<&TaskId, String> = all_tasks
.iter()
.map(|t| (&t.id, normalized_state_name(t.state.as_str(), machines.for_task(&t.id))))
.collect();
let mut ready = Vec::new();
for task in &all_tasks {
let task = *task;
match supervision_verdict_for(task, &index, machines, rhei.metadata.as_ref(), spawned) {
SupervisionVerdict::Held { .. } | SupervisionVerdict::SupervisorWaiting => continue,
SupervisionVerdict::SupervisorReady => {}
SupervisionVerdict::Unsupervised => {
if !descendants_are_terminal(task, machines) {
continue;
}
}
}
let machine = machines.for_task(&task.id);
let current_state = task.state.as_str();
let normalized_state = normalized_state_name(current_state, machine);
if is_terminal_state(current_state, machine)
|| machine.states.get(&normalized_state).map(|def| def.gating).unwrap_or(false)
{
continue;
}
if machine.states.get(&normalized_state).and_then(|def| def.poll.as_ref()).is_some()
&& poll_next_attempt_at(rhei.metadata.as_ref(), &task.id, &normalized_state)
.is_some_and(|deadline| deadline > current_unix_secs())
{
continue;
}
let all_priors_done = task.prior.iter().all(|dep_id| {
state_map
.get(dep_id)
.map(|s| dependency_is_satisfied(s, machines.for_task(dep_id)))
.unwrap_or(false)
});
let task_id = task.id.to_string();
let artifact_root = task_roots.get(&task_id).map_or(workspace_root, |root| root.as_path());
if all_priors_done
&& state_inputs_exist_for_ready_set(
workspace_root,
artifact_root,
rhei,
machine,
task,
&normalized_state,
)
{
ready.push(task);
}
}
ready
}
fn find_runnable_tasks<'a>(
rhei: &'a rhei_core::ast::Rhei,
machines: &rhei_validator::MachineSet,
workspace_root: &Path,
spawned: &HashSet<String>,
) -> Vec<&'a rhei_core::ast::Task> {
find_ready_tasks(rhei, machines, workspace_root, &std::collections::HashMap::new(), spawned)
.into_iter()
.filter(|task| task.assignee.is_none())
.collect()
}
fn find_held_tasks<'a>(
rhei: &'a rhei_core::ast::Rhei,
machines: &rhei_validator::MachineSet,
workspace_root: &Path,
) -> Vec<&'a rhei_core::ast::Task> {
find_ready_tasks(
rhei,
machines,
workspace_root,
&std::collections::HashMap::new(),
&HashSet::new(),
)
.into_iter()
.filter(|task| task.assignee.is_some())
.collect()
}
fn format_supervisor_holds(
rhei: &rhei_core::ast::Rhei,
machines: &rhei_validator::MachineSet,
scope: &RheiScope,
) -> Vec<String> {
let mut all = Vec::new();
collect_plan_tasks(&rhei.tasks, &mut all);
let mut counts: Vec<(String, usize)> = Vec::new();
for task in all
.iter()
.copied()
.filter(|task| task_in_rhei_scope(scope, &task.id.to_string()))
.filter(|task| !is_terminal_state(task.state.as_str(), machines.for_task(&task.id)))
{
let Some(hold) = held_by_supervisor(task, rhei, machines) else { continue };
let key = hold.supervisor.to_string();
match counts.iter_mut().find(|(id, _)| *id == key) {
Some((_, count)) => *count += 1,
None => counts.push((key, 1)),
}
}
counts
.into_iter()
.map(|(supervisor, count)| {
format!("{count} ticket(s) held by supervisor Task {supervisor}")
})
.collect()
}
fn format_held_tasks(held: &[&rhei_core::ast::Task]) -> String {
held.iter()
.map(|task| {
format!("Task {} (assignee {})", task.id, task.assignee.as_deref().unwrap_or("?"))
})
.collect::<Vec<_>>()
.join(", ")
}
fn find_claimable_tasks<'a>(
rhei: &'a rhei_core::ast::Rhei,
machines: &rhei_validator::MachineSet,
workspace_root: &Path,
task_roots: &std::collections::HashMap<String, std::path::PathBuf>,
) -> Vec<&'a rhei_core::ast::Task> {
find_ready_tasks(rhei, machines, workspace_root, task_roots, &HashSet::new())
.into_iter()
.filter(|task| task.assignee.is_none())
.filter(|task| {
let machine = machines.for_task(&task.id);
let state = normalized_state_name(task.state.as_str(), machine);
task_is_in_initial_state(task, &state, machine)
})
.collect()
}
fn task_is_in_initial_state(
task: &rhei_core::ast::Task,
normalized_state: &str,
machine: &rhei_validator::StateMachine,
) -> bool {
machine
.profile_for_node(task.kind.as_str(), task.profile_level())
.map(|profile| profile.initial == normalized_state)
.unwrap_or_else(|| machine.states.get(normalized_state).map(|def| def.initial).unwrap_or(false))
}
fn collect_plan_tasks<'a>(
tasks: &'a [rhei_core::ast::Task],
out: &mut Vec<&'a rhei_core::ast::Task>,
) {
for task in tasks {
out.push(task);
collect_plan_tasks(&task.children, out);
}
}
fn plan_state_map<'a>(
tasks: &[&'a rhei_core::ast::Task],
machines: &rhei_validator::MachineSet,
) -> std::collections::HashMap<&'a TaskId, String> {
tasks
.iter()
.map(|task| {
(&task.id, normalized_state_name(task.state.as_str(), machines.for_task(&task.id)))
})
.collect()
}
fn blocking_priors(
task: &rhei_core::ast::Task,
state_map: &std::collections::HashMap<&TaskId, String>,
machines: &rhei_validator::MachineSet,
) -> Vec<String> {
task.prior
.iter()
.filter_map(|dep_id| match state_map.get(dep_id) {
Some(state) if !dependency_is_satisfied(state, machines.for_task(dep_id)) => {
Some(format!("Task {} ({})", dep_id, state))
}
None => Some(format!("Task {} (missing)", dep_id)),
_ => None,
})
.collect()
}
fn first_blocking_prior(
task: &rhei_core::ast::Task,
state_map: &std::collections::HashMap<&TaskId, String>,
machines: &rhei_validator::MachineSet,
scope: &RheiScope,
) -> Option<String> {
task.prior.iter().find_map(|dep_id| match state_map.get(dep_id) {
Some(state) if !dependency_is_satisfied(state, machines.for_task(dep_id)) => {
let outside = if task_in_rhei_scope(scope, &dep_id.to_string()) {
""
} else {
", outside the --rhei scope"
};
Some(format!("Task {} ({}{})", dep_id, state, outside))
}
None => Some(format!("Task {} (missing)", dep_id)),
_ => None,
})
}
fn is_terminal_state(state: &str, machine: &rhei_validator::StateMachine) -> bool {
let normalized = normalized_state_name(state, machine);
machine.states.get(&normalized).map(|def| def.terminal).unwrap_or(false)
}