1use std::collections::HashMap;
34use std::process::{Child, Command, Stdio};
35use std::sync::{Arc, Mutex, MutexGuard};
36use std::time::{Duration, Instant};
37
38use super::*;
39use onlyne_config::layout::RoleWorkspace;
40
41const TERMINATE_GRACE: Duration = Duration::from_secs(5);
44const REAP_POLL: Duration = Duration::from_millis(25);
46const OUTPUT_TAIL_LINES: usize = 200;
50const OUTPUT_TAIL_BYTES: usize = 16 * 1024;
52
53#[derive(Clone, Default)]
54pub struct ExecBackend {
55 children: Arc<Mutex<HashMap<String, Child>>>,
56}
57
58fn guard(mutex: &Mutex<HashMap<String, Child>>) -> Result<MutexGuard<'_, HashMap<String, Child>>> {
60 mutex
61 .lock()
62 .map_err(|_| anyhow::anyhow!("exec backend child map lock poisoned"))
63}
64
65impl ExecBackend {
66 pub fn new() -> Self {
67 Self::default()
68 }
69
70 fn pid_alive(pid: u32) -> bool {
76 #[cfg(unix)]
77 {
78 pid_alive_unix(pid)
79 }
80 #[cfg(windows)]
81 {
82 pid_alive_windows(pid)
83 }
84 }
85
86 fn pid_of(session: &SessionRef) -> Result<u32> {
87 session
88 .backend_ref
89 .get("pid")
90 .and_then(Value::as_u64)
91 .map(|pid| pid as u32)
92 .ok_or_else(|| anyhow::anyhow!("exec session ref missing pid"))
93 }
94
95 fn stop(child: &mut Child, force: bool) -> Result<()> {
106 if let Ok(Some(_)) = child.try_wait() {
107 return Ok(());
108 }
109 let pid = child.id();
110 if !force {
111 if !signal_group(pid, "TERM") {
112 signal_pid(pid, "TERM");
113 }
114 let deadline = Instant::now() + TERMINATE_GRACE;
115 while Instant::now() < deadline {
116 if let Ok(Some(_)) = child.try_wait() {
117 return Ok(());
118 }
119 std::thread::sleep(REAP_POLL);
120 }
121 }
122 signal_group(pid, "KILL");
125 let _ = child.kill();
126 let _ = child.wait();
127 Ok(())
128 }
129}
130
131#[cfg(unix)]
138fn signal_group(pid: u32, signal: &str) -> bool {
139 let Some(pgid) = unix_pid(pid) else {
140 return false;
141 };
142 let Some(sig) = unix_sig(signal) else {
143 return false;
144 };
145 send_signal(-pgid, sig)
146}
147
148#[cfg(windows)]
153fn signal_group(pid: u32, signal: &str) -> bool {
154 if signal != "TERM" {
155 return false;
156 }
157 generate_ctrl_break(pid)
158}
159
160#[cfg(unix)]
162fn signal_pid(pid: u32, signal: &str) {
163 let Some(pid) = unix_pid(pid) else {
164 return;
165 };
166 let Some(sig) = unix_sig(signal) else {
167 return;
168 };
169 let _ = send_signal(pid, sig);
170}
171
172#[cfg(unix)]
175fn unix_pid(pid: u32) -> Option<i32> {
176 i32::try_from(pid).ok().filter(|&pid| pid > 1)
177}
178
179#[cfg(unix)]
180fn unix_sig(signal: &str) -> Option<i32> {
181 match signal {
182 "TERM" => Some(libc::SIGTERM),
183 "KILL" => Some(libc::SIGKILL),
184 _ => None,
185 }
186}
187
188#[cfg(unix)]
189fn send_signal(pid: i32, sig: i32) -> bool {
190 if pid == 0 || pid == -1 {
191 return false;
192 }
193 unsafe { libc::kill(pid, sig) == 0 }
194}
195
196#[cfg(unix)]
197fn pid_alive_unix(pid: u32) -> bool {
198 let Some(pid) = unix_pid(pid) else {
199 return false;
200 };
201 unsafe { libc::kill(pid, 0) == 0 }
202}
203
204#[cfg(windows)]
205fn signal_pid(pid: u32, signal: &str) {
206 if signal == "TERM" {
207 let _ = generate_ctrl_break(pid);
208 }
209}
210
211#[cfg(windows)]
212fn generate_ctrl_break(pid: u32) -> bool {
213 use windows_sys::Win32::System::Console::{CTRL_BREAK_EVENT, GenerateConsoleCtrlEvent};
214 if pid == 0 {
216 return false;
217 }
218 unsafe { GenerateConsoleCtrlEvent(CTRL_BREAK_EVENT, pid) != 0 }
219}
220
221#[cfg(windows)]
222fn pid_alive_windows(pid: u32) -> bool {
223 use windows_sys::Win32::Foundation::{
224 CloseHandle, GetLastError, INVALID_HANDLE_VALUE, STILL_ACTIVE,
225 };
226 use windows_sys::Win32::System::Threading::{
227 GetExitCodeProcess, OpenProcess, PROCESS_QUERY_LIMITED_INFORMATION,
228 };
229 if pid == 0 {
230 return false;
231 }
232 unsafe {
233 let handle = OpenProcess(PROCESS_QUERY_LIMITED_INFORMATION, 0, pid);
234 if handle.is_null() || handle == INVALID_HANDLE_VALUE {
235 return false;
236 }
237 let mut code = 0u32;
238 let ok = GetExitCodeProcess(handle, &mut code);
239 let err = GetLastError();
240 CloseHandle(handle);
241 if ok == 0 {
242 let _ = err;
244 return false;
245 }
246 code == STILL_ACTIVE as u32
248 }
249}
250
251fn read_output_tail(path: &str) -> Option<String> {
255 let data = std::fs::read(path).ok()?;
256 let start = data.len().saturating_sub(OUTPUT_TAIL_BYTES);
257 let slice = if start == 0 {
258 data.as_slice()
259 } else {
260 match data[start..].iter().position(|&b| b == b'\n') {
261 Some(offset) => &data[start + offset + 1..],
262 None => &data[start..],
263 }
264 };
265 let text = String::from_utf8_lossy(slice);
266 let lines: Vec<&str> = text.lines().collect();
267 let skip = lines.len().saturating_sub(OUTPUT_TAIL_LINES);
268 Some(lines[skip..].join("\n"))
269}
270
271impl SessionBackend for ExecBackend {
272 fn name(&self) -> &'static str {
273 "exec"
274 }
275 fn capabilities(&self) -> Capabilities {
276 Capabilities {
277 spawn: true,
278 attach: true,
279 probe: true,
280 close: true,
281 focus: false,
282 rename: false,
283 }
284 }
285 fn available(&self) -> Result<bool> {
286 Ok(true)
287 }
288 fn spawn(&self, spec: SpawnSpec) -> Result<SessionRef> {
289 let Some(program) = spec.command.first() else {
290 return Err(anyhow::anyhow!(
291 "exec: the role's `[client.runtime] command` is empty; there is nothing to run \
292 for task {}",
293 spec.task_id
294 ));
295 };
296 let layout = RoleWorkspace::resolve(&spec.cwd);
297 let logs = layout.logs_dir();
298 std::fs::create_dir_all(&logs)
299 .map_err(|error| anyhow::anyhow!("create {}: {error}", logs.display()))?;
300 let log_path = layout.session_log_path(&spec.task_id);
301 let log = std::fs::OpenOptions::new()
302 .create(true)
303 .append(true)
304 .open(&log_path)
305 .map_err(|error| anyhow::anyhow!("open {}: {error}", log_path.display()))?;
306 let errors = log
307 .try_clone()
308 .map_err(|error| anyhow::anyhow!("clone {}: {error}", log_path.display()))?;
309 let mut command = Command::new(program);
310 command
311 .args(&spec.command[1..])
312 .current_dir(&spec.cwd)
313 .envs(&spec.env)
314 .stdin(Stdio::piped())
316 .stdout(Stdio::from(log))
317 .stderr(Stdio::from(errors));
318 #[cfg(unix)]
319 {
320 use std::os::unix::process::CommandExt;
321 command.process_group(0);
322 }
323 #[cfg(windows)]
324 {
325 use std::os::windows::process::CommandExt;
326 command.creation_flags(0x0000_0200);
328 }
329 let child = command
330 .spawn()
331 .map_err(|error| anyhow::anyhow!("spawn {}: {error}", spec.command.join(" ")))?;
332 let pid = child.id();
333 tracing::info!(task = %spec.task_id, pid, log = %log_path.display(), "exec session started");
334 guard(&self.children)?.insert(spec.task_id.clone(), child);
335 Ok(SessionRef {
336 task_id: spec.task_id.clone(),
337 backend: self.name().into(),
338 backend_ref: serde_json::json!({
339 "id": spec.task_id,
340 "pid": pid,
341 "pgid": pid,
342 "log": log_path.to_string_lossy(),
343 }),
344 generation: 1,
345 })
346 }
347 fn attach(&self, session: &SessionRef) -> Result<SessionRef> {
348 let pid = Self::pid_of(session)?;
349 if Self::pid_alive(pid) {
350 Ok(session.clone())
351 } else {
352 anyhow::bail!("exec session {} is gone (pid {pid})", session.task_id)
353 }
354 }
355 fn probe(&self, session: &SessionRef) -> Result<ResourceProbe> {
356 let mut children = guard(&self.children)?;
357 if let Some(child) = children.get_mut(&session.task_id) {
358 return match child.try_wait() {
361 Ok(Some(status)) => {
362 let mut detail = serde_json::json!({"exit": status.code()});
363 if let Some(path) = session.backend_ref.get("log").and_then(Value::as_str) {
364 if let Some(tail) = read_output_tail(path) {
365 detail["output_tail"] = Value::String(tail);
366 }
367 }
368 Ok(ResourceProbe {
369 alive: false,
370 attached: false,
371 detail: Some(detail),
372 })
373 }
374 Ok(None) => Ok(ResourceProbe {
375 alive: true,
376 attached: true,
377 detail: Some(serde_json::json!({"pid": child.id()})),
378 }),
379 Err(error) => Ok(ResourceProbe {
380 alive: false,
381 attached: false,
382 detail: Some(serde_json::json!({"error": error.to_string()})),
383 }),
384 };
385 }
386 drop(children);
387 let pid = Self::pid_of(session)?;
388 let alive = Self::pid_alive(pid);
389 Ok(ResourceProbe {
390 alive,
391 attached: alive,
392 detail: Some(serde_json::json!({"pid": pid, "reattached": true})),
393 })
394 }
395 fn close(&self, session: &SessionRef, reason: CloseReason, force: bool) -> Result<()> {
396 let child = guard(&self.children)?.remove(&session.task_id);
397 match child {
398 Some(mut child) => {
399 Self::stop(&mut child, force)?;
400 tracing::info!(task = %session.task_id, ?reason, "exec session closed");
401 Ok(())
402 }
403 None => {
404 let pid = Self::pid_of(session)?;
408 if Self::pid_alive(pid) {
409 anyhow::bail!(
410 "exec session {} has no handle; pid {pid} is still alive",
411 session.task_id
412 )
413 }
414 Ok(())
415 }
416 }
417 }
418}