use std::collections::HashMap;
use std::process::{Child, Command, Stdio};
use std::sync::{Arc, Mutex, MutexGuard};
use std::time::{Duration, Instant};
use super::*;
use onlyne_config::layout::RoleWorkspace;
const TERMINATE_GRACE: Duration = Duration::from_secs(5);
const REAP_POLL: Duration = Duration::from_millis(25);
const OUTPUT_TAIL_LINES: usize = 200;
const OUTPUT_TAIL_BYTES: usize = 16 * 1024;
#[derive(Clone, Default)]
pub struct ExecBackend {
children: Arc<Mutex<HashMap<String, Child>>>,
}
fn guard(mutex: &Mutex<HashMap<String, Child>>) -> Result<MutexGuard<'_, HashMap<String, Child>>> {
mutex
.lock()
.map_err(|_| anyhow::anyhow!("exec backend child map lock poisoned"))
}
impl ExecBackend {
pub fn new() -> Self {
Self::default()
}
fn pid_alive(pid: u32) -> bool {
#[cfg(unix)]
{
pid_alive_unix(pid)
}
#[cfg(windows)]
{
pid_alive_windows(pid)
}
}
fn pid_of(session: &SessionRef) -> Result<u32> {
session
.backend_ref
.get("pid")
.and_then(Value::as_u64)
.map(|pid| pid as u32)
.ok_or_else(|| anyhow::anyhow!("exec session ref missing pid"))
}
fn stop(child: &mut Child, force: bool) -> Result<()> {
if let Ok(Some(_)) = child.try_wait() {
return Ok(());
}
let pid = child.id();
if !force {
if !signal_group(pid, "TERM") {
signal_pid(pid, "TERM");
}
let deadline = Instant::now() + TERMINATE_GRACE;
while Instant::now() < deadline {
if let Ok(Some(_)) = child.try_wait() {
return Ok(());
}
std::thread::sleep(REAP_POLL);
}
}
signal_group(pid, "KILL");
let _ = child.kill();
let _ = child.wait();
Ok(())
}
}
#[cfg(unix)]
fn signal_group(pid: u32, signal: &str) -> bool {
let Some(pgid) = unix_pid(pid) else {
return false;
};
let Some(sig) = unix_sig(signal) else {
return false;
};
send_signal(-pgid, sig)
}
#[cfg(windows)]
fn signal_group(pid: u32, signal: &str) -> bool {
if signal != "TERM" {
return false;
}
generate_ctrl_break(pid)
}
#[cfg(unix)]
fn signal_pid(pid: u32, signal: &str) {
let Some(pid) = unix_pid(pid) else {
return;
};
let Some(sig) = unix_sig(signal) else {
return;
};
let _ = send_signal(pid, sig);
}
#[cfg(unix)]
fn unix_pid(pid: u32) -> Option<i32> {
i32::try_from(pid).ok().filter(|&pid| pid > 1)
}
#[cfg(unix)]
fn unix_sig(signal: &str) -> Option<i32> {
match signal {
"TERM" => Some(libc::SIGTERM),
"KILL" => Some(libc::SIGKILL),
_ => None,
}
}
#[cfg(unix)]
fn send_signal(pid: i32, sig: i32) -> bool {
if pid == 0 || pid == -1 {
return false;
}
unsafe { libc::kill(pid, sig) == 0 }
}
#[cfg(unix)]
fn pid_alive_unix(pid: u32) -> bool {
let Some(pid) = unix_pid(pid) else {
return false;
};
unsafe { libc::kill(pid, 0) == 0 }
}
#[cfg(windows)]
fn signal_pid(pid: u32, signal: &str) {
if signal == "TERM" {
let _ = generate_ctrl_break(pid);
}
}
#[cfg(windows)]
fn generate_ctrl_break(pid: u32) -> bool {
use windows_sys::Win32::System::Console::{CTRL_BREAK_EVENT, GenerateConsoleCtrlEvent};
if pid == 0 {
return false;
}
unsafe { GenerateConsoleCtrlEvent(CTRL_BREAK_EVENT, pid) != 0 }
}
#[cfg(windows)]
fn pid_alive_windows(pid: u32) -> bool {
use windows_sys::Win32::Foundation::{
CloseHandle, GetLastError, INVALID_HANDLE_VALUE, STILL_ACTIVE,
};
use windows_sys::Win32::System::Threading::{
GetExitCodeProcess, OpenProcess, PROCESS_QUERY_LIMITED_INFORMATION,
};
if pid == 0 {
return false;
}
unsafe {
let handle = OpenProcess(PROCESS_QUERY_LIMITED_INFORMATION, 0, pid);
if handle.is_null() || handle == INVALID_HANDLE_VALUE {
return false;
}
let mut code = 0u32;
let ok = GetExitCodeProcess(handle, &mut code);
let err = GetLastError();
CloseHandle(handle);
if ok == 0 {
let _ = err;
return false;
}
code == STILL_ACTIVE as u32
}
}
fn read_output_tail(path: &str) -> Option<String> {
let data = std::fs::read(path).ok()?;
let start = data.len().saturating_sub(OUTPUT_TAIL_BYTES);
let slice = if start == 0 {
data.as_slice()
} else {
match data[start..].iter().position(|&b| b == b'\n') {
Some(offset) => &data[start + offset + 1..],
None => &data[start..],
}
};
let text = String::from_utf8_lossy(slice);
let lines: Vec<&str> = text.lines().collect();
let skip = lines.len().saturating_sub(OUTPUT_TAIL_LINES);
Some(lines[skip..].join("\n"))
}
impl SessionBackend for ExecBackend {
fn name(&self) -> &'static str {
"exec"
}
fn capabilities(&self) -> Capabilities {
Capabilities {
spawn: true,
attach: true,
probe: true,
close: true,
focus: false,
rename: false,
}
}
fn available(&self) -> Result<bool> {
Ok(true)
}
fn spawn(&self, spec: SpawnSpec) -> Result<SessionRef> {
let Some(program) = spec.command.first() else {
return Err(anyhow::anyhow!(
"exec: the role's `[client.runtime] command` is empty; there is nothing to run \
for task {}",
spec.task_id
));
};
let layout = RoleWorkspace::resolve(&spec.cwd);
let logs = layout.logs_dir();
std::fs::create_dir_all(&logs)
.map_err(|error| anyhow::anyhow!("create {}: {error}", logs.display()))?;
let log_path = layout.session_log_path(&spec.task_id);
let log = std::fs::OpenOptions::new()
.create(true)
.append(true)
.open(&log_path)
.map_err(|error| anyhow::anyhow!("open {}: {error}", log_path.display()))?;
let errors = log
.try_clone()
.map_err(|error| anyhow::anyhow!("clone {}: {error}", log_path.display()))?;
let mut command = Command::new(program);
command
.args(&spec.command[1..])
.current_dir(&spec.cwd)
.envs(&spec.env)
.stdin(Stdio::piped())
.stdout(Stdio::from(log))
.stderr(Stdio::from(errors));
#[cfg(unix)]
{
use std::os::unix::process::CommandExt;
command.process_group(0);
}
#[cfg(windows)]
{
use std::os::windows::process::CommandExt;
command.creation_flags(0x0000_0200);
}
let child = command
.spawn()
.map_err(|error| anyhow::anyhow!("spawn {}: {error}", spec.command.join(" ")))?;
let pid = child.id();
tracing::info!(task = %spec.task_id, pid, log = %log_path.display(), "exec session started");
guard(&self.children)?.insert(spec.task_id.clone(), child);
Ok(SessionRef {
task_id: spec.task_id.clone(),
backend: self.name().into(),
backend_ref: serde_json::json!({
"id": spec.task_id,
"pid": pid,
"pgid": pid,
"log": log_path.to_string_lossy(),
}),
generation: 1,
})
}
fn attach(&self, session: &SessionRef) -> Result<SessionRef> {
let pid = Self::pid_of(session)?;
if Self::pid_alive(pid) {
Ok(session.clone())
} else {
anyhow::bail!("exec session {} is gone (pid {pid})", session.task_id)
}
}
fn probe(&self, session: &SessionRef) -> Result<ResourceProbe> {
let mut children = guard(&self.children)?;
if let Some(child) = children.get_mut(&session.task_id) {
return match child.try_wait() {
Ok(Some(status)) => {
let mut detail = serde_json::json!({"exit": status.code()});
if let Some(path) = session.backend_ref.get("log").and_then(Value::as_str) {
if let Some(tail) = read_output_tail(path) {
detail["output_tail"] = Value::String(tail);
}
}
Ok(ResourceProbe {
alive: false,
attached: false,
detail: Some(detail),
})
}
Ok(None) => Ok(ResourceProbe {
alive: true,
attached: true,
detail: Some(serde_json::json!({"pid": child.id()})),
}),
Err(error) => Ok(ResourceProbe {
alive: false,
attached: false,
detail: Some(serde_json::json!({"error": error.to_string()})),
}),
};
}
drop(children);
let pid = Self::pid_of(session)?;
let alive = Self::pid_alive(pid);
Ok(ResourceProbe {
alive,
attached: alive,
detail: Some(serde_json::json!({"pid": pid, "reattached": true})),
})
}
fn close(&self, session: &SessionRef, reason: CloseReason, force: bool) -> Result<()> {
let child = guard(&self.children)?.remove(&session.task_id);
match child {
Some(mut child) => {
Self::stop(&mut child, force)?;
tracing::info!(task = %session.task_id, ?reason, "exec session closed");
Ok(())
}
None => {
let pid = Self::pid_of(session)?;
if Self::pid_alive(pid) {
anyhow::bail!(
"exec session {} has no handle; pid {pid} is still alive",
session.task_id
)
}
Ok(())
}
}
}
}