#[derive(Args, Clone, Debug, Default)]
#[command(next_help_heading = "Standalone Execution")]
struct StandaloneExecutionFlags {
#[arg(long)]
dry_run: bool,
#[arg(long)]
no_callbacks: bool,
#[arg(long)]
continue_on_error: bool,
#[arg(long, default_value_t = 1, add = ArgValueCompleter::new(complete_parallel))]
parallel: usize,
#[arg(long = "rhei", value_name = "RHEI_ID", add = ArgValueCompleter::new(complete_rhei_id))]
rhei: Vec<String>,
#[arg(long, conflicts_with = "no_tui")]
tui: bool,
#[arg(long)]
no_tui: bool,
#[arg(long, conflicts_with = "no_dashboard")]
dashboard: bool,
#[arg(long)]
no_dashboard: bool,
}
#[derive(Args, Clone, Debug, Default)]
#[command(next_help_heading = "Agent Execution")]
struct AgentExecutionFlags {
#[arg(long)]
no_agent: bool,
#[arg(long, value_name = "AGENT", add = ArgValueCompleter::new(complete_agent_name))]
agent: Option<String>,
#[arg(long, value_name = "MODE", add = ArgValueCompleter::new(complete_agent_mode))]
agent_mode: Option<String>,
#[arg(long, value_name = "MODEL", add = ArgValueCompleter::new(complete_model_name))]
model: Option<String>,
}
#[derive(Args, Clone, Debug, Default)]
#[command(next_help_heading = "Program Execution")]
struct ProgramExecutionFlags {
#[arg(long)]
no_program: bool,
#[arg(long, value_name = "DURATION", add = ArgValueCompleter::new(complete_duration))]
program_timeout: Option<String>,
}
#[derive(Args, Clone, Debug, Default)]
#[command(next_help_heading = "Snapshots")]
struct SnapshotExecutionFlags {
#[arg(long, value_name = "REF")]
from_snapshot: Option<String>,
#[arg(long, requires = "from_snapshot")]
override_inherit: bool,
#[arg(long = "task", value_name = "TASK_ID", add = ArgValueCompleter::new(complete_task_id))]
snapshot_task: Option<String>,
#[arg(long = "target", value_name = "SLUG")]
snapshot_target: Option<String>,
}
struct RunOptions {
standalone: StandaloneExecutionFlags,
agent: AgentExecutionFlags,
program: ProgramExecutionFlags,
snapshot: SnapshotExecutionFlags,
}
impl RunOptions {
fn dry_run(&self) -> bool {
self.standalone.dry_run
}
fn no_callbacks(&self) -> bool {
self.standalone.no_callbacks
}
fn continue_on_error(&self) -> bool {
self.standalone.continue_on_error
}
fn parallel(&self) -> usize {
self.standalone.parallel
}
fn rhei_scope(&self) -> &[String] {
&self.standalone.rhei
}
fn narrow_to(&mut self, scope: Vec<String>) {
self.standalone.rhei = scope;
}
fn frontend_kind(&self) -> rhei_tui::FrontendKind {
if self.standalone.tui {
rhei_tui::FrontendKind::Tui
} else if self.standalone.no_tui {
rhei_tui::FrontendKind::Stdout
} else {
rhei_tui::FrontendKind::Auto
}
}
fn dashboard_enabled(&self, frontend_is_tui: bool) -> bool {
if self.standalone.dashboard {
true
} else if self.standalone.no_dashboard {
false
} else {
frontend_is_tui
}
}
fn no_agent(&self) -> bool {
self.agent.no_agent
}
fn agent_override(&self) -> Option<&str> {
self.agent.agent.as_deref()
}
fn agent_mode_override(&self) -> Option<&str> {
self.agent.agent_mode.as_deref()
}
fn model_override(&self) -> Option<&str> {
self.agent.model.as_deref()
}
fn no_program(&self) -> bool {
self.program.no_program
}
fn program_timeout_override(&self) -> Option<&str> {
self.program.program_timeout.as_deref()
}
fn snapshot_override_ref(&self) -> Option<&str> {
self.snapshot.from_snapshot.as_deref()
}
fn override_inherit(&self) -> bool {
self.snapshot.override_inherit
}
fn snapshot_task_selector(&self) -> Option<&str> {
self.snapshot.snapshot_task.as_deref()
}
fn snapshot_target_selector(&self) -> Option<&str> {
self.snapshot.snapshot_target.as_deref()
}
}
struct ActiveRunFrontend {
sink: Arc<dyn rhei_tui::EventSink>,
is_tui: bool,
dashboard: Option<Arc<rhei_tui::DashboardSink>>,
summary: Arc<SummarySink>,
intervene: Option<Arc<RunInterveneSink>>,
_frontend: Option<rhei_tui::Frontend>,
}
struct RunGateTransitionSink {
input: PathBuf,
machines: ExecutionMachines,
no_callbacks: bool,
}
impl RunGateTransitionSink {
fn new(input: PathBuf, machines: ExecutionMachines, no_callbacks: bool) -> Self {
Self { input, machines, no_callbacks }
}
}
impl rhei_tui::GateTransitionSink for RunGateTransitionSink {
fn transition_gate(
&self,
task_id: &str,
from: &str,
to: &str,
result: Option<&str>,
) -> Result<String, String> {
transition_dashboard_gate(
&self.input,
self.machines.for_task_str(task_id),
self.machines.callbacks_for_str(task_id),
task_id,
from,
to,
result,
self.no_callbacks,
)
.map_err(|err| err.to_string())
}
}
impl ActiveRunFrontend {
fn announce_dashboard(&self) {
if let Some(dashboard) = &self.dashboard {
self.sink.emit(rhei_tui::RunEvent::RunLink {
label: "Dashboard".to_string(),
url: dashboard.url().to_string(),
});
}
}
fn write_frozen_dashboard(&self) {
let Some(dashboard) = &self.dashboard else {
return;
};
match dashboard.write_frozen_dashboard() {
Ok(path) => self.sink.emit(rhei_tui::RunEvent::Message {
level: rhei_tui::MessageLevel::Info,
text: format!("Final dashboard: {}", path.display()),
}),
Err(err) => self.sink.emit(rhei_tui::RunEvent::Message {
level: rhei_tui::MessageLevel::Warn,
text: format!("warning: could not write final dashboard: {err}"),
}),
}
}
}
fn start_run_frontend(
workspace_root: &Path,
plan_input: &Path,
machines: &ExecutionMachines,
opts: &RunOptions,
parallel: u16,
total_tasks: usize,
) -> ActiveRunFrontend {
if opts.dry_run() {
return ActiveRunFrontend {
sink: Arc::new(rhei_tui::StdoutSink::new()),
is_tui: false,
dashboard: None,
summary: Arc::new(SummarySink::new()),
intervene: None,
_frontend: None,
};
}
let plan_path = plan_input.to_path_buf();
let loader_machines = machines.set.clone();
let loader: rhei_tui::PlanLoader =
Arc::new(move || load_plan_for_dashboard(&plan_path, &loader_machines));
let registry = Arc::new(RunInterveneSink::new(workspace_root.join("runtime")));
let gate = Arc::new(RunGateTransitionSink::new(
plan_input.to_path_buf(),
machines.clone(),
opts.no_callbacks(),
));
let tui_context = rhei_tui::TuiContext {
workspace: workspace_root.to_path_buf(),
plan_loader: Some(loader.clone()),
intervene: Some(registry.clone() as Arc<dyn rhei_tui::InterveneSink>),
gate: Some(gate.clone() as Arc<dyn rhei_tui::GateTransitionSink>),
};
let frontend = rhei_tui::select_frontend(
workspace_root,
opts.frontend_kind(),
parallel,
total_tasks,
tui_context,
);
let dashboard = if opts.dashboard_enabled(frontend.is_tui) {
match rhei_tui::DashboardSink::start_with_plan_intervene_and_gate(
workspace_root.to_path_buf(),
parallel,
total_tasks,
Some(loader.clone()),
Some(registry.clone() as Arc<dyn rhei_tui::InterveneSink>),
Some(gate.clone() as Arc<dyn rhei_tui::GateTransitionSink>),
) {
Ok(sink) => Some(Arc::new(sink)),
Err(err) => {
frontend.sink.emit(rhei_tui::RunEvent::Message {
level: rhei_tui::MessageLevel::Warn,
text: format!("warning: could not start dashboard: {err}"),
});
None
}
}
} else {
None
};
let intervene: Option<Arc<RunInterveneSink>> =
(frontend.is_tui || dashboard.is_some()).then(|| registry.clone());
let summary = Arc::new(SummarySink::new());
let mut inner: Vec<Arc<dyn rhei_tui::EventSink>> = vec![frontend.sink.clone(), summary.clone()];
if let Some(dashboard) = &dashboard {
inner.push(dashboard.clone());
}
let sink: Arc<dyn rhei_tui::EventSink> = Arc::new(rhei_tui::Tee::new(inner));
let is_tui = frontend.is_tui;
ActiveRunFrontend { sink, is_tui, dashboard, summary, intervene, _frontend: Some(frontend) }
}
#[allow(clippy::too_many_arguments)]
fn transition_dashboard_gate(
input: &Path,
machine: &rhei_validator::StateMachine,
callback_paths: &CallbackPaths,
task_id_str: &str,
from: &str,
to: &str,
result: Option<&str>,
no_callbacks: bool,
) -> MietteResult<String> {
let loaded = load_plan(input)?;
let task = find_task_by_id_str(&loaded.rhei.tasks, task_id_str)
.ok_or_else(|| {
miette!(
help = format!(
"list the task ids in this plan with: rhei list {}",
shell_quote(&input.display().to_string())
),
"task '{}' not found in the plan",
task_id_str
)
})?;
let current_state = normalized_state_name(task.state.as_str(), machine);
if current_state != from {
return Err(miette!(
help = format!(
"someone moved the task since you looked. Re-read its current state with: \
rhei list {}",
shell_quote(&input.display().to_string())
),
"conflict: Task {} is in state '{}', expected '{}'",
task_id_str,
task.state,
from
));
}
if !machine.states.get(¤t_state).map(|def| def.gating).unwrap_or(false) {
return Err(miette!(
help = "only human-gate states are released this way. Advance a non-gating state \
with `rhei transition` or let `rhei run` drive it. See which states gate \
with: rhei states",
"Task {} is in state '{}', which is not a gating state",
task_id_str,
current_state
));
}
let explicit_transition =
machine.transitions().iter().any(|rule| rule.from.0 == from && rule.to.0 == to);
if !explicit_transition {
return Err(miette!(
help = format!(
"the state machine declares no '{from}' -> '{to}' edge. List the edges \
leaving '{from}' with: rhei states"
),
"transition from '{}' to '{}' is not an explicit human-gate transition",
from,
to
));
}
let result = result.map(str::trim).filter(|message| !message.is_empty());
let route = loaded.task_route(task_id_str, input);
execute_transition(
TransitionFiles {
task_file: &route.task_file,
metadata_file: &route.metadata_file,
metadata_id: &route.metadata_id,
artifact_root: &route.execution_root,
artifact_id: task_id_str,
},
callback_paths,
machine,
&route.local_id,
from,
to,
result,
no_callbacks,
)
}
fn load_plan_for_dashboard(
plan_path: &Path,
machines: &rhei_validator::MachineSet,
) -> Option<rhei_viz_model::VizModel> {
let loaded = load_plan(plan_path).ok()?;
let default_root = execution_workspace_root(plan_path);
Some(rhei_viz::build_set_with_history_roots(
&loaded.rhei,
machines,
&default_root,
&loaded.task_roots,
))
}
impl
From<(
StandaloneExecutionFlags,
AgentExecutionFlags,
ProgramExecutionFlags,
SnapshotExecutionFlags,
)> for RunOptions
{
fn from(
(standalone, agent, program, snapshot): (
StandaloneExecutionFlags,
AgentExecutionFlags,
ProgramExecutionFlags,
SnapshotExecutionFlags,
),
) -> Self {
Self { standalone, agent, program, snapshot }
}
}