Skip to main content

scv_tools/
process.rs

1//! Running a child process for a tool: spawn it in its own process group,
2//! drain its output while it runs, and always finish the whole group when
3//! it exits, times out, or is cancelled.
4
5use std::{
6    ffi::OsString, os::unix::process::CommandExt as _, path::PathBuf, sync::Arc, time::Duration,
7};
8
9use scv_core::{ToolError, ToolFailure, ToolOutput};
10use serde_json::json;
11use tokio::{
12    io::AsyncReadExt,
13    process::Command,
14    sync::Mutex,
15    task::JoinHandle,
16    time::{Instant, sleep, sleep_until, timeout, timeout_at},
17};
18
19use crate::delegate::{adapters, records};
20
21/// Give a native agent command its adapter environment: remove every
22/// inherited credential, endpoint, and state-location variable any adapter
23/// declares, then set `environment` (such as the relocated config home).
24pub fn apply_agent_environment(
25    command: &mut std::process::Command,
26    environment: &[(OsString, OsString)],
27) {
28    apply_agent_environment_from(
29        command,
30        std::env::vars_os().map(|(variable, _)| variable),
31        environment,
32    );
33}
34
35fn apply_agent_environment_from(
36    command: &mut std::process::Command,
37    inherited: impl IntoIterator<Item = OsString>,
38    environment: &[(OsString, OsString)],
39) {
40    for variable in inherited {
41        if adapters::is_removed_agent_variable(&variable) {
42            command.env_remove(variable);
43        }
44    }
45    command.envs(environment.iter().map(|(key, value)| (key, value)));
46}
47
48/// What to run for a tool, and its bounds.
49pub(crate) struct ProcessSpec {
50    pub(crate) executable: OsString,
51    pub(crate) args: Vec<OsString>,
52    pub(crate) cwd: PathBuf,
53    pub(crate) environment: Vec<(OsString, OsString)>,
54    /// Use only the supplied environment, as for reproducible tool probes.
55    pub(crate) clear_environment: bool,
56    /// Strip inherited agent credentials and state locations first
57    /// ([`apply_agent_environment`]), as for a delegated agent CLI.
58    pub(crate) sanitize_scv_environment: bool,
59    pub(crate) timeout: Duration,
60    pub(crate) output_limit: usize,
61}
62
63pub(crate) async fn execute_process(
64    spec: ProcessSpec,
65    cancellation: tokio_util::sync::CancellationToken,
66) -> Result<ToolOutput, ToolError> {
67    let deadline = Instant::now() + spec.timeout;
68    let mut child = spawn_process(&spec)?;
69    let pid = child_pid(&child)?;
70    let output = Arc::new(Mutex::new(BoundedOutput::new(spec.output_limit)));
71    let stdout_task = child
72        .stdout
73        .take()
74        .map(|stdout| tokio::spawn(drain_output(stdout, Arc::clone(&output))));
75    let stderr_task = child
76        .stderr
77        .take()
78        .map(|stderr| tokio::spawn(drain_output(stderr, Arc::clone(&output))));
79    let finished = supervise(
80        &mut child,
81        pid,
82        deadline,
83        cancellation,
84        stdout_task,
85        stderr_task,
86    )
87    .await;
88    records::untrack_spawned(pid as u32);
89    let finished = finished?;
90    let collected = output.lock().await;
91    let text = String::from_utf8_lossy(&collected.bytes).into_owned();
92    let content = json!({
93        "exit_code": finished.status.code(),
94        "timed_out": finished.timed_out,
95        "output": text,
96        "truncated": collected.truncated
97    })
98    .to_string();
99    let failure = if finished.timed_out {
100        Some(ToolFailure::Limit)
101    } else {
102        (!finished.status.success()).then_some(ToolFailure::Failed)
103    };
104    Ok(ToolOutput {
105        content,
106        failure,
107        truncated: collected.truncated,
108    })
109}
110
111pub(crate) fn spawn_process(spec: &ProcessSpec) -> Result<tokio::process::Child, ToolError> {
112    let mut command = Command::new(&spec.executable);
113    if spec.clear_environment {
114        command.env_clear();
115    }
116    if spec.sanitize_scv_environment {
117        apply_agent_environment(command.as_std_mut(), &spec.environment);
118    } else {
119        command.envs(spec.environment.iter().map(|(key, value)| (key, value)));
120    }
121    command
122        .args(&spec.args)
123        .current_dir(&spec.cwd)
124        .stdin(std::process::Stdio::null())
125        .stdout(std::process::Stdio::piped())
126        .stderr(std::process::Stdio::piped())
127        .kill_on_drop(true);
128    command.as_std_mut().process_group(0);
129    let child = command.spawn().map_err(|error| {
130        ToolError::unavailable(format!("launch {:?}: {error}", spec.executable))
131    })?;
132    if let Some(pid) = child.id() {
133        records::track_spawned(pid);
134    }
135    Ok(child)
136}
137
138pub(crate) fn child_pid(child: &tokio::process::Child) -> Result<i32, ToolError> {
139    child
140        .id()
141        .and_then(|pid| i32::try_from(pid).ok())
142        .ok_or_else(|| ToolError::failed("child process has no pid"))
143}
144
145pub(crate) struct Finished {
146    pub(crate) status: std::process::ExitStatus,
147    pub(crate) timed_out: bool,
148}
149
150/// Wait for a spawned process group until it exits, times out, or is
151/// cancelled, always finishing the whole group and draining its output.
152pub(crate) async fn supervise(
153    child: &mut tokio::process::Child,
154    pid: i32,
155    deadline: Instant,
156    cancellation: tokio_util::sync::CancellationToken,
157    stdout_task: Option<JoinHandle<()>>,
158    stderr_task: Option<JoinHandle<()>>,
159) -> Result<Finished, ToolError> {
160    enum Completion {
161        Exited(std::process::ExitStatus),
162        TimedOut,
163        Cancelled,
164    }
165    let completion = tokio::select! {
166        status = child.wait() => Completion::Exited(status.map_err(|error| ToolError::failed(format!("wait for child: {error}")))?),
167        () = cancellation.cancelled() => {
168            Completion::Cancelled
169        },
170        () = sleep_until(deadline) => Completion::TimedOut,
171    };
172
173    let (status, timed_out, drain_deadline) = match completion {
174        Completion::Exited(status) => {
175            let cleanup_deadline = deadline.min(Instant::now() + Duration::from_secs(2));
176            let status = terminate_group(pid, child, Some(status), cleanup_deadline, true).await?;
177            (
178                status,
179                false,
180                deadline.min(Instant::now() + Duration::from_millis(250)),
181            )
182        }
183        Completion::TimedOut => {
184            let status = terminate_group(pid, child, None, Instant::now(), false).await?;
185            (status, true, Instant::now() + Duration::from_millis(250))
186        }
187        Completion::Cancelled => {
188            let cleanup_deadline = Instant::now() + Duration::from_secs(2);
189            let _ = terminate_group(pid, child, None, cleanup_deadline, true).await;
190            finish_drain(stdout_task, Instant::now() + Duration::from_millis(250)).await;
191            finish_drain(stderr_task, Instant::now() + Duration::from_millis(250)).await;
192            return Err(ToolError::cancelled("process cancelled"));
193        }
194    };
195    finish_drain(stdout_task, drain_deadline).await;
196    finish_drain(stderr_task, drain_deadline).await;
197    Ok(Finished { status, timed_out })
198}
199
200async fn terminate_group(
201    pid: i32,
202    child: &mut tokio::process::Child,
203    mut status: Option<std::process::ExitStatus>,
204    deadline: Instant,
205    graceful: bool,
206) -> Result<std::process::ExitStatus, ToolError> {
207    // The child was spawned as the leader of its own group.
208    let group = u32::try_from(pid).ok().and_then(ProcessGroup::new);
209    if let Some(group) = group {
210        group.signal(if graceful {
211            libc::SIGTERM
212        } else {
213            libc::SIGKILL
214        });
215    }
216    while Instant::now() < deadline {
217        if status.is_none() {
218            status = child
219                .try_wait()
220                .map_err(|error| ToolError::failed(format!("wait for child: {error}")))?;
221        }
222        if !group.is_some_and(ProcessGroup::is_signalable)
223            && let Some(status) = status
224        {
225            return Ok(status);
226        }
227        sleep(Duration::from_millis(20)).await;
228    }
229    // Always finish the process group, even if its original leader already exited.
230    if let Some(group) = group {
231        group.signal(libc::SIGKILL);
232    }
233    if let Some(status) = status {
234        return Ok(status);
235    }
236    timeout(Duration::from_secs(1), child.wait())
237        .await
238        .map_err(|_| ToolError::failed("child did not exit after process-group kill"))?
239        .map_err(|error| ToolError::failed(format!("wait after KILL: {error}")))
240}
241
242/// A process group SCV may signal. It is never group 0 or 1 (this process's
243/// own group, or init's), which an unset or corrupt ID would otherwise
244/// address, so every signal SCV sends to a group goes through this type.
245#[derive(Debug, Clone, Copy, PartialEq, Eq)]
246pub(crate) struct ProcessGroup(i32);
247
248impl ProcessGroup {
249    /// The group with ID `pgid`, or `None` for an ID SCV must never signal.
250    pub(crate) fn new(pgid: u32) -> Option<Self> {
251        i32::try_from(pgid).ok().filter(|&id| id > 1).map(Self)
252    }
253
254    /// Send `signal` to every member of the group.
255    pub(crate) fn signal(self, signal: i32) {
256        // SAFETY: kill(2) takes plain integers and touches no memory of
257        // ours; the negative ID addresses the group, never 0 or 1 (`new`).
258        unsafe {
259            libc::kill(-self.0, signal);
260        }
261    }
262
263    /// Whether any member can still be signalled, zombies included.
264    pub(crate) fn is_signalable(self) -> bool {
265        // SAFETY: as in `signal`; signal 0 only checks existence and access.
266        let result = unsafe { libc::kill(-self.0, 0) };
267        result == 0 || std::io::Error::last_os_error().raw_os_error() == Some(libc::EPERM)
268    }
269}
270
271async fn finish_drain(task: Option<JoinHandle<()>>, deadline: Instant) {
272    let Some(mut task) = task else { return };
273    if timeout_at(deadline, &mut task).await.is_err() {
274        task.abort();
275        let _ = task.await;
276    }
277}
278
279/// Where a child's output goes as it is read.
280pub(crate) trait OutputSink: Send + 'static {
281    fn push(&mut self, bytes: &[u8]);
282}
283
284impl OutputSink for BoundedOutput {
285    fn push(&mut self, bytes: &[u8]) {
286        BoundedOutput::push(self, bytes);
287    }
288}
289
290pub(crate) async fn drain_output<R, S>(mut reader: R, output: Arc<Mutex<S>>)
291where
292    R: tokio::io::AsyncRead + Unpin,
293    S: OutputSink,
294{
295    let mut chunk = [0u8; 8192];
296    loop {
297        match reader.read(&mut chunk).await {
298            Ok(0) | Err(_) => break,
299            Ok(read) => output.lock().await.push(&chunk[..read]),
300        }
301    }
302}
303
304struct BoundedOutput {
305    bytes: Vec<u8>,
306    limit: usize,
307    truncated: bool,
308}
309
310impl BoundedOutput {
311    fn new(limit: usize) -> Self {
312        Self {
313            bytes: Vec::with_capacity(limit.min(8192)),
314            limit,
315            truncated: false,
316        }
317    }
318
319    fn push(&mut self, bytes: &[u8]) {
320        let remaining = self.limit.saturating_sub(self.bytes.len());
321        self.bytes
322            .extend_from_slice(&bytes[..bytes.len().min(remaining)]);
323        self.truncated |= bytes.len() > remaining;
324    }
325}
326
327#[cfg(test)]
328mod tests;