Skip to main content

sim_platform_ubuntu_pc/
process.rs

1use sim_lib_exec::{
2    BindingValue, DispatchEvidence, PrivateArtifactRef, ProcResult, ProcessAttempt,
3    ProcessCancellation, ProcessPort, ProcessReceipt, ProcessRefusal, ProcessRequest, ProgramRef,
4    ProjectRootRef, StopReceipt,
5};
6use std::{
7    collections::BTreeMap,
8    io::{Read, Write},
9    os::unix::process::CommandExt as _,
10    process::{Child, Command, Stdio},
11    sync::{Arc, Mutex, mpsc},
12    thread,
13    time::{Duration, Instant},
14};
15
16/// Ubuntu process adapter. All native process mechanics are confined here.
17#[derive(Clone, Debug, Default)]
18pub struct UbuntuProcess {
19    programs: BTreeMap<ProgramRef, std::path::PathBuf>,
20    roots: BTreeMap<ProjectRootRef, std::path::PathBuf>,
21    artifacts: BTreeMap<PrivateArtifactRef, std::path::PathBuf>,
22}
23impl UbuntuProcess {
24    /// Creates a capsule from boot-trusted native resource mappings.
25    #[must_use]
26    pub fn new(
27        programs: BTreeMap<ProgramRef, std::path::PathBuf>,
28        roots: BTreeMap<ProjectRootRef, std::path::PathBuf>,
29        artifacts: BTreeMap<PrivateArtifactRef, std::path::PathBuf>,
30    ) -> Self {
31        Self {
32            programs,
33            roots,
34            artifacts,
35        }
36    }
37}
38impl ProcessPort for UbuntuProcess {
39    fn run(&self, request: &ProcessRequest, cancellation: &ProcessCancellation) -> ProcessAttempt {
40        let Some(program) = self.programs.get(&request.program) else {
41            return ProcessAttempt::NotDispatched {
42                refusal: ProcessRefusal::Refused("program reference is not boot-authorized".into()),
43            };
44        };
45        let Some(root) = self.roots.get(&request.root) else {
46            return ProcessAttempt::NotDispatched {
47                refusal: ProcessRefusal::Refused(
48                    "project-root reference is not boot-authorized".into(),
49                ),
50            };
51        };
52        let root = match root.canonicalize() {
53            Ok(v) => v,
54            Err(e) => {
55                return ProcessAttempt::NotDispatched {
56                    refusal: ProcessRefusal::Refused(format!("project root unavailable: {e}")),
57                };
58            }
59        };
60        let mut environment = BTreeMap::new();
61        for (name, value) in request.environment.iter() {
62            let rendered = match value {
63                BindingValue::Literal(v) => v.clone(),
64                BindingValue::ProjectRoot(reference) => match self.roots.get(reference) {
65                    Some(v) => v.to_string_lossy().into_owned(),
66                    None => {
67                        return ProcessAttempt::NotDispatched {
68                            refusal: ProcessRefusal::Refused(
69                                "binding references wrong or unavailable resource kind".into(),
70                            ),
71                        };
72                    }
73                },
74                BindingValue::PrivateArtifact(reference) => match self.artifacts.get(reference) {
75                    Some(v) => v.to_string_lossy().into_owned(),
76                    None => {
77                        return ProcessAttempt::NotDispatched {
78                            refusal: ProcessRefusal::Refused(
79                                "binding references wrong or unavailable resource kind".into(),
80                            ),
81                        };
82                    }
83                },
84            };
85            environment.insert(name, rendered);
86        }
87        if request
88            .private_artifacts
89            .iter()
90            .any(|v| !self.artifacts.contains_key(v))
91        {
92            return ProcessAttempt::NotDispatched {
93                refusal: ProcessRefusal::Refused("private artifact is not boot-authorized".into()),
94            };
95        }
96        let mut command = Command::new(program);
97        command
98            .args(request.argv.iter().map(sim_lib_exec::ArgAtom::as_str))
99            .current_dir(root)
100            .env_clear()
101            .envs(environment)
102            .stdin(Stdio::piped())
103            .stdout(Stdio::piped())
104            .stderr(Stdio::piped())
105            .process_group(0);
106        let mut child = match command.spawn() {
107            Ok(v) => v,
108            Err(e) => {
109                return ProcessAttempt::NotDispatched {
110                    refusal: ProcessRefusal::SpawnFailed(e.to_string()),
111                };
112            }
113        };
114        run_child(&mut child, request, cancellation)
115    }
116}
117
118enum Event {
119    Stdout(Result<Vec<u8>, String>),
120    Stderr(Result<Vec<u8>, String>),
121    Stdin(Result<(), String>),
122}
123struct Budget {
124    left: usize,
125    truncated: bool,
126}
127pub(crate) fn run_child(
128    child: &mut Child,
129    request: &ProcessRequest,
130    cancellation: &ProcessCancellation,
131) -> ProcessAttempt {
132    let started = Instant::now();
133    let deadline = started.checked_add(Duration::from_millis(request.budget.timeout_ms));
134    let Some(deadline) = deadline else {
135        return unknown(started, "deadline", "timeout overflow");
136    };
137    let Some(stdout) = child.stdout.take() else {
138        return unknown(started, "capture", "stdout pipe missing");
139    };
140    let Some(stderr) = child.stderr.take() else {
141        return unknown(started, "capture", "stderr pipe missing");
142    };
143    let Some(stdin) = child.stdin.take() else {
144        return unknown(started, "capture", "stdin pipe missing");
145    };
146    let budget = Arc::new(Mutex::new(Budget {
147        left: request.budget.max_output_bytes,
148        truncated: false,
149    }));
150    let (tx, rx) = mpsc::channel();
151    reader(stdout, Arc::clone(&budget), tx.clone(), true);
152    reader(stderr, Arc::clone(&budget), tx.clone(), false);
153    writer(stdin, request.budget.stdin.clone(), tx.clone());
154    drop(tx);
155    let (mut status, mut out, mut err, mut input) = (None, None, None, None);
156    loop {
157        if status.is_none() {
158            status = match child.try_wait() {
159                Ok(v) => v,
160                Err(e) => return unknown(started, "reap", &e.to_string()),
161            };
162        }
163        while let Ok(event) = rx.try_recv() {
164            match event {
165                Event::Stdout(v) => out = Some(v),
166                Event::Stderr(v) => err = Some(v),
167                Event::Stdin(v) => input = Some(v),
168            }
169        }
170        if status.is_some() && out.is_some() && err.is_some() && input.is_some() {
171            let Some(status) = status else {
172                unreachable!("checked above")
173            };
174            let Some(out) = out.take() else {
175                unreachable!("checked above")
176            };
177            let Some(err) = err.take() else {
178                unreachable!("checked above")
179            };
180            let Some(input) = input.take() else {
181                unreachable!("checked above")
182            };
183            return completed(child, started, status, out, err, input, &budget);
184        }
185        let cancelled = cancellation.is_cancelled();
186        if cancelled || Instant::now() >= deadline {
187            let cleanup = terminate_tree(child);
188            let elapsed = u64::try_from(started.elapsed().as_nanos()).unwrap_or(u64::MAX);
189            return match cleanup {
190                Ok(detail) => {
191                    let receipt = StopReceipt {
192                        provider: "platform/site/ubuntu-pc".into(),
193                        elapsed_mono_ns: elapsed,
194                        cleanup: detail,
195                    };
196                    if cancelled {
197                        ProcessAttempt::StoppedAfterCancel { receipt }
198                    } else {
199                        ProcessAttempt::StoppedAfterTimeout { receipt }
200                    }
201                }
202                Err(detail) => ProcessAttempt::UnknownAfterDispatch {
203                    evidence: DispatchEvidence {
204                        provider: "platform/site/ubuntu-pc".into(),
205                        stage: "cleanup".into(),
206                        detail,
207                    },
208                },
209            };
210        }
211        thread::sleep(Duration::from_millis(2));
212    }
213}
214
215fn completed(
216    child: &mut Child,
217    started: Instant,
218    status: std::process::ExitStatus,
219    out: Result<Vec<u8>, String>,
220    err: Result<Vec<u8>, String>,
221    input: Result<(), String>,
222    budget: &Mutex<Budget>,
223) -> ProcessAttempt {
224    if let Err(error) = input {
225        return unknown(started, "stdin", &error);
226    }
227    let out = match out {
228        Ok(value) => value,
229        Err(error) => return unknown(started, "capture", &error),
230    };
231    let err = match err {
232        Ok(value) => value,
233        Err(error) => return unknown(started, "capture", &error),
234    };
235    let truncated = match budget.lock() {
236        Ok(value) => value.truncated,
237        Err(_) => return unknown(started, "capture", "output budget poisoned"),
238    };
239    if group_exists(child.id())
240        && let Err(detail) = terminate_tree(child)
241    {
242        return unknown(started, "cleanup", &detail);
243    }
244    ProcessAttempt::Completed {
245        receipt: ProcessReceipt {
246            provider: "platform/site/ubuntu-pc".into(),
247            elapsed_mono_ns: u64::try_from(started.elapsed().as_nanos()).unwrap_or(u64::MAX),
248            result: ProcResult {
249                stdout: String::from_utf8_lossy(&out).into_owned(),
250                stderr: String::from_utf8_lossy(&err).into_owned(),
251                exit_code: status.code().unwrap_or(-1),
252                truncated,
253            },
254        },
255    }
256}
257fn reader<R: Read + Send + 'static>(
258    mut stream: R,
259    budget: Arc<Mutex<Budget>>,
260    tx: mpsc::Sender<Event>,
261    stdout: bool,
262) {
263    thread::spawn(move || {
264        let mut result = Vec::new();
265        let mut chunk = [0; 4096];
266        let value = loop {
267            match stream.read(&mut chunk) {
268                Ok(0) => break Ok(result),
269                Ok(n) => {
270                    let Ok(mut b) = budget.lock() else {
271                        break Err("output budget poisoned".into());
272                    };
273                    let keep = n.min(b.left);
274                    result.extend_from_slice(&chunk[..keep]);
275                    b.left -= keep;
276                    b.truncated |= keep < n;
277                }
278                Err(e) => break Err(format!("capture: {e}")),
279            }
280        };
281        let _ = tx.send(if stdout {
282            Event::Stdout(value)
283        } else {
284            Event::Stderr(value)
285        });
286    });
287}
288fn writer(mut stream: std::process::ChildStdin, input: Option<Vec<u8>>, tx: mpsc::Sender<Event>) {
289    thread::spawn(move || {
290        let value = input.map_or(Ok(()), |v| {
291            stream.write_all(&v).map_err(|e| format!("stdin: {e}"))
292        });
293        drop(stream);
294        let _ = tx.send(Event::Stdin(value));
295    });
296}
297fn unknown(started: Instant, stage: &str, detail: &str) -> ProcessAttempt {
298    ProcessAttempt::UnknownAfterDispatch {
299        evidence: DispatchEvidence {
300            provider: "platform/site/ubuntu-pc".into(),
301            stage: stage.into(),
302            detail: format!("{}; elapsed_ns={}", detail, started.elapsed().as_nanos()),
303        },
304    }
305}
306fn terminate_tree(child: &mut Child) -> Result<String, String> {
307    let pgid = child.id();
308    let term = signal_group(pgid, "TERM");
309    let until = Instant::now() + Duration::from_millis(100);
310    while Instant::now() < until {
311        if child.try_wait().ok().flatten().is_some() {
312            break;
313        }
314        thread::sleep(Duration::from_millis(5));
315    }
316    let escalated = child.try_wait().ok().flatten().is_none() || group_exists(pgid);
317    let kill = escalated.then(|| signal_group(pgid, "KILL"));
318    let wait = child.wait().map(|_| ()).map_err(|e| format!("wait: {e}"));
319    let cleanup_deadline = Instant::now() + Duration::from_millis(500);
320    while group_exists(pgid) && Instant::now() < cleanup_deadline {
321        thread::sleep(Duration::from_millis(5));
322    }
323    let leaked = group_exists(pgid);
324    if leaked {
325        return Err(term
326            .err()
327            .or_else(|| kill.and_then(Result::err))
328            .or_else(|| wait.err())
329            .unwrap_or_else(|| "descendants remained after bounded cleanup".into()));
330    }
331    wait.map(|()| {
332        if escalated {
333            "process group received TERM, escalated to KILL, and was reaped".into()
334        } else {
335            "process group received TERM and was reaped without escalation".into()
336        }
337    })
338}
339fn signal_group(pgid: u32, signal: &str) -> Result<(), String> {
340    let status = Command::new("kill")
341        .args([format!("-{signal}"), "--".into(), format!("-{pgid}")])
342        .stdout(Stdio::null())
343        .stderr(Stdio::null())
344        .status()
345        .map_err(|e| e.to_string())?;
346    if status.success() {
347        Ok(())
348    } else {
349        Err(format!("kill {signal} group {pgid}: {status}"))
350    }
351}
352fn group_exists(pgid: u32) -> bool {
353    let Ok(entries) = std::fs::read_dir("/proc") else {
354        return true;
355    };
356    entries
357        .flatten()
358        .filter_map(|entry| std::fs::read_to_string(entry.path().join("stat")).ok())
359        .any(|stat| {
360            let Some((_, fields)) = stat.rsplit_once(')') else {
361                return false;
362            };
363            let mut fields = fields.split_whitespace();
364            let state = fields.next();
365            let _ppid = fields.next();
366            let group = fields.next().and_then(|v| v.parse::<u32>().ok());
367            group == Some(pgid) && !matches!(state, Some("Z" | "X"))
368        })
369}
370
371#[cfg(test)]
372mod tests {
373    use super::*;
374    use sim_lib_exec::{ArgAtom, ProcessBudget, SealedBindings};
375    use std::{
376        fs,
377        path::PathBuf,
378        time::{SystemTime, UNIX_EPOCH},
379    };
380    fn root(label: &str) -> PathBuf {
381        let path = std::env::temp_dir().join(format!(
382            "sim-platform-process-{label}-{}",
383            SystemTime::now()
384                .duration_since(UNIX_EPOCH)
385                .unwrap()
386                .as_nanos()
387        ));
388        fs::create_dir_all(&path).unwrap();
389        path
390    }
391    fn fixture(root: &std::path::Path, program: &str) -> (UbuntuProcess, ProcessRequest) {
392        let program_ref = ProgramRef::new("shell").unwrap();
393        let root_ref = ProjectRootRef::new("project").unwrap();
394        let process = UbuntuProcess::new(
395            BTreeMap::from([(program_ref.clone(), PathBuf::from(program))]),
396            BTreeMap::from([(root_ref.clone(), root.to_owned())]),
397            BTreeMap::new(),
398        );
399        let request = ProcessRequest {
400            program: program_ref,
401            argv: vec![],
402            root: root_ref,
403            environment: SealedBindings::literals([("SIM_PROCESS_VISIBLE".into(), "yes".into())])
404                .unwrap(),
405            private_artifacts: vec![],
406            budget: ProcessBudget {
407                timeout_ms: 1_000,
408                max_output_bytes: 192,
409                stdin: None,
410            },
411        };
412        (process, request)
413    }
414    #[test]
415    fn native_conformance_clears_environment_confines_cwd_and_caps_output() {
416        let root = root("success");
417        let (process, mut request) = fixture(&root, "/bin/sh");
418        request.argv=["-c","printf '%s:%s:' \"$SIM_PROCESS_VISIBLE\" \"${HOME-unset}\"; pwd; printf 12345678901234567890123456789012345678901234567890123456789012345678901234567890123456789012345678901234567890123456789012345678901234567890123456789012345678901234567890123456789012345678901234567890"].into_iter().map(|v|ArgAtom::new(v).unwrap()).collect();
419        let ProcessAttempt::Completed { receipt } =
420            process.run(&request, &ProcessCancellation::default())
421        else {
422            panic!("expected completion")
423        };
424        assert!(receipt.result.stdout.starts_with("yes:unset:"));
425        assert!(receipt.result.stdout.contains(root.to_str().unwrap()));
426        assert!(receipt.result.truncated);
427        fs::remove_dir_all(root).unwrap();
428    }
429    #[test]
430    fn native_conformance_refuses_unknown_refs_and_cleans_timed_out_tree() {
431        let root = root("cleanup");
432        let (process, mut timed) = fixture(&root, "/bin/sh");
433        let mut unknown = timed.clone();
434        unknown.program = ProgramRef::new("unknown").unwrap();
435        assert!(matches!(
436            process.run(&unknown, &ProcessCancellation::default()),
437            ProcessAttempt::NotDispatched { .. }
438        ));
439        timed.argv = ["-c", "sleep 5 & wait"]
440            .into_iter()
441            .map(|v| ArgAtom::new(v).unwrap())
442            .collect();
443        timed.budget.timeout_ms = 20;
444        let outcome = process.run(&timed, &ProcessCancellation::default());
445        assert!(
446            matches!(outcome, ProcessAttempt::StoppedAfterTimeout { .. }),
447            "{outcome:?}"
448        );
449        fs::remove_dir_all(root).unwrap();
450    }
451    #[test]
452    fn native_conformance_cancellation_wins_and_cleans_tree() {
453        let root = root("cancel");
454        let (process, mut request) = fixture(&root, "/bin/sleep");
455        request.argv.push(ArgAtom::new("5").unwrap());
456        let token = ProcessCancellation::default();
457        token.cancel();
458        assert!(matches!(
459            process.run(&request, &token),
460            ProcessAttempt::StoppedAfterCancel { .. }
461        ));
462        fs::remove_dir_all(root).unwrap();
463    }
464    #[test]
465    fn native_conformance_escalates_ignored_term_and_proves_group_reaped() {
466        let root = root("signal-escalation");
467        let (process, mut request) = fixture(&root, "/bin/sh");
468        request.argv = ["-c", "trap '' TERM; while :; do :; done"]
469            .into_iter()
470            .map(|value| ArgAtom::new(value).unwrap())
471            .collect();
472        request.budget.timeout_ms = 20;
473        let ProcessAttempt::StoppedAfterTimeout { receipt } =
474            process.run(&request, &ProcessCancellation::default())
475        else {
476            panic!("expected timeout with proven cleanup")
477        };
478        assert!(
479            receipt.cleanup.contains("escalated to KILL"),
480            "{}",
481            receipt.cleanup
482        );
483        fs::remove_dir_all(root).unwrap();
484    }
485}