use std::os::unix::process::{CommandExt, ExitStatusExt};
use std::process::Command;
use clap::Args as ClapArgs;
use serde_json::{json, Value};
use octl_core::{append_and_apply_event, read_manifest_opt, read_node_opt, NodeId, RunPaths};
use crate::error::CliError;
use crate::run::{from_core, parse_node_id, parse_run_id, run_paths_exact};
const SHIM_IGNORED_SIGNALS: [libc::c_int; 4] =
[libc::SIGINT, libc::SIGTERM, libc::SIGHUP, libc::SIGQUIT];
const SPAWN_FAILURE_EXIT_CODE: i32 = 127;
#[derive(ClapArgs, Debug)]
pub struct RunWorkerArgs {
pub run_id: String,
pub node_id: String,
#[arg(last = true, required = true)]
pub command: Vec<String>,
}
pub fn dispatch(args: RunWorkerArgs) -> Result<(), CliError> {
let run_id = parse_run_id(&args.run_id)?;
let node_id = parse_node_id(&args.node_id)?;
let root = crate::home::root_dir()?;
let paths = run_paths_exact(&root, &run_id)?;
if read_manifest_opt(&paths).map_err(from_core)?.is_none() {
return Err(
CliError::user("run_not_found", format!("no run with id {}", args.run_id))
.with_invalid_value(&args.run_id),
);
}
if read_node_opt(&paths, &node_id)
.map_err(from_core)?
.is_none()
{
return Err(CliError::user(
"node_not_found",
format!("no node {} in run {}", args.node_id, args.run_id),
)
.with_invalid_value(&args.node_id));
}
let (program, prog_args) = args
.command
.split_first()
.ok_or_else(|| CliError::user("missing_worker_command", "no worker command after `--`"))?;
for sig in SHIM_IGNORED_SIGNALS {
unsafe { libc::signal(sig, libc::SIG_IGN) };
}
let mut cmd = Command::new(program);
cmd.args(prog_args);
unsafe {
cmd.pre_exec(|| {
for sig in SHIM_IGNORED_SIGNALS {
libc::signal(sig, libc::SIG_DFL);
}
Ok(())
});
}
let status = match cmd.status() {
Ok(s) => s,
Err(e) => {
let mut data = serde_json::Map::new();
data.insert("exit_code".into(), json!(SPAWN_FAILURE_EXIT_CODE));
record_worker_exit(&paths, &node_id, &args, Value::Object(data));
return Err(CliError::system(
"worker_spawn_failed",
format!("could not launch worker `{program}`: {e}"),
));
}
};
let exit_code = status.code();
let signal = status.signal();
let mut data = serde_json::Map::new();
match (signal, exit_code) {
(Some(s), _) => {
data.insert("signal".into(), json!(s));
}
(None, Some(c)) => {
data.insert("exit_code".into(), json!(c));
}
(None, None) => {
data.insert("exit_code".into(), json!(-1));
}
}
record_worker_exit(&paths, &node_id, &args, Value::Object(data));
let code = exit_code.unwrap_or_else(|| 128 + signal.unwrap_or(1));
crate::cli::flush_logs();
std::process::exit(code);
}
fn record_worker_exit(paths: &RunPaths, node_id: &NodeId, args: &RunWorkerArgs, data: Value) {
if let Err(e) = append_and_apply_event(paths, "worker.exited", Some(node_id), None, data) {
let msg = from_core(e).message;
eprintln!(
"orchestratectl run-worker: failed to record worker.exited for {}/{}: {msg}",
args.run_id, args.node_id
);
tracing::warn!(
target: "orchestratectl::run_worker",
run_id = %args.run_id,
node_id = %args.node_id,
error = %msg,
"failed to record worker.exited event; relying on the crash backstop"
);
}
}