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#[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 #[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}