const ATTACH_POLL: Duration = Duration::from_millis(150);
pub(crate) fn attach_command(
reference: Option<&str>,
json: bool,
since: u64,
wait: bool,
) -> MietteResult<()> {
let descriptor = resolve_run(reference)?;
let events_path = resolve_workspace_relative(&descriptor, &descriptor.events);
if json {
require_event_log(&descriptor, &events_path)?;
return stream_run_json(&descriptor, &events_path, since, wait);
}
if wait {
return wait_for_run_end(&descriptor);
}
match descriptor.liveness() {
Liveness::Ended | Liveness::Gone => {
report_finished_run(&descriptor);
Ok(())
}
Liveness::Live => {
require_event_log(&descriptor, &events_path)?;
attach_surface(&descriptor, &events_path, None)
}
Liveness::Unknown(reason) => {
require_event_log(&descriptor, &events_path)?;
attach_surface(&descriptor, &events_path, Some(reason))
}
}
}
fn require_event_log(descriptor: &RunDescriptor, events_path: &Path) -> MietteResult<()> {
if descriptor.liveness().has_ended() || events_path.is_file() {
return Ok(());
}
Err(miette!(
help = "the run is live but is not publishing an event log; watch it in its own \
terminal, or open the browser dashboard",
"run {} has no event log at {} to follow",
descriptor.id,
events_path.display()
))
}
fn wait_for_run_end(descriptor: &RunDescriptor) -> MietteResult<()> {
poll_until_run_ends(descriptor)?;
let final_state = settled_end(descriptor);
report_finished_run(&final_state);
exit_with_run_status(&final_state)
}
fn poll_until_run_ends(descriptor: &RunDescriptor) -> MietteResult<()> {
let mut undecided = UndecidedWatch::default();
loop {
match descriptor.liveness() {
Liveness::Ended | Liveness::Gone => return Ok(()),
Liveness::Live => undecided.decided(),
Liveness::Unknown(reason) => {
if undecided.exhausted(&reason) {
return Err(undecided_run(descriptor, undecided.reason()));
}
}
}
std::thread::sleep(ATTACH_POLL);
}
}
fn settled_end(descriptor: &RunDescriptor) -> RunDescriptor {
let deadline = Instant::now() + UNDECIDED_GRACE;
loop {
let current = recorded_end(descriptor);
if current.exit_code.is_some() || Instant::now() >= deadline {
return current;
}
std::thread::sleep(ATTACH_POLL);
}
}
fn undecided_run(descriptor: &RunDescriptor, reason: &str) -> miette::Report {
miette!(
help = "fix what is in the way and ask again; the run itself is untouched",
"could not tell whether run {} is still running: {reason}",
descriptor.id
)
}
fn recorded_end(descriptor: &RunDescriptor) -> RunDescriptor {
read_descriptor(&run_descriptor_path(&descriptor.workspace))
.filter(|current| current.id == descriptor.id)
.unwrap_or_else(|| descriptor.clone())
}
fn resolve_workspace_relative(descriptor: &RunDescriptor, path: &Path) -> PathBuf {
if path.is_absolute() {
path.to_path_buf()
} else {
descriptor.workspace.join(path)
}
}
fn report_finished_run(descriptor: &RunDescriptor) {
println!("Run {} has ended.", descriptor.id);
match descriptor.exit_code {
Some(0) => println!(" It exited 0."),
Some(code) => println!(" It exited {code}."),
None => println!(" It recorded no exit status."),
}
for (label, relative) in
[("Report", "runtime/run-report.md"), ("Dashboard", "runtime/dashboard.html")]
{
let path = descriptor.workspace.join(relative);
if path.is_file() {
println!(" {label}: {}", path.display());
}
}
}
fn exit_with_run_status(descriptor: &RunDescriptor) -> MietteResult<()> {
match recorded_end(descriptor).exit_code.or(descriptor.exit_code) {
Some(0) => Ok(()),
Some(code) => std::process::exit(code),
None => Err(miette!(
help = "read what it managed to write: runtime/run.log and runtime/run-report.md",
"run {} recorded no exit status, so it did not end on its own",
descriptor.id
)),
}
}
fn stream_run_json(
descriptor: &RunDescriptor,
events_path: &Path,
since: u64,
wait: bool,
) -> MietteResult<()> {
let mut reader = rhei_tui::EventLogReader::open(events_path);
let mut ending = EndingWatch::default();
loop {
for record in reader.poll() {
if record.seq.is_some_and(|seq| seq <= since) {
continue;
}
ending.note(&record.event);
println!(
"{}",
rhei_tui::encode_event(
record.seq,
&record.event,
record.ts,
Some(&descriptor.workspace)
)
);
}
match ending.verdict(descriptor) {
Ending::Ended => break,
Ending::Undecided(reason) => return Err(undecided_run(descriptor, &reason)),
Ending::Following => {}
}
std::thread::sleep(ATTACH_POLL);
}
if wait {
poll_until_run_ends(descriptor)?;
return exit_with_run_status(&settled_end(descriptor));
}
Ok(())
}
fn attach_surface(
descriptor: &RunDescriptor,
events_path: &Path,
undecided: Option<String>,
) -> MietteResult<()> {
use rhei_tui::EventSink as _;
install_interrupt_handlers();
let machines = attach_machines(descriptor)?;
let plan_path = descriptor.plan.clone();
let loader: rhei_tui::PlanLoader =
Arc::new(move || load_plan_for_dashboard(&plan_path, &machines));
let context = rhei_tui::TuiContext {
workspace: descriptor.workspace.clone(),
plan_loader: Some(loader),
intervene: Some(Arc::new(ControlInterveneSink::new(descriptor.control_url.clone()))),
gate: Some(Arc::new(ControlGateSink::new(descriptor.control_url.clone()))),
stop_requested: Arc::new(interrupt_requested),
attached: true,
};
let parallel = descriptor.parallel.clamp(1, usize::from(u16::MAX)) as u16;
let tui = rhei_tui::TuiSink::start(parallel, 0, context).map_err(|err| {
miette!(
help = "attaching needs an interactive terminal; follow the run's records \
instead with: rhei attach --json",
"could not start the attached surface: {err}"
)
})?;
tui.emit(rhei_tui::RunEvent::Message {
level: rhei_tui::MessageLevel::Info,
text: format!(
"Attached to run {} (pid {}). Ctrl+C or q detaches; `rhei stop {}` stops the run.",
descriptor.id, descriptor.pid, descriptor.id
),
});
if descriptor.control_url.is_none() {
tui.emit(rhei_tui::RunEvent::Message {
level: rhei_tui::MessageLevel::Warn,
text: "This run serves no control endpoint: intervene and gate release are \
unavailable from here."
.to_string(),
});
}
if let Some(reason) = undecided {
tui.emit(rhei_tui::RunEvent::Message {
level: rhei_tui::MessageLevel::Warn,
text: format!(
"Could not confirm this run is still live ({reason}); attaching anyway."
),
});
}
follow_run(descriptor, events_path, &tui);
tui.finish();
Ok(())
}
enum Ending {
Following,
Ended,
Undecided(String),
}
#[derive(Default)]
struct EndingWatch {
saw_finish: bool,
drained_after_end: bool,
undecided: UndecidedWatch,
}
impl EndingWatch {
fn note(&mut self, event: &rhei_tui::RunEvent) {
self.saw_finish |= matches!(event, rhei_tui::RunEvent::RunFinished { .. });
}
fn verdict(&mut self, descriptor: &RunDescriptor) -> Ending {
if !self.saw_finish {
match descriptor.liveness() {
Liveness::Ended | Liveness::Gone => {}
Liveness::Live => {
self.undecided.decided();
return Ending::Following;
}
Liveness::Unknown(reason) => {
return if self.undecided.exhausted(&reason) {
Ending::Undecided(self.undecided.reason().to_string())
} else {
Ending::Following
};
}
}
}
if self.drained_after_end {
return Ending::Ended;
}
self.drained_after_end = true;
Ending::Following
}
}
fn follow_run(descriptor: &RunDescriptor, events_path: &Path, tui: &rhei_tui::TuiSink) {
use rhei_tui::EventSink as _;
let mut reader = rhei_tui::EventLogReader::open(events_path);
let mut tailer = AgentLogTailer::default();
let mut ending = EndingWatch::default();
loop {
for record in reader.poll() {
let event = record.event;
match &event {
rhei_tui::RunEvent::SlotAssigned { task, slot, log_path, .. } => {
tailer.follow(&descriptor.workspace, task, *slot, log_path);
}
rhei_tui::RunEvent::SlotReleased { task, slot, .. } => {
for output in tailer.release(task, *slot) {
tui.emit(output);
}
}
_ => {}
}
ending.note(&event);
tui.emit(event);
}
for output in tailer.poll() {
tui.emit(output);
}
if tui.screen_restored() || interrupt_requested() {
return;
}
match ending.verdict(descriptor) {
Ending::Ended => return,
Ending::Undecided(reason) => tui.emit(rhei_tui::RunEvent::Message {
level: rhei_tui::MessageLevel::Warn,
text: format!("Still cannot confirm this run is live ({reason})."),
}),
Ending::Following => {}
}
std::thread::sleep(ATTACH_POLL);
}
}
fn attach_machines(descriptor: &RunDescriptor) -> MietteResult<rhei_validator::MachineSet> {
let loaded = load_plan(&descriptor.plan).map_err(|err| {
miette!(
help = format!(
"the run's plan must still be readable to attach to it: {}",
descriptor.plan.display()
),
"could not load the plan of run {}: {err}",
descriptor.id
)
})?;
let resolved = resolve_state_machines_for_loaded_plan(
&descriptor.plan,
&loaded,
descriptor.state_machine.as_deref(),
)?;
Ok(ExecutionMachines::build(&resolved, &descriptor.plan)?.set)
}