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 kill = signal_group(pgid, "KILL");
317    let wait = child.wait().map(|_| ()).map_err(|e| format!("wait: {e}"));
318    let cleanup_deadline = Instant::now() + Duration::from_millis(500);
319    while group_exists(pgid) && Instant::now() < cleanup_deadline {
320        thread::sleep(Duration::from_millis(5));
321    }
322    let leaked = group_exists(pgid);
323    if leaked {
324        return Err(term
325            .err()
326            .or_else(|| kill.err())
327            .or_else(|| wait.err())
328            .unwrap_or_else(|| "descendants remained after bounded cleanup".into()));
329    }
330    wait.map(|()| "process group killed and reaped".into())
331}
332fn signal_group(pgid: u32, signal: &str) -> Result<(), String> {
333    let status = Command::new("kill")
334        .args([format!("-{signal}"), "--".into(), format!("-{pgid}")])
335        .stdout(Stdio::null())
336        .stderr(Stdio::null())
337        .status()
338        .map_err(|e| e.to_string())?;
339    if status.success() {
340        Ok(())
341    } else {
342        Err(format!("kill {signal} group {pgid}: {status}"))
343    }
344}
345fn group_exists(pgid: u32) -> bool {
346    let Ok(entries) = std::fs::read_dir("/proc") else {
347        return true;
348    };
349    entries
350        .flatten()
351        .filter_map(|entry| std::fs::read_to_string(entry.path().join("stat")).ok())
352        .any(|stat| {
353            let Some((_, fields)) = stat.rsplit_once(')') else {
354                return false;
355            };
356            let mut fields = fields.split_whitespace();
357            let state = fields.next();
358            let _ppid = fields.next();
359            let group = fields.next().and_then(|v| v.parse::<u32>().ok());
360            group == Some(pgid) && !matches!(state, Some("Z" | "X"))
361        })
362}
363
364#[cfg(test)]
365mod tests {
366    use super::*;
367    use sim_lib_exec::{ArgAtom, ProcessBudget, SealedBindings};
368    use std::{
369        fs,
370        path::PathBuf,
371        time::{SystemTime, UNIX_EPOCH},
372    };
373    fn root(label: &str) -> PathBuf {
374        let path = std::env::temp_dir().join(format!(
375            "sim-platform-process-{label}-{}",
376            SystemTime::now()
377                .duration_since(UNIX_EPOCH)
378                .unwrap()
379                .as_nanos()
380        ));
381        fs::create_dir_all(&path).unwrap();
382        path
383    }
384    fn fixture(root: &std::path::Path, program: &str) -> (UbuntuProcess, ProcessRequest) {
385        let program_ref = ProgramRef::new("shell").unwrap();
386        let root_ref = ProjectRootRef::new("project").unwrap();
387        let process = UbuntuProcess::new(
388            BTreeMap::from([(program_ref.clone(), PathBuf::from(program))]),
389            BTreeMap::from([(root_ref.clone(), root.to_owned())]),
390            BTreeMap::new(),
391        );
392        let request = ProcessRequest {
393            program: program_ref,
394            argv: vec![],
395            root: root_ref,
396            environment: SealedBindings::literals([("SIM_PROCESS_VISIBLE".into(), "yes".into())])
397                .unwrap(),
398            private_artifacts: vec![],
399            budget: ProcessBudget {
400                timeout_ms: 1_000,
401                max_output_bytes: 192,
402                stdin: None,
403            },
404        };
405        (process, request)
406    }
407    #[test]
408    fn native_conformance_clears_environment_confines_cwd_and_caps_output() {
409        let root = root("success");
410        let (process, mut request) = fixture(&root, "/bin/sh");
411        request.argv=["-c","printf '%s:%s:' \"$SIM_PROCESS_VISIBLE\" \"${HOME-unset}\"; pwd; printf 12345678901234567890123456789012345678901234567890123456789012345678901234567890123456789012345678901234567890123456789012345678901234567890123456789012345678901234567890123456789012345678901234567890"].into_iter().map(|v|ArgAtom::new(v).unwrap()).collect();
412        let ProcessAttempt::Completed { receipt } =
413            process.run(&request, &ProcessCancellation::default())
414        else {
415            panic!("expected completion")
416        };
417        assert!(receipt.result.stdout.starts_with("yes:unset:"));
418        assert!(receipt.result.stdout.contains(root.to_str().unwrap()));
419        assert!(receipt.result.truncated);
420        fs::remove_dir_all(root).unwrap();
421    }
422    #[test]
423    fn native_conformance_refuses_unknown_refs_and_cleans_timed_out_tree() {
424        let root = root("cleanup");
425        let (process, mut timed) = fixture(&root, "/bin/sh");
426        let mut unknown = timed.clone();
427        unknown.program = ProgramRef::new("unknown").unwrap();
428        assert!(matches!(
429            process.run(&unknown, &ProcessCancellation::default()),
430            ProcessAttempt::NotDispatched { .. }
431        ));
432        timed.argv = ["-c", "sleep 5 & wait"]
433            .into_iter()
434            .map(|v| ArgAtom::new(v).unwrap())
435            .collect();
436        timed.budget.timeout_ms = 20;
437        let outcome = process.run(&timed, &ProcessCancellation::default());
438        assert!(
439            matches!(outcome, ProcessAttempt::StoppedAfterTimeout { .. }),
440            "{outcome:?}"
441        );
442        fs::remove_dir_all(root).unwrap();
443    }
444    #[test]
445    fn native_conformance_cancellation_wins_and_cleans_tree() {
446        let root = root("cancel");
447        let (process, mut request) = fixture(&root, "/bin/sleep");
448        request.argv.push(ArgAtom::new("5").unwrap());
449        let token = ProcessCancellation::default();
450        token.cancel();
451        assert!(matches!(
452            process.run(&request, &token),
453            ProcessAttempt::StoppedAfterCancel { .. }
454        ));
455        fs::remove_dir_all(root).unwrap();
456    }
457}