返回 CodeWhale
shell.rs
根目录 / crates / tui / src / tools / shell.rs
1 //! Advanced shell execution with background process support and sandboxing.
2 //!
3 //! Provides:
4 //! - Synchronous command execution with timeout
5 //! - Background process execution
6 //! - Process output retrieval
7 //! - Process termination
8 //! - Sandbox support (macOS Seatbelt and opt-in Linux bubblewrap)
9 //! - Streaming output (future)
10
11 use anyhow::{Context, Result, anyhow};
12 use base64::Engine as _;
13 use serde::{Deserialize, Serialize};
14 use std::collections::HashMap;
15 use std::fs::File;
16 use std::io;
17 use std::io::{Read, Write};
18 use std::path::{Path, PathBuf};
19 use std::process::{Child, ChildStdin, Command, Stdio};
20 use std::sync::{Arc, Mutex};
21 use std::time::{Duration, Instant};
22 use uuid::Uuid;
23 use wait_timeout::ChildExt;
24
25 #[cfg(unix)]
26 use std::os::fd::FromRawFd;
27 #[cfg(unix)]
28 use std::os::unix::process::CommandExt;
29 #[cfg(windows)]
30 use std::os::windows::io::AsRawHandle;
31 #[cfg(windows)]
32 use std::os::windows::io::FromRawHandle;
33 #[cfg(windows)]
34 use windows::Win32::Foundation::{CloseHandle, HANDLE};
35 #[cfg(windows)]
36 use windows::Win32::System::JobObjects::{
37 AssignProcessToJobObject, CreateJobObjectW, JOB_OBJECT_LIMIT_KILL_ON_JOB_CLOSE,
38 JOBOBJECT_EXTENDED_LIMIT_INFORMATION, JobObjectExtendedLimitInformation,
39 SetInformationJobObject, TerminateJobObject,
40 };
41 #[cfg(windows)]
42 use windows::core::PCWSTR;
43
44 #[cfg(not(target_env = "ohos"))]
45 use portable_pty::{CommandBuilder, PtySize, native_pty_system};
46
47 mod guidance;
48 mod output;
49
50 use super::shell_output::{summarize_output, truncate_with_meta};
51 use crate::child_env;
52 use crate::sandbox::{
53 CommandSpec,
54 ExecEnv,
55 SandboxManager,
56 SandboxPolicy as ExecutionSandboxPolicy, // Rename to avoid conflict with spec::SandboxPolicy
57 SandboxType,
58 };
59 use crate::tools::resource_admission::{
60 CommandExpense, HeavyCommandPermit, MemoryPressure, acquire_heavy_command_permit,
61 infer_command_expense,
62 };
63 use crate::work_graph::{
64 EvidenceKind, EvidenceRef, OperationIntent, OperationOwnerSnapshot, OwnerState,
65 SharedWorkRuntime,
66 };
67 use crate::worker_profile::ShellPolicy;
68 use output::{
69 BoundedOutputAccumulator, BoundedOutputSnapshot, RAW_STREAM_SETTLED_TAIL_BYTES,
70 RawOutputBuffer, SharedRawOutput, new_shared_raw_output, tail_from_buffer, tail_text,
71 take_delta_from_buffer,
72 };
73
74 const READONLY_ENV_MARKER: &str = "CODEWHALE_INTERNAL_READONLY_ARGV";
75
76 #[cfg(unix)]
77 static PENDING_PERSISTENT_PROCESS_GROUPS: std::sync::OnceLock<
78 Mutex<std::collections::HashSet<u32>>,
79 > = std::sync::OnceLock::new();
80
81 #[cfg(unix)]
82 fn pending_persistent_process_groups() -> &'static Mutex<std::collections::HashSet<u32>> {
83 PENDING_PERSISTENT_PROCESS_GROUPS.get_or_init(|| Mutex::new(std::collections::HashSet::new()))
84 }
85
86 #[cfg(unix)]
87 fn register_pending_persistent_process_group(process_group_id: u32) {
88 let mut groups = pending_persistent_process_groups()
89 .lock()
90 .unwrap_or_else(std::sync::PoisonError::into_inner);
91 groups.insert(process_group_id);
92 }
93
94 #[cfg(unix)]
95 fn unregister_pending_persistent_process_group(process_group_id: u32) {
96 let mut groups = pending_persistent_process_groups()
97 .lock()
98 .unwrap_or_else(std::sync::PoisonError::into_inner);
99 groups.remove(&process_group_id);
100 }
101
102 /// Kill services that were staged for ownership transfer but have not yet
103 /// been released. The process-wide signal path calls this immediately before
104 /// `process::exit`, where Rust destructors cannot run.
105 #[cfg(unix)]
106 pub(crate) fn abort_pending_persistent_process_groups_for_exit() {
107 let groups = {
108 let mut groups = pending_persistent_process_groups()
109 .lock()
110 .unwrap_or_else(std::sync::PoisonError::into_inner);
111 groups.drain().collect::<Vec<_>>()
112 };
113 for process_group_id in groups {
114 if let Ok(process_group_id) = i32::try_from(process_group_id) {
115 // SAFETY: the id was captured from a child spawned with
116 // `process_group(0)`. A negative pid targets that child's process
117 // group, never Codewhale's own group.
118 unsafe {
119 libc::kill(-process_group_id, libc::SIGKILL);
120 }
121 }
122 }
123 }
124
125 fn validate_shell_working_dir(path: &Path, inherited_session_workspace: bool) -> Result<()> {
126 let metadata = std::fs::metadata(path).with_context(|| {
127 let source = if inherited_session_workspace {
128 "saved session workspace"
129 } else {
130 "requested working directory"
131 };
132 format!(
133 "{source} is unavailable: {}. Restore or remap that directory, resume/fork the session from an existing workspace, or pass an explicit `working_dir`/`cwd` to exec_shell",
134 path.display()
135 )
136 })?;
137 if !metadata.is_dir() {
138 let source = if inherited_session_workspace {
139 "saved session workspace"
140 } else {
141 "requested working directory"
142 };
143 return Err(anyhow!(
144 "{source} is not a directory: {}. Resume/fork from an existing workspace or pass an explicit `working_dir`/`cwd`",
145 path.display()
146 ));
147 }
148 Ok(())
149 }
150
151 /// Status of a shell process.
152 #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
153 pub enum ShellStatus {
154 Running,
155 Completed,
156 Failed,
157 Killed,
158 TimedOut,
159 }
160
161 /// Result from a shell command execution.
162 #[derive(Debug, Clone, Serialize, Deserialize)]
163 pub struct ShellResult {
164 pub task_id: Option<String>,
165 pub status: ShellStatus,
166 /// Lossless process exit status. Windows exception/NTSTATUS values use
167 /// the full unsigned 32-bit range, so an i32 would corrupt them.
168 pub exit_code: Option<i64>,
169 pub stdout: String,
170 pub stderr: String,
171 pub duration_ms: u64,
172 /// Original stdout length in bytes.
173 #[serde(default)]
174 pub stdout_len: usize,
175 /// Original stderr length in bytes.
176 #[serde(default)]
177 pub stderr_len: usize,
178 /// Bytes omitted from stdout due to truncation.
179 #[serde(default)]
180 pub stdout_omitted: usize,
181 /// Bytes omitted from stderr due to truncation.
182 #[serde(default)]
183 pub stderr_omitted: usize,
184 /// Whether stdout was truncated.
185 #[serde(default)]
186 pub stdout_truncated: bool,
187 /// Whether stderr was truncated.
188 #[serde(default)]
189 pub stderr_truncated: bool,
190 /// Whether the command was executed in a sandbox.
191 #[serde(default)]
192 pub sandboxed: bool,
193 /// Type of sandbox used (if any).
194 #[serde(skip_serializing_if = "Option::is_none")]
195 pub sandbox_type: Option<String>,
196 /// Whether the command was blocked by sandbox restrictions.
197 #[serde(default)]
198 pub sandbox_denied: bool,
199 }
200
201 /// Compact, UI-oriented view of a tracked background shell job.
202 #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
203 pub struct ShellJobSnapshot {
204 pub id: String,
205 pub job_id: String,
206 pub command: String,
207 pub cwd: PathBuf,
208 pub status: ShellStatus,
209 /// Explicitly backgrounded or detached from the foreground waiter.
210 #[serde(default)]
211 pub background: bool,
212 /// First observed terminal time; legacy snapshots cannot trigger notices.
213 #[serde(default, skip_serializing_if = "Option::is_none")]
214 pub finished_at: Option<chrono::DateTime<chrono::Utc>>,
215 pub exit_code: Option<i64>,
216 pub elapsed_ms: u64,
217 pub stdout_tail: String,
218 pub stderr_tail: String,
219 pub stdout_len: usize,
220 pub stderr_len: usize,
221 pub stdin_available: bool,
222 pub stale: bool,
223 #[serde(default, skip_serializing_if = "Option::is_none")]
224 pub elapsed_since_output_ms: Option<u64>,
225 pub linked_task_id: Option<String>,
226 #[serde(default, skip_serializing_if = "Option::is_none")]
227 pub owner_agent_id: Option<String>,
228 #[serde(default, skip_serializing_if = "Option::is_none")]
229 pub owner_agent_name: Option<String>,
230 #[serde(default, skip_serializing_if = "Option::is_none")]
231 pub origin_tool_call_id: Option<String>,
232 #[serde(default, skip_serializing_if = "Option::is_none")]
233 pub origin_turn_id: Option<String>,
234 /// Immutable root session that launched the job. Empty legacy records are
235 /// intentionally hidden from session-scoped completion drains.
236 #[serde(default, skip_serializing_if = "String::is_empty")]
237 pub owner_session_id: String,
238 }
239
240 /// Once-only completion event for a tracked background shell job.
241 #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
242 pub struct ShellCompletionEvent {
243 pub task_id: String,
244 pub command: String,
245 pub status: ShellStatus,
246 pub exit_code: Option<i64>,
247 pub duration_ms: u64,
248 pub stdout_tail: String,
249 pub stderr_tail: String,
250 #[serde(default)]
251 pub stdout_len: usize,
252 #[serde(default)]
253 pub stderr_len: usize,
254 #[serde(default, skip_serializing_if = "Option::is_none")]
255 pub evidence_ref: Option<String>,
256 pub linked_task_id: Option<String>,
257 #[serde(default, skip_serializing_if = "Option::is_none")]
258 pub owner_agent_id: Option<String>,
259 #[serde(default, skip_serializing_if = "Option::is_none")]
260 pub owner_agent_name: Option<String>,
261 #[serde(default, skip_serializing_if = "Option::is_none")]
262 pub origin_tool_call_id: Option<String>,
263 #[serde(default, skip_serializing_if = "Option::is_none")]
264 pub origin_turn_id: Option<String>,
265 #[serde(default, skip_serializing_if = "String::is_empty")]
266 pub owner_session_id: String,
267 }
268
269 /// Byte evidence captured alongside a bounded completion event. Exact unless
270 /// the stream exceeded the in-memory retention ceiling, in which case the
271 /// omission is declared per stream rather than presented as complete.
272 #[derive(Debug, Clone)]
273 pub(crate) struct ShellCompletionEvidence {
274 pub event: ShellCompletionEvent,
275 stdout: Vec<u8>,
276 stderr: Vec<u8>,
277 stdout_omitted: usize,
278 stderr_omitted: usize,
279 }
280
281 impl ShellCompletionEvidence {
282 /// Encode each stream losslessly. UTF-8 remains readable; arbitrary bytes
283 /// use base64 so `retrieve_tool_result` can still recover exact output.
284 pub(crate) fn artifact_bytes(&self) -> Vec<u8> {
285 fn stream(bytes: &[u8], omitted: usize) -> serde_json::Value {
286 let mut value = match std::str::from_utf8(bytes) {
287 Ok(content) => serde_json::json!({
288 "encoding": "utf-8",
289 "byte_length": bytes.len(),
290 "content": content,
291 }),
292 Err(_) => serde_json::json!({
293 "encoding": "base64",
294 "byte_length": bytes.len(),
295 "content": base64::engine::general_purpose::STANDARD.encode(bytes),
296 }),
297 };
298 // Additive and only present when something was actually dropped, so
299 // the common case stays byte-identical to the v1 artifact readers
300 // already parse.
301 if omitted > 0
302 && let Some(object) = value.as_object_mut()
303 {
304 object.insert("leading_bytes_omitted".into(), omitted.into());
305 object.insert(
306 "total_byte_length".into(),
307 bytes.len().saturating_add(omitted).into(),
308 );
309 }
310 value
311 }
312
313 serde_json::json!({
314 "schema": "codewhale.shell_completion.evidence.v1",
315 "task_id": self.event.task_id,
316 "command": self.event.command,
317 "status": format!("{:?}", self.event.status),
318 "exit_code": self.event.exit_code,
319 "duration_ms": self.event.duration_ms,
320 "origin_tool_call_id": self.event.origin_tool_call_id,
321 "origin_turn_id": self.event.origin_turn_id,
322 "stdout": stream(&self.stdout, self.stdout_omitted),
323 "stderr": stream(&self.stderr, self.stderr_omitted),
324 })
325 .to_string()
326 .into_bytes()
327 }
328 }
329
330 // Keep the two inline streams at a 2 KiB combined hard ceiling. The durable
331 // artifact carries the exact bytes beyond these diagnostic tails.
332 const SHELL_COMPLETION_TAIL_BYTES: usize = 1_024;
333
334 /// How long a finished shell record stays listed in `/jobs`.
335 const FINISHED_SHELL_MAX_AGE: Duration = Duration::from_secs(3600);
336 /// Ceiling on finished records kept for the jobs panel. A long automation run
337 /// makes hundreds of `Bash` calls per hour; the panel is only useful for the
338 /// recent ones (#5472).
339 const MAX_FINISHED_SHELL_RECORDS: usize = 128;
340 /// Ceiling on bytes still held across all finished records. Each settled record
341 /// releases down to a 64 KiB tail, so this only binds when many large outputs
342 /// finish inside the same window.
343 const MAX_FINISHED_SHELL_BYTES: usize = 8 * 1024 * 1024;
344
345 fn bounded_completion_tail(buffer: &SharedRawOutput, max_bytes: usize) -> (usize, String) {
346 let (total, candidate) = tail_from_buffer(buffer, max_bytes);
347 if candidate.len() <= max_bytes {
348 return (total, candidate);
349 }
350 let content_budget = max_bytes.saturating_sub(3);
351 let mut start = candidate.len().saturating_sub(content_budget);
352 while start < candidate.len() && !candidate.is_char_boundary(start) {
353 start += 1;
354 }
355 (total, format!("...{}", &candidate[start..]))
356 }
357
358 /// Optional owner attribution for background shell work.
359 #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
360 pub struct ShellJobOwner {
361 pub agent_id: String,
362 pub agent_name: String,
363 }
364
365 /// Full output view used by `/jobs show <id>`.
366 #[derive(Debug, Clone, Serialize, Deserialize)]
367 pub struct ShellJobDetail {
368 pub snapshot: ShellJobSnapshot,
369 pub stdout: String,
370 pub stderr: String,
371 }
372
373 pub struct ShellDeltaResult {
374 pub command: String,
375 pub result: ShellResult,
376 pub stdout_total_len: usize,
377 pub stderr_total_len: usize,
378 }
379
380 /// Which of a job's raw output streams to read. Stderr is a separate stream
381 /// only for piped jobs; PTY and merged modes fold it into stdout.
382 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
383 pub enum ShellOutputStream {
384 Stdout,
385 Stderr,
386 }
387
388 /// A non-consuming window of a job's raw output stream at absolute byte
389 /// offsets. Unlike [`ShellManager::get_output_delta`], reading a chunk never
390 /// advances anyone else's cursor, so several HTTP clients can follow the same
391 /// job without splitting the stream.
392 pub struct ShellOutputChunk {
393 /// Absolute offset of `bytes[0]`. Exceeds the requested cursor when the
394 /// bounded buffer already discarded that prefix — the gap is reported via
395 /// `dropped`, never silently re-sent.
396 pub offset: usize,
397 /// Raw stream bytes. Output is arbitrary bytes, not guaranteed UTF-8.
398 pub bytes: Vec<u8>,
399 /// Absolute offset just past the last returned byte; the next cursor.
400 pub next_offset: usize,
401 /// Total bytes this stream has produced, including discarded bytes.
402 pub total: usize,
403 /// Leading bytes permanently discarded by the in-flight bound.
404 pub dropped: usize,
405 pub status: ShellStatus,
406 pub exit_code: Option<i64>,
407 }
408
409 enum ShellChild {
410 Process(Child),
411 #[cfg(not(target_env = "ohos"))]
412 Pty(Box<dyn portable_pty::Child + Send>),
413 }
414 #[cfg(unix)]
415 impl ShellChild {
416 fn process_id(&self) -> Option<u32> {
417 match self {
418 Self::Process(child) => Some(child.id()),
419 #[cfg(not(target_env = "ohos"))]
420 Self::Pty(_) => None,
421 }
422 }
423 }
424
425 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
426 enum ShellOwnership {
427 Managed,
428 PersistPending,
429 Released,
430 }
431
432 #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
433 pub struct PersistentServiceReceipt {
434 pub task_id: String,
435 pub pid: u32,
436 pub process_group_id: u32,
437 pub ownership: String,
438 }
439
440 #[cfg(unix)]
441 fn signal_child_process_group(child: &Child, signal: libc::c_int) -> std::io::Result<()> {
442 let pgid = child.id() as libc::pid_t;
443 if pgid <= 0 {
444 return Ok(());
445 }
446
447 // SAFETY: kill(2) dereferences no pointers.
448 let result = unsafe { libc::kill(-pgid, signal) };
449 if result == 0 {
450 Ok(())
451 } else {
452 let err = std::io::Error::last_os_error();
453 if err.raw_os_error() == Some(libc::ESRCH) {
454 // The group is already gone (or never formed); nothing to signal.
455 Ok(())
456 } else {
457 Err(err)
458 }
459 }
460 }
461
462 #[cfg(unix)]
463 fn kill_child_process_group(child: &mut Child) -> std::io::Result<()> {
464 let pgid = child.id() as libc::pid_t;
465 if pgid <= 0 {
466 return child.kill();
467 }
468
469 signal_child_process_group(child, libc::SIGKILL).or_else(|_| child.kill())
470 }
471
472 /// Bounded wait for the direct child to exit. Returns true once the child was
473 /// reaped (or the wait errored), false when the grace elapsed first. Unlike
474 /// `Child::wait`, this can never wedge the caller behind a child stuck in
475 /// uninterruptible sleep.
476 #[cfg(unix)]
477 fn wait_child_bounded(child: &mut Child, grace: Duration) -> bool {
478 let deadline = Instant::now() + grace;
479 loop {
480 match child.try_wait() {
481 Ok(Some(_)) | Err(_) => return true,
482 Ok(None) => {}
483 }
484 if Instant::now() >= deadline {
485 return false;
486 }
487 std::thread::sleep(Duration::from_millis(10));
488 }
489 }
490
491 /// Terminate a shell's whole process group with a bounded SIGTERM → SIGKILL
492 /// escalation (#52). The previous kill path SIGKILLed only the direct child
493 /// and then joined output-reader threads with no timeout, so the tool
494 /// returned whenever the command's descendants felt like exiting — observed
495 /// as a 120s foreground timeout returning after 300s. Every step here is
496 /// bounded: the tool returns at ~timeout + grace.
497 #[cfg(unix)]
498 fn terminate_child_process_group(child: &mut Child) -> std::io::Result<()> {
499 // Cooperative stop first so shells and their children can run traps and
500 // clean up; bounded so a SIGTERM-ignoring command cannot stall the caller.
501 let _ = signal_child_process_group(child, libc::SIGTERM);
502 if wait_child_bounded(child, KILL_TERM_GRACE) {
503 // The leader exited on SIGTERM; descendants may linger, so SIGKILL
504 // the rest of the group (ESRCH when it is already empty).
505 kill_child_process_group(child)?;
506 return Ok(());
507 }
508 kill_child_process_group(child)?;
509 let _ = wait_child_bounded(child, KILL_REAP_GRACE);
510 Ok(())
511 }
512
513 /// Configure parent-death signaling so shell-spawned children are reaped when
514 /// the TUI dies abnormally (#421). On Linux this installs
515 /// `PR_SET_PDEATHSIG(SIGTERM)` via `pre_exec` — the kernel then sends SIGTERM
516 /// to the child the moment the parent process exits, even on SIGKILL of the
517 /// TUI. The cancellation path already SIGKILLs the whole process group, so
518 /// this only fires when the parent dies without running its drop / cleanup
519 /// code (panic during shutdown, OOM, hardware crash, etc.).
520 ///
521 /// On macOS / Windows there's no kernel equivalent. The existing graceful
522 /// path (`kill_child_process_group` from the cancellation token) still
523 /// handles normal shutdown; abnormal exit can leak children — tracked as a
524 /// follow-up watchdog item per the original issue's acceptance criteria.
525 #[cfg(all(target_os = "linux", not(target_env = "ohos")))]
526 fn install_parent_death_signal(cmd: &mut Command) {
527 use std::os::unix::process::CommandExt;
528 // Captured before the fork so the child can tell whether the TUI already
529 // died in the fork→prctl window, where the signal would never arrive.
530 let parent_pid = std::process::id() as libc::pid_t;
531 // SAFETY: `pre_exec` runs in the child between fork and exec. The closure
532 // only calls `libc::prctl` / `libc::getppid` with stack-allocated
533 // arguments and does not touch heap memory or the parent's locks. Both
534 // requirements (async-signal-safe + no allocation in the post-fork
535 // window) are met.
536 unsafe {
537 cmd.pre_exec(move || {
538 let result = libc::prctl(libc::PR_SET_PDEATHSIG, libc::SIGTERM, 0, 0, 0);
539 if result == -1 {
540 // Surface the errno but do not abort the spawn — the child
541 // will simply lose the parent-death cleanup safety net.
542 return Err(std::io::Error::last_os_error());
543 }
544 if libc::getppid() != parent_pid {
545 // The TUI exited before the signal was armed: do not exec an
546 // orphan nobody will ever reap.
547 return Err(std::io::Error::from_raw_os_error(libc::ESRCH));
548 }
549 Ok(())
550 });
551 }
552 }
553
554 /// Watcher script for [`watch_process_group_for_parent_death`]. `$1` is the
555 /// process group id. The outer `sh` only forks the watcher and exits, so the
556 /// watcher is reparented away from the TUI and never needs reaping. The
557 /// watcher blocks reading the pipe the TUI holds the other end of: the read
558 /// ends at EOF when every write end closes, which happens however the TUI
559 /// dies. It then SIGTERMs the group, waits briefly and SIGKILLs it. It sits in
560 /// the job's own process group, so the group id cannot be reused while it
561 /// waits, and the TUI's normal group kills take it down with the job.
562 #[cfg(unix)]
563 const PROCESS_GROUP_WATCHER_SCRIPT: &str = r#"exec 3<&0
564 (
565 trap '' TERM
566 read -r _ <&3 || {
567 kill -TERM -"$1" 2>/dev/null
568 sleep 2
569 kill -KILL -"$1" 2>/dev/null
570 }
571 ) >/dev/null 2>&1 &"#;
572
573 /// Kill a `Managed` background shell's whole process group when the TUI dies
574 /// (#6654).
575 ///
576 /// `PR_SET_PDEATHSIG` is not enough here: it signals only the direct child
577 /// (the `sh -c` wrapper), a fork clears it in grandchildren, so a command that
578 /// does not `exec` its last step (`cd app && npm run dev; echo done`,
579 /// `a | b`) keeps its real workload alive; it also fires when the forking
580 /// *thread* exits, and background spawns run on threads that retire. So a
581 /// small watcher joins the job's process group and blocks on a pipe whose
582 /// write end only the TUI holds (it is close-on-exec, so no child inherits
583 /// it). The returned write end lives in the `BackgroundShell`; closing it —
584 /// by dropping the shell or by the TUI dying in any way, SIGKILL included —
585 /// makes the watcher terminate the group. This works on Linux and macOS alike.
586 ///
587 /// Best effort: if the watcher cannot start, the job still runs, just
588 /// without this cleanup. A TUI that dies between the job spawn and the
589 /// watcher spawn also leaves the job running.
590 #[cfg(unix)]
591 fn watch_process_group_for_parent_death(
592 process_group_id: u32,
593 ) -> std::io::Result<std::io::PipeWriter> {
594 let (reader, writer) = std::io::pipe()?;
595 let pgid = i32::try_from(process_group_id)
596 .map_err(|_| std::io::Error::from_raw_os_error(libc::EINVAL))?;
597 let mut cmd = Command::new("/bin/sh");
598 cmd.args([
599 "-c",
600 PROCESS_GROUP_WATCHER_SCRIPT,
601 "codewhale-bg-watch",
602 &process_group_id.to_string(),
603 ])
604 .stdin(Stdio::from(reader))
605 .stdout(Stdio::null())
606 .stderr(Stdio::null())
607 .process_group(pgid);
608 // The outer `sh` exits as soon as it has forked the watcher.
609 let status = cmd.spawn()?.wait()?;
610 if status.success() {
611 Ok(writer)
612 } else {
613 Err(std::io::Error::other(format!(
614 "process-group watcher exited with {status}"
615 )))
616 }
617 }
618
619 /// Attach `args` to a `std::process::Command`, honoring shell-quoting on
620 /// Windows.
621 ///
622 /// Issue #1691: on Windows the shell command is invoked as
623 /// `cmd /C "chcp 65001 >NUL & <command>"`. Rust's `Command::arg` applies
624 /// MSVCRT (`CommandLineToArgvW`) escaping, turning the embedded `"` in a
625 /// quoted argument (e.g. `git commit -m "feat: complete sub-pages"`) into
626 /// `\"`. `cmd.exe` does NOT use MSVCRT parsing — it treats `\` literally and
627 /// `"` as a bare quote toggle — so the escaped payload is mis-tokenized and
628 /// `git` receives `feat:`, `complete`, `sub-pages"` as separate pathspecs
629 /// (the reported `pathspec 'sub-pages"' did not match` symptom). Passing the
630 /// `cmd /C` payload through `CommandExt::raw_arg` suppresses std's escaping so
631 /// the string reaches `cmd.exe` verbatim, exactly as a terminal would.
632 #[cfg(windows)]
633 fn push_shell_args(cmd: &mut Command, program: &str, args: &[String]) {
634 use std::os::windows::process::CommandExt;
635 // The `cmd /C <payload>` shape is the only place std's per-arg escaping
636 // corrupts a quoted command. Pass `/C` and the payload raw so the quotes
637 // survive; any other program keeps normal (correct) escaping. Match `cmd`
638 // by file stem so a full path (`C:\Windows\System32\cmd.exe`) or `.exe`
639 // suffix still triggers the raw-arg path.
640 let is_cmd = std::path::Path::new(program)
641 .file_stem()
642 .and_then(|s| s.to_str())
643 .map(|s| s.eq_ignore_ascii_case("cmd"))
644 .unwrap_or(false);
645 if is_cmd && args.len() == 2 && args[0].eq_ignore_ascii_case("/C") {
646 cmd.raw_arg(&args[0]);
647 cmd.raw_arg(&args[1]);
648 } else {
649 cmd.args(args);
650 }
651 }
652
653 #[cfg(not(windows))]
654 fn push_shell_args(cmd: &mut Command, _program: &str, args: &[String]) {
655 // Unix delegates tokenization entirely to `sh -c <command>`; the command
656 // string is passed as a single argv entry and never split by us.
657 cmd.args(args);
658 }
659
660 #[cfg(not(all(target_os = "linux", not(target_env = "ohos"))))]
661 fn install_parent_death_signal(_cmd: &mut Command) {
662 // No kernel-level equivalent on macOS / Windows. The cooperative
663 // cancellation + process_group SIGKILL path covers normal shutdown;
664 // abnormal exit (panic without unwind, SIGKILL of the TUI) can still
665 // leak children on those platforms — tracked as a follow-up. Pipe-backed
666 // background shells are covered on every platform: Unix by the
667 // process-group watcher (`watch_process_group_for_parent_death`), Windows
668 // by their KILL_ON_JOB_CLOSE job object.
669 //
670 // Known limitations on every platform (#6654): `tty: true` background
671 // shells spawn through `portable_pty`, which never goes through
672 // `std::process::Command`, so they get no watcher (or job object); and
673 // staged persistent services (`persist_pending`) stay deliberately
674 // unwatched because the user can take ownership of the service.
675 }
676
677 #[cfg(windows)]
678 #[derive(Debug)]
679 struct WindowsJob {
680 handle: HANDLE,
681 }
682
683 #[cfg(windows)]
684 // SAFETY: Windows job handles are process-wide kernel handles. Moving the
685 // wrapper between threads does not invalidate the handle, and access is
686 // externally synchronized by ShellManager's mutex.
687 unsafe impl Send for WindowsJob {}
688 #[cfg(windows)]
689 // SAFETY: The wrapper exposes only terminate/drop operations around a kernel
690 // handle; concurrent use is guarded by ShellManager.
691 unsafe impl Sync for WindowsJob {}
692
693 #[cfg(windows)]
694 impl WindowsJob {
695 fn attach_to_child(child: &Child) -> std::io::Result<Self> {
696 // SAFETY: returned handle is owned by the new wrapper.
697 let handle = unsafe { CreateJobObjectW(None, PCWSTR::null()).map_err(windows_io_error)? };
698 let job = Self { handle };
699
700 let mut limits = JOBOBJECT_EXTENDED_LIMIT_INFORMATION::default();
701 limits.BasicLimitInformation.LimitFlags = JOB_OBJECT_LIMIT_KILL_ON_JOB_CLOSE;
702
703 // SAFETY: `limits` is live with matching size; both handles are live.
704 unsafe {
705 SetInformationJobObject(
706 job.handle,
707 JobObjectExtendedLimitInformation,
708 &limits as *const _ as *const core::ffi::c_void,
709 std::mem::size_of::<JOBOBJECT_EXTENDED_LIMIT_INFORMATION>() as u32,
710 )
711 .map_err(windows_io_error)?;
712
713 let process_handle = HANDLE(child.as_raw_handle());
714 AssignProcessToJobObject(job.handle, process_handle).map_err(windows_io_error)?;
715 }
716
717 Ok(job)
718 }
719
720 fn terminate(&self) -> std::io::Result<()> {
721 // SAFETY: `self.handle` is a live owned job handle.
722 unsafe { TerminateJobObject(self.handle, 1).map_err(windows_io_error) }
723 }
724 }
725
726 #[cfg(windows)]
727 impl Drop for WindowsJob {
728 fn drop(&mut self) {
729 // SAFETY: `self.handle` is owned here; Drop runs once.
730 unsafe {
731 let _ = CloseHandle(self.handle);
732 }
733 }
734 }
735
736 #[cfg(windows)]
737 fn windows_io_error(error: windows::core::Error) -> std::io::Error {
738 std::io::Error::other(error)
739 }
740
741 #[cfg(windows)]
742 fn terminate_windows_job(job: Option<&WindowsJob>, child: &mut Child) -> std::io::Result<()> {
743 if let Some(job) = job {
744 match job.terminate() {
745 Ok(()) => return Ok(()),
746 Err(error) => {
747 tracing::warn!(
748 ?error,
749 "failed to terminate Windows job object; falling back to immediate child kill"
750 );
751 }
752 }
753 }
754 child.kill()
755 }
756
757 #[cfg(windows)]
758 fn terminate_and_close_windows_job(windows_job: Option<WindowsJob>) {
759 if let Some(job) = windows_job.as_ref()
760 && let Err(err) = job.terminate()
761 {
762 tracing::warn!(
763 ?err,
764 "failed to terminate Windows shell job before closing job handle"
765 );
766 }
767 drop(windows_job);
768 }
769
770 #[cfg(windows)]
771 fn terminate_child_and_close_windows_job(
772 windows_job: Option<WindowsJob>,
773 child: &mut Child,
774 ) -> std::io::Result<()> {
775 let result = terminate_windows_job(windows_job.as_ref(), child);
776 drop(windows_job);
777 result
778 }
779
780 #[cfg(windows)]
781 fn attach_windows_job(child: &Child, command: &str) -> Option<WindowsJob> {
782 match WindowsJob::attach_to_child(child) {
783 Ok(job) => Some(job),
784 Err(error) => {
785 tracing::warn!(
786 ?error,
787 command,
788 "failed to attach Windows shell process to job object; descendant cleanup degraded"
789 );
790 None
791 }
792 }
793 }
794
795 #[cfg(windows)]
796 fn terminate_unregistered_process(child: &mut Child, job: Option<&WindowsJob>) {
797 let _ = terminate_windows_job(job, child);
798 let _ = child.wait();
799 }
800
801 #[cfg(not(windows))]
802 fn terminate_unregistered_process(child: &mut Child) {
803 #[cfg(unix)]
804 {
805 let _ = kill_child_process_group(child);
806 let _ = wait_child_bounded(child, KILL_REAP_GRACE);
807 }
808 #[cfg(not(unix))]
809 {
810 let _ = child.kill();
811 let _ = child.wait();
812 }
813 }
814
815 #[derive(Clone, Copy, Debug)]
816 struct ShellExitStatus {
817 code: Option<i64>,
818 success: bool,
819 }
820
821 impl ShellExitStatus {
822 fn from_std(status: std::process::ExitStatus) -> Self {
823 Self {
824 code: status.code().map(std_exit_code_i64),
825 success: status.success(),
826 }
827 }
828
829 #[cfg(not(target_env = "ohos"))]
830 fn from_pty(status: portable_pty::ExitStatus) -> Self {
831 Self {
832 code: Some(i64::from(status.exit_code())),
833 success: status.success(),
834 }
835 }
836 }
837
838 #[cfg(windows)]
839 fn std_exit_code_i64(code: i32) -> i64 {
840 // std exposes Windows DWORD process statuses through i32. Reinterpret
841 // negative values as their original unsigned bit pattern so codes such
842 // as 0xC0000005 survive JSON, persistence, and diagnostics unchanged.
843 i64::from(code as u32)
844 }
845
846 #[cfg(not(windows))]
847 fn std_exit_code_i64(code: i32) -> i64 {
848 i64::from(code)
849 }
850
851 impl ShellChild {
852 fn try_wait(&mut self) -> std::io::Result<Option<ShellExitStatus>> {
853 match self {
854 ShellChild::Process(child) => child
855 .try_wait()
856 .map(|status| status.map(ShellExitStatus::from_std)),
857 #[cfg(not(target_env = "ohos"))]
858 ShellChild::Pty(child) => child
859 .try_wait()
860 .map(|status| status.map(ShellExitStatus::from_pty)),
861 }
862 }
863
864 #[cfg(not(windows))]
865 fn kill(&mut self) -> std::io::Result<()> {
866 match self {
867 #[cfg(unix)]
868 ShellChild::Process(child) => kill_child_process_group(child),
869 #[cfg(not(unix))]
870 ShellChild::Process(child) => child.kill(),
871 #[cfg(not(target_env = "ohos"))]
872 ShellChild::Pty(child) => child.kill(),
873 }
874 }
875 }
876
877 enum StdinWriter {
878 Pipe(ChildStdin),
879 #[cfg(not(target_env = "ohos"))]
880 Pty(Box<dyn Write + Send>),
881 }
882
883 impl StdinWriter {
884 fn write_all(&mut self, data: &[u8]) -> std::io::Result<()> {
885 match self {
886 StdinWriter::Pipe(stdin) => stdin.write_all(data),
887 #[cfg(not(target_env = "ohos"))]
888 StdinWriter::Pty(writer) => writer.write_all(data),
889 }
890 }
891
892 fn flush(&mut self) -> std::io::Result<()> {
893 match self {
894 StdinWriter::Pipe(stdin) => stdin.flush(),
895 #[cfg(not(target_env = "ohos"))]
896 StdinWriter::Pty(writer) => writer.flush(),
897 }
898 }
899 }
900
901 fn spawn_reader_thread<R: Read + Send + 'static>(
902 mut reader: R,
903 buffer: SharedRawOutput,
904 ) -> std::thread::JoinHandle<()> {
905 std::thread::spawn(move || {
906 let mut chunk = [0u8; 4096];
907 loop {
908 match reader.read(&mut chunk) {
909 Ok(0) => break,
910 Ok(n) => {
911 // `RawOutputBuffer::append` enforces the in-flight ceiling
912 // here, at the only writer, so a chatty command cannot grow
913 // the process without bound while it runs (#5472). It
914 // returns false once the stream has been abandoned, which
915 // is this thread's only exit when a descendant holds the
916 // pipe open and EOF never arrives.
917 let keep_reading = buffer
918 .lock()
919 .unwrap_or_else(|e| e.into_inner())
920 .append(&chunk[..n]);
921 if !keep_reading {
922 break;
923 }
924 }
925 Err(_) => break,
926 }
927 }
928 })
929 }
930
931 fn spawn_bounded_reader_thread<R: Read + Send + 'static>(
932 mut reader: R,
933 output: Arc<Mutex<BoundedOutputAccumulator>>,
934 ) -> std::thread::JoinHandle<()> {
935 std::thread::spawn(move || {
936 let mut chunk = [0u8; 4096];
937 loop {
938 match reader.read(&mut chunk) {
939 Ok(0) => break,
940 Ok(n) => {
941 let mut guard = output.lock().unwrap_or_else(|error| error.into_inner());
942 if let Err(error) = guard.append(&chunk[..n]) {
943 guard.record_error(&error);
944 return;
945 }
946 }
947 Err(error) => {
948 output
949 .lock()
950 .unwrap_or_else(|poison| poison.into_inner())
951 .record_error(&error);
952 return;
953 }
954 }
955 }
956 let mut guard = output.lock().unwrap_or_else(|error| error.into_inner());
957 if let Err(error) = guard.finish() {
958 guard.record_error(&error);
959 }
960 })
961 }
962
963 #[cfg(unix)]
964 fn shared_output_pipe() -> io::Result<(File, File, File)> {
965 let mut descriptors = [0; 2];
966 // SAFETY: `pipe` initializes both descriptors on success. Each descriptor
967 // is immediately transferred into exactly one owned `File`.
968 if unsafe { libc::pipe(descriptors.as_mut_ptr()) } != 0 {
969 return Err(io::Error::last_os_error());
970 }
971 // SAFETY: successful `pipe` returned two live, uniquely owned descriptors.
972 let reader = unsafe { File::from_raw_fd(descriptors[0]) };
973 let writer = unsafe { File::from_raw_fd(descriptors[1]) };
974 let stderr_writer = writer.try_clone()?;
975 Ok((reader, writer, stderr_writer))
976 }
977
978 #[cfg(windows)]
979 fn shared_output_pipe() -> io::Result<(File, File, File)> {
980 let mut read_handle = std::ptr::null_mut();
981 let mut write_handle = std::ptr::null_mut();
982 // SAFETY: CreatePipe initializes both handles on success; ownership is
983 // transferred to `File` immediately below.
984 if unsafe {
985 windows_sys::Win32::System::Pipes::CreatePipe(
986 &mut read_handle,
987 &mut write_handle,
988 std::ptr::null(),
989 0,
990 )
991 } == 0
992 {
993 return Err(io::Error::last_os_error());
994 }
995 // SAFETY: successful CreatePipe returned two live, uniquely owned handles.
996 let reader = unsafe { File::from_raw_handle(read_handle.cast()) };
997 let writer = unsafe { File::from_raw_handle(write_handle.cast()) };
998 let stderr_writer = writer.try_clone()?;
999 Ok((reader, writer, stderr_writer))
1000 }
1001
1002 const SYNC_READER_DRAIN_TIMEOUT: Duration = Duration::from_secs(5);
1003 const STALE_NO_OUTPUT_AFTER: Duration = Duration::from_secs(60);
1004
1005 /// Grace between SIGTERM and SIGKILL on the shell kill path (timeout,
1006 /// cancel, drop). Bounded so a SIGTERM-ignoring command is force-killed
1007 /// instead of stalling the tool (#52).
1008 #[cfg(unix)]
1009 const KILL_TERM_GRACE: Duration = Duration::from_millis(500);
1010 /// Bounded reap wait after SIGKILL; a child stuck in uninterruptible sleep
1011 /// must not wedge the caller behind an unbounded `wait`.
1012 #[cfg(unix)]
1013 const KILL_REAP_GRACE: Duration = Duration::from_millis(1_000);
1014 /// Bounded join for output-reader threads after the process group is killed.
1015 /// A descendant that escaped the group (its own session/process group) keeps
1016 /// its inherited pipe write-end open, so the reader cannot see EOF until that
1017 /// descendant exits on its own — an unbounded join held the shell-manager
1018 /// lock for minutes and overshot the tool timeout (#52).
1019 const READER_JOIN_GRACE: Duration = Duration::from_millis(2_000);
1020
1021 fn spawn_sync_reader_thread<R: Read + Send + 'static>(
1022 mut reader: R,
1023 ) -> std::sync::mpsc::Receiver<Vec<u8>> {
1024 let (tx, rx) = std::sync::mpsc::channel();
1025 std::thread::spawn(move || {
1026 // Bounded, unlike the `read_to_end` this replaces (#5472 finding 2).
1027 // `recv_sync_reader_output` gives up after 5 s, but the thread lives as
1028 // long as the pipe does — an interactive command that keeps printing
1029 // grew this Vec without limit, for a result nobody was still waiting
1030 // for. The tail is what the caller renders, so keep the tail.
1031 let mut buf = RawOutputBuffer::new();
1032 let mut chunk = [0u8; 4096];
1033 loop {
1034 match reader.read(&mut chunk) {
1035 Ok(0) => break,
1036 Ok(n) => {
1037 if !buf.append(&chunk[..n]) {
1038 break;
1039 }
1040 }
1041 Err(_) => break,
1042 }
1043 }
1044 tx.send(buf.retained().to_vec()).ok();
1045 });
1046 rx
1047 }
1048
1049 fn recv_sync_reader_output(rx: &std::sync::mpsc::Receiver<Vec<u8>>) -> Vec<u8> {
1050 rx.recv_timeout(SYNC_READER_DRAIN_TIMEOUT)
1051 .unwrap_or_default()
1052 }
1053
1054 /// Cell dimensions accepted by the existing PTY owner. Pixel sizes remain
1055 /// unspecified; callers must not allocate an unbounded terminal grid.
1056 #[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
1057 #[serde(deny_unknown_fields)]
1058 pub struct PtyDimensions {
1059 pub rows: u16,
1060 pub cols: u16,
1061 }
1062
1063 impl Default for PtyDimensions {
1064 fn default() -> Self {
1065 Self { rows: 24, cols: 80 }
1066 }
1067 }
1068
1069 impl PtyDimensions {
1070 pub fn validate(self) -> Result<Self> {
1071 anyhow::ensure!(
1072 (1..=1000).contains(&self.rows) && (1..=1000).contains(&self.cols),
1073 "PTY rows and columns must each be between 1 and 1000"
1074 );
1075 Ok(self)
1076 }
1077 }
1078
1079 /// A background shell process being tracked
1080 pub struct BackgroundShell {
1081 pub id: String,
1082 pub command: String,
1083 pub working_dir: PathBuf,
1084 pub status: ShellStatus,
1085 background: bool,
1086 pub exit_code: Option<i64>,
1087 pub started_at: Instant,
1088 /// When the job reached a terminal status. A finished job reports the
1089 /// duration it finished with; without this, `started_at.elapsed()` kept
1090 /// growing and `/jobs` showed "2m 07s" for a 12-second command (#5478).
1091 finished_at: Option<Instant>,
1092 finished_at_utc: Option<chrono::DateTime<chrono::Utc>>,
1093 last_output_at: Instant,
1094 last_observed_output_len: usize,
1095 pub sandbox_type: SandboxType,
1096 pub linked_task_id: Option<String>,
1097 pub owner_agent: Option<ShellJobOwner>,
1098 owner_session_id: String,
1099 origin_tool_call_id: Option<String>,
1100 origin_turn_id: Option<String>,
1101 ownership: ShellOwnership,
1102 stdout_buffer: SharedRawOutput,
1103 stderr_buffer: Option<SharedRawOutput>,
1104 /// Lowercase `bash` streams one combined process pipe through a bounded
1105 /// small-contract-compatible accumulator while persisting the complete output.
1106 bounded_output: Option<Arc<Mutex<BoundedOutputAccumulator>>>,
1107 heavy_permit: Option<HeavyCommandPermit>,
1108 stdout_cursor: usize,
1109 stderr_cursor: usize,
1110 completion_reported: bool,
1111 stdin: Option<StdinWriter>,
1112 /// Retain the existing PTY owner for resize; never create another session.
1113 #[cfg(not(target_env = "ohos"))]
1114 pty_master: Option<Box<dyn portable_pty::MasterPty + Send>>,
1115 terminal_size: Option<PtyDimensions>,
1116 child: Option<ShellChild>,
1117 #[cfg(windows)]
1118 windows_job: Option<WindowsJob>,
1119 /// Write end of the process-group watcher's pipe for `Managed` shells;
1120 /// closing it (drop, or the TUI dying) kills the job's group (#6654).
1121 #[cfg(unix)]
1122 #[allow(dead_code, reason = "held only so its Drop closes the pipe")]
1123 parent_death_watch: Option<std::io::PipeWriter>,
1124 stdout_thread: Option<std::thread::JoinHandle<()>>,
1125 stderr_thread: Option<std::thread::JoinHandle<()>>,
1126 work_lifecycle: Option<ShellWorkLifecycle>,
1127 lifecycle_seq: u64,
1128 last_lifecycle_status: Option<ShellStatus>,
1129 last_lifecycle_bytes: usize,
1130 /// The terminal status came from a failed `try_wait`, not an observed
1131 /// exit or signal. Such a run gets no execution receipt (#6689).
1132 wait_failed: bool,
1133 }
1134
1135 #[derive(Clone)]
1136 struct ShellWorkLifecycle {
1137 work: SharedWorkRuntime,
1138 session_id: String,
1139 }
1140
1141 impl ShellWorkLifecycle {
1142 fn register(&self, id: &str, command: &str) -> Result<()> {
1143 self.work
1144 .register_operation(
1145 &self.session_id,
1146 OperationIntent::new(
1147 format!("shell:{id}"),
1148 format!("Shell · {command}"),
1149 false,
1150 "exec_shell",
1151 id,
1152 ),
1153 )
1154 .map(|_| ())
1155 .map_err(anyhow::Error::msg)
1156 }
1157
1158 fn observe(&self, id: &str, status: &ShellStatus, seq: u64, raw_bytes: usize) -> Result<()> {
1159 let owner_state = match status {
1160 ShellStatus::Running => OwnerState::Running,
1161 ShellStatus::Completed => OwnerState::Completed,
1162 ShellStatus::Failed | ShellStatus::TimedOut => OwnerState::Failed,
1163 ShellStatus::Killed => OwnerState::Cancelled,
1164 };
1165 let raw_bytes = u64::try_from(raw_bytes).unwrap_or(u64::MAX);
1166 let output = EvidenceRef::new(
1167 EvidenceKind::Receipt {
1168 owner: "shell".to_string(),
1169 },
1170 format!("shell:{id}:output"),
1171 Some(raw_bytes),
1172 false,
1173 )
1174 .map_err(|err| anyhow!(err.to_string()))?;
1175 self.work
1176 .reconcile_operation(
1177 &self.session_id,
1178 OperationOwnerSnapshot::new(
1179 format!("shell:{id}"),
1180 owner_state,
1181 seq,
1182 lifecycle_now_ms(),
1183 )
1184 .with_output(output),
1185 )
1186 .map(|_| ())
1187 .map_err(anyhow::Error::msg)
1188 }
1189 }
1190
1191 struct ShellSpawnIntentGuard {
1192 lifecycle: Option<ShellWorkLifecycle>,
1193 id: String,
1194 armed: bool,
1195 }
1196
1197 struct ShellSpawnContext {
1198 owner_agent: Option<ShellJobOwner>,
1199 owner_session_id: String,
1200 origin_tool_call_id: Option<String>,
1201 origin_turn_id: Option<String>,
1202 work_lifecycle: Option<ShellWorkLifecycle>,
1203 }
1204
1205 impl ShellSpawnIntentGuard {
1206 /// Register the spawn intent with the Work graph.
1207 ///
1208 /// Registration is observability bookkeeping — the same subsystem already
1209 /// treats the `observe` half as best-effort (a graph-write failure must
1210 /// not relabel a completed command) — so a transiently busy To-do/Plan
1211 /// state must not veto the command itself. Live sessions hit this: a
1212 /// shell call issued right after another tool call failed outright with
1213 /// "To-do state is busy; operation was not registered" because
1214 /// `register_operation` gives the lock only a short try-lock spin.
1215 ///
1216 /// On failure the guard goes inert: the shell still runs, and no later
1217 /// `observe` pretends the operation was bound.
1218 fn new(lifecycle: Option<ShellWorkLifecycle>, id: &str, command: &str) -> Self {
1219 let lifecycle = lifecycle.and_then(|lifecycle| match lifecycle.register(id, command) {
1220 Ok(()) => Some(lifecycle),
1221 Err(err) => {
1222 tracing::warn!(
1223 shell_id = %id,
1224 error = %err,
1225 "shell work-graph registration skipped; running without a bound operation"
1226 );
1227 None
1228 }
1229 });
1230 Self {
1231 lifecycle,
1232 id: id.to_string(),
1233 armed: true,
1234 }
1235 }
1236
1237 fn disarm(&mut self) {
1238 self.armed = false;
1239 }
1240 }
1241
1242 impl Drop for ShellSpawnIntentGuard {
1243 fn drop(&mut self) {
1244 if self.armed
1245 && let Some(lifecycle) = self.lifecycle.as_ref()
1246 && let Err(err) = lifecycle.observe(&self.id, &ShellStatus::Failed, 1, 0)
1247 {
1248 tracing::warn!(shell_id = %self.id, error = %err, "failed to record shell spawn failure");
1249 }
1250 }
1251 }
1252
1253 impl BackgroundShell {
1254 /// Wall time to report: elapsed while running, frozen once finished.
1255 fn wall_duration(&self) -> Duration {
1256 self.finished_at
1257 .unwrap_or_else(Instant::now)
1258 .saturating_duration_since(self.started_at)
1259 }
1260
1261 fn wall_millis(&self) -> u64 {
1262 u64::try_from(self.wall_duration().as_millis()).unwrap_or(u64::MAX)
1263 }
1264
1265 /// Stamp the finish instant the first time a terminal status is observed.
1266 /// Idempotent: a later poll must not restate when the job ended.
1267 fn mark_finished(&mut self) {
1268 if self.finished_at.is_none() && self.status != ShellStatus::Running {
1269 self.finished_at = Some(Instant::now());
1270 self.finished_at_utc = Some(chrono::Utc::now());
1271 }
1272 }
1273
1274 /// Check if the process has completed and update status
1275 fn poll(&mut self) -> bool {
1276 self.refresh_output_activity();
1277 if self.status != ShellStatus::Running {
1278 self.mark_finished();
1279 self.publish_lifecycle_best_effort();
1280 return true;
1281 }
1282
1283 #[cfg(unix)]
1284 let pending_process_group = (self.ownership == ShellOwnership::PersistPending)
1285 .then(|| self.child.as_ref().and_then(ShellChild::process_id))
1286 .flatten();
1287 let completed = if let Some(ref mut child) = self.child {
1288 match child.try_wait() {
1289 Ok(Some(status)) => {
1290 self.exit_code = status.code;
1291 self.status = if status.success {
1292 ShellStatus::Completed
1293 } else {
1294 ShellStatus::Failed
1295 };
1296 self.heavy_permit.take();
1297 self.collect_output();
1298 true
1299 }
1300 Ok(None) => false, // Still running
1301 Err(_) => {
1302 self.status = ShellStatus::Failed;
1303 self.wait_failed = true;
1304 self.heavy_permit.take();
1305 self.collect_output();
1306 true
1307 }
1308 }
1309 } else {
1310 true
1311 };
1312 #[cfg(unix)]
1313 if completed && let Some(process_group_id) = pending_process_group {
1314 unregister_pending_persistent_process_group(process_group_id);
1315 }
1316 self.mark_finished();
1317 self.publish_lifecycle_best_effort();
1318 completed
1319 }
1320
1321 fn publish_lifecycle(&mut self) -> Result<()> {
1322 let bytes = self.observed_output_len();
1323 if self.last_lifecycle_status.as_ref() == Some(&self.status)
1324 && self.last_lifecycle_bytes == bytes
1325 {
1326 return Ok(());
1327 }
1328 let next_seq = self.lifecycle_seq.saturating_add(1);
1329 if let Some(lifecycle) = self.work_lifecycle.as_ref() {
1330 lifecycle.observe(&self.id, &self.status, next_seq, bytes)?;
1331 }
1332 self.lifecycle_seq = next_seq;
1333 self.last_lifecycle_status = Some(self.status.clone());
1334 self.last_lifecycle_bytes = bytes;
1335 Ok(())
1336 }
1337
1338 fn publish_lifecycle_best_effort(&mut self) {
1339 if let Err(err) = self.publish_lifecycle() {
1340 tracing::warn!(shell_id = %self.id, error = %err, "failed to reconcile shell lifecycle");
1341 }
1342 }
1343
1344 fn refresh_output_activity(&mut self) {
1345 let observed_len = self.observed_output_len();
1346 if observed_len != self.last_observed_output_len {
1347 self.last_observed_output_len = observed_len;
1348 self.last_output_at = Instant::now();
1349 }
1350 }
1351
1352 fn observed_output_len(&self) -> usize {
1353 if let Some(output) = self.bounded_output.as_ref() {
1354 return output
1355 .lock()
1356 .map(|output| output.total_bytes())
1357 .unwrap_or(0);
1358 }
1359 let stdout_len = self
1360 .stdout_buffer
1361 .lock()
1362 .map(|data| data.total_len())
1363 .unwrap_or(0);
1364 let stderr_len = self
1365 .stderr_buffer
1366 .as_ref()
1367 .and_then(|buffer| buffer.lock().ok().map(|data| data.total_len()))
1368 .unwrap_or(0);
1369 stdout_len.saturating_add(stderr_len)
1370 }
1371
1372 /// Drop everything but a bounded tail of both raw streams.
1373 ///
1374 /// Called only once a job is terminal **and** its bytes have already been
1375 /// delivered — either returned as the foreground tool result, or written to
1376 /// its durable completion artifact. Before that the full bytes are a
1377 /// contract; after it they are pure residency for the up-to-1 h the
1378 /// finished record stays listed, which is what took the owner's host to
1379 /// 11 GB of swap (#5472 finding 1).
1380 fn release_delivered_output(&mut self) {
1381 if self.status == ShellStatus::Running {
1382 return;
1383 }
1384 for buffer in [Some(&self.stdout_buffer), self.stderr_buffer.as_ref()]
1385 .into_iter()
1386 .flatten()
1387 {
1388 buffer
1389 .lock()
1390 .unwrap_or_else(|poison| poison.into_inner())
1391 .release_to_tail(RAW_STREAM_SETTLED_TAIL_BYTES);
1392 }
1393 }
1394
1395 /// Bytes still held in memory for this job, for the eviction accounting in
1396 /// [`ShellManager::cleanup`] and for tests that assert the bound.
1397 fn retained_output_bytes(&self) -> usize {
1398 let stdout = self
1399 .stdout_buffer
1400 .lock()
1401 .map(|data| data.retained().len())
1402 .unwrap_or(0);
1403 let stderr = self
1404 .stderr_buffer
1405 .as_ref()
1406 .and_then(|buffer| buffer.lock().ok().map(|data| data.retained().len()))
1407 .unwrap_or(0);
1408 stdout.saturating_add(stderr)
1409 }
1410
1411 /// Collect output from the background threads
1412 fn collect_output(&mut self) {
1413 // Kill the whole process group before joining reader threads.
1414 // When the shell spawned persistent background jobs (e.g. `nohup curl`),
1415 // those subprocesses keep the pipe write-ends open after the shell exits.
1416 // Without this kill, the reader join would block until the descendant
1417 // exits, freezing the UI event loop that calls list_jobs() → poll() →
1418 // collect_output(). The joins themselves are additionally bounded
1419 // (READER_JOIN_GRACE) because a descendant in its own session/process
1420 // group escapes even the group kill (#52).
1421 #[cfg(unix)]
1422 if let Some(child) = self.child.as_mut() {
1423 match child {
1424 ShellChild::Process(proc) => {
1425 let _ = kill_child_process_group(proc);
1426 }
1427 #[cfg(not(target_env = "ohos"))]
1428 ShellChild::Pty(_) => {}
1429 }
1430 }
1431 #[cfg(windows)]
1432 terminate_and_close_windows_job(self.windows_job.take());
1433 if let Some(handle) = self.stdout_thread.take() {
1434 finish_background_reader(handle, &self.status, Some(&self.stdout_buffer));
1435 }
1436 if let Some(handle) = self.stderr_thread.take() {
1437 finish_background_reader(handle, &self.status, self.stderr_buffer.as_ref());
1438 }
1439 self.stdin = None;
1440 #[cfg(not(target_env = "ohos"))]
1441 {
1442 self.pty_master = None;
1443 }
1444 self.child = None;
1445 }
1446
1447 fn write_stdin(&mut self, input: &str, close: bool) -> Result<()> {
1448 self.write_stdin_bytes(input.as_bytes(), close)
1449 }
1450
1451 fn write_stdin_bytes(&mut self, input: &[u8], close: bool) -> Result<()> {
1452 if let Some(stdin) = self.stdin.as_mut() {
1453 if !input.is_empty() {
1454 stdin.write_all(input).context("Failed to write to stdin")?;
1455 stdin.flush().context("Failed to flush stdin")?;
1456 }
1457 if close {
1458 self.stdin = None;
1459 }
1460 return Ok(());
1461 }
1462
1463 if input.is_empty() && close {
1464 return Ok(());
1465 }
1466
1467 Err(anyhow!("stdin is not available for task {}", self.id))
1468 }
1469
1470 fn full_output(&self) -> (String, String, usize, usize) {
1471 if let Some(snapshot) = self.bounded_output_snapshot(false).ok().flatten() {
1472 return (snapshot.content, String::new(), snapshot.total_bytes, 0);
1473 }
1474 let (stdout_bytes, stderr_bytes, stdout_omitted, stderr_omitted) =
1475 self.retained_output_bytes_with_omissions();
1476 // Report what the stream produced, not what is still held.
1477 let stdout_len = stdout_bytes.len().saturating_add(stdout_omitted);
1478 let stderr_len = stderr_bytes.len().saturating_add(stderr_omitted);
1479
1480 (
1481 String::from_utf8_lossy(&stdout_bytes).to_string(),
1482 String::from_utf8_lossy(&stderr_bytes).to_string(),
1483 stdout_len,
1484 stderr_len,
1485 )
1486 }
1487
1488 /// Retained bytes for both streams plus how many leading bytes the memory
1489 /// bound discarded. Callers that publish these bytes as evidence must
1490 /// declare the omission rather than presenting a clipped stream as exact.
1491 fn retained_output_bytes_with_omissions(&self) -> (Vec<u8>, Vec<u8>, usize, usize) {
1492 if let Some(snapshot) = self.bounded_output_snapshot(false).ok().flatten() {
1493 let omitted = snapshot.total_bytes.saturating_sub(snapshot.retained_bytes);
1494 return (snapshot.content.into_bytes(), Vec::new(), omitted, 0);
1495 }
1496 let (stdout_bytes, stdout_omitted) = self
1497 .stdout_buffer
1498 .lock()
1499 .map(|data| (data.retained().to_vec(), data.dropped()))
1500 .unwrap_or_default();
1501 let (stderr_bytes, stderr_omitted) = self
1502 .stderr_buffer
1503 .as_ref()
1504 .and_then(|buffer| {
1505 buffer
1506 .lock()
1507 .ok()
1508 .map(|data| (data.retained().to_vec(), data.dropped()))
1509 })
1510 .unwrap_or_default();
1511 (stdout_bytes, stderr_bytes, stdout_omitted, stderr_omitted)
1512 }
1513
1514 fn take_delta(&mut self) -> (String, String, usize, usize, usize, usize) {
1515 // Only the bytes after the cursor: returning the whole retained tail
1516 // made every `wait` poll repeat output the caller already holds.
1517 let bounded_delta = self.bounded_output.as_ref().and_then(|output| {
1518 let output = output.lock().unwrap_or_else(|error| error.into_inner());
1519 let total = output.total_bytes();
1520 output
1521 .delta_since(self.stdout_cursor)
1522 .ok()
1523 .map(|delta| (delta, total))
1524 });
1525 if let Some(((mut delta, omitted), total_bytes)) = bounded_delta {
1526 let delta_len = total_bytes.saturating_sub(self.stdout_cursor);
1527 self.stdout_cursor = total_bytes;
1528 if delta_len == 0 {
1529 return (String::new(), String::new(), 0, 0, total_bytes, 0);
1530 }
1531 self.last_output_at = Instant::now();
1532 self.last_observed_output_len = total_bytes;
1533 if omitted > 0
1534 && let Some(output) = self.bounded_output.as_ref()
1535 {
1536 let notice = output
1537 .lock()
1538 .unwrap_or_else(|error| error.into_inner())
1539 .omitted_notice(omitted);
1540 delta.insert_str(0, &notice);
1541 }
1542 return (delta, String::new(), delta_len, 0, total_bytes, 0);
1543 }
1544 let (stdout_delta, stdout_total) =
1545 take_delta_from_buffer(&self.stdout_buffer, &mut self.stdout_cursor);
1546 let (stderr_delta, stderr_total) = if let Some(buffer) = self.stderr_buffer.as_ref() {
1547 take_delta_from_buffer(buffer, &mut self.stderr_cursor)
1548 } else {
1549 (Vec::new(), 0)
1550 };
1551
1552 let stdout_delta_len = stdout_delta.len();
1553 let stderr_delta_len = stderr_delta.len();
1554
1555 if stdout_delta_len > 0 || stderr_delta_len > 0 {
1556 self.last_output_at = Instant::now();
1557 self.last_observed_output_len = stdout_total.saturating_add(stderr_total);
1558 }
1559
1560 (
1561 String::from_utf8_lossy(&stdout_delta).to_string(),
1562 String::from_utf8_lossy(&stderr_delta).to_string(),
1563 stdout_delta_len,
1564 stderr_delta_len,
1565 stdout_total,
1566 stderr_total,
1567 )
1568 }
1569
1570 fn sandbox_denied(&self) -> bool {
1571 if matches!(self.status, ShellStatus::Running) {
1572 return false;
1573 }
1574 let (_, stderr_full, _, _) = self.full_output();
1575 SandboxManager::was_denied(
1576 self.sandbox_type,
1577 self.exit_code
1578 .and_then(|code| i32::try_from(code).ok())
1579 .unwrap_or(-1),
1580 &stderr_full,
1581 )
1582 }
1583
1584 /// Kill the process
1585 fn kill(&mut self) -> Result<()> {
1586 #[cfg(unix)]
1587 if self.ownership == ShellOwnership::PersistPending
1588 && let Some(process_group_id) = self.child.as_ref().and_then(ShellChild::process_id)
1589 {
1590 unregister_pending_persistent_process_group(process_group_id);
1591 }
1592 if let Some(ref mut child) = self.child {
1593 match child {
1594 ShellChild::Process(proc) => {
1595 #[cfg(windows)]
1596 {
1597 terminate_windows_job(self.windows_job.as_ref(), proc)
1598 .context("Failed to kill process tree")?;
1599 let _ = proc.wait();
1600 }
1601 #[cfg(all(not(windows), unix))]
1602 {
1603 // Bounded SIGTERM → SIGKILL escalation against the
1604 // whole process group; returns within ~grace even if
1605 // the command ignores SIGTERM (#52).
1606 terminate_child_process_group(proc).context("Failed to kill process")?;
1607 }
1608 #[cfg(all(not(windows), not(unix)))]
1609 {
1610 proc.kill().context("Failed to kill process")?;
1611 let _ = proc.wait();
1612 }
1613 }
1614 #[cfg(not(target_env = "ohos"))]
1615 ShellChild::Pty(child) => {
1616 child.kill().context("Failed to kill process")?;
1617 let _ = child.wait();
1618 }
1619 }
1620 }
1621 self.status = ShellStatus::Killed;
1622 self.mark_finished();
1623 self.heavy_permit.take();
1624 self.collect_output();
1625 self.publish_lifecycle_best_effort();
1626 Ok(())
1627 }
1628
1629 /// Get a snapshot of the current state
1630 pub fn snapshot(&self) -> Result<ShellResult> {
1631 let sandboxed = !matches!(self.sandbox_type, SandboxType::None);
1632 if let Some(snapshot) = self.bounded_output_snapshot(self.status != ShellStatus::Running)? {
1633 return Ok(ShellResult {
1634 task_id: Some(self.id.clone()),
1635 status: self.status.clone(),
1636 exit_code: self.exit_code,
1637 stdout: snapshot.content,
1638 stderr: String::new(),
1639 duration_ms: self.wall_millis(),
1640 stdout_len: snapshot.total_bytes,
1641 stderr_len: 0,
1642 stdout_omitted: snapshot.total_bytes.saturating_sub(snapshot.retained_bytes),
1643 stderr_omitted: 0,
1644 stdout_truncated: snapshot.truncated,
1645 stderr_truncated: false,
1646 sandboxed,
1647 sandbox_type: sandboxed.then(|| self.sandbox_type.to_string()),
1648 sandbox_denied: false,
1649 });
1650 }
1651 let (stdout_full, stderr_full, stdout_total, stderr_total) = self.full_output();
1652 let (stdout, stdout_meta) = truncate_with_meta(&stdout_full);
1653 let (stderr, stderr_meta) = truncate_with_meta(&stderr_full);
1654 // `truncate_with_meta` can only see the bytes still held. Fold in what
1655 // the in-memory bound dropped so a >16 MiB stream reports its real
1656 // length and its real omission instead of silently shrinking (#5472).
1657 let stdout_dropped = stdout_total.saturating_sub(stdout_meta.original_len);
1658 let stderr_dropped = stderr_total.saturating_sub(stderr_meta.original_len);
1659 Ok(ShellResult {
1660 task_id: Some(self.id.clone()),
1661 status: self.status.clone(),
1662 exit_code: self.exit_code,
1663 stdout,
1664 stderr,
1665 duration_ms: self.wall_millis(),
1666 stdout_len: stdout_total,
1667 stderr_len: stderr_total,
1668 stdout_omitted: stdout_meta.omitted.saturating_add(stdout_dropped),
1669 stderr_omitted: stderr_meta.omitted.saturating_add(stderr_dropped),
1670 stdout_truncated: stdout_meta.truncated || stdout_dropped > 0,
1671 stderr_truncated: stderr_meta.truncated || stderr_dropped > 0,
1672 sandboxed,
1673 sandbox_type: if sandboxed {
1674 Some(self.sandbox_type.to_string())
1675 } else {
1676 None
1677 },
1678 sandbox_denied: self.sandbox_denied(),
1679 })
1680 }
1681
1682 fn bounded_output_snapshot(&self, finalize: bool) -> Result<Option<BoundedOutputSnapshot>> {
1683 self.bounded_output
1684 .as_ref()
1685 .map(|output| {
1686 output
1687 .lock()
1688 .unwrap_or_else(|error| error.into_inner())
1689 .snapshot(finalize)
1690 .map_err(anyhow::Error::from)
1691 })
1692 .transpose()
1693 }
1694
1695 fn job_snapshot(&self) -> ShellJobSnapshot {
1696 // Use tail_from_buffer instead of full_output so we never clone the
1697 // entire accumulated stdout/stderr for display purposes. full_output
1698 // is O(total_bytes_written), which caused the ShellManager mutex to be
1699 // held for an arbitrarily long time during list_jobs() calls from the
1700 // TUI event loop — freezing input handling on long automation runs.
1701 let (stdout_len, stdout_tail) =
1702 if let Some(snapshot) = self.bounded_output_snapshot(false).ok().flatten() {
1703 (snapshot.total_bytes, tail_text(&snapshot.content, 1_200))
1704 } else {
1705 tail_from_buffer(&self.stdout_buffer, 1200)
1706 };
1707 let (stderr_len, stderr_tail) = self
1708 .stderr_buffer
1709 .as_ref()
1710 .map(|buf| tail_from_buffer(buf, 1200))
1711 .unwrap_or((0, String::new()));
1712 let elapsed_since_output_ms = (self.status == ShellStatus::Running)
1713 .then(|| u64::try_from(self.last_output_at.elapsed().as_millis()).unwrap_or(u64::MAX));
1714 // A live PTY can wait indefinitely for input. Silence does not mean it
1715 // lost its process owner; marking it stale hides its transport from API
1716 // clients and incorrectly prevents reconnect/resize after one minute.
1717 let stale = self.terminal_size.is_none()
1718 && elapsed_since_output_ms.is_some_and(|elapsed| {
1719 elapsed >= u64::try_from(STALE_NO_OUTPUT_AFTER.as_millis()).unwrap_or(u64::MAX)
1720 });
1721 ShellJobSnapshot {
1722 id: self.id.clone(),
1723 job_id: self.id.clone(),
1724 command: self.command.clone(),
1725 cwd: self.working_dir.clone(),
1726 status: self.status.clone(),
1727 background: self.background,
1728 finished_at: self.finished_at_utc,
1729 exit_code: self.exit_code,
1730 elapsed_ms: self.wall_millis(),
1731 stdout_tail,
1732 stderr_tail,
1733 stdout_len,
1734 stderr_len,
1735 stdin_available: self.stdin.is_some() && self.status == ShellStatus::Running,
1736 stale,
1737 elapsed_since_output_ms,
1738 linked_task_id: self.linked_task_id.clone(),
1739 owner_agent_id: self
1740 .owner_agent
1741 .as_ref()
1742 .map(|owner| owner.agent_id.clone()),
1743 owner_agent_name: self
1744 .owner_agent
1745 .as_ref()
1746 .map(|owner| owner.agent_name.clone()),
1747 origin_tool_call_id: self.origin_tool_call_id.clone(),
1748 origin_turn_id: self.origin_turn_id.clone(),
1749 owner_session_id: self.owner_session_id.clone(),
1750 }
1751 }
1752
1753 fn completion_event(&self) -> ShellCompletionEvent {
1754 let snapshot = self.job_snapshot();
1755 let (stdout_len, stdout_tail) =
1756 if let Some(output) = self.bounded_output_snapshot(false).ok().flatten() {
1757 (
1758 output.total_bytes,
1759 tail_text(&output.content, SHELL_COMPLETION_TAIL_BYTES),
1760 )
1761 } else {
1762 bounded_completion_tail(&self.stdout_buffer, SHELL_COMPLETION_TAIL_BYTES)
1763 };
1764 let (stderr_len, stderr_tail) = self
1765 .stderr_buffer
1766 .as_ref()
1767 .map(|buffer| bounded_completion_tail(buffer, SHELL_COMPLETION_TAIL_BYTES))
1768 .unwrap_or((0, String::new()));
1769 ShellCompletionEvent {
1770 task_id: snapshot.id,
1771 command: snapshot.command,
1772 status: snapshot.status,
1773 exit_code: snapshot.exit_code,
1774 duration_ms: snapshot.elapsed_ms,
1775 stdout_tail,
1776 stderr_tail,
1777 stdout_len,
1778 stderr_len,
1779 evidence_ref: None,
1780 linked_task_id: snapshot.linked_task_id,
1781 owner_agent_id: snapshot.owner_agent_id,
1782 owner_agent_name: snapshot.owner_agent_name,
1783 origin_tool_call_id: snapshot.origin_tool_call_id,
1784 origin_turn_id: snapshot.origin_turn_id,
1785 owner_session_id: snapshot.owner_session_id,
1786 }
1787 }
1788
1789 fn completion_evidence(&self) -> ShellCompletionEvidence {
1790 let event = self.completion_event();
1791 let (stdout, stderr, stdout_omitted, stderr_omitted) =
1792 self.retained_output_bytes_with_omissions();
1793 ShellCompletionEvidence {
1794 event,
1795 stdout,
1796 stderr,
1797 stdout_omitted,
1798 stderr_omitted,
1799 }
1800 }
1801
1802 fn job_detail(&self) -> ShellJobDetail {
1803 let (stdout, stderr, _, _) = self.full_output();
1804 ShellJobDetail {
1805 snapshot: self.job_snapshot(),
1806 stdout,
1807 stderr,
1808 }
1809 }
1810 }
1811
1812 fn finish_background_reader(
1813 handle: std::thread::JoinHandle<()>,
1814 status: &ShellStatus,
1815 buffer: Option<&SharedRawOutput>,
1816 ) {
1817 // A killed Windows process can leave a pipe reader blocked even after its
1818 // Job Object has been closed. Cancellation must return promptly instead of
1819 // waiting for that reader to observe EOF. Other terminal states still join
1820 // so their final output is collected before the shell is discarded.
1821 #[cfg(windows)]
1822 if *status == ShellStatus::Killed {
1823 drop(handle);
1824 return;
1825 }
1826
1827 #[cfg(not(windows))]
1828 let _ = status;
1829
1830 // Bounded join (#52): after the process group is killed the reader
1831 // normally sees EOF immediately, but a descendant that escaped the group
1832 // (its own session/process group) keeps its inherited pipe write-end
1833 // open, so the reader stays blocked until that descendant exits on its
1834 // own. Joining unboundedly froze the foreground shell — and, through the
1835 // shell-manager lock, every other shell — for minutes. On timeout the
1836 // join is handed to a helper thread and we return; the reader thread
1837 // still finishes on its own once the pipe finally closes.
1838 let (done_tx, done_rx) = std::sync::mpsc::channel();
1839 std::thread::spawn(move || {
1840 let _ = handle.join();
1841 let _ = done_tx.send(());
1842 });
1843 if done_rx.recv_timeout(READER_JOIN_GRACE).is_ok() {
1844 return;
1845 }
1846 // The reader is still blocked in `read()` on a pipe a descendant refuses to
1847 // close. Previously both it and the helper thread above stayed alive for the
1848 // life of the process, the reader still appending into a buffer nobody would
1849 // ever read (#5472 finding 2). Abandoning the stream releases what it holds
1850 // and gives the reader an exit on its next wakeup, which also lets the
1851 // helper's `join` return.
1852 if let Some(buffer) = buffer {
1853 buffer
1854 .lock()
1855 .unwrap_or_else(|poison| poison.into_inner())
1856 .abandon();
1857 }
1858 }
1859
1860 impl Drop for BackgroundShell {
1861 fn drop(&mut self) {
1862 #[cfg(unix)]
1863 if self.ownership == ShellOwnership::PersistPending
1864 && let Some(process_group_id) = self.child.as_ref().and_then(ShellChild::process_id)
1865 {
1866 unregister_pending_persistent_process_group(process_group_id);
1867 }
1868 if self.ownership != ShellOwnership::Released
1869 && self.status == ShellStatus::Running
1870 && let Some(ref mut child) = self.child
1871 {
1872 #[cfg(windows)]
1873 match child {
1874 ShellChild::Process(proc) => {
1875 let _ = terminate_windows_job(self.windows_job.as_ref(), proc);
1876 }
1877 #[cfg(not(target_env = "ohos"))]
1878 ShellChild::Pty(child) => {
1879 let _ = child.kill();
1880 }
1881 }
1882 #[cfg(all(not(windows), unix))]
1883 {
1884 let _ = child.kill();
1885 match child {
1886 ShellChild::Process(proc) => {
1887 let _ = wait_child_bounded(proc, KILL_REAP_GRACE);
1888 }
1889 #[cfg(not(target_env = "ohos"))]
1890 ShellChild::Pty(child) => {
1891 let _ = child.wait();
1892 }
1893 }
1894 }
1895 #[cfg(all(not(windows), not(unix)))]
1896 {
1897 let _ = child.kill();
1898 let _ = child.wait();
1899 }
1900 }
1901 }
1902 }
1903
1904 #[cfg(all(unix, not(target_env = "ohos")))]
1905 pub(crate) fn inherited_interactive_terminal_refusal() -> Option<&'static str> {
1906 Some(
1907 "Inherited interactive terminal takeover is unavailable on Unix because foreground TTY \
1908 ownership cannot be transferred safely. Use Bash with `background: true, tty: true`, \
1909 then continue it with `action: \"interact\"` and the returned `task_id`; alternatively \
1910 use `terminal/run` and `terminal/send`, or launch the command in a new terminal. For \
1911 non-interactive work, omit `interactive: true`.",
1912 )
1913 }
1914
1915 #[cfg(all(unix, target_env = "ohos"))]
1916 pub(crate) fn inherited_interactive_terminal_refusal() -> Option<&'static str> {
1917 Some(
1918 "Inherited interactive terminal takeover is unavailable on Unix because foreground TTY \
1919 ownership cannot be transferred safely. Launch the command in a new terminal, or omit \
1920 `interactive: true` for non-interactive work.",
1921 )
1922 }
1923
1924 #[cfg(not(unix))]
1925 pub(crate) fn inherited_interactive_terminal_refusal() -> Option<&'static str> {
1926 None
1927 }
1928
1929 /// Manages background shell processes with optional sandboxing.
1930 pub struct ShellManager {
1931 processes: HashMap<String, BackgroundShell>,
1932 stale_jobs: HashMap<String, ShellJobSnapshot>,
1933 default_workspace: PathBuf,
1934 sandbox_manager: SandboxManager,
1935 sandbox_policy: ExecutionSandboxPolicy,
1936 foreground_background_requested: bool,
1937 /// Directory for lowercase-`bash` complete-output spill files
1938 /// (`None` = process temp dir). Overridable so tests can fault-inject a
1939 /// missing/unwritable spill location.
1940 output_spill_dir: Option<PathBuf>,
1941 }
1942
1943 impl std::fmt::Debug for ShellManager {
1944 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1945 f.debug_struct("ShellManager")
1946 .field("processes", &self.processes.len())
1947 .field("stale_jobs", &self.stale_jobs.len())
1948 .field("default_workspace", &self.default_workspace)
1949 .field("sandbox_policy", &self.sandbox_policy)
1950 .field(
1951 "foreground_background_requested",
1952 &self.foreground_background_requested,
1953 )
1954 .finish()
1955 }
1956 }
1957
1958 impl ShellManager {
1959 fn require_session_owner(&self, task_id: &str, active_session_id: &str) -> Result<()> {
1960 let owned = self.processes.get(task_id).is_some_and(|shell| {
1961 !active_session_id.is_empty() && shell.owner_session_id == active_session_id
1962 }) || self.stale_jobs.get(task_id).is_some_and(|job| {
1963 !active_session_id.is_empty() && job.owner_session_id == active_session_id
1964 });
1965 if owned {
1966 Ok(())
1967 } else {
1968 // Do not disclose whether the id exists in another session.
1969 Err(anyhow!("Job {task_id} not found"))
1970 }
1971 }
1972
1973 /// Create a new `ShellManager` with default (no sandbox) policy.
1974 pub fn new(workspace: PathBuf) -> Self {
1975 Self {
1976 processes: HashMap::new(),
1977 stale_jobs: HashMap::new(),
1978 default_workspace: workspace,
1979 sandbox_manager: SandboxManager::new(),
1980 sandbox_policy: ExecutionSandboxPolicy::default(),
1981 foreground_background_requested: false,
1982 output_spill_dir: None,
1983 }
1984 }
1985
1986 /// Point lowercase-`bash` complete-output spill files at `dir` instead of
1987 /// the process temp dir. Tests use a nonexistent dir to simulate a full or
1988 /// broken temp volume. (unix-only: the regression test that uses it drives
1989 /// a POSIX shell loop.)
1990 #[cfg(all(test, unix))]
1991 pub(crate) fn set_output_spill_dir_for_test(&mut self, dir: Option<PathBuf>) {
1992 self.output_spill_dir = dir;
1993 }
1994
1995 /// Insert a finished job without spawning. Count-bound tests would
1996 /// otherwise pay for 100+ live shells.
1997 #[cfg(test)]
1998 pub(crate) fn seed_finished_record_for_test(&mut self, id: impl Into<String>, age: Duration) {
1999 let now = Instant::now();
2000 let started_at = now.checked_sub(age).unwrap_or(now);
2001 let id = id.into();
2002 self.processes.insert(
2003 id.clone(),
2004 BackgroundShell {
2005 id,
2006 command: String::new(),
2007 working_dir: self.default_workspace.clone(),
2008 status: ShellStatus::Completed,
2009 background: true,
2010 exit_code: Some(0),
2011 started_at,
2012 finished_at: Some(now),
2013 finished_at_utc: Some(chrono::Utc::now()),
2014 last_output_at: now,
2015 last_observed_output_len: 0,
2016 sandbox_type: SandboxType::None,
2017 linked_task_id: None,
2018 owner_agent: None,
2019 owner_session_id: String::new(),
2020 origin_tool_call_id: None,
2021 origin_turn_id: None,
2022 ownership: ShellOwnership::Managed,
2023 stdout_buffer: new_shared_raw_output(),
2024 stderr_buffer: Some(new_shared_raw_output()),
2025 bounded_output: None,
2026 heavy_permit: None,
2027 stdout_cursor: 0,
2028 stderr_cursor: 0,
2029 completion_reported: false,
2030 stdin: None,
2031 #[cfg(not(target_env = "ohos"))]
2032 pty_master: None,
2033 terminal_size: None,
2034 child: None,
2035 #[cfg(windows)]
2036 windows_job: None,
2037 #[cfg(unix)]
2038 parent_death_watch: None,
2039 stdout_thread: None,
2040 stderr_thread: None,
2041 work_lifecycle: None,
2042 lifecycle_seq: 0,
2043 last_lifecycle_status: None,
2044 last_lifecycle_bytes: 0,
2045 wait_failed: false,
2046 },
2047 );
2048 }
2049
2050 /// Test-only observation of the workspace selected by runtime rebuilds.
2051 #[cfg(test)]
2052 pub(crate) fn default_workspace(&self) -> &Path {
2053 &self.default_workspace
2054 }
2055
2056 /// Enable or disable bubblewrap passthrough (#2184).
2057 ///
2058 /// When enabled and `/usr/bin/bwrap` is executable on Linux, exec_shell
2059 /// commands are routed through bubblewrap for filesystem isolation.
2060 pub fn set_prefer_bwrap(&mut self, prefer: bool) {
2061 self.sandbox_manager.set_prefer_bwrap(prefer);
2062 }
2063
2064 /// Move the fallback working directory. Callers that pass an explicit
2065 /// `working_dir` are unaffected; this only keeps `None` honest when a
2066 /// thread's workspace changes while its jobs are still tracked here.
2067 pub fn set_default_workspace(&mut self, workspace: PathBuf) {
2068 self.default_workspace = workspace;
2069 }
2070
2071 /// Set user-configured bwrap mount extensions (#5410): extra read-only
2072 /// roots and writable device nodes such as `/dev/null`.
2073 pub fn set_bwrap_extensions(&mut self, extensions: crate::sandbox::BwrapMountExtensions) {
2074 self.sandbox_manager.set_bwrap_extensions(extensions);
2075 }
2076
2077 /// Forward the opt-in sandbox read deny-list (S1, #5568) to the sandbox
2078 /// manager; `~` prefixes expand there.
2079 pub fn set_denied_read_subpaths(&mut self, paths: Vec<std::path::PathBuf>) {
2080 self.sandbox_manager.set_denied_read_subpaths(paths);
2081 }
2082
2083 /// Return the OS sandbox wrapper this shell manager is configured and able
2084 /// to apply to commands.
2085 pub fn configured_sandbox_type(&self) -> Option<SandboxType> {
2086 self.sandbox_manager.configured_sandbox()
2087 }
2088
2089 /// Prepare a program for a tool that runs workspace code outside
2090 /// `exec_shell`, under the policy and sandbox configuration a shell
2091 /// command would get in this session.
2092 fn prepare_runner(
2093 &self,
2094 program: &str,
2095 args: Vec<String>,
2096 cwd: &Path,
2097 timeout: Duration,
2098 policy_override: Option<ExecutionSandboxPolicy>,
2099 ) -> ExecEnv {
2100 let policy = policy_override.unwrap_or_else(|| self.sandbox_policy.clone());
2101 let spec =
2102 CommandSpec::program(program, args, cwd.to_path_buf(), timeout).with_policy(policy);
2103 self.sandbox_manager.prepare(&spec)
2104 }
2105
2106 /// Request that the active foreground shell wait detach and leave its
2107 /// process running in the background job table.
2108 pub fn request_foreground_background(&mut self) {
2109 self.foreground_background_requested = true;
2110 }
2111
2112 #[cfg(test)]
2113 pub(crate) fn foreground_background_requested_for_test(&self) -> bool {
2114 self.foreground_background_requested
2115 }
2116
2117 fn clear_foreground_background_request(&mut self) {
2118 self.foreground_background_requested = false;
2119 }
2120
2121 fn take_foreground_background_request(&mut self, task_id: &str) -> bool {
2122 let requested = self.foreground_background_requested;
2123 self.foreground_background_requested = false;
2124 if requested && let Some(shell) = self.processes.get_mut(task_id) {
2125 shell.background = true;
2126 }
2127 requested
2128 }
2129
2130 /// Execute a shell command with stdin/TTY options plus an extra env-var map
2131 /// that is merged into the spawned process environment. Used by the
2132 /// `shell_env` hook injection path (#456).
2133 #[allow(clippy::too_many_arguments)]
2134 #[cfg(test)]
2135 pub fn execute_with_options_env(
2136 &mut self,
2137 command: &str,
2138 working_dir: Option<&str>,
2139 timeout_ms: u64,
2140 background: bool,
2141 stdin_data: Option<&str>,
2142 tty: bool,
2143 policy_override: Option<ExecutionSandboxPolicy>,
2144 extra_env: HashMap<String, String>,
2145 ) -> Result<ShellResult> {
2146 self.execute_with_options_env_for_owner(
2147 command,
2148 working_dir,
2149 timeout_ms,
2150 background,
2151 stdin_data,
2152 tty,
2153 policy_override,
2154 extra_env,
2155 None,
2156 )
2157 }
2158
2159 /// Launch a parent-owned job stamped to the immutable root session.
2160 #[allow(clippy::too_many_arguments)]
2161 pub fn execute_with_options_env_for_session(
2162 &mut self,
2163 command: &str,
2164 working_dir: Option<&str>,
2165 timeout_ms: u64,
2166 background: bool,
2167 stdin_data: Option<&str>,
2168 tty: bool,
2169 policy_override: Option<ExecutionSandboxPolicy>,
2170 extra_env: HashMap<String, String>,
2171 owner_session_id: &str,
2172 ) -> Result<ShellResult> {
2173 self.execute_with_options_env_for_owner_and_work(
2174 command,
2175 working_dir,
2176 timeout_ms,
2177 background,
2178 stdin_data,
2179 tty,
2180 policy_override,
2181 extra_env,
2182 None,
2183 owner_session_id.to_string(),
2184 None,
2185 None,
2186 None,
2187 None,
2188 false,
2189 (1_000, 600_000),
2190 )
2191 }
2192
2193 /// Same as `execute_with_options_env`, with optional background-job owner
2194 /// attribution for sub-agent launched jobs.
2195 #[allow(clippy::too_many_arguments)]
2196 #[cfg(test)]
2197 pub fn execute_with_options_env_for_owner(
2198 &mut self,
2199 command: &str,
2200 working_dir: Option<&str>,
2201 timeout_ms: u64,
2202 background: bool,
2203 stdin_data: Option<&str>,
2204 tty: bool,
2205 policy_override: Option<ExecutionSandboxPolicy>,
2206 extra_env: HashMap<String, String>,
2207 owner_agent: Option<ShellJobOwner>,
2208 ) -> Result<ShellResult> {
2209 self.execute_with_options_env_for_owner_and_work(
2210 command,
2211 working_dir,
2212 timeout_ms,
2213 background,
2214 stdin_data,
2215 tty,
2216 policy_override,
2217 extra_env,
2218 owner_agent,
2219 String::new(),
2220 None,
2221 None,
2222 None,
2223 None,
2224 false,
2225 (1_000, 600_000),
2226 )
2227 }
2228
2229 /// Test-only owner-aware launch with an explicit immutable session owner.
2230 #[allow(clippy::too_many_arguments)]
2231 #[cfg(test)]
2232 pub fn execute_with_options_env_for_owner_and_session(
2233 &mut self,
2234 command: &str,
2235 working_dir: Option<&str>,
2236 timeout_ms: u64,
2237 background: bool,
2238 stdin_data: Option<&str>,
2239 tty: bool,
2240 policy_override: Option<ExecutionSandboxPolicy>,
2241 extra_env: HashMap<String, String>,
2242 owner_agent: Option<ShellJobOwner>,
2243 owner_session_id: &str,
2244 ) -> Result<ShellResult> {
2245 self.execute_with_options_env_for_owner_and_work(
2246 command,
2247 working_dir,
2248 timeout_ms,
2249 background,
2250 stdin_data,
2251 tty,
2252 policy_override,
2253 extra_env,
2254 owner_agent,
2255 owner_session_id.to_string(),
2256 None,
2257 None,
2258 None,
2259 None,
2260 false,
2261 (1_000, 600_000),
2262 )
2263 }
2264
2265 /// Owner-aware execution with an optional Work Graph lifecycle sink.
2266 #[allow(clippy::too_many_arguments)]
2267 fn execute_with_options_env_for_owner_and_work(
2268 &mut self,
2269 command: &str,
2270 working_dir: Option<&str>,
2271 timeout_ms: u64,
2272 background: bool,
2273 stdin_data: Option<&str>,
2274 tty: bool,
2275 policy_override: Option<ExecutionSandboxPolicy>,
2276 extra_env: HashMap<String, String>,
2277 owner_agent: Option<ShellJobOwner>,
2278 owner_session_id: String,
2279 origin_tool_call_id: Option<String>,
2280 origin_turn_id: Option<String>,
2281 work_lifecycle: Option<ShellWorkLifecycle>,
2282 readonly_workspace: Option<&std::path::Path>,
2283 persist_pending: bool,
2284 timeout_bounds_ms: (u64, u64),
2285 ) -> Result<ShellResult> {
2286 // Log execution via ShellDispatcher when SHELL_DISPATCHER_LOG is set.
2287 crate::shell_dispatcher::ShellDispatcher::log_exec(command);
2288
2289 let work_dir = working_dir.map_or_else(|| self.default_workspace.clone(), PathBuf::from);
2290 validate_shell_working_dir(&work_dir, working_dir.is_none())?;
2291
2292 let timeout_ms = timeout_ms.clamp(timeout_bounds_ms.0, timeout_bounds_ms.1);
2293
2294 // Use override policy if provided, otherwise use the manager's policy
2295 let policy = policy_override.unwrap_or_else(|| self.sandbox_policy.clone());
2296
2297 // Create command spec and prepare sandboxed environment
2298 let spec = if let Some(workspace) = readonly_workspace {
2299 if readonly_command_needs_shell(command) {
2300 let script = hardened_readonly_script(command, workspace)?;
2301 CommandSpec::shell(&script, work_dir.clone(), Duration::from_millis(timeout_ms))
2302 } else {
2303 let (program, args) = hardened_readonly_argv(command)?;
2304 let program = resolve_readonly_program(&program, workspace)?;
2305 CommandSpec::program(
2306 program
2307 .to_str()
2308 .ok_or_else(|| anyhow!("read-only executable path is not valid UTF-8"))?,
2309 args,
2310 work_dir.clone(),
2311 Duration::from_millis(timeout_ms),
2312 )
2313 }
2314 } else {
2315 CommandSpec::shell(command, work_dir.clone(), Duration::from_millis(timeout_ms))
2316 };
2317 let spec = spec.with_policy(policy).with_env(extra_env);
2318 let exec_env = self.sandbox_manager.prepare(&spec);
2319 if matches!(spec.sandbox_policy, ExecutionSandboxPolicy::ReadOnly)
2320 && readonly_workspace.is_none()
2321 {
2322 // Arbitrary code with a read-only policy needs kernel enforcement;
2323 // only the separately hardened argv subset may run without it.
2324 require_native_readonly_execution(&exec_env)?;
2325 }
2326
2327 if background {
2328 let bounded_output = timeout_bounds_ms == (1, BASH_MAX_TIMEOUT_MS);
2329 self.spawn_background_sandboxed(
2330 command,
2331 &work_dir,
2332 &exec_env,
2333 None,
2334 stdin_data,
2335 tty,
2336 ShellSpawnContext {
2337 owner_agent,
2338 owner_session_id,
2339 origin_tool_call_id,
2340 origin_turn_id,
2341 work_lifecycle,
2342 },
2343 persist_pending,
2344 bounded_output,
2345 )
2346 } else {
2347 if tty {
2348 return Err(anyhow!(
2349 "TTY mode requires background execution (set background: true)."
2350 ));
2351 }
2352 Self::execute_sync_sandboxed(command, &work_dir, timeout_ms, stdin_data, &exec_env)
2353 }
2354 }
2355
2356 /// Interactive variant that accepts extra env vars (#456 shell_env hook).
2357 pub fn execute_interactive_with_policy_env(
2358 &mut self,
2359 command: &str,
2360 working_dir: Option<&str>,
2361 timeout_ms: u64,
2362 policy_override: Option<ExecutionSandboxPolicy>,
2363 extra_env: HashMap<String, String>,
2364 ) -> Result<ShellResult> {
2365 crate::shell_dispatcher::ShellDispatcher::log_exec(command);
2366
2367 // A new Unix process group that inherits the terminal is not its
2368 // foreground owner. Letting it read stdin triggers SIGTTIN; sharing
2369 // Codewhale's group instead would make cooked-mode Ctrl+C terminate
2370 // both parent and child. Until this path owns a complete POSIX job-
2371 // control lease, fail closed before spawning. Persistent PTY tools
2372 // already provide a safe interactive lane without taking over the
2373 // operator's live terminal.
2374 if let Some(message) = inherited_interactive_terminal_refusal() {
2375 return Err(anyhow!(message));
2376 }
2377
2378 let work_dir = working_dir.map_or_else(|| self.default_workspace.clone(), PathBuf::from);
2379 validate_shell_working_dir(&work_dir, working_dir.is_none())?;
2380
2381 let timeout_ms = timeout_ms.clamp(1000, 600_000);
2382 let policy = policy_override.unwrap_or_else(|| self.sandbox_policy.clone());
2383
2384 let spec = CommandSpec::shell(command, work_dir.clone(), Duration::from_millis(timeout_ms))
2385 .with_policy(policy)
2386 .with_env(extra_env);
2387 let exec_env = self.sandbox_manager.prepare(&spec);
2388
2389 Self::execute_interactive_sandboxed(command, &work_dir, timeout_ms, &exec_env)
2390 }
2391
2392 /// Execute command synchronously with timeout (sandboxed).
2393 fn execute_sync_sandboxed(
2394 original_command: &str,
2395 working_dir: &std::path::Path,
2396 timeout_ms: u64,
2397 stdin_data: Option<&str>,
2398 exec_env: &ExecEnv,
2399 ) -> Result<ShellResult> {
2400 let started = Instant::now();
2401 let timeout = Duration::from_millis(timeout_ms);
2402 let sandbox_type = exec_env.sandbox_type;
2403 let sandboxed = exec_env.is_sandboxed();
2404
2405 // Build the command from ExecEnv
2406 let program = exec_env.program();
2407 let args = exec_env.args();
2408
2409 let mut cmd = Command::new(program);
2410 crate::utils::suppress_console_window(&mut cmd);
2411 push_shell_args(&mut cmd, program, args);
2412 cmd.current_dir(working_dir)
2413 .stdout(Stdio::piped())
2414 .stderr(Stdio::piped());
2415 #[cfg(unix)]
2416 {
2417 cmd.process_group(0);
2418 }
2419 install_parent_death_signal(&mut cmd);
2420
2421 // Without input, stdin is closed rather than inherited: an unexpected
2422 // read (`cat`, a prompt) gets EOF instead of blocking on, or reading,
2423 // the operator's terminal until the timeout.
2424 cmd.stdin(if stdin_data.is_some() {
2425 Stdio::piped()
2426 } else {
2427 Stdio::null()
2428 });
2429
2430 child_env::apply_to_command(&mut cmd, child_env::string_map_env(&exec_env.env));
2431 remove_readonly_redirect_env(&mut cmd, &exec_env.env);
2432
2433 // Leave raw mode before spawn; the guard restores it only if raw mode
2434 // was active on entry (issue #1690).
2435 let _raw_mode = crate::host_terminal::suspend_raw_mode();
2436
2437 let mut child = cmd
2438 .spawn()
2439 .with_context(|| format!("Failed to execute: {original_command}"))?;
2440 #[cfg(windows)]
2441 let windows_job = attach_windows_job(&child, original_command);
2442
2443 if let Some(input) = stdin_data
2444 && let Some(mut stdin) = child.stdin.take()
2445 {
2446 stdin
2447 .write_all(input.as_bytes())
2448 .context("Failed to write to stdin")?;
2449 stdin.flush().ok();
2450 }
2451
2452 let stdout_handle = child.stdout.take().context("Failed to capture stdout")?;
2453 let stderr_handle = child.stderr.take().context("Failed to capture stderr")?;
2454
2455 // Spawn threads to read output. Use bounded receives below so a killed
2456 // or detached descendant that keeps pipe handles open cannot wedge the
2457 // foreground shell path while the global tool lock is held (#2571).
2458 let stdout_rx = spawn_sync_reader_thread(stdout_handle);
2459 let stderr_rx = spawn_sync_reader_thread(stderr_handle);
2460
2461 // Wait with timeout
2462 if let Some(status) = child.wait_timeout(timeout)? {
2463 let status = ShellExitStatus::from_std(status);
2464 #[cfg(unix)]
2465 let _ = kill_child_process_group(&mut child);
2466 #[cfg(windows)]
2467 terminate_and_close_windows_job(windows_job);
2468 let stdout = recv_sync_reader_output(&stdout_rx);
2469 let stderr = recv_sync_reader_output(&stderr_rx);
2470 let stdout_str = String::from_utf8_lossy(&stdout).to_string();
2471 let stderr_str = String::from_utf8_lossy(&stderr).to_string();
2472 let exit_code = status
2473 .code
2474 .and_then(|code| i32::try_from(code).ok())
2475 .unwrap_or(-1);
2476
2477 // Check if sandbox denied the operation
2478 let sandbox_denied = SandboxManager::was_denied(sandbox_type, exit_code, &stderr_str);
2479 let (stdout, stdout_meta) = truncate_with_meta(&stdout_str);
2480 let (stderr, stderr_meta) = truncate_with_meta(&stderr_str);
2481
2482 Ok(ShellResult {
2483 task_id: None,
2484 status: if status.success {
2485 ShellStatus::Completed
2486 } else {
2487 ShellStatus::Failed
2488 },
2489 exit_code: status.code,
2490 stdout,
2491 stderr,
2492 duration_ms: u64::try_from(started.elapsed().as_millis()).unwrap_or(u64::MAX),
2493 stdout_len: stdout_meta.original_len,
2494 stderr_len: stderr_meta.original_len,
2495 stdout_omitted: stdout_meta.omitted,
2496 stderr_omitted: stderr_meta.omitted,
2497 stdout_truncated: stdout_meta.truncated,
2498 stderr_truncated: stderr_meta.truncated,
2499 sandboxed,
2500 sandbox_type: if sandboxed {
2501 Some(sandbox_type.to_string())
2502 } else {
2503 None
2504 },
2505 sandbox_denied,
2506 })
2507 } else {
2508 // Timeout - kill the process
2509 #[cfg(unix)]
2510 let _ = kill_child_process_group(&mut child);
2511 #[cfg(windows)]
2512 let _ = terminate_child_and_close_windows_job(windows_job, &mut child);
2513 #[cfg(all(not(unix), not(windows)))]
2514 let _ = child.kill();
2515 let status = child.wait().ok();
2516 let stdout = recv_sync_reader_output(&stdout_rx);
2517 let stderr = recv_sync_reader_output(&stderr_rx);
2518 let stdout_str = String::from_utf8_lossy(&stdout).to_string();
2519 let stderr_str = String::from_utf8_lossy(&stderr).to_string();
2520 let (stdout, stdout_meta) = truncate_with_meta(&stdout_str);
2521 let (stderr, stderr_meta) = truncate_with_meta(&stderr_str);
2522
2523 Ok(ShellResult {
2524 task_id: None,
2525 status: ShellStatus::TimedOut,
2526 exit_code: status
2527 .map(ShellExitStatus::from_std)
2528 .and_then(|status| status.code),
2529 stdout,
2530 stderr,
2531 duration_ms: u64::try_from(started.elapsed().as_millis()).unwrap_or(u64::MAX),
2532 stdout_len: stdout_meta.original_len,
2533 stderr_len: stderr_meta.original_len,
2534 stdout_omitted: stdout_meta.omitted,
2535 stderr_omitted: stderr_meta.omitted,
2536 stdout_truncated: stdout_meta.truncated,
2537 stderr_truncated: stderr_meta.truncated,
2538 sandboxed,
2539 sandbox_type: if sandboxed {
2540 Some(sandbox_type.to_string())
2541 } else {
2542 None
2543 },
2544 sandbox_denied: false,
2545 })
2546 }
2547 }
2548
2549 /// Execute command interactively with timeout (sandboxed).
2550 fn execute_interactive_sandboxed(
2551 original_command: &str,
2552 working_dir: &std::path::Path,
2553 timeout_ms: u64,
2554 exec_env: &ExecEnv,
2555 ) -> Result<ShellResult> {
2556 let started = Instant::now();
2557 let timeout = Duration::from_millis(timeout_ms);
2558 let sandbox_type = exec_env.sandbox_type;
2559 let sandboxed = exec_env.is_sandboxed();
2560
2561 let program = exec_env.program();
2562 let args = exec_env.args();
2563
2564 let mut cmd = Command::new(program);
2565 crate::utils::suppress_console_window(&mut cmd);
2566 push_shell_args(&mut cmd, program, args);
2567 cmd.current_dir(working_dir)
2568 .stdin(Stdio::inherit())
2569 .stdout(Stdio::inherit())
2570 .stderr(Stdio::inherit());
2571 #[cfg(unix)]
2572 {
2573 cmd.process_group(0);
2574 }
2575 install_parent_death_signal(&mut cmd);
2576
2577 // Leave raw mode before spawn; the guard restores it only if raw mode
2578 // was active on entry (issue #1690).
2579 let _raw_mode = crate::host_terminal::suspend_raw_mode();
2580
2581 child_env::apply_to_command(&mut cmd, child_env::string_map_env(&exec_env.env));
2582
2583 let mut child = cmd
2584 .spawn()
2585 .with_context(|| format!("Failed to execute: {original_command}"))?;
2586 #[cfg(windows)]
2587 let windows_job = attach_windows_job(&child, original_command);
2588
2589 if let Some(status) = child.wait_timeout(timeout)? {
2590 let status = ShellExitStatus::from_std(status);
2591 #[cfg(windows)]
2592 terminate_and_close_windows_job(windows_job);
2593 Ok(ShellResult {
2594 task_id: None,
2595 status: if status.success {
2596 ShellStatus::Completed
2597 } else {
2598 ShellStatus::Failed
2599 },
2600 exit_code: status.code,
2601 stdout: String::new(),
2602 stderr: String::new(),
2603 duration_ms: u64::try_from(started.elapsed().as_millis()).unwrap_or(u64::MAX),
2604 stdout_len: 0,
2605 stderr_len: 0,
2606 stdout_omitted: 0,
2607 stderr_omitted: 0,
2608 stdout_truncated: false,
2609 stderr_truncated: false,
2610 sandboxed,
2611 sandbox_type: if sandboxed {
2612 Some(sandbox_type.to_string())
2613 } else {
2614 None
2615 },
2616 sandbox_denied: false,
2617 })
2618 } else {
2619 #[cfg(unix)]
2620 let _ = kill_child_process_group(&mut child);
2621 #[cfg(windows)]
2622 let _ = terminate_child_and_close_windows_job(windows_job, &mut child);
2623 #[cfg(all(not(unix), not(windows)))]
2624 let _ = child.kill();
2625 let status = child.wait().ok();
2626
2627 Ok(ShellResult {
2628 task_id: None,
2629 status: ShellStatus::TimedOut,
2630 exit_code: status
2631 .map(ShellExitStatus::from_std)
2632 .and_then(|status| status.code),
2633 stdout: String::new(),
2634 stderr: String::new(),
2635 duration_ms: u64::try_from(started.elapsed().as_millis()).unwrap_or(u64::MAX),
2636 stdout_len: 0,
2637 stderr_len: 0,
2638 stdout_omitted: 0,
2639 stderr_omitted: 0,
2640 stdout_truncated: false,
2641 stderr_truncated: false,
2642 sandboxed,
2643 sandbox_type: if sandboxed {
2644 Some(sandbox_type.to_string())
2645 } else {
2646 None
2647 },
2648 sandbox_denied: false,
2649 })
2650 }
2651 }
2652
2653 /// Spawn a background process (sandboxed).
2654 #[allow(clippy::too_many_arguments)]
2655 fn spawn_background_sandboxed(
2656 &mut self,
2657 original_command: &str,
2658 working_dir: &std::path::Path,
2659 exec_env: &ExecEnv,
2660 heavy_permit: Option<HeavyCommandPermit>,
2661 stdin_data: Option<&str>,
2662 tty: bool,
2663 spawn_context: ShellSpawnContext,
2664 persist_pending: bool,
2665 small_contract_mode: bool,
2666 ) -> Result<ShellResult> {
2667 let ShellSpawnContext {
2668 owner_agent,
2669 owner_session_id,
2670 origin_tool_call_id,
2671 origin_turn_id,
2672 work_lifecycle,
2673 } = spawn_context;
2674 let task_id = format!("shell_{}", &Uuid::new_v4().to_string()[..8]);
2675 let mut spawn_guard =
2676 ShellSpawnIntentGuard::new(work_lifecycle, &task_id, original_command);
2677 // The guard owns the registration outcome: when bookkeeping could not
2678 // bind this operation, nothing below may publish to it (#6435).
2679 let work_lifecycle = spawn_guard.lifecycle.clone();
2680 let started = Instant::now();
2681 let sandbox_type = exec_env.sandbox_type;
2682 let sandboxed = exec_env.is_sandboxed();
2683
2684 // Build the command from ExecEnv
2685 let program = exec_env.program();
2686 let args = exec_env.args();
2687
2688 #[cfg(target_env = "ohos")]
2689 if tty {
2690 return Err(anyhow!(
2691 "TTY shell mode is not supported on HarmonyOS/OpenHarmony yet."
2692 ));
2693 }
2694
2695 let stdout_buffer = new_shared_raw_output();
2696 let stderr_buffer = if tty || persist_pending || small_contract_mode {
2697 None
2698 } else {
2699 Some(new_shared_raw_output())
2700 };
2701 // The spill file is best-effort: a full disk or exhausted descriptor
2702 // table must not make `echo ok` unrunnable (that is exactly how the
2703 // owner's session got wedged under swap exhaustion).
2704 let bounded_output = small_contract_mode.then(|| {
2705 Arc::new(Mutex::new(BoundedOutputAccumulator::new_in(
2706 self.output_spill_dir.as_deref(),
2707 )))
2708 });
2709
2710 #[cfg(windows)]
2711 let mut windows_job = None;
2712 #[cfg(unix)]
2713 let mut parent_death_watch = None;
2714
2715 #[cfg(not(target_env = "ohos"))]
2716 let mut pty_master = None;
2717 let (child, stdin, stdout_thread, stderr_thread) = if tty {
2718 #[cfg(target_env = "ohos")]
2719 unreachable!("OHOS TTY mode returns before PTY setup");
2720
2721 #[cfg(not(target_env = "ohos"))]
2722 {
2723 let pty_system = native_pty_system();
2724 let pair = pty_system
2725 .openpty(PtySize {
2726 rows: 24,
2727 cols: 80,
2728 pixel_width: 0,
2729 pixel_height: 0,
2730 })
2731 .context("Failed to open PTY")?;
2732
2733 let mut cmd = CommandBuilder::new(program);
2734 for arg in args {
2735 cmd.arg(arg);
2736 }
2737 cmd.cwd(working_dir);
2738 child_env::apply_to_pty_command(&mut cmd, child_env::string_map_env(&exec_env.env));
2739
2740 let mut child = pair
2741 .slave
2742 .spawn_command(cmd)
2743 .with_context(|| format!("Failed to spawn PTY command: {original_command}"))?;
2744 drop(pair.slave);
2745
2746 let reader = match pair.master.try_clone_reader() {
2747 Ok(reader) => reader,
2748 Err(err) => {
2749 let _ = child.kill();
2750 let _ = child.wait();
2751 return Err(err).context("Failed to clone PTY reader");
2752 }
2753 };
2754 let writer = match pair.master.take_writer() {
2755 Ok(writer) => writer,
2756 Err(err) => {
2757 let _ = child.kill();
2758 let _ = child.wait();
2759 return Err(err).context("Failed to take PTY writer");
2760 }
2761 };
2762 let stdout_thread = Some(spawn_reader_thread(reader, Arc::clone(&stdout_buffer)));
2763 pty_master = Some(pair.master);
2764
2765 (
2766 ShellChild::Pty(child),
2767 Some(StdinWriter::Pty(writer)),
2768 stdout_thread,
2769 None,
2770 )
2771 }
2772 } else if persist_pending {
2773 let mut cmd = Command::new(program);
2774 crate::utils::suppress_console_window(&mut cmd);
2775 push_shell_args(&mut cmd, program, args);
2776 cmd.current_dir(working_dir)
2777 .stdin(Stdio::null())
2778 .stdout(Stdio::null())
2779 .stderr(Stdio::null());
2780 #[cfg(unix)]
2781 {
2782 cmd.process_group(0);
2783 }
2784 // Deliberately no parent-death cleanup: a staged service can be
2785 // handed to the user (`ShellOwnership::Released`) and must then
2786 // outlive the TUI (#6654).
2787
2788 child_env::apply_to_command(&mut cmd, child_env::string_map_env(&exec_env.env));
2789 remove_readonly_redirect_env(&mut cmd, &exec_env.env);
2790
2791 let child = cmd.spawn().with_context(|| {
2792 format!("Failed to spawn persistent service: {original_command}")
2793 })?;
2794 (ShellChild::Process(child), None, None, None)
2795 } else {
2796 let mut cmd = Command::new(program);
2797 crate::utils::suppress_console_window(&mut cmd);
2798 push_shell_args(&mut cmd, program, args);
2799 cmd.current_dir(working_dir).stdin(Stdio::piped());
2800 let combined_reader = if small_contract_mode {
2801 let (reader, stdout, stderr) =
2802 shared_output_pipe().context("Failed to create combined shell output pipe")?;
2803 cmd.stdout(Stdio::from(stdout)).stderr(Stdio::from(stderr));
2804 Some(reader)
2805 } else {
2806 cmd.stdout(Stdio::piped()).stderr(Stdio::piped());
2807 None
2808 };
2809 #[cfg(unix)]
2810 {
2811 cmd.process_group(0);
2812 }
2813
2814 child_env::apply_to_command(&mut cmd, child_env::string_map_env(&exec_env.env));
2815 remove_readonly_redirect_env(&mut cmd, &exec_env.env);
2816
2817 let mut child = cmd
2818 .spawn()
2819 .with_context(|| format!("Failed to spawn background: {original_command}"))?;
2820 // Managed children die with the TUI (#6654); unlike the
2821 // persistent branch above, nobody can take ownership of them.
2822 #[cfg(unix)]
2823 {
2824 parent_death_watch = watch_process_group_for_parent_death(child.id())
2825 .inspect_err(|err| {
2826 tracing::warn!(
2827 "background shell {task_id} has no parent-death cleanup: {err}"
2828 );
2829 })
2830 .ok();
2831 }
2832 #[cfg(windows)]
2833 {
2834 windows_job = attach_windows_job(&child, original_command);
2835 }
2836
2837 let stdin_handle = child.stdin.take().map(StdinWriter::Pipe);
2838
2839 let (stdout_thread, stderr_thread) =
2840 if let (Some(reader), Some(output)) = (combined_reader, bounded_output.as_ref()) {
2841 (
2842 Some(spawn_bounded_reader_thread(reader, Arc::clone(output))),
2843 None,
2844 )
2845 } else {
2846 let stdout_handle = child.stdout.take().ok_or_else(|| {
2847 #[cfg(windows)]
2848 terminate_unregistered_process(&mut child, windows_job.as_ref());
2849 #[cfg(not(windows))]
2850 terminate_unregistered_process(&mut child);
2851 anyhow!("Failed to capture stdout")
2852 })?;
2853 let stderr_handle = child.stderr.take().ok_or_else(|| {
2854 #[cfg(windows)]
2855 terminate_unregistered_process(&mut child, windows_job.as_ref());
2856 #[cfg(not(windows))]
2857 terminate_unregistered_process(&mut child);
2858 anyhow!("Failed to capture stderr")
2859 })?;
2860 (
2861 Some(spawn_reader_thread(
2862 stdout_handle,
2863 Arc::clone(&stdout_buffer),
2864 )),
2865 stderr_buffer
2866 .as_ref()
2867 .map(|buffer| spawn_reader_thread(stderr_handle, Arc::clone(buffer))),
2868 )
2869 };
2870
2871 (
2872 ShellChild::Process(child),
2873 stdin_handle,
2874 stdout_thread,
2875 stderr_thread,
2876 )
2877 };
2878
2879 let mut bg_shell = BackgroundShell {
2880 id: task_id.clone(),
2881 command: original_command.to_string(),
2882 working_dir: working_dir.to_path_buf(),
2883 status: ShellStatus::Running,
2884 background: true,
2885 exit_code: None,
2886 started_at: started,
2887 finished_at: None,
2888 finished_at_utc: None,
2889 last_output_at: started,
2890 last_observed_output_len: 0,
2891 sandbox_type,
2892 linked_task_id: None,
2893 owner_agent,
2894 owner_session_id,
2895 origin_tool_call_id,
2896 origin_turn_id,
2897 ownership: if persist_pending {
2898 ShellOwnership::PersistPending
2899 } else {
2900 ShellOwnership::Managed
2901 },
2902 stdout_buffer,
2903 stderr_buffer,
2904 bounded_output,
2905 heavy_permit,
2906 stdout_cursor: 0,
2907 stderr_cursor: 0,
2908 completion_reported: false,
2909 stdin,
2910 #[cfg(not(target_env = "ohos"))]
2911 pty_master,
2912 terminal_size: tty.then(PtyDimensions::default),
2913 child: Some(child),
2914 #[cfg(windows)]
2915 windows_job,
2916 #[cfg(unix)]
2917 parent_death_watch,
2918 stdout_thread,
2919 stderr_thread,
2920 work_lifecycle,
2921 lifecycle_seq: 0,
2922 last_lifecycle_status: None,
2923 last_lifecycle_bytes: 0,
2924 wait_failed: false,
2925 };
2926
2927 #[cfg(unix)]
2928 if persist_pending {
2929 let process_group_id = bg_shell
2930 .child
2931 .as_ref()
2932 .and_then(ShellChild::process_id)
2933 .ok_or_else(|| anyhow!("Persistent service has no process group id"))?;
2934 register_pending_persistent_process_group(process_group_id);
2935 }
2936
2937 if let Some(input) = stdin_data
2938 && let Err(err) = bg_shell.write_stdin(input, false)
2939 {
2940 let _ = bg_shell.kill();
2941 return Err(err);
2942 }
2943
2944 // Work-graph publication is bookkeeping: a graph write that fails
2945 // now must not kill a command that already started (#6435).
2946 bg_shell.publish_lifecycle_best_effort();
2947
2948 self.processes.insert(task_id.clone(), bg_shell);
2949 spawn_guard.disarm();
2950 // Evict here, not only from `list_jobs()`: retention must not depend on
2951 // the user opening the jobs panel, and every spawn is exactly the moment
2952 // the previous calls' records became one call staler (#5472).
2953 self.cleanup(FINISHED_SHELL_MAX_AGE);
2954
2955 Ok(ShellResult {
2956 task_id: Some(task_id),
2957 status: ShellStatus::Running,
2958 exit_code: None,
2959 stdout: String::new(),
2960 stderr: String::new(),
2961 duration_ms: 0,
2962 stdout_len: 0,
2963 stderr_len: 0,
2964 stdout_omitted: 0,
2965 stderr_omitted: 0,
2966 stdout_truncated: false,
2967 stderr_truncated: false,
2968 sandboxed,
2969 sandbox_type: if sandboxed {
2970 Some(sandbox_type.to_string())
2971 } else {
2972 None
2973 },
2974 sandbox_denied: false,
2975 })
2976 }
2977
2978 /// Get output from a background process
2979 pub fn get_output(
2980 &mut self,
2981 task_id: &str,
2982 block: bool,
2983 timeout_ms: u64,
2984 ) -> Result<ShellResult> {
2985 let shell = self
2986 .processes
2987 .get_mut(task_id)
2988 .ok_or_else(|| anyhow!("Task {task_id} not found"))?;
2989
2990 if block && shell.status == ShellStatus::Running {
2991 let timeout = Duration::from_millis(timeout_ms.clamp(1000, 600_000));
2992 let deadline = Instant::now() + timeout;
2993
2994 while shell.status == ShellStatus::Running && Instant::now() < deadline {
2995 if shell.poll() {
2996 break;
2997 }
2998 std::thread::sleep(Duration::from_millis(100));
2999 }
3000
3001 // If still running after timeout
3002 if shell.status == ShellStatus::Running {
3003 return shell.snapshot();
3004 }
3005 } else {
3006 shell.poll();
3007 }
3008
3009 shell.snapshot()
3010 }
3011
3012 /// Poll a job and return only its status.
3013 ///
3014 /// The foreground wait loop ticks every 100 ms and discards the snapshot
3015 /// unless the job is terminal, but `get_output` → `snapshot` clones both
3016 /// raw buffers to build it. On a command printing 50 MB that was ~1.5 GB of
3017 /// allocate-and-drop churn per wait, all of it thrown away (#5472 finding 1,
3018 /// the transient term). Status is what the loop actually needs.
3019 fn poll_status(&mut self, task_id: &str) -> Result<ShellStatus> {
3020 let shell = self
3021 .processes
3022 .get_mut(task_id)
3023 .ok_or_else(|| anyhow!("Task {task_id} not found"))?;
3024 shell.poll();
3025 Ok(shell.status.clone())
3026 }
3027
3028 /// Write data to stdin of a background process.
3029 pub fn write_stdin(&mut self, task_id: &str, input: &str, close: bool) -> Result<()> {
3030 self.write_stdin_bytes(task_id, input.as_bytes(), close)
3031 }
3032
3033 /// Exact input bytes from an authenticated client; never round-trip through UTF-8.
3034 pub fn write_stdin_bytes(&mut self, task_id: &str, input: &[u8], close: bool) -> Result<()> {
3035 let shell = self
3036 .processes
3037 .get_mut(task_id)
3038 .ok_or_else(|| anyhow!("Task {task_id} not found"))?;
3039 shell.write_stdin_bytes(input, close)?;
3040 Ok(())
3041 }
3042
3043 /// Historical evicted records have no terminal-size evidence.
3044 pub fn job_terminal_size(&self, task_id: &str) -> Option<PtyDimensions> {
3045 self.processes
3046 .get(task_id)
3047 .and_then(|shell| shell.terminal_size)
3048 }
3049
3050 pub fn resize_pty(&mut self, task_id: &str, size: PtyDimensions) -> Result<()> {
3051 let size = size.validate()?;
3052 let shell = self
3053 .processes
3054 .get_mut(task_id)
3055 .ok_or_else(|| anyhow!("Job {task_id} not found"))?;
3056 shell.poll();
3057 anyhow::ensure!(
3058 shell.status == ShellStatus::Running,
3059 "PTY job is no longer running"
3060 );
3061 #[cfg(not(target_env = "ohos"))]
3062 {
3063 let master = shell.pty_master.as_ref().context("Job is not a PTY")?;
3064 master
3065 .resize(PtySize {
3066 rows: size.rows,
3067 cols: size.cols,
3068 pixel_width: 0,
3069 pixel_height: 0,
3070 })
3071 .context("Failed to resize PTY")?;
3072 shell.terminal_size = Some(size);
3073 Ok(())
3074 }
3075 #[cfg(target_env = "ohos")]
3076 {
3077 let _ = size;
3078 Err(anyhow!("PTY resize is unavailable on this platform"))
3079 }
3080 }
3081
3082 pub fn write_stdin_for_session(
3083 &mut self,
3084 active_session_id: &str,
3085 task_id: &str,
3086 input: &str,
3087 close: bool,
3088 ) -> Result<()> {
3089 self.require_session_owner(task_id, active_session_id)?;
3090 self.write_stdin(task_id, input, close)
3091 }
3092
3093 /// Get incremental output from a background process, consuming any new output.
3094 fn get_output_delta(
3095 &mut self,
3096 task_id: &str,
3097 wait: bool,
3098 timeout_ms: u64,
3099 ) -> Result<ShellDeltaResult> {
3100 let shell = self
3101 .processes
3102 .get_mut(task_id)
3103 .ok_or_else(|| anyhow!("Task {task_id} not found"))?;
3104
3105 if wait && shell.status == ShellStatus::Running {
3106 let timeout = Duration::from_millis(timeout_ms.clamp(1000, 600_000));
3107 let deadline = Instant::now() + timeout;
3108
3109 while shell.status == ShellStatus::Running && Instant::now() < deadline {
3110 if shell.poll() {
3111 break;
3112 }
3113 std::thread::sleep(Duration::from_millis(100));
3114 }
3115 } else {
3116 shell.poll();
3117 }
3118
3119 let (
3120 stdout_delta,
3121 stderr_delta,
3122 stdout_delta_len,
3123 stderr_delta_len,
3124 stdout_total,
3125 stderr_total,
3126 ) = shell.take_delta();
3127 let (stdout, stdout_meta) = truncate_with_meta(&stdout_delta);
3128 let (stderr, stderr_meta) = truncate_with_meta(&stderr_delta);
3129 let sandboxed = !matches!(shell.sandbox_type, SandboxType::None);
3130
3131 let command = shell.command.clone();
3132 let result = ShellResult {
3133 task_id: Some(shell.id.clone()),
3134 status: shell.status.clone(),
3135 exit_code: shell.exit_code,
3136 stdout,
3137 stderr,
3138 duration_ms: u64::try_from(shell.started_at.elapsed().as_millis()).unwrap_or(u64::MAX),
3139 stdout_len: stdout_meta.original_len.max(stdout_delta_len),
3140 stderr_len: stderr_meta.original_len.max(stderr_delta_len),
3141 stdout_omitted: stdout_meta.omitted,
3142 stderr_omitted: stderr_meta.omitted,
3143 stdout_truncated: stdout_meta.truncated,
3144 stderr_truncated: stderr_meta.truncated,
3145 sandboxed,
3146 sandbox_type: if sandboxed {
3147 Some(shell.sandbox_type.to_string())
3148 } else {
3149 None
3150 },
3151 sandbox_denied: shell.sandbox_denied(),
3152 };
3153
3154 Ok(ShellDeltaResult {
3155 command,
3156 result,
3157 stdout_total_len: stdout_total,
3158 stderr_total_len: stderr_total,
3159 })
3160 }
3161
3162 fn attach_heavy_permit(&mut self, task_id: &str, permit: HeavyCommandPermit) -> Result<()> {
3163 let shell = self
3164 .processes
3165 .get_mut(task_id)
3166 .ok_or_else(|| anyhow!("Task {task_id} not found"))?;
3167 shell.heavy_permit = Some(permit);
3168 Ok(())
3169 }
3170
3171 /// Kill a running background process
3172 pub fn kill(&mut self, task_id: &str) -> Result<ShellResult> {
3173 let shell = self
3174 .processes
3175 .get_mut(task_id)
3176 .ok_or_else(|| anyhow!("Task {task_id} not found"))?;
3177
3178 shell.kill()?;
3179 shell.snapshot()
3180 }
3181
3182 pub fn kill_for_session(
3183 &mut self,
3184 active_session_id: &str,
3185 task_id: &str,
3186 ) -> Result<ShellResult> {
3187 self.require_session_owner(task_id, active_session_id)?;
3188 self.kill(task_id)
3189 }
3190
3191 /// Kill every currently running background shell process.
3192 #[cfg(test)]
3193 pub fn kill_running(&mut self) -> Result<Vec<ShellResult>> {
3194 let ids = self
3195 .processes
3196 .iter()
3197 .filter(|(_, shell)| shell.status == ShellStatus::Running)
3198 .map(|(id, _)| id.clone())
3199 .collect::<Vec<_>>();
3200
3201 let mut results = Vec::with_capacity(ids.len());
3202 for id in ids {
3203 results.push(self.kill(&id)?);
3204 }
3205 Ok(results)
3206 }
3207
3208 pub fn kill_running_for_session(
3209 &mut self,
3210 active_session_id: &str,
3211 ) -> Result<Vec<ShellResult>> {
3212 let ids = self
3213 .processes
3214 .iter()
3215 .filter(|(_, shell)| {
3216 shell.status == ShellStatus::Running
3217 && !active_session_id.is_empty()
3218 && shell.owner_session_id == active_session_id
3219 })
3220 .map(|(id, _)| id.clone())
3221 .collect::<Vec<_>>();
3222 let mut results = Vec::with_capacity(ids.len());
3223 for id in ids {
3224 results.push(self.kill(&id)?);
3225 }
3226 Ok(results)
3227 }
3228
3229 /// Transfer every still-running `persist:true` process out of Codewhale's
3230 /// ownership. This is called only by the real headless exec host after the
3231 /// enclosing turn has completed successfully.
3232 #[cfg(unix)]
3233 pub fn commit_persistent_services(&mut self) -> Result<Vec<PersistentServiceReceipt>> {
3234 let mut ids = self
3235 .processes
3236 .iter()
3237 .filter(|(_, shell)| shell.ownership == ShellOwnership::PersistPending)
3238 .map(|(id, _)| id.clone())
3239 .collect::<Vec<_>>();
3240 ids.sort();
3241
3242 for id in &ids {
3243 let shell = self
3244 .processes
3245 .get_mut(id)
3246 .ok_or_else(|| anyhow!("Persistent service {id} disappeared before commit"))?;
3247 shell.poll();
3248 if shell.status != ShellStatus::Running {
3249 return Err(anyhow!(
3250 "Persistent service {id} exited before ownership transfer (status {:?}, exit code {:?})",
3251 shell.status,
3252 shell.exit_code
3253 ));
3254 }
3255 if shell
3256 .child
3257 .as_ref()
3258 .and_then(ShellChild::process_id)
3259 .is_none()
3260 {
3261 return Err(anyhow!(
3262 "Persistent service {id} has no releasable process id"
3263 ));
3264 }
3265 }
3266
3267 let mut receipts = Vec::with_capacity(ids.len());
3268 for id in ids {
3269 let mut shell = self
3270 .processes
3271 .remove(&id)
3272 .ok_or_else(|| anyhow!("Persistent service {id} disappeared during commit"))?;
3273 let pid = shell
3274 .child
3275 .as_ref()
3276 .and_then(ShellChild::process_id)
3277 .ok_or_else(|| anyhow!("Persistent service {id} lost its process id"))?;
3278 unregister_pending_persistent_process_group(pid);
3279 shell.ownership = ShellOwnership::Released;
3280 shell.stdin = None;
3281 shell.heavy_permit.take();
3282 shell.work_lifecycle = None;
3283 receipts.push(PersistentServiceReceipt {
3284 task_id: id,
3285 pid,
3286 process_group_id: pid,
3287 ownership: "external".to_string(),
3288 });
3289 }
3290 Ok(receipts)
3291 }
3292
3293 /// Kill only services waiting for a successful exec ownership transfer.
3294 /// Ordinary background jobs retain their existing manager lifetime.
3295 pub fn abort_persistent_services(&mut self) {
3296 let ids = self
3297 .processes
3298 .iter()
3299 .filter(|(_, shell)| shell.ownership == ShellOwnership::PersistPending)
3300 .map(|(id, _)| id.clone())
3301 .collect::<Vec<_>>();
3302 for id in ids {
3303 if let Err(error) = self.kill(&id) {
3304 tracing::warn!(shell_id = %id, %error, "failed to abort pending persistent service");
3305 }
3306 }
3307 }
3308
3309 /// Poll a background process and return incremental output.
3310 #[cfg(test)]
3311 pub fn poll_delta(
3312 &mut self,
3313 task_id: &str,
3314 wait: bool,
3315 timeout_ms: u64,
3316 ) -> Result<ShellDeltaResult> {
3317 self.get_output_delta(task_id, wait, timeout_ms)
3318 }
3319
3320 pub fn poll_delta_for_session(
3321 &mut self,
3322 active_session_id: &str,
3323 task_id: &str,
3324 wait: bool,
3325 timeout_ms: u64,
3326 ) -> Result<ShellDeltaResult> {
3327 self.require_session_owner(task_id, active_session_id)?;
3328 self.get_output_delta(task_id, wait, timeout_ms)
3329 }
3330
3331 fn get_output_delta_for_session(
3332 &mut self,
3333 active_session_id: &str,
3334 task_id: &str,
3335 wait: bool,
3336 timeout_ms: u64,
3337 ) -> Result<ShellDeltaResult> {
3338 self.require_session_owner(task_id, active_session_id)?;
3339 self.get_output_delta(task_id, wait, timeout_ms)
3340 }
3341
3342 /// Read a job's raw stream at an absolute byte offset without consuming
3343 /// anything. This is the `/v1/jobs` byte-stream contract: HTTP clients hold
3344 /// the cursor, so reads must not disturb the engine's own delta consumer.
3345 ///
3346 /// `cursor` is a byte offset into the stream's lifetime output (matching
3347 /// `total`). When the bounded buffer has already discarded `[0, dropped)`,
3348 /// the window starts at `dropped` instead and the caller sees the gap in
3349 /// the response rather than a replayed tail. With `wait_ms > 0` on a
3350 /// running job, polls up to that bound for new bytes past `cursor` before
3351 /// answering — long-poll instead of a hot loop.
3352 pub fn read_output_chunk(
3353 &mut self,
3354 task_id: &str,
3355 stream: ShellOutputStream,
3356 cursor: usize,
3357 max_bytes: usize,
3358 wait_ms: u64,
3359 ) -> Result<ShellOutputChunk> {
3360 let Some(shell) = self.processes.get_mut(task_id) else {
3361 // Evicted jobs retain only their snapshot tails. Serve that tail as
3362 // the final retained window so a late reader still gets the ending
3363 // of the stream instead of a bare not-found.
3364 let snapshot = self
3365 .stale_jobs
3366 .get(task_id)
3367 .ok_or_else(|| anyhow!("Job {task_id} not found"))?;
3368 let (tail, total) = match stream {
3369 ShellOutputStream::Stdout => (&snapshot.stdout_tail, snapshot.stdout_len),
3370 ShellOutputStream::Stderr => (&snapshot.stderr_tail, snapshot.stderr_len),
3371 };
3372 let tail_start = total.saturating_sub(tail.len());
3373 let offset = cursor.max(tail_start).min(total);
3374 let next_offset = offset.saturating_add(max_bytes).min(total);
3375 return Ok(ShellOutputChunk {
3376 offset,
3377 bytes: tail.as_bytes()[offset - tail_start..next_offset - tail_start].to_vec(),
3378 next_offset,
3379 total,
3380 dropped: tail_start,
3381 status: snapshot.status.clone(),
3382 exit_code: snapshot.exit_code,
3383 });
3384 };
3385 let buffer = match stream {
3386 ShellOutputStream::Stdout => shell.stdout_buffer.clone(),
3387 ShellOutputStream::Stderr => shell
3388 .stderr_buffer
3389 .clone()
3390 .ok_or_else(|| anyhow!("Job {task_id} merges stderr into stdout"))?,
3391 };
3392
3393 let wait_deadline = (wait_ms > 0 && shell.status == ShellStatus::Running)
3394 .then(|| Instant::now() + Duration::from_millis(wait_ms.clamp(50, 30_000)));
3395 loop {
3396 shell.poll();
3397 let total = buffer.lock().map(|guard| guard.total_len()).unwrap_or(0);
3398 let done_waiting = total > cursor
3399 || shell.status != ShellStatus::Running
3400 || wait_deadline.is_none_or(|deadline| Instant::now() >= deadline);
3401 if done_waiting {
3402 break;
3403 }
3404 std::thread::sleep(Duration::from_millis(50));
3405 }
3406
3407 let (bytes, offset, next_offset, total, dropped) = {
3408 let guard = buffer.lock().unwrap_or_else(|e| e.into_inner());
3409 let total = guard.total_len();
3410 let dropped = guard.dropped();
3411 let offset = cursor.max(dropped).min(total);
3412 let next_offset = offset.saturating_add(max_bytes).min(total);
3413 let retained = guard.retained();
3414 let bytes = retained[offset - dropped..next_offset - dropped].to_vec();
3415 (bytes, offset, next_offset, total, dropped)
3416 };
3417 Ok(ShellOutputChunk {
3418 offset,
3419 bytes,
3420 next_offset,
3421 total,
3422 dropped,
3423 status: shell.status.clone(),
3424 exit_code: shell.exit_code,
3425 })
3426 }
3427
3428 /// Attach durable task context to a live shell job.
3429 pub fn tag_linked_task(&mut self, task_id: &str, linked_task_id: Option<String>) -> Result<()> {
3430 let shell = self
3431 .processes
3432 .get_mut(task_id)
3433 .ok_or_else(|| anyhow!("Task {task_id} not found"))?;
3434 shell.linked_task_id = linked_task_id;
3435 Ok(())
3436 }
3437
3438 /// Inspect full output for a live or stale job.
3439 pub fn inspect_job(&mut self, task_id: &str) -> Result<ShellJobDetail> {
3440 if let Some(shell) = self.processes.get_mut(task_id) {
3441 shell.poll();
3442 return Ok(shell.job_detail());
3443 }
3444 if let Some(snapshot) = self.stale_jobs.get(task_id) {
3445 return Ok(ShellJobDetail {
3446 snapshot: snapshot.clone(),
3447 stdout: snapshot.stdout_tail.clone(),
3448 stderr: snapshot.stderr_tail.clone(),
3449 });
3450 }
3451 Err(anyhow!("Job {task_id} not found"))
3452 }
3453
3454 pub fn inspect_job_for_session(
3455 &mut self,
3456 active_session_id: &str,
3457 task_id: &str,
3458 ) -> Result<ShellJobDetail> {
3459 self.require_session_owner(task_id, active_session_id)?;
3460 self.inspect_job(task_id)
3461 }
3462
3463 /// List all live and known-stale background shell jobs for the TUI.
3464 pub fn list_jobs(&mut self) -> Vec<ShellJobSnapshot> {
3465 for shell in self.processes.values_mut() {
3466 shell.poll();
3467 }
3468 // Evict completed processes older than 1 hour to bound memory growth.
3469 self.cleanup(FINISHED_SHELL_MAX_AGE);
3470
3471 let mut jobs = self
3472 .processes
3473 .values()
3474 .map(BackgroundShell::job_snapshot)
3475 .collect::<Vec<_>>();
3476 jobs.extend(self.stale_jobs.values().cloned());
3477 jobs.sort_by(|a, b| {
3478 job_status_rank(&a.status, a.stale)
3479 .cmp(&job_status_rank(&b.status, b.stale))
3480 .then_with(|| a.id.cmp(&b.id))
3481 });
3482 jobs
3483 }
3484
3485 pub fn list_jobs_for_session(&mut self, active_session_id: &str) -> Vec<ShellJobSnapshot> {
3486 if active_session_id.is_empty() {
3487 return Vec::new();
3488 }
3489 self.list_jobs()
3490 .into_iter()
3491 .filter(|job| job.owner_session_id == active_session_id)
3492 .collect()
3493 }
3494
3495 /// Whether a finished parent-owned job's completion is waiting to be
3496 /// claimed. Unlike
3497 /// [`Self::may_have_undelivered_completion`] this polls, so it reports
3498 /// readiness the moment the process exits; the engine's idle shell wake
3499 /// uses it to fire exactly when evidence exists.
3500 #[cfg(test)]
3501 pub(crate) fn has_finished_unreported_jobs(&mut self) -> bool {
3502 self.processes.values_mut().any(|shell| {
3503 shell.poll();
3504 shell.owner_agent.is_none()
3505 && shell.status != ShellStatus::Running
3506 && !shell.completion_reported
3507 })
3508 }
3509
3510 pub(crate) fn has_finished_unreported_jobs_for_session(
3511 &mut self,
3512 active_session_id: &str,
3513 ) -> bool {
3514 !active_session_id.is_empty()
3515 && self.processes.values_mut().any(|shell| {
3516 shell.poll();
3517 shell.owner_session_id == active_session_id
3518 && shell.owner_agent.is_none()
3519 && shell.status != ShellStatus::Running
3520 && !shell.completion_reported
3521 })
3522 }
3523
3524 /// Drain once-only completion events together with lossless stream bytes.
3525 /// The engine publishes the bytes outside this manager's mutex and puts
3526 /// only the bounded event plus resulting handle into model context.
3527 #[cfg(test)]
3528 pub(crate) fn drain_finished_jobs_with_evidence(&mut self) -> Vec<ShellCompletionEvidence> {
3529 self.drain_finished_jobs_with_evidence_inner(None)
3530 }
3531
3532 pub(crate) fn drain_finished_jobs_with_evidence_for_session(
3533 &mut self,
3534 active_session_id: &str,
3535 ) -> Vec<ShellCompletionEvidence> {
3536 self.drain_finished_jobs_with_evidence_inner(Some(active_session_id))
3537 }
3538
3539 fn drain_finished_jobs_with_evidence_inner(
3540 &mut self,
3541 active_session_id: Option<&str>,
3542 ) -> Vec<ShellCompletionEvidence> {
3543 let mut completions = Vec::new();
3544 for shell in self.processes.values_mut() {
3545 shell.poll();
3546 let owned = active_session_id.is_none_or(|session_id| {
3547 !session_id.is_empty() && shell.owner_session_id == session_id
3548 });
3549 if owned && shell.status != ShellStatus::Running && !shell.completion_reported {
3550 shell.completion_reported = true;
3551 completions.push(shell.completion_evidence());
3552 // The bytes are now in the caller's hands (they become a durable
3553 // session artifact). Holding a second copy here for the rest of
3554 // the retention hour is what #5472 measured.
3555 shell.release_delivered_output();
3556 }
3557 }
3558 completions.sort_by(|a, b| a.event.task_id.cmp(&b.event.task_id));
3559 completions
3560 }
3561
3562 /// A terminal foreground result is already returned as the tool result;
3563 /// do not emit it again through the background-completion channel.
3564 fn acknowledge_foreground_completion(&mut self, task_id: &str) {
3565 if let Some(shell) = self.processes.get_mut(task_id) {
3566 shell.completion_reported = true;
3567 // The caller already holds this job's `ShellResult`; the record only
3568 // stays listed so `/jobs` can show it. A 1,200-char tail is all any
3569 // remaining consumer reads, so the rest is released now instead of
3570 // at the 1 h `cleanup` (#5472 finding 1 — the dominant term: every
3571 // uppercase `Bash` call, foreground included, went through here).
3572 shell.release_delivered_output();
3573 }
3574 }
3575
3576 /// Whether the next production turn may inject a parent-owned shell
3577 /// completion event.
3578 ///
3579 /// This deliberately does not poll processes or flip
3580 /// `completion_reported`: preview is read-only. A running job counts as
3581 /// pending because it can finish before production drains completions; in
3582 /// that race an exact request body cannot be proved without mutation.
3583 #[cfg(test)]
3584 pub fn may_have_undelivered_completion(&self) -> bool {
3585 self.processes
3586 .values()
3587 .any(|shell| shell.owner_agent.is_none() && !shell.completion_reported)
3588 }
3589
3590 pub fn may_have_undelivered_completion_for_session(&self, active_session_id: &str) -> bool {
3591 !active_session_id.is_empty()
3592 && self.processes.values().any(|shell| {
3593 shell.owner_session_id == active_session_id
3594 && shell.owner_agent.is_none()
3595 && !shell.completion_reported
3596 })
3597 }
3598
3599 /// Return agent owners whose tracked shell work is still running. The
3600 /// engine uses this to keep a worker's heartbeat alive while its only
3601 /// pending work is an explicitly tracked background shell task.
3602 #[cfg(test)]
3603 pub fn running_owner_agent_ids(&mut self) -> Vec<String> {
3604 self.running_owner_agent_ids_inner(None)
3605 }
3606
3607 pub fn running_owner_agent_ids_for_session(&mut self, active_session_id: &str) -> Vec<String> {
3608 self.running_owner_agent_ids_inner(Some(active_session_id))
3609 }
3610
3611 fn running_owner_agent_ids_inner(&mut self, active_session_id: Option<&str>) -> Vec<String> {
3612 let mut owners = self
3613 .processes
3614 .values_mut()
3615 .filter_map(|shell| {
3616 shell.poll();
3617 (shell.status == ShellStatus::Running
3618 && active_session_id.is_none_or(|session_id| {
3619 !session_id.is_empty() && shell.owner_session_id == session_id
3620 }))
3621 .then(|| {
3622 shell
3623 .owner_agent
3624 .as_ref()
3625 .map(|owner| owner.agent_id.clone())
3626 })
3627 .flatten()
3628 })
3629 .collect::<Vec<_>>();
3630 owners.sort();
3631 owners.dedup();
3632 owners
3633 }
3634
3635 /// Remember a restart-stale job so the UI can show it instead of hiding it.
3636 #[cfg(test)]
3637 pub fn remember_stale_job(
3638 &mut self,
3639 id: impl Into<String>,
3640 command: impl Into<String>,
3641 cwd: PathBuf,
3642 linked_task_id: Option<String>,
3643 ) {
3644 let id = id.into();
3645 self.stale_jobs.insert(
3646 id.clone(),
3647 ShellJobSnapshot {
3648 id: id.clone(),
3649 job_id: id,
3650 command: command.into(),
3651 cwd,
3652 status: ShellStatus::Killed,
3653 background: true,
3654 finished_at: None,
3655 exit_code: None,
3656 elapsed_ms: 0,
3657 stdout_tail: String::new(),
3658 stderr_tail: "Process is no longer attached to this TUI session.".to_string(),
3659 stdout_len: 0,
3660 stderr_len: 0,
3661 stdin_available: false,
3662 stale: true,
3663 elapsed_since_output_ms: None,
3664 linked_task_id,
3665 owner_agent_id: None,
3666 owner_agent_name: None,
3667 origin_tool_call_id: None,
3668 origin_turn_id: None,
3669 owner_session_id: String::new(),
3670 },
3671 );
3672 }
3673
3674 /// Clean up completed processes older than the given duration, then enforce
3675 /// the count and byte ceilings on what is left.
3676 ///
3677 /// Age alone is not a bound: it only fired from `list_jobs()`, so a session
3678 /// that never opened the jobs panel evicted nothing, and 500 finished
3679 /// records inside one hour were all retained regardless of size (#5472).
3680 ///
3681 /// Age counts from when a job finished, not when it started: a job that
3682 /// ran longer than `max_age` would otherwise be dropped the moment it
3683 /// exited, before its completion was delivered. Undelivered completions
3684 /// age out too: jobs launched over the Runtime API, or owned by a session
3685 /// that is no longer active, are never drained.
3686 pub fn cleanup(&mut self, max_age: Duration) {
3687 self.processes.retain(|_, shell| {
3688 if shell.status == ShellStatus::Running {
3689 return true;
3690 }
3691 shell.mark_finished();
3692 shell
3693 .finished_at
3694 .is_none_or(|finished| finished.elapsed() < max_age)
3695 });
3696 self.enforce_finished_job_bounds();
3697 }
3698
3699 /// Bytes still held across every tracked job. The retention bound in #5472
3700 /// is stated in these terms, so the tests assert on them directly.
3701 #[cfg(all(test, unix))]
3702 pub(crate) fn retained_output_bytes_total(&self) -> usize {
3703 self.processes
3704 .values()
3705 .map(BackgroundShell::retained_output_bytes)
3706 .fold(0, usize::saturating_add)
3707 }
3708
3709 #[cfg(test)]
3710 pub(crate) fn tracked_job_count(&self) -> usize {
3711 self.processes.len()
3712 }
3713
3714 /// Drop the oldest finished records once they exceed either ceiling.
3715 /// Running jobs are never evicted — killing the handle would orphan a live
3716 /// process — and a job whose completion has not been delivered yet is
3717 /// evicted last, since dropping it loses the only copy of its result.
3718 fn enforce_finished_job_bounds(&mut self) {
3719 let mut finished = self
3720 .processes
3721 .iter()
3722 .filter(|(_, shell)| shell.status != ShellStatus::Running)
3723 .map(|(id, shell)| {
3724 (
3725 id.clone(),
3726 shell.completion_reported,
3727 // Oldest by finish, like the age rule in `cleanup`.
3728 shell.finished_at.unwrap_or(shell.started_at),
3729 shell.retained_output_bytes(),
3730 )
3731 })
3732 .collect::<Vec<_>>();
3733 let total_bytes: usize = finished
3734 .iter()
3735 .map(|(_, _, _, bytes)| *bytes)
3736 .fold(0, usize::saturating_add);
3737 if finished.len() <= MAX_FINISHED_SHELL_RECORDS && total_bytes <= MAX_FINISHED_SHELL_BYTES {
3738 return;
3739 }
3740 // Undelivered completions last, then oldest first.
3741 finished.sort_by(|a, b| a.1.cmp(&b.1).reverse().then_with(|| a.2.cmp(&b.2)));
3742 let mut remaining_count = finished.len();
3743 let mut remaining_bytes = total_bytes;
3744 for (id, _, _, bytes) in finished {
3745 if remaining_count <= MAX_FINISHED_SHELL_RECORDS
3746 && remaining_bytes <= MAX_FINISHED_SHELL_BYTES
3747 {
3748 break;
3749 }
3750 self.processes.remove(&id);
3751 remaining_count -= 1;
3752 remaining_bytes = remaining_bytes.saturating_sub(bytes);
3753 }
3754 }
3755 }
3756
3757 fn job_status_rank(status: &ShellStatus, stale: bool) -> u8 {
3758 if stale {
3759 return 4;
3760 }
3761 match status {
3762 ShellStatus::Running => 0,
3763 ShellStatus::Failed | ShellStatus::TimedOut => 1,
3764 ShellStatus::Killed => 2,
3765 ShellStatus::Completed => 3,
3766 }
3767 }
3768
3769 /// Thread-safe wrapper for `ShellManager`
3770 pub type SharedShellManager = Arc<Mutex<ShellManager>>;
3771
3772 /// Create a new shared shell manager with default sandbox policy.
3773 pub fn new_shared_shell_manager(workspace: PathBuf) -> SharedShellManager {
3774 Arc::new(Mutex::new(ShellManager::new(workspace)))
3775 }
3776
3777 // === ToolSpec Implementations ===
3778
3779 use crate::features::Feature;
3780 use crate::tools::cargo_failure_summary::summarize_cargo_failure;
3781 use crate::tools::spec::{
3782 ApprovalRequirement, ToolCapability, ToolContext, ToolError, ToolResult, ToolSpec,
3783 optional_bool, optional_str, optional_u64, required_str, type_mismatch,
3784 };
3785 use async_trait::async_trait;
3786 use codewhale_execpolicy::command_safety::{
3787 NetworkRead, ReadonlyRejection, SafetyLevel, agent_readonly_verdict, analyze_command,
3788 extract_primary_command, is_parallel_readonly_command, normalize_windows_command_paths,
3789 readonly_network_reads, split_leading_cd,
3790 };
3791 use codewhale_execpolicy::toml_rules::{ExecPolicyConfig, RuleDecision};
3792 use serde_json::json;
3793
3794 /// The TOML execpolicy file lives in the user config home; the rule engine
3795 /// itself is `codewhale_execpolicy::toml_rules`.
3796 fn default_execpolicy_path() -> Option<std::path::PathBuf> {
3797 crate::config::effective_home_dir().map(|home| home.join(".deepseek").join("execpolicy.toml"))
3798 }
3799
3800 fn load_default_policy() -> anyhow::Result<Option<ExecPolicyConfig>> {
3801 /// A parsed rules file, tagged with the identity it was parsed from.
3802 type PolicyKey = (std::path::PathBuf, u64, Option<std::time::SystemTime>);
3803
3804 let Some(path) = default_execpolicy_path() else {
3805 return Ok(None);
3806 };
3807 // An unreadable or missing file (including a permissions error, which
3808 // `exists()` also swallows) means "no file rules" — the same answer as
3809 // before, just reached with one `stat` instead of an existence check plus a
3810 // full read.
3811 let Ok(metadata) = std::fs::metadata(&path) else {
3812 return Ok(None);
3813 };
3814 let key: PolicyKey = (path.clone(), metadata.len(), metadata.modified().ok());
3815
3816 // #6208: this runs on every shell execution, so the read and TOML parse
3817 // happen only when the file's identity changes. Length joins the timestamp
3818 // because a coarse-mtime filesystem can hand back the same instant for two
3819 // different revisions.
3820 static CACHE: std::sync::OnceLock<std::sync::Mutex<Option<(PolicyKey, ExecPolicyConfig)>>> =
3821 std::sync::OnceLock::new();
3822 let cache = CACHE.get_or_init(|| std::sync::Mutex::new(None));
3823 let mut cache = cache
3824 .lock()
3825 .unwrap_or_else(|poisoned| poisoned.into_inner());
3826 if let Some((cached_key, config)) = cache.as_ref()
3827 && *cached_key == key
3828 {
3829 return Ok(Some(config.clone()));
3830 }
3831
3832 let config = ExecPolicyConfig::from_path(&path)?;
3833 *cache = Some((key, config.clone()));
3834 Ok(Some(config))
3835 }
3836
3837 const FOREGROUND_TIMEOUT_RECOVERY_HINT: &str = "Foreground Bash is for bounded commands. \
3838 The timed-out process was killed; rerun long work as Bash action=\"run\" background=true, \
3839 then poll with Bash action=\"wait\" task_id=\"<id>\".";
3840
3841 const MACOS_PROVENANCE_HINT: &str = "Docker buildx failed to update its activity file due to a macOS \
3842 com.apple.provenance restriction. Files created by Docker Desktop's signed process carry a \
3843 kernel-enforced provenance tag that blocks writes from child processes (including the TUI \
3844 shell sandbox). Workarounds: (1) run the Docker build from a regular terminal outside the \
3845 TUI, or (2) disable BuildKit with DOCKER_BUILDKIT=0 (only works if your Dockerfiles do not \
3846 use RUN --mount directives).";
3847
3848 /// Human-readable exit status for a shell result: the numeric code when the
3849 /// process returned one, or "terminated by signal" when it did not (rather
3850 /// than leaking `Some(127)` / `None` Debug output to the user).
3851 fn exit_code_label(code: Option<i64>) -> String {
3852 match (code, exit_code_hex(code)) {
3853 (Some(code), Some(hex)) => format!("exit code {code} ({hex})"),
3854 (Some(code), None) => format!("exit code {code}"),
3855 (None, _) => "terminated by signal".to_string(),
3856 }
3857 }
3858
3859 fn exit_code_hex(code: Option<i64>) -> Option<String> {
3860 code.filter(|code| *code > i64::from(i32::MAX) && *code <= i64::from(u32::MAX))
3861 .map(|code| format!("0x{code:08X}"))
3862 }
3863 const PYTHON_BUILD_DEPENDENCY_HINT: &str = "Python build dependency missing: setuptools is not \
3864 available in the active environment. Install the declared build requirements first, for example \
3865 `python -m pip install -U pip setuptools wheel build`, then rerun the build command.";
3866
3867 fn attach_cargo_failure_summary(
3868 metadata: &mut serde_json::Value,
3869 command: &str,
3870 result: &ShellResult,
3871 ) {
3872 if let Some(summary) = summarize_cargo_failure(
3873 command,
3874 &result.stdout,
3875 &result.stderr,
3876 result.exit_code.and_then(|code| i32::try_from(code).ok()),
3877 ) {
3878 metadata["cargo_failure_summary"] = summary.to_metadata_value();
3879 }
3880 }
3881
3882 fn attach_python_build_dependency_hint(
3883 metadata: &mut serde_json::Value,
3884 hint: Option<&'static str>,
3885 ) {
3886 if let Some(hint) = hint {
3887 metadata["python_build_dependency_hint"] = json!({
3888 "kind": "missing_setuptools",
3889 "hint": hint,
3890 "recommended_first_step": "python -m pip install -U pip setuptools wheel build",
3891 });
3892 }
3893 }
3894
3895 pub(crate) fn looks_like_macos_provenance_failure(result: &ShellResult) -> bool {
3896 if matches!(result.status, ShellStatus::Completed) && result.exit_code == Some(0) {
3897 return false;
3898 }
3899 let combined = format!("{}\n{}", result.stdout, result.stderr).to_ascii_lowercase();
3900 combined.contains("com.apple.provenance")
3901 || combined.contains("update builder last activity")
3902 || (combined.contains("buildx/activity") && combined.contains("operation not permitted"))
3903 }
3904
3905 fn macos_provenance_hint(result: &ShellResult) -> Option<&'static str> {
3906 if looks_like_macos_provenance_failure(result) {
3907 Some(MACOS_PROVENANCE_HINT)
3908 } else {
3909 None
3910 }
3911 }
3912
3913 fn python_build_dependency_hint(command: &str, result: &ShellResult) -> Option<&'static str> {
3914 if matches!(result.status, ShellStatus::Completed) && result.exit_code == Some(0) {
3915 return None;
3916 }
3917
3918 let command = command.to_ascii_lowercase();
3919 let combined = format!("{}\n{}", result.stdout, result.stderr).to_ascii_lowercase();
3920 let mentions_missing_setuptools = [
3921 "no module named 'setuptools'",
3922 "no module named \"setuptools\"",
3923 "setuptools is not available",
3924 "cannot import 'setuptools",
3925 "cannot import \"setuptools",
3926 "missing dependencies",
3927 ]
3928 .iter()
3929 .any(|needle| combined.contains(needle))
3930 && combined.contains("setuptools");
3931 if !mentions_missing_setuptools {
3932 return None;
3933 }
3934
3935 let pythonish_command = [
3936 "python",
3937 "pip",
3938 "pytest",
3939 "tox",
3940 "nox",
3941 "cython",
3942 "setup.py",
3943 "build_ext",
3944 ]
3945 .iter()
3946 .any(|needle| command.contains(needle));
3947 let pythonish_output = [
3948 "setup.py",
3949 "pyproject.toml",
3950 "build_meta",
3951 "build_ext",
3952 "pep 517",
3953 "cython",
3954 ]
3955 .iter()
3956 .any(|needle| combined.contains(needle));
3957
3958 if pythonish_command || pythonish_output {
3959 Some(PYTHON_BUILD_DEPENDENCY_HINT)
3960 } else {
3961 None
3962 }
3963 }
3964
3965 fn command_likely_needs_network(command: &str) -> bool {
3966 let normalized = command.to_ascii_lowercase();
3967 let Some(primary) = extract_primary_command(&normalized) else {
3968 return false;
3969 };
3970 let primary = primary.rsplit(['/', '\\']).next().unwrap_or(primary);
3971
3972 match primary {
3973 "curl" | "wget" | "fetch" | "nc" | "netcat" | "ncat" | "ssh" | "scp" | "sftp" | "rsync"
3974 | "ftp" | "ping" | "traceroute" | "nslookup" | "dig" | "host" | "nmap" | "gh" | "hub" => {
3975 true
3976 }
3977 "git" => [
3978 " fetch",
3979 " pull",
3980 " clone",
3981 " ls-remote",
3982 " submodule",
3983 " push",
3984 ]
3985 .iter()
3986 .any(|needle| normalized.contains(needle)),
3987 "cargo" => [" install", " fetch", " update", " publish", " search"]
3988 .iter()
3989 .any(|needle| normalized.contains(needle)),
3990 "npm" | "pnpm" | "yarn" => [" install", " i", " add", " update", " publish"]
3991 .iter()
3992 .any(|needle| normalized.contains(needle)),
3993 "pip" | "pip3" | "uv" | "poetry" => [" install", " add", " sync", " update"]
3994 .iter()
3995 .any(|needle| normalized.contains(needle)),
3996 "brew" | "apt" | "apt-get" | "yum" | "dnf" | "pacman" => true,
3997 "go" => [" get", " install", " mod download"]
3998 .iter()
3999 .any(|needle| normalized.contains(needle)),
4000 _ => false,
4001 }
4002 }
4003
4004 fn looks_like_network_blocked_failure(result: &ShellResult) -> bool {
4005 if matches!(result.status, ShellStatus::Completed | ShellStatus::Running)
4006 || result.exit_code == Some(0)
4007 {
4008 return false;
4009 }
4010
4011 if result.stdout.trim() == "000" {
4012 return true;
4013 }
4014 if result.sandboxed && result.stdout.is_empty() && result.stderr.is_empty() {
4015 return true;
4016 }
4017
4018 let output = format!("{}\n{}", result.stdout, result.stderr).to_ascii_lowercase();
4019 [
4020 "operation not permitted",
4021 "network is unreachable",
4022 "could not resolve host",
4023 "couldn't resolve host",
4024 "failed to resolve",
4025 "temporary failure in name resolution",
4026 "name or service not known",
4027 "nodename nor servname provided",
4028 "no address associated",
4029 "failed to connect",
4030 "couldn't connect",
4031 "connection timed out",
4032 "connection reset",
4033 ]
4034 .iter()
4035 .any(|pattern| output.contains(pattern))
4036 }
4037
4038 fn shell_network_restricted_hint<'a>(
4039 context: &'a ToolContext,
4040 command: &str,
4041 result: &ShellResult,
4042 ) -> Option<&'a str> {
4043 let hint = context.shell_network_denied_hint.as_deref()?;
4044 let policy_blocks_network = context
4045 .elevated_sandbox_policy
4046 .as_ref()
4047 .is_some_and(|policy| !policy.has_network_access());
4048 if !policy_blocks_network || !command_likely_needs_network(command) {
4049 return None;
4050 }
4051 if result.sandbox_denied || looks_like_network_blocked_failure(result) {
4052 Some(hint)
4053 } else {
4054 None
4055 }
4056 }
4057
4058 /// Coaching line when the execution sandbox denied a command and the
4059 /// Plan-mode network hint did not already explain it. Most often a write under
4060 /// a read-only posture: name the effective posture and the Ask-only retry shape
4061 /// so other postures do not mistake it for autonomous authority.
4062 fn shell_sandbox_denied_hint(context: &ToolContext, result: &ShellResult) -> Option<String> {
4063 if !result.sandbox_denied {
4064 return None;
4065 }
4066 let policy = context.elevated_sandbox_policy.as_ref()?;
4067 Some(format!(
4068 "The execution sandbox blocked this command. Effective sandbox posture: {}. [sandbox: Ask-only escalation — retry this exact command once with sandbox_permissions (the narrowest wider mode that suffices) + justification; the approval prompt asks the user]",
4069 policy.posture_label()
4070 ))
4071 }
4072
4073 fn shell_job_owner_from_context(context: &ToolContext) -> Option<ShellJobOwner> {
4074 let agent_id = context
4075 .owner_agent_id
4076 .as_deref()
4077 .map(str::trim)
4078 .filter(|value| !value.is_empty())?;
4079 let agent_name = context
4080 .owner_agent_name
4081 .as_deref()
4082 .map(str::trim)
4083 .filter(|value| !value.is_empty())
4084 .unwrap_or(agent_id);
4085 Some(ShellJobOwner {
4086 agent_id: agent_id.to_string(),
4087 agent_name: agent_name.to_string(),
4088 })
4089 }
4090
4091 fn shell_work_lifecycle_from_context(context: &ToolContext) -> Option<ShellWorkLifecycle> {
4092 context
4093 .runtime
4094 .work
4095 .as_ref()
4096 .map(|work| ShellWorkLifecycle {
4097 work: work.clone(),
4098 session_id: context.state_namespace.clone(),
4099 })
4100 }
4101
4102 fn lifecycle_now_ms() -> i64 {
4103 std::time::SystemTime::now()
4104 .duration_since(std::time::UNIX_EPOCH)
4105 .unwrap_or_default()
4106 .as_millis()
4107 .try_into()
4108 .unwrap_or(i64::MAX)
4109 }
4110
4111 fn attach_shell_owner_metadata(metadata: &mut serde_json::Value, context: &ToolContext) {
4112 let Some(owner) = shell_job_owner_from_context(context) else {
4113 return;
4114 };
4115 metadata["owner_agent_id"] = json!(owner.agent_id);
4116 metadata["owner_agent_name"] = json!(owner.agent_name);
4117 }
4118
4119 /// NUL bytes cannot cross the `exec` boundary: `Command` panics on them.
4120 /// Refuse with the byte offset before anything spawns (#5529).
4121 fn require_no_nul<'a>(value: &'a str, field: &str) -> Result<&'a str, ToolError> {
4122 if let Some(offset) = value.find('\0') {
4123 return Err(ToolError::invalid_input(format!(
4124 "Shell {field} contains a NUL byte at byte offset {offset}; it cannot cross the exec boundary. Remove it (usually a truncated heredoc or binary paste) and re-send."
4125 )));
4126 }
4127 Ok(value)
4128 }
4129
4130 /// Apply the network policy to every network read inside a read-only
4131 /// command, segment by segment, so a pipeline or chain cannot hide one.
4132 /// A full shell keeps its historical scope here: only a lone `gh` read is
4133 /// judged, because its other network use is governed elsewhere.
4134 fn enforce_readonly_network_reads(command: &str, context: &ToolContext) -> Result<(), ToolError> {
4135 let Some(decider) = context.network_policy.as_ref() else {
4136 return Ok(());
4137 };
4138 let mut reads = readonly_network_reads(command);
4139 if context.shell_policy != ShellPolicy::ReadOnly {
4140 let lone = agent_readonly_verdict(command).is_ok_and(|segments| segments.len() == 1);
4141 reads.retain(|read| lone && *read == NetworkRead::GitHub);
4142 }
4143 use crate::network_policy::Decision;
4144 for read in reads {
4145 let host = read.host();
4146 match decider.evaluate(host, "Bash") {
4147 Decision::Allow => {}
4148 Decision::Deny => {
4149 return Err(ToolError::permission_denied(format!(
4150 "Read-only network access to '{host}' is blocked by the active network policy."
4151 )));
4152 }
4153 Decision::Prompt => {
4154 return Err(ToolError::permission_denied(format!(
4155 "Read-only network access to '{host}' requires network approval; allow that host in the parent session or network policy first."
4156 )));
4157 }
4158 }
4159 }
4160 Ok(())
4161 }
4162
4163 /// This is a request for mandatory filesystem/network isolation, not a claim
4164 /// about what the command does. The executor refuses it without a native
4165 /// enforcing sandbox and checks the prepared environment again before spawn.
4166 fn enforced_readonly_input(input: &serde_json::Value) -> bool {
4167 let Some(fields) = input.as_object() else {
4168 return false;
4169 };
4170 fields.get("read_only").and_then(serde_json::Value::as_bool) == Some(true)
4171 && fields.keys().all(|key| {
4172 matches!(
4173 key.as_str(),
4174 "action" | "command" | "cwd" | "timeout_ms" | "read_only"
4175 )
4176 })
4177 && fields
4178 .get("action")
4179 .is_none_or(|action| action.as_str() == Some("run"))
4180 && fields
4181 .get("command")
4182 .and_then(serde_json::Value::as_str)
4183 .is_some_and(|command| !command.trim().is_empty())
4184 }
4185
4186 fn is_native_readonly_sandbox(sandbox_type: SandboxType) -> bool {
4187 match sandbox_type {
4188 #[cfg(target_os = "macos")]
4189 SandboxType::MacosSeatbelt => true,
4190 #[cfg(all(target_os = "linux", not(target_env = "ohos")))]
4191 SandboxType::LinuxBubblewrap => true,
4192 _ => false,
4193 }
4194 }
4195
4196 /// Build the child for a tool that runs workspace code outside `exec_shell`
4197 /// (gate commands, the cargo test runner). It gets exactly the confinement a
4198 /// shell command gets in this session: the same policy, the same OS sandbox
4199 /// wrapper, and the sanitized child environment. Callers keep their own
4200 /// process-tree containment, timeout and cancellation.
4201 ///
4202 /// Known limitations: where the platform has no enforcing sandbox configured
4203 /// (Linux without bubblewrap, Windows) a workspace-write policy runs
4204 /// unsandboxed, as it does for `exec_shell`; a read-only policy refuses there.
4205 /// A session whose commands run in an external sandbox backend cannot run
4206 /// these local children at all.
4207 pub(crate) fn sandboxed_runner_command(
4208 context: &crate::tools::spec::ToolContext,
4209 program: &str,
4210 args: Vec<String>,
4211 cwd: &Path,
4212 timeout: Duration,
4213 ) -> std::result::Result<tokio::process::Command, crate::tools::spec::ToolError> {
4214 use crate::tools::spec::ToolError;
4215 if context.sandbox_backend.is_some() {
4216 return Err(ToolError::not_available(
4217 "this tool starts a local process, and this session runs commands in an external sandbox; run the command through the shell tool instead",
4218 ));
4219 }
4220 let exec_env = context
4221 .shell_manager
4222 .lock()
4223 .unwrap_or_else(std::sync::PoisonError::into_inner)
4224 .prepare_runner(
4225 program,
4226 args,
4227 cwd,
4228 timeout,
4229 context.elevated_sandbox_policy.clone(),
4230 );
4231 if matches!(exec_env.policy, ExecutionSandboxPolicy::ReadOnly) {
4232 require_native_readonly_execution(&exec_env)
4233 .map_err(|error| ToolError::permission_denied(error.to_string()))?;
4234 }
4235 let mut cmd = tokio::process::Command::new(exec_env.program());
4236 crate::utils::suppress_tokio_console_window(&mut cmd);
4237 cmd.args(exec_env.args())
4238 .current_dir(&exec_env.cwd)
4239 .stdout(std::process::Stdio::piped())
4240 .stderr(std::process::Stdio::piped());
4241 // Workspace code never inherits parent credentials.
4242 crate::child_env::apply_to_tokio_command(
4243 &mut cmd,
4244 crate::child_env::string_map_env(&exec_env.env),
4245 );
4246 Ok(cmd)
4247 }
4248
4249 fn require_native_readonly_execution(exec_env: &ExecEnv) -> Result<()> {
4250 if matches!(exec_env.policy, ExecutionSandboxPolicy::ReadOnly)
4251 && is_native_readonly_sandbox(exec_env.sandbox_type)
4252 {
4253 Ok(())
4254 } else {
4255 Err(anyhow!(
4256 "read_only execution requires an enforcing native read-only sandbox; nothing was run"
4257 ))
4258 }
4259 }
4260
4261 /// `exec_shell_input_is_parallel_readonly` with the agent-posture classifier:
4262 /// same input-shape restrictions (run action only, no background/tty/stdin),
4263 /// but commands are judged by
4264 /// [`codewhale_execpolicy::command_safety::agent_readonly_verdict`] so
4265 /// `ShellPolicy::ReadOnly` agents keep a usable inspection surface
4266 /// (pipelines and chains of reads, globs, `git -C`, `find`, `sed -n`,
4267 /// `npm view`). Callers pass input already through [`normalize_readonly_cd`].
4268 fn exec_shell_input_agent_readonly_verdict(
4269 input: &serde_json::Value,
4270 ) -> Result<(), ReadonlyRejection> {
4271 if enforced_readonly_input(input) {
4272 return Ok(());
4273 }
4274 if !exec_shell_input_is_parallel_readonly_shape(input) {
4275 return Err(ReadonlyRejection::new(
4276 "shape",
4277 "read-only shell accepts only a foreground `command` with optional `cwd` and timeout; background, stdin, interactive and TTY modes are not admitted",
4278 ));
4279 }
4280 let command = input
4281 .get("command")
4282 .and_then(serde_json::Value::as_str)
4283 .expect("shape check established a command string");
4284 agent_readonly_verdict(command).map(|_| ())
4285 }
4286
4287 /// Move a leading `cd <dir> &&` into the `cwd` field, resolved against any
4288 /// `cwd` already present, so no gate has to interpret `cd`: the ordinary
4289 /// working-directory workspace check then judges the directory, and every
4290 /// gate and the executor see the same rewritten input. Only the shape
4291 /// [`split_leading_cd`] accepts is rewritten; `read_only: true` input runs
4292 /// its command unchanged under the enforced lane.
4293 pub(crate) fn normalize_readonly_cd(input: &serde_json::Value) -> serde_json::Value {
4294 let mut input = input.clone();
4295 if enforced_readonly_input(&input) {
4296 return input;
4297 }
4298 // Bounded: each pass removes one leading `cd`.
4299 for _ in 0..4 {
4300 let Some((dir, rest)) = input
4301 .get("command")
4302 .and_then(serde_json::Value::as_str)
4303 .and_then(split_leading_cd)
4304 else {
4305 break;
4306 };
4307 let cwd = match input.get("cwd") {
4308 None | Some(serde_json::Value::Null) => dir,
4309 Some(serde_json::Value::String(existing)) => std::path::Path::new(existing)
4310 .join(&dir)
4311 .to_string_lossy()
4312 .into_owned(),
4313 // A wrong type is reported by the executor's own type check.
4314 Some(_) => break,
4315 };
4316 input["command"] = json!(rest);
4317 input["cwd"] = json!(cwd);
4318 }
4319 input
4320 }
4321
4322 /// The refusal a read-only agent sees, shared by every gate that judges the
4323 /// agent grammar (subagent posture, durable authority, and the executor), so
4324 /// the rule text is byte-identical wherever a command is refused. The next
4325 /// steps name only tools such an agent has.
4326 pub(crate) fn readonly_refusal(rejection: &ReadonlyRejection, enforced_lane: bool) -> String {
4327 let lane = if enforced_lane {
4328 ", rerun it with `read_only: true` to run it in an enforced read-only sandbox"
4329 } else {
4330 ""
4331 };
4332 format!(
4333 "{rejection}. Next: use admitted reads only (the shell tool description lists the grammar), read or search files with the File tool{lane}, or, if the probe is essential, return your findings and the blocked probe to the parent."
4334 )
4335 }
4336
4337 /// Whether `read_only: true` can run here: a native enforcing sandbox and
4338 /// no external backend.
4339 pub(crate) fn readonly_enforced_lane_available(context: &ToolContext) -> bool {
4340 context.sandbox_backend.is_none()
4341 && context.shell_manager.lock().is_ok_and(|manager| {
4342 manager
4343 .configured_sandbox_type()
4344 .is_some_and(is_native_readonly_sandbox)
4345 })
4346 }
4347
4348 /// `exec_shell_input_agent_readonly_verdict` is also the gate-side predicate for the
4349 /// subagent posture check (#5426): the catalog carve-out that admits
4350 /// canonical `bash` to Scout/Reviewer/Planner must judge the same call the
4351 /// `BashTool::execute` `ShellPolicy::ReadOnly` branch will judge, so the
4352 /// posture gate can admit a proven-readonly call without ever widening past
4353 /// the execute-time refusal.
4354 pub(crate) fn agent_readonly_bash_input(input: &serde_json::Value) -> bool {
4355 agent_readonly_bash_verdict(input).is_ok()
4356 }
4357
4358 /// [`agent_readonly_bash_input`] with the rule that refused the call.
4359 pub(crate) fn agent_readonly_bash_verdict(
4360 input: &serde_json::Value,
4361 ) -> Result<(), ReadonlyRejection> {
4362 // Canonical lowercase `bash` advertises `timeout` in seconds, while the
4363 // internal executor contract uses `timeout_ms`. Normalize through the same
4364 // translator the concrete tool uses so posture, session approval, envelope,
4365 // and execute judge one input (#5595). Legacy/internal shapes fall back to
4366 // their existing direct classification. A leading `cd` is then moved into
4367 // `cwd` exactly as `BashTool::execute` does.
4368 let translated = contract_bash_legacy_input(input).unwrap_or_else(|_| input.clone());
4369 exec_shell_input_agent_readonly_verdict(&normalize_readonly_cd(&translated))
4370 }
4371
4372 fn exec_shell_input_is_parallel_readonly(input: &serde_json::Value) -> bool {
4373 if enforced_readonly_input(input) {
4374 return true;
4375 }
4376 if !exec_shell_input_is_parallel_readonly_shape(input) {
4377 return false;
4378 }
4379 let command = input
4380 .get("command")
4381 .and_then(serde_json::Value::as_str)
4382 .expect("shape check established a command string");
4383 is_parallel_readonly_command(command)
4384 }
4385
4386 fn exec_shell_input_is_parallel_readonly_shape(input: &serde_json::Value) -> bool {
4387 let Some(fields) = input.as_object() else {
4388 return false;
4389 };
4390 if fields
4391 .keys()
4392 .any(|key| !matches!(key.as_str(), "action" | "command" | "cwd" | "timeout_ms"))
4393 {
4394 return false;
4395 }
4396 match input.get("action") {
4397 None | Some(serde_json::Value::Null) => {}
4398 Some(serde_json::Value::String(action)) if action == "run" => {}
4399 Some(_) => return false,
4400 }
4401 if ["background", "interactive", "tty", "combined_output"]
4402 .iter()
4403 .any(|key| {
4404 !matches!(
4405 input.get(*key),
4406 None | Some(serde_json::Value::Null | serde_json::Value::Bool(false))
4407 )
4408 })
4409 {
4410 return false;
4411 }
4412 if ["stdin", "input", "data"]
4413 .iter()
4414 .any(|key| input.get(*key).is_some())
4415 {
4416 return false;
4417 }
4418 if ["task_id", "id", "wait", "block", "close_stdin", "all"]
4419 .iter()
4420 .any(|key| input.get(*key).is_some())
4421 {
4422 return false;
4423 }
4424
4425 input
4426 .get("command")
4427 .and_then(serde_json::Value::as_str)
4428 .is_some()
4429 }
4430
4431 /// Whether a classifier-approved read needs a shell to join its segments or
4432 /// apply an admitted redirect; a lone plain command runs as direct argv.
4433 fn readonly_command_needs_shell(command: &str) -> bool {
4434 agent_readonly_verdict(command).map_or(true, |segments| {
4435 segments.len() > 1 || segments.iter().any(|segment| !segment.redirects.is_empty())
4436 })
4437 }
4438
4439 fn hardened_readonly_script(command: &str, workspace: &std::path::Path) -> Result<String> {
4440 use crate::shell_dispatcher::ShellKind;
4441 // POSIX quoting must never be passed to a different command interpreter.
4442 let supported = match crate::shell_dispatcher::global_dispatcher().kind() {
4443 ShellKind::Bash => true,
4444 ShellKind::Custom { binary, .. } => matches!(
4445 std::path::Path::new(binary)
4446 .file_name()
4447 .and_then(|name| name.to_str()),
4448 Some("bash" | "zsh")
4449 ),
4450 _ => false,
4451 };
4452 if !supported {
4453 return Err(anyhow!(
4454 "read-only pipelines and chains require bash or zsh; run each read separately"
4455 ));
4456 }
4457 let segments = agent_readonly_verdict(command)
4458 .map_err(|rejection| anyhow!("command is outside the read-only policy: {rejection}"))?;
4459 let mut script = String::from("set -o pipefail;");
4460 for segment in &segments {
4461 let (program, args) = hardened_readonly_argv(&segment.command)?;
4462 let program = resolve_readonly_program(&program, workspace)?;
4463 let program = program
4464 .to_str()
4465 .ok_or_else(|| anyhow!("read-only executable path is not valid UTF-8"))?;
4466 for word in std::iter::once(program).chain(args.iter().map(String::as_str)) {
4467 script.push(' ');
4468 script.push_str(&shell_words::quote(word));
4469 }
4470 for redirect in &segment.redirects {
4471 script.push(' ');
4472 script.push_str(redirect);
4473 }
4474 if let Some(join) = segment.join {
4475 script.push(' ');
4476 script.push_str(join.as_str());
4477 }
4478 }
4479 // The only shell syntax is the admitted operators and redirect words.
4480 // Filenames cannot expand into options or unvalidated symlinks; every Git
4481 // stage retains helper guards.
4482 Ok(script)
4483 }
4484
4485 fn hardened_readonly_argv(command: &str) -> Result<(String, Vec<String>)> {
4486 let mut argv = shell_words::split(&normalize_windows_command_paths(command))
4487 .map_err(|error| anyhow!("could not parse classifier-approved read command: {error}"))?;
4488 if argv.is_empty() {
4489 return Err(anyhow!("classifier-approved read command was empty"));
4490 }
4491
4492 // Even when repository/user configuration names a diff or signature
4493 // helper, these flags make Git keep the read inside its own process.
4494 if argv.first().is_some_and(|program| program == "git") {
4495 // The agent read-only classifier admits `git -C <dir>` and
4496 // `git --no-pager` before the subcommand; keep the preamble but
4497 // locate the subcommand after it so the hardening flags splice in
4498 // the right place. `-C` targets were already workspace-checked by
4499 // `enforce_readonly_workspace_operands`.
4500 let mut subcommand_index = 1;
4501 while let Some(flag) = argv.get(subcommand_index) {
4502 match flag.as_str() {
4503 "--no-pager" => subcommand_index += 1,
4504 "-C" => subcommand_index += 2,
4505 _ => break,
4506 }
4507 }
4508 let subcommand = argv
4509 .get(subcommand_index)
4510 .map(String::as_str)
4511 .ok_or_else(|| {
4512 anyhow!("classifier-approved Git read was missing its literal subcommand")
4513 })?;
4514 match subcommand {
4515 "diff" => {
4516 let at = subcommand_index + 1;
4517 argv.splice(
4518 at..at,
4519 crate::dependencies::Git::REVIEW_DIFF_ARGS.map(String::from),
4520 );
4521 }
4522 "log" | "show" => {
4523 let at = subcommand_index + 1;
4524 argv.splice(
4525 at..at,
4526 crate::dependencies::Git::REVIEW_DIFF_ARGS
4527 .into_iter()
4528 .chain(["--no-show-signature"])
4529 .map(String::from),
4530 );
4531 }
4532 "blame" => {
4533 let at = subcommand_index + 1;
4534 argv.insert(at, "--no-textconv".to_string());
4535 }
4536 "status" | "ls-files" | "grep" => {}
4537 _ => {
4538 return Err(anyhow!(
4539 "classifier-approved Git read did not keep its subcommand in argv[1]"
4540 ));
4541 }
4542 }
4543 }
4544
4545 let program = argv.remove(0);
4546 Ok((program, argv))
4547 }
4548
4549 /// The directory each `git` segment of a read-only command runs in: `cwd`
4550 /// followed through any leading `-C` hops, as git itself resolves them.
4551 fn readonly_git_dirs(command: &str, cwd: &std::path::Path) -> Vec<std::path::PathBuf> {
4552 // Split with the same lexer the classifier and executor use (#6637), so a
4553 // `git` segment after `&&`, `||` or `;` gets its filter hardening too.
4554 let segments = agent_readonly_verdict(command).map_or_else(
4555 |_| command.split('|').map(str::to_string).collect::<Vec<_>>(),
4556 |segments| {
4557 segments
4558 .into_iter()
4559 .map(|segment| segment.command)
4560 .collect()
4561 },
4562 );
4563 segments
4564 .iter()
4565 .filter_map(|segment| shell_words::split(&normalize_windows_command_paths(segment)).ok())
4566 .filter(|argv| argv.first().is_some_and(|program| program == "git"))
4567 .map(|argv| {
4568 let mut dir = cwd.to_path_buf();
4569 let mut args = argv.iter().skip(1);
4570 while let Some(flag) = args.next() {
4571 match flag.as_str() {
4572 "--no-pager" => {}
4573 "-C" => match args.next() {
4574 Some(target) => dir = dir.join(target),
4575 None => break,
4576 },
4577 _ => break,
4578 }
4579 }
4580 dir
4581 })
4582 .collect()
4583 }
4584
4585 fn enforce_readonly_workspace_operands(
4586 command: &str,
4587 workspace: &std::path::Path,
4588 effective_cwd: &std::path::Path,
4589 ) -> Result<(), ToolError> {
4590 // Judge each segment of a pipeline or chain on its own, so the `gh`
4591 // exemption below covers only the `gh` segment itself.
4592 let segments = agent_readonly_verdict(command).map_or_else(
4593 |_| vec![command.to_string()],
4594 |segments| {
4595 segments
4596 .into_iter()
4597 .map(|segment| segment.command)
4598 .collect()
4599 },
4600 );
4601 for segment in &segments {
4602 enforce_readonly_segment_operands(segment, workspace, effective_cwd)?;
4603 }
4604 Ok(())
4605 }
4606
4607 fn enforce_readonly_segment_operands(
4608 command: &str,
4609 workspace: &std::path::Path,
4610 effective_cwd: &std::path::Path,
4611 ) -> Result<(), ToolError> {
4612 let argv = shell_words::split(&normalize_windows_command_paths(command)).map_err(|error| {
4613 ToolError::invalid_input(format!(
4614 "Could not parse read-only command arguments: {error}"
4615 ))
4616 })?;
4617 if argv.first().is_some_and(|program| program == "gh") {
4618 // High-level gh reads do not consume local path operands. Their host
4619 // is pinned and evaluated separately by the network-policy guard.
4620 return Ok(());
4621 }
4622 let workspace = workspace.canonicalize().map_err(|error| {
4623 ToolError::execution_failed(format!(
4624 "Could not resolve the Scout workspace before shell dispatch: {error}"
4625 ))
4626 })?;
4627 let effective_cwd = effective_cwd.canonicalize().map_err(|error| {
4628 ToolError::permission_denied(format!(
4629 "Could not prove the read-only shell working directory stays in the workspace: {error}"
4630 ))
4631 })?;
4632 if !effective_cwd.starts_with(&workspace) {
4633 return Err(ToolError::permission_denied(
4634 "[shell.readonly.cwd.outside_workspace] Read-only Scout shell working directory resolves outside the workspace.",
4635 ));
4636 }
4637
4638 for token in argv.iter().skip(1) {
4639 if token.starts_with('-') && (token.contains('/') || token.contains('\\')) {
4640 return Err(ToolError::permission_denied(format!(
4641 "[shell.readonly.option.attached_path] Read-only Scout shell options may not carry attached paths; refused {token:?}. Use the bounded File read/search actions for project evidence."
4642 )));
4643 }
4644 let value = token
4645 .split_once('=')
4646 .map_or(token.as_str(), |(_, value)| value)
4647 .trim();
4648 // Windows `Path::canonicalize` returns verbatim device paths
4649 // (`\\?\C:\...`), which `Path::is_absolute` does not recognize; without
4650 // stripping the prefix the operand falls into the shape refusal and
4651 // every canonical absolute read is denied on Windows. The stripped
4652 // form resolves to the same location, so location-based judgement is
4653 // unchanged. On unix hosts such spellings stay fail-closed below.
4654 let value = value
4655 .strip_prefix(r"\\?\")
4656 .or_else(|| value.strip_prefix(r"\\.\"))
4657 .unwrap_or(value);
4658 if value.is_empty() || value == "-" {
4659 continue;
4660 }
4661 let candidate = std::path::Path::new(value);
4662 let bytes = value.as_bytes();
4663 let windows_prefixed =
4664 (bytes.len() >= 2 && bytes[0].is_ascii_alphabetic() && bytes[1] == b':')
4665 || value.starts_with("\\\\")
4666 || candidate
4667 .components()
4668 .any(|component| matches!(component, std::path::Component::Prefix(_)));
4669
4670 // #5595: an absolute operand is safe by resolved location, not by
4671 // spelling. This is the canonical `git -C /absolute/workspace log`
4672 // case. Requiring canonicalization also keeps nonexistent output-like
4673 // targets fail-closed, while symlinks are judged by their destination.
4674 if candidate.is_absolute() {
4675 let resolved = candidate.canonicalize().map_err(|error| {
4676 ToolError::permission_denied(format!(
4677 "[shell.readonly.operand.unresolved] Could not prove absolute read-only operand {value:?} stays inside the workspace because it could not be resolved: {error}"
4678 ))
4679 })?;
4680 if !resolved.starts_with(&workspace) {
4681 return Err(ToolError::permission_denied(format!(
4682 "[shell.readonly.operand.outside_workspace] Read-only Scout shell operand {value:?} resolves outside the workspace. Use the bounded File read/search actions for project evidence."
4683 )));
4684 }
4685 continue;
4686 }
4687
4688 if value.starts_with('~')
4689 || value.contains('\\')
4690 || windows_prefixed
4691 || candidate.has_root()
4692 || candidate
4693 .components()
4694 .any(|component| matches!(component, std::path::Component::ParentDir))
4695 {
4696 return Err(ToolError::permission_denied(format!(
4697 "[shell.readonly.operand.shape] Read-only Scout shell operands must stay inside the workspace; refused {value:?}. Use the bounded File read/search actions for project evidence."
4698 )));
4699 }
4700
4701 let joined = effective_cwd.join(candidate);
4702 if joined.exists() {
4703 let resolved = joined.canonicalize().map_err(|error| {
4704 ToolError::permission_denied(format!(
4705 "[shell.readonly.operand.unresolved] Could not prove read-only operand {value:?} stays in the workspace: {error}"
4706 ))
4707 })?;
4708 if !resolved.starts_with(&workspace) {
4709 return Err(ToolError::permission_denied(format!(
4710 "[shell.readonly.operand.outside_workspace] Read-only Scout shell operand {value:?} resolves outside the workspace. Use the bounded File read/search actions for project evidence."
4711 )));
4712 }
4713 }
4714 }
4715 Ok(())
4716 }
4717
4718 fn readonly_sanitized_path_from(
4719 workspace: &std::path::Path,
4720 path: &std::ffi::OsStr,
4721 ) -> Option<std::ffi::OsString> {
4722 let workspace = workspace.canonicalize().ok()?;
4723 let safe = std::env::split_paths(path).filter_map(|entry| {
4724 if !entry.is_absolute() {
4725 return None;
4726 }
4727 let resolved = entry.canonicalize().ok()?;
4728 (!resolved.starts_with(&workspace)).then_some(resolved)
4729 });
4730 std::env::join_paths(safe).ok()
4731 }
4732
4733 fn readonly_sanitized_path(workspace: &std::path::Path) -> Option<String> {
4734 let path = std::env::var_os("PATH")?;
4735 readonly_sanitized_path_from(workspace, &path).map(|value| value.to_string_lossy().into_owned())
4736 }
4737
4738 fn resolve_readonly_program(program: &str, workspace: &std::path::Path) -> Result<PathBuf> {
4739 let path = std::env::var_os("PATH")
4740 .ok_or_else(|| anyhow!("no executable search path is configured"))?;
4741 resolve_readonly_program_from_path(program, workspace, &path)
4742 }
4743
4744 fn resolve_readonly_program_from_path(
4745 program: &str,
4746 workspace: &std::path::Path,
4747 path: &std::ffi::OsStr,
4748 ) -> Result<PathBuf> {
4749 let workspace = workspace.canonicalize()?;
4750 if std::path::Path::new(program).components().count() != 1 {
4751 return Err(anyhow!(
4752 "read-only command must name a bare allowlisted executable"
4753 ));
4754 }
4755 let safe_path = readonly_sanitized_path_from(&workspace, path).ok_or_else(|| {
4756 anyhow!("no trusted executable search path remains outside the workspace")
4757 })?;
4758 let names = if cfg!(windows) {
4759 vec![format!("{program}.exe"), format!("{program}.com")]
4760 } else {
4761 vec![program.to_string()]
4762 };
4763 for directory in std::env::split_paths(&safe_path) {
4764 for name in &names {
4765 let candidate = directory.join(name);
4766 if !candidate.is_file() {
4767 continue;
4768 }
4769 #[cfg(unix)]
4770 {
4771 use std::os::unix::fs::PermissionsExt as _;
4772 if candidate.metadata()?.permissions().mode() & 0o111 == 0 {
4773 continue;
4774 }
4775 }
4776 let resolved = candidate.canonicalize()?;
4777 if resolved.is_absolute() && !resolved.starts_with(&workspace) {
4778 return Ok(resolved);
4779 }
4780 }
4781 }
4782 Err(anyhow!(
4783 "allowlisted read-only executable {program:?} was not found at a canonical path outside the workspace"
4784 ))
4785 }
4786
4787 fn remove_readonly_redirect_env(cmd: &mut Command, env: &HashMap<String, String>) {
4788 if env.get(READONLY_ENV_MARKER).map(String::as_str) != Some("1") {
4789 return;
4790 }
4791 cmd.env_remove(READONLY_ENV_MARKER);
4792 let removals = cmd
4793 .get_envs()
4794 .filter_map(|(key, _)| {
4795 let upper = key.to_string_lossy().to_ascii_uppercase();
4796 let guarded = upper.starts_with("GIT_")
4797 || upper.starts_with("GH_")
4798 || upper.starts_with("GITHUB_");
4799 let safe = matches!(
4800 upper.as_str(),
4801 "GIT_OPTIONAL_LOCKS"
4802 | "GIT_NO_LAZY_FETCH"
4803 | "GIT_PAGER"
4804 | "GIT_CONFIG_NOSYSTEM"
4805 | "GIT_CONFIG_GLOBAL"
4806 | "GIT_CONFIG_PARAMETERS"
4807 | "GIT_EXTERNAL_DIFF"
4808 | "GIT_ATTR_NOSYSTEM"
4809 | "GIT_CONFIG_COUNT"
4810 | "GH_PAGER"
4811 | "GH_PROMPT_DISABLED"
4812 | "GH_NO_UPDATE_NOTIFIER"
4813 | "GH_HOST"
4814 | "GH_REPO"
4815 ) || upper.starts_with("GIT_CONFIG_KEY_")
4816 || upper.starts_with("GIT_CONFIG_VALUE_");
4817 (guarded && !safe).then(|| key.to_os_string())
4818 })
4819 .collect::<Vec<_>>();
4820 for key in removals {
4821 cmd.env_remove(key);
4822 }
4823 }
4824
4825 fn exec_shell_input_starts_detached(input: &serde_json::Value) -> bool {
4826 input
4827 .get("command")
4828 .and_then(serde_json::Value::as_str)
4829 .is_some()
4830 && input
4831 .get("interactive")
4832 .and_then(serde_json::Value::as_bool)
4833 != Some(true)
4834 && (input.get("background").and_then(serde_json::Value::as_bool) == Some(true)
4835 || input.get("tty").and_then(serde_json::Value::as_bool) == Some(true))
4836 }
4837
4838 fn persistent_services_enabled_for(context: &ToolContext) -> bool {
4839 #[cfg(unix)]
4840 {
4841 context.persist_services_enabled
4842 && context.owner_agent_id.is_none()
4843 && context.tool_authority.is_none()
4844 && context.sandbox_backend.is_none()
4845 && matches!(context.shell_policy, ShellPolicy::Full)
4846 && matches!(
4847 context.elevated_sandbox_policy,
4848 Some(ExecutionSandboxPolicy::DangerFullAccess)
4849 )
4850 }
4851 #[cfg(not(unix))]
4852 {
4853 let _ = context;
4854 false
4855 }
4856 }
4857
4858 #[allow(clippy::too_many_arguments)]
4859 async fn execute_foreground_via_background(
4860 context: &ToolContext,
4861 command: &str,
4862 heavy_permit: Option<HeavyCommandPermit>,
4863 working_dir: Option<String>,
4864 timeout_ms: Option<u64>,
4865 stdin_data: Option<&str>,
4866 tty: bool,
4867 policy_override: Option<ExecutionSandboxPolicy>,
4868 extra_env: HashMap<String, String>,
4869 direct_argv: bool,
4870 timeout_bounds_ms: (u64, u64),
4871 wants_receipt: bool,
4872 receipt_identity: &mut Option<ShellExecutionIdentity>,
4873 ) -> Result<ShellResult> {
4874 let timeout_ms =
4875 timeout_ms.map(|timeout| timeout.clamp(timeout_bounds_ms.0, timeout_bounds_ms.1));
4876 let spawn_timeout_ms = timeout_ms.unwrap_or(timeout_bounds_ms.1);
4877 // Freeze the receipt's directory before execution and hand that same
4878 // resolved spelling to the OS. Looking up the original symlink after
4879 // the command runs can name a different directory than the one it used.
4880 // Resolution is asynchronous; failure leaves the ordinary execution
4881 // path intact but cannot produce an exact receipt.
4882 let receipt_cwd = if wants_receipt {
4883 match working_dir.as_deref() {
4884 Some(cwd) => tokio::fs::canonicalize(cwd)
4885 .await
4886 .ok()
4887 .and_then(|path| path.into_os_string().into_string().ok()),
4888 None => None,
4889 }
4890 } else {
4891 None
4892 };
4893 let working_dir = receipt_cwd.clone().or(working_dir);
4894 let task_id = {
4895 let mut manager = context
4896 .shell_manager
4897 .lock()
4898 .map_err(|_| anyhow!("shell manager lock poisoned"))?;
4899 manager.clear_foreground_background_request();
4900 let owner = shell_job_owner_from_context(context);
4901 let lifecycle = shell_work_lifecycle_from_context(context);
4902 let spawned = manager.execute_with_options_env_for_owner_and_work(
4903 command,
4904 working_dir.as_deref(),
4905 spawn_timeout_ms,
4906 true,
4907 stdin_data,
4908 tty,
4909 policy_override,
4910 extra_env,
4911 owner,
4912 context.state_namespace.clone(),
4913 context.origin_tool_call_id.clone(),
4914 context.origin_turn_id.clone(),
4915 lifecycle,
4916 direct_argv.then_some(context.workspace.as_path()),
4917 false,
4918 timeout_bounds_ms,
4919 )?;
4920 let task_id = spawned
4921 .task_id
4922 .ok_or_else(|| anyhow!("foreground shell did not return a process id"))?;
4923 // Classify before releasing the manager lock: even an immediately
4924 // completed foreground command must never look like background work.
4925 let process = manager
4926 .processes
4927 .get_mut(&task_id)
4928 .ok_or_else(|| anyhow!("foreground shell {task_id} is not tracked"))?;
4929 process.background = false;
4930 // #6689: the receipt identity is what the manager recorded for this
4931 // spawn — the admitted command and the directory handed to the OS —
4932 // never the before-hook request. Only pipe-backed, unsandboxed local
4933 // runs through a POSIX-style `<shell> <flag> <source>` dispatcher
4934 // qualify: a PTY, the hardened read-only argv rewrite, an OS sandbox
4935 // wrapper, Windows shell prefixes, and a PowerShell dispatcher (even
4936 // `$SHELL=pwsh` on Unix, which wraps the source or runs it from a temp
4937 // `-File`) all change what the process executes relative to this string.
4938 if receipt_cwd.is_some()
4939 && !cfg!(windows)
4940 && !tty
4941 && !direct_argv
4942 && !spawned.sandboxed
4943 && !spawned.sandbox_denied
4944 && !crate::shell_dispatcher::global_dispatcher()
4945 .kind()
4946 .is_powershell()
4947 {
4948 *receipt_identity = Some(ShellExecutionIdentity {
4949 command: process.command.clone(),
4950 cwd: process.working_dir.clone(),
4951 });
4952 }
4953 task_id
4954 };
4955 let mut foreground = ForegroundShellGuard {
4956 manager: context.shell_manager.clone(),
4957 task_id: task_id.clone(),
4958 armed: true,
4959 };
4960 if let Some(permit) = heavy_permit {
4961 let mut manager = context
4962 .shell_manager
4963 .lock()
4964 .map_err(|_| anyhow!("shell manager lock poisoned"))?;
4965 manager.attach_heavy_permit(&task_id, permit)?;
4966 }
4967
4968 // A foreground pipe command gets EOF on stdin: an unexpected read (`cat`,
4969 // `read x`, a confirmation prompt) then fails at once instead of blocking
4970 // until the timeout kills it. A TTY keeps its terminal input. A command
4971 // later moved to /jobs keeps the closed stdin; interactive input needs
4972 // `background: true` from the start.
4973 if stdin_data.is_some() || !tty {
4974 let mut manager = context
4975 .shell_manager
4976 .lock()
4977 .map_err(|_| anyhow!("shell manager lock poisoned"))?;
4978 manager.write_stdin(&task_id, "", true)?;
4979 }
4980
4981 let deadline = timeout_ms.map(|timeout| Instant::now() + Duration::from_millis(timeout));
4982 // Adaptive poll cadence: fast commands (the common case — grep, wc, echo)
4983 // finish in single-digit milliseconds, and a fixed 100ms tick made every
4984 // foreground call pay that full quantum before completion was noticed.
4985 // Start fine-grained and back off to the 100ms cap for long-running work.
4986 let mut poll_tick_ms: u64 = FOREGROUND_POLL_INITIAL_MS;
4987 loop {
4988 if context
4989 .cancel_token
4990 .as_ref()
4991 .is_some_and(|token| token.is_cancelled())
4992 {
4993 let mut manager = context
4994 .shell_manager
4995 .lock()
4996 .map_err(|_| anyhow!("shell manager lock poisoned"))?;
4997 let result = manager.kill(&task_id);
4998 if result.is_ok() {
4999 manager.acknowledge_foreground_completion(&task_id);
5000 foreground.armed = false;
5001 }
5002 return result;
5003 }
5004
5005 // Poll status only. The snapshot — and the buffer clones behind it — is
5006 // built once, when there is actually a result to return (#5472). Both
5007 // happen under one lock acquisition so the record cannot be evicted
5008 // between observing that it finished and reading its result.
5009 let finished = {
5010 let mut manager = context
5011 .shell_manager
5012 .lock()
5013 .map_err(|_| anyhow!("shell manager lock poisoned"))?;
5014 if manager.take_foreground_background_request(&task_id) {
5015 let snapshot = manager.get_output(&task_id, false, 0)?;
5016 foreground.armed = false;
5017 return Ok(snapshot);
5018 }
5019 if manager.poll_status(&task_id)? == ShellStatus::Running {
5020 None
5021 } else {
5022 // #6689: a status from a failed wait is neither an observed
5023 // exit nor an interruption, so the receipt is left out rather
5024 // than reporting a guessed state.
5025 if manager
5026 .processes
5027 .get(&task_id)
5028 .is_none_or(|shell| shell.wait_failed)
5029 {
5030 *receipt_identity = None;
5031 }
5032 let snapshot = manager.get_output(&task_id, false, 0)?;
5033 // Ordering matters: the snapshot is taken before the
5034 // acknowledgement releases the retained bytes.
5035 manager.acknowledge_foreground_completion(&task_id);
5036 Some(snapshot)
5037 }
5038 };
5039
5040 if let Some(snapshot) = finished {
5041 foreground.armed = false;
5042 return Ok(snapshot);
5043 }
5044
5045 if deadline.is_some_and(|deadline| Instant::now() >= deadline) {
5046 let mut manager = context
5047 .shell_manager
5048 .lock()
5049 .map_err(|_| anyhow!("shell manager lock poisoned"))?;
5050 let mut result = manager.kill(&task_id)?;
5051 manager.acknowledge_foreground_completion(&task_id);
5052 result.status = ShellStatus::TimedOut;
5053 foreground.armed = false;
5054 return Ok(result);
5055 }
5056
5057 tokio::time::sleep(Duration::from_millis(poll_tick_ms)).await;
5058 poll_tick_ms = (poll_tick_ms * 2).min(FOREGROUND_POLL_MAX_MS);
5059 }
5060 }
5061
5062 /// A foreground wait owns its process even if its caller drops the future
5063 /// before cooperative cancellation can be polled. Only an explicit transfer
5064 /// to /jobs releases that ownership while the process is still running.
5065 struct ForegroundShellGuard {
5066 manager: SharedShellManager,
5067 task_id: String,
5068 armed: bool,
5069 }
5070
5071 impl Drop for ForegroundShellGuard {
5072 fn drop(&mut self) {
5073 if !self.armed {
5074 return;
5075 }
5076 let mut manager = self
5077 .manager
5078 .lock()
5079 .unwrap_or_else(std::sync::PoisonError::into_inner);
5080 let result = manager.poll_status(&self.task_id).and_then(|status| {
5081 if status == ShellStatus::Running {
5082 manager.kill(&self.task_id).map(|_| ())
5083 } else {
5084 Ok(())
5085 }
5086 });
5087 if let Err(error) = result {
5088 tracing::warn!(shell_id = %self.task_id, %error, "foreground shell cleanup failed");
5089 }
5090 if let Some(shell) = manager.processes.get_mut(&self.task_id) {
5091 // No result received these bytes. Retain them for /jobs inspection,
5092 // but do not let an abandoned foreground wait wake a new model turn.
5093 shell.completion_reported = true;
5094 }
5095 }
5096 }
5097
5098 /// The admitted command and working directory a foreground spawn recorded,
5099 /// for the `tool_call_after` execution receipt (#6689).
5100 struct ShellExecutionIdentity {
5101 command: String,
5102 cwd: PathBuf,
5103 }
5104
5105 /// Largest command or working directory a receipt carries. Identities are
5106 /// exact or absent, never truncated.
5107 const EXECUTION_RECEIPT_IDENTITY_MAX_BYTES: usize = 8 * 1024;
5108
5109 /// Starting per-stream output preview in an execution receipt. Halved until
5110 /// the serialized receipt fits `HOOK_EXECUTION_RECEIPT_MAX_BYTES`.
5111 const EXECUTION_RECEIPT_PREVIEW_MAX_BYTES: usize = 8 * 1024;
5112
5113 /// Keep both ends of an output stream — the first lines and the final
5114 /// diagnostic — cut on UTF-8 boundaries with a visible marker.
5115 fn execution_receipt_preview(text: &str, budget: usize) -> (String, bool) {
5116 if text.len() <= budget {
5117 return (text.to_owned(), false);
5118 }
5119 let mut head = budget / 2;
5120 while !text.is_char_boundary(head) {
5121 head -= 1;
5122 }
5123 let mut tail = text.len() - budget / 2;
5124 while !text.is_char_boundary(tail) {
5125 tail += 1;
5126 }
5127 (
5128 format!(
5129 "{}\n[receipt preview truncated]\n{}",
5130 &text[..head],
5131 &text[tail..]
5132 ),
5133 true,
5134 )
5135 }
5136
5137 /// Build the schema-1 execution receipt for a settled foreground shell run
5138 /// (#6689), exported to `tool_call_after` hooks as
5139 /// `DEEPSEEK_TOOL_EXECUTION_RECEIPT`.
5140 ///
5141 /// Returns `None` — absence, which implies neither success nor failure — for
5142 /// a run still in flight, a sandboxed run, or an identity that is not exact:
5143 /// empty, over the identity bound, containing NUL, or a relative or non-UTF-8
5144 /// directory. `output_kind` is `"combined"` when stdout and stderr shared one
5145 /// pipe, so `stdout` holds the combined preview.
5146 fn shell_execution_receipt(
5147 identity: &ShellExecutionIdentity,
5148 result: &ShellResult,
5149 output_kind: &'static str,
5150 ) -> Option<serde_json::Value> {
5151 let command = identity.command.as_str();
5152 let cwd = identity.cwd.to_str()?;
5153 if command.is_empty()
5154 || command.len() > EXECUTION_RECEIPT_IDENTITY_MAX_BYTES
5155 || command.contains('\0')
5156 || cwd.len() > EXECUTION_RECEIPT_IDENTITY_MAX_BYTES
5157 || cwd.contains('\0')
5158 || !identity.cwd.is_absolute()
5159 || result.sandboxed
5160 || result.sandbox_denied
5161 {
5162 return None;
5163 }
5164 // A nonzero exit is still a completed run; `interrupted` is a signal,
5165 // kill, cancel, or timeout. The exit code is only ever the observed one.
5166 let state = match result.status {
5167 ShellStatus::Completed => "completed",
5168 ShellStatus::Failed if result.exit_code.is_some() => "completed",
5169 ShellStatus::Failed | ShellStatus::Killed | ShellStatus::TimedOut => "interrupted",
5170 ShellStatus::Running => return None,
5171 };
5172 let mut budget = EXECUTION_RECEIPT_PREVIEW_MAX_BYTES;
5173 loop {
5174 let (stdout, stdout_clipped) = execution_receipt_preview(&result.stdout, budget);
5175 let (stderr, stderr_clipped) = execution_receipt_preview(&result.stderr, budget);
5176 let receipt = json!({
5177 "schema_version": 1,
5178 "command": command,
5179 "cwd": cwd,
5180 "state": state,
5181 "scope": "local",
5182 "exit_code": result.exit_code,
5183 "stdout": stdout,
5184 "stderr": stderr,
5185 "stdout_truncated": result.stdout_truncated || stdout_clipped,
5186 "stderr_truncated": result.stderr_truncated || stderr_clipped,
5187 "output_kind": output_kind,
5188 });
5189 // The bound is on the serialized form, so JSON escaping counts.
5190 if serde_json::to_vec(&receipt).ok()?.len()
5191 <= crate::hooks::HOOK_EXECUTION_RECEIPT_MAX_BYTES
5192 {
5193 return Some(receipt);
5194 }
5195 if budget == 0 {
5196 return None;
5197 }
5198 budget /= 2;
5199 }
5200 }
5201
5202 const BASH_MAX_TIMEOUT_MS: u64 = i32::MAX as u64;
5203
5204 /// Initial cadence for foreground-via-background completion polling. Fast
5205 /// commands dominate real agent traffic; detection latency on `true`-class
5206 /// commands drops from ~100ms to ~10ms, while long commands reach the
5207 /// 100ms cap within one doubling step.
5208 const FOREGROUND_POLL_INITIAL_MS: u64 = 10;
5209 /// Poll-cadence ceiling; matches the previous fixed tick so long-running
5210 /// command overhead is unchanged.
5211 const FOREGROUND_POLL_MAX_MS: u64 = 100;
5212
5213 /// Default foreground lifetime for a contract-`bash` `action=run` that names
5214 /// no `timeout_ms`. Matches the value the tool's own input schema advertises;
5215 /// before this existed the omitted case fell through to
5216 /// `BASH_MAX_TIMEOUT_MS`.
5217 const CONTRACT_BASH_FOREGROUND_DEFAULT_TIMEOUT_MS: u64 = 120_000;
5218
5219 /// Resolve the lifetime for one `bash` run.
5220 ///
5221 /// A foreground contract-`bash` run that names no timeout used to inherit
5222 /// `BASH_MAX_TIMEOUT_MS` (~24.8 days), so a command that blocked on an
5223 /// interactive prompt or a hung network call pinned the turn indefinitely —
5224 /// the tool row just counted seconds while the model waited. The tool's own
5225 /// schema already promises `action=run 120000`, and its description already
5226 /// says foreground is for bounded commands, so honor that: an omitted
5227 /// timeout takes the advertised default, which lets
5228 /// `FOREGROUND_TIMEOUT_RECOVERY_HINT` kill the process and tell the model to
5229 /// rerun with `background=true`.
5230 ///
5231 /// An explicit `timeout_ms` is still honored up to the full contract ceiling,
5232 /// and background and interactive runs keep their own lifetimes: their
5233 /// processes are meant to outlive the call, so bounding them here would kill
5234 /// long-lived jobs the model deliberately detached.
5235 fn contract_bash_timeout_ms(
5236 optional_timeout: bool,
5237 requested_ms: Option<u64>,
5238 background: bool,
5239 interactive: bool,
5240 ) -> Option<u64> {
5241 if optional_timeout && requested_ms.is_none() && !background && !interactive {
5242 return Some(CONTRACT_BASH_FOREGROUND_DEFAULT_TIMEOUT_MS);
5243 }
5244 requested_ms
5245 }
5246
5247 fn contract_bash_error_status(result: &ShellResult, timeout_ms: Option<u64>) -> String {
5248 match result.status {
5249 ShellStatus::TimedOut => {
5250 let millis = timeout_ms.unwrap_or(BASH_MAX_TIMEOUT_MS);
5251 let seconds = if millis.is_multiple_of(1_000) {
5252 (millis / 1_000).to_string()
5253 } else {
5254 format!("{}", millis as f64 / 1_000.0)
5255 };
5256 format!("Command timed out after {seconds} seconds")
5257 }
5258 ShellStatus::Killed => "Command aborted".to_string(),
5259 ShellStatus::Failed | ShellStatus::Completed | ShellStatus::Running => format!(
5260 "Command exited with code {}",
5261 result.exit_code.unwrap_or(-1)
5262 ),
5263 }
5264 }
5265
5266 fn finish_contract_bash_result(
5267 result: ShellResult,
5268 timeout_ms: Option<u64>,
5269 context: &ToolContext,
5270 execution_receipt: Option<serde_json::Value>,
5271 ) -> Result<ToolResult, ToolError> {
5272 let sandbox_denied_hint = shell_sandbox_denied_hint(context, &result);
5273 let mut output = result.stdout.clone();
5274 output.push_str(&result.stderr);
5275 if let Some(hint) = sandbox_denied_hint {
5276 output = if output.is_empty() {
5277 hint
5278 } else {
5279 format!("{hint}\n\n{output}")
5280 };
5281 }
5282 let mut metadata = json!({
5283 "evidence_routing": "inline", "exit_code": result.exit_code,
5284 "status": format!("{:?}", result.status), "duration_ms": result.duration_ms,
5285 "sandboxed": result.sandboxed, "sandbox_type": result.sandbox_type,
5286 "task_id": result.task_id, "backgrounded": result.status == ShellStatus::Running,
5287 });
5288 if let Some(receipt) = execution_receipt {
5289 metadata["execution_receipt"] = receipt;
5290 }
5291 if result.status == ShellStatus::Running {
5292 let task_id = result.task_id.as_deref().unwrap_or("unknown");
5293 let partial = (!output.is_empty()).then(|| format!("\n\nOutput so far:\n{output}"));
5294 return Ok(ToolResult::success(format!(
5295 "Foreground shell wait moved to /jobs: {task_id}{}\n\nThe command is still running; completion will appear as a runtime event.",
5296 partial.as_deref().unwrap_or_default()
5297 )).with_metadata(metadata));
5298 }
5299 if result.status != ShellStatus::Completed {
5300 // A nonzero exit, timeout, or kill stays a failed tool call, but it
5301 // keeps the same metadata a success carries: hooks read `exit_code`
5302 // and `status` from it, and an error with no metadata left them blind
5303 // to exactly the commands that failed.
5304 let status = contract_bash_error_status(&result, timeout_ms);
5305 return Err(ToolError::execution_failed_with_metadata(
5306 if output.is_empty() {
5307 status
5308 } else {
5309 format!("{output}\n\n{status}")
5310 },
5311 metadata,
5312 ));
5313 }
5314
5315 Ok(ToolResult::success(if output.is_empty() {
5316 "(no output)".to_string()
5317 } else {
5318 output
5319 })
5320 .with_metadata(metadata))
5321 }
5322
5323 /// Small foreground-only shell surface shown to new model turns.
5324 pub struct LowercaseBashTool;
5325
5326 #[async_trait]
5327 impl ToolSpec for LowercaseBashTool {
5328 fn name(&self) -> &'static str {
5329 "bash"
5330 }
5331
5332 fn description(&self) -> &'static str {
5333 guidance::foreground_description()
5334 }
5335
5336 fn input_schema(&self) -> serde_json::Value {
5337 json!({
5338 "type": "object",
5339 "properties": {
5340 "command": { "type": "string", "description": guidance::runtime_command_guidance() },
5341 "timeout": { "type": "number", "description": "Optional timeout in seconds; when omitted the command is killed after 120 seconds." },
5342 "read_only": { "type": "boolean", "description": "Set true to run analysis code (including Python/SQLite) with mandatory native filesystem read-only isolation and no network. Available during peer writes. Refused when native enforcement is unavailable; no background, stdin, external backend, or sandbox escalation." },
5343 "sandbox_permissions": {
5344 "type": "string",
5345 "enum": ["workspace-write", "danger-full-access"],
5346 "description": "The wider sandbox mode this exact command needs. Use only as a one-shot retry after a sandbox denial; requires justification and user approval in Ask."
5347 },
5348 "justification": {
5349 "type": "string",
5350 "description": "Required with sandbox_permissions: one sentence explaining why this exact command needs wider access."
5351 }
5352 },
5353 "required": ["command"],
5354 "additionalProperties": false
5355 })
5356 }
5357
5358 fn capabilities(&self) -> Vec<ToolCapability> {
5359 BashTool::contract_delegate().capabilities()
5360 }
5361
5362 fn approval_requirement(&self) -> ApprovalRequirement {
5363 ApprovalRequirement::Required
5364 }
5365
5366 fn approval_requirement_for(&self, input: &serde_json::Value) -> ApprovalRequirement {
5367 let translated = contract_bash_legacy_input(input).unwrap_or_else(|_| input.clone());
5368 BashTool::contract_delegate().approval_requirement_for(&translated)
5369 }
5370
5371 fn is_read_only_for(&self, input: &serde_json::Value) -> bool {
5372 contract_bash_legacy_input(input)
5373 .is_ok_and(|translated| BashTool::contract_delegate().is_read_only_for(&translated))
5374 }
5375
5376 fn supports_parallel_for(&self, input: &serde_json::Value) -> bool {
5377 self.is_read_only_for(input)
5378 }
5379
5380 async fn execute(
5381 &self,
5382 input: serde_json::Value,
5383 context: &ToolContext,
5384 ) -> Result<ToolResult, ToolError> {
5385 let translated = contract_bash_legacy_input(&input)?;
5386 BashTool::contract_delegate()
5387 .execute(translated, context)
5388 .await
5389 }
5390 }
5391
5392 fn contract_bash_legacy_input(input: &serde_json::Value) -> Result<serde_json::Value, ToolError> {
5393 let object = input
5394 .as_object()
5395 .ok_or_else(|| ToolError::invalid_input("bash input must be an object"))?;
5396 let unexpected = object
5397 .keys()
5398 .filter(|key| {
5399 !matches!(
5400 key.as_str(),
5401 "command" | "timeout" | "read_only" | "sandbox_permissions" | "justification"
5402 )
5403 })
5404 .cloned()
5405 .collect::<Vec<_>>();
5406 if !unexpected.is_empty() {
5407 return Err(ToolError::invalid_input(format!(
5408 "unexpected bash parameter(s): {}",
5409 unexpected.join(", ")
5410 )));
5411 }
5412 let command = required_str(input, "command")?;
5413 let mut translated = json!({"command": command});
5414 if let Some(timeout) = input.get("timeout") {
5415 let seconds = timeout.as_f64().ok_or_else(|| {
5416 ToolError::invalid_input("Invalid timeout: expected a finite number of seconds")
5417 })?;
5418 if !seconds.is_finite() || seconds <= 0.0 {
5419 return Err(ToolError::invalid_input(
5420 "Invalid timeout: expected a positive finite number of seconds",
5421 ));
5422 }
5423 let millis = seconds * 1000.0;
5424 if millis > BASH_MAX_TIMEOUT_MS as f64 {
5425 return Err(ToolError::invalid_input(format!(
5426 "Invalid timeout: maximum is {} seconds",
5427 BASH_MAX_TIMEOUT_MS as f64 / 1000.0
5428 )));
5429 }
5430 translated["timeout_ms"] = json!((millis as u64).max(1));
5431 }
5432 if let Some(value) = input.get("read_only") {
5433 if !value.is_boolean() {
5434 return Err(type_mismatch("read_only", value, "a boolean"));
5435 }
5436 translated["read_only"] = value.clone();
5437 }
5438 for field in ["sandbox_permissions", "justification"] {
5439 if let Some(value) = input.get(field) {
5440 translated[field] = value.clone();
5441 }
5442 }
5443 Ok(translated)
5444 }
5445
5446 /// Compatibility shell tool retained for saved v0.9.x transcripts and the
5447 /// background/session control surface. It is hidden from new model catalogs.
5448 pub struct BashTool {
5449 name: &'static str,
5450 forced_action: Option<&'static str>,
5451 read_only: bool,
5452 optional_timeout: bool,
5453 }
5454
5455 pub(crate) fn readonly_bash_input_schema() -> serde_json::Value {
5456 json!({
5457 "type": "object",
5458 "properties": {
5459 "action": { "type": "string", "enum": ["run"] },
5460 "command": { "type": "string", "description": "A classifier-approved read command, or analysis code with read_only=true" },
5461 "read_only": { "type": "boolean", "description": "Require native filesystem read-only and no-network enforcement for analysis code; unavailable sandboxes fail closed." },
5462 "cwd": { "type": "string", "description": "Workspace-relative working directory" },
5463 "timeout_ms": { "type": "integer", "description": "Timeout in milliseconds (1000-600000)" }
5464 },
5465 "required": ["command"],
5466 "additionalProperties": false
5467 })
5468 }
5469
5470 impl BashTool {
5471 pub const fn new(name: &'static str) -> Self {
5472 Self {
5473 name,
5474 forced_action: None,
5475 read_only: false,
5476 optional_timeout: false,
5477 }
5478 }
5479
5480 pub const fn read_only(name: &'static str) -> Self {
5481 Self {
5482 name,
5483 forced_action: None,
5484 read_only: true,
5485 optional_timeout: false,
5486 }
5487 }
5488
5489 pub const fn alias(name: &'static str, action: &'static str) -> Self {
5490 Self {
5491 name,
5492 forced_action: Some(action),
5493 read_only: false,
5494 optional_timeout: false,
5495 }
5496 }
5497
5498 const fn contract_delegate() -> Self {
5499 Self {
5500 name: "bash",
5501 forced_action: Some("run"),
5502 read_only: false,
5503 optional_timeout: true,
5504 }
5505 }
5506 }
5507
5508 #[async_trait]
5509 impl ToolSpec for BashTool {
5510 fn name(&self) -> &'static str {
5511 self.name
5512 }
5513
5514 fn model_visible(&self) -> bool {
5515 false
5516 }
5517
5518 fn description(&self) -> &'static str {
5519 if self.read_only {
5520 "Inspect with classifier-bounded commands run directly as argv, never through a shell. For analysis code, read_only=true instead requires a native filesystem read-only and no-network sandbox. Only foreground action=run with command, cwd, timeout_ms, and read_only is accepted."
5521 } else {
5522 guidance::description()
5523 }
5524 }
5525
5526 fn input_schema(&self) -> serde_json::Value {
5527 if self.read_only {
5528 return readonly_bash_input_schema();
5529 }
5530 json!({
5531 "type": "object",
5532 "properties": {
5533 "action": {
5534 "type": "string",
5535 "enum": ["run", "wait", "interact", "cancel"],
5536 "description": "Action to perform (default: run)"
5537 },
5538 "command": {
5539 "type": "string",
5540 "description": guidance::runtime_command_guidance()
5541 },
5542 "read_only": { "type": "boolean", "description": "Set true to require native filesystem read-only and no-network execution. Only foreground run with command, cwd and timeout_ms; unavailable enforcement fails closed." },
5543 "timeout_ms": {
5544 "type": "integer",
5545 "description": "Timeout in milliseconds. The default depends on the action: action=run 120000 (the standalone Bash tool caps it at 600000), action=wait 30000, action=interact 1000. A foreground action=run that omits this is bounded by that default and killed with a background-rerun hint; pass an explicit value for longer foreground work, or background=true. For action=wait, `timeout_secs` (seconds) and `timeout` (milliseconds) are accepted aliases."
5546 },
5547 "background": {
5548 "type": "boolean",
5549 "description": "Temporary background; killed at session exit. Surviving headless services need background:true,persist:true. It is not killed at timeout_ms; plan to poll it with action=wait or stop it with action=cancel."
5550 },
5551 "interactive": {
5552 "type": "boolean",
5553 "description": "Run interactively with terminal IO (default: false)"
5554 },
5555 "stdin": {
5556 "type": "string",
5557 "description": "Stdin data to send (action=run: before waiting; action=interact: to the background task). Also accepted as `input` or `data` — send only one."
5558 },
5559 "input": {
5560 "type": "string",
5561 "description": "Alias for `stdin`."
5562 },
5563 "data": {
5564 "type": "string",
5565 "description": "Alias for `stdin`."
5566 },
5567 "cwd": {
5568 "type": "string",
5569 "description": "Optional working directory for the command"
5570 },
5571 "tty": {
5572 "type": "boolean",
5573 "description": "Allocate a pseudo-terminal for interactive programs (implies background)"
5574 },
5575 "combined_output": {
5576 "type": "boolean",
5577 "description": "Capture stdout and stderr as one chronological PTY stream (default false)"
5578 },
5579 "task_id": {
5580 "type": "string",
5581 "description": "Task ID for action=wait/interact/cancel. Also accepted as `id`."
5582 },
5583 "task_ids": {
5584 "type": "array",
5585 "items": { "type": "string" },
5586 "description": "For action=wait: wait on several background task ids at once (alternative to `task_id`)."
5587 },
5588 "until": {
5589 "type": "string",
5590 "enum": ["any", "all"],
5591 "description": "For action=wait with `task_ids`: return as soon as any task is terminal (`any`) or when every task is terminal (`all`, default)."
5592 },
5593 "id": {
5594 "type": "string",
5595 "description": "Alias for `task_id`."
5596 },
5597 "wait": {
5598 "type": "boolean",
5599 "description": "For action=wait, block until the task completes or timeout elapses (default: true). Pass false for a nonblocking snapshot; `block` is an accepted alias."
5600 },
5601 "close_stdin": {
5602 "type": "boolean",
5603 "description": "Close stdin after sending (action=interact)"
5604 },
5605 "all": {
5606 "type": "boolean",
5607 "description": "Cancel all running background tasks (action=cancel)"
5608 },
5609 "persist": {
5610 "type": "boolean",
5611 "description": "Keep this background service running after a successful headless exec (default: false). Requires background:true and explicit danger-full-access. Run the service itself in the foreground; do not use nohup or a trailing `&`."
5612 },
5613 "sandbox_permissions": {
5614 "type": "string",
5615 "enum": ["workspace-write", "danger-full-access"],
5616 "description": "The wider sandbox mode this exact command needs. Use only as a one-shot retry after a sandbox denial; requires justification and user approval in Ask."
5617 },
5618 "justification": {
5619 "type": "string",
5620 "description": "Required with sandbox_permissions: one sentence explaining why this exact command needs wider access."
5621 }
5622 },
5623 // The schema used to declare nothing required at all, so
5624 // `Bash{}` was schema-valid for the tool that runs shell
5625 // commands. What is required is per-action and cannot be spelled
5626 // as a flat `required` list: `run` needs `command`,
5627 // `wait`/`interact`/`cancel` need `task_id` (or its `id` alias),
5628 // and `cancel` needs `all` instead when cancelling everything.
5629 // A root `anyOf` of `required` groups is how this repo already
5630 // spells that (`finance`, `apply_patch`), and `schema_sanitize`
5631 // knows the shape: providers that reject root composition get the
5632 // groups merged and the constraint restated as a description note
5633 // (`root_composition_constraint_note`). The cost is that
5634 // `strict_schema_supported` rejects a root `anyOf`, so `Bash`
5635 // opts out of DeepSeek strict mode — as `finance` already does on
5636 // the same default agent surface, which turns strict mode off for
5637 // the whole tool set regardless.
5638 "anyOf": [
5639 { "required": ["command"] },
5640 { "required": ["task_id"] },
5641 { "required": ["id"] },
5642 { "required": ["all"] }
5643 ]
5644 })
5645 }
5646
5647 fn capabilities(&self) -> Vec<ToolCapability> {
5648 vec![
5649 ToolCapability::ExecutesCode,
5650 ToolCapability::Sandboxable,
5651 ToolCapability::RequiresApproval,
5652 ]
5653 }
5654
5655 fn approval_requirement(&self) -> ApprovalRequirement {
5656 ApprovalRequirement::Required
5657 }
5658
5659 fn approval_requirement_for(&self, input: &serde_json::Value) -> ApprovalRequirement {
5660 if exec_shell_input_is_parallel_readonly(input) {
5661 ApprovalRequirement::Auto
5662 } else {
5663 self.approval_requirement()
5664 }
5665 }
5666
5667 fn is_read_only_for(&self, input: &serde_json::Value) -> bool {
5668 exec_shell_input_is_parallel_readonly(input)
5669 }
5670
5671 fn supports_parallel_for(&self, input: &serde_json::Value) -> bool {
5672 exec_shell_input_is_parallel_readonly(input)
5673 }
5674
5675 fn starts_detached_for(&self, input: &serde_json::Value) -> bool {
5676 exec_shell_input_starts_detached(input)
5677 }
5678
5679 async fn execute(
5680 &self,
5681 input: serde_json::Value,
5682 context: &ToolContext,
5683 ) -> Result<ToolResult, ToolError> {
5684 // A leading `cd <dir> &&` becomes the working directory before any
5685 // check, exactly as the read-only gates judged it.
5686 let input = if context.shell_policy == ShellPolicy::ReadOnly {
5687 normalize_readonly_cd(&input)
5688 } else {
5689 input
5690 };
5691 // `and_then(as_str).unwrap_or("run")` treated *any* non-string
5692 // `action` as absent and fell through to the branch that runs
5693 // arbitrary code: `Bash{action: 3, command: "…"}` executed the
5694 // command. Every sibling family refuses a non-string action
5695 // (`canonical_action::required_action`), and `Bash` cannot be the
5696 // lenient one. `optional_str` is the type-strictness lane's extractor:
5697 // absent or `null` takes the documented `run` default, anything else
5698 // is a `type_mismatch` naming the field and the type it needed.
5699 let action = match self.forced_action {
5700 Some(forced) => forced,
5701 None => optional_str(&input, "action")?.unwrap_or("run"),
5702 };
5703 let enforced_readonly = match input.get("read_only") {
5704 None => false,
5705 Some(serde_json::Value::Bool(enabled)) => *enabled,
5706 Some(value) => return Err(type_mismatch("read_only", value, "a boolean")),
5707 };
5708 if enforced_readonly && (!enforced_readonly_input(&input) || action != "run") {
5709 return Err(ToolError::invalid_input(
5710 "read_only=true accepts only foreground run with command, cwd and timeout_ms; background, input, interactive modes and sandbox escalation are incompatible",
5711 ));
5712 }
5713 if enforced_readonly {
5714 if context.sandbox_backend.is_some() {
5715 return Err(ToolError::permission_denied(
5716 "read_only execution requires a native enforcing sandbox; external backends cannot attest this policy",
5717 ));
5718 }
5719 let manager = context
5720 .shell_manager
5721 .lock()
5722 .map_err(|_| ToolError::execution_failed("shell manager lock poisoned"))?;
5723 if !manager
5724 .configured_sandbox_type()
5725 .is_some_and(is_native_readonly_sandbox)
5726 {
5727 return Err(ToolError::not_available(
5728 "read_only execution requires native read-only enforcement (macOS Seatbelt or configured Linux bubblewrap); nothing was run",
5729 ));
5730 }
5731 }
5732 let mut policy_input = input.clone();
5733 if let Some(object) = policy_input.as_object_mut() {
5734 object.insert("action".into(), json!(action));
5735 }
5736 crate::core::engine::tool_catalog::enforce_tool_denial(
5737 context,
5738 self.name(),
5739 &policy_input,
5740 )?;
5741 if action == "interact" && context.shell_policy != ShellPolicy::Full {
5742 return Err(ToolError::permission_denied(
5743 "Sending shell input requires full shell permission.",
5744 ));
5745 }
5746 match action {
5747 "wait" => return self.execute_wait(&input, context).await,
5748 "interact" => return self.execute_interact(&input, context).await,
5749 "cancel" => return self.execute_cancel(&input, context).await,
5750 "run" => {}
5751 // Bash was the only action wrapper whose catch-all fell through to
5752 // its most dangerous branch: `{"action":"kill", "command":…}` ran
5753 // the command instead of cancelling, and a mis-cased "Cancel" did
5754 // the same. Every sibling (`File`, `Git`, `Web`, `Run`) already
5755 // refuses an unknown action; the tool that executes arbitrary code
5756 // should not be the lenient one.
5757 other => {
5758 return Err(ToolError::invalid_input(format!(
5759 "Unknown Bash action \"{other}\"; nothing was run. Pass one of: run, wait, interact, cancel."
5760 )));
5761 }
5762 }
5763 let command = require_no_nul(required_str(&input, "command")?, "command")?;
5764 match context.shell_policy {
5765 ShellPolicy::None => {
5766 return Ok(ToolResult::error(
5767 "Shell tools are disabled by the active permission profile.",
5768 ));
5769 }
5770 ShellPolicy::ReadOnly => {
5771 if let Err(rejection) = exec_shell_input_agent_readonly_verdict(&input) {
5772 // A typed denial, so a Fleet worker's no-progress guard
5773 // counts it. #6298: an agent has no mode to switch to,
5774 // so it gets the same next steps as the other read-only
5775 // gates; only a parent session is pointed at Work mode.
5776 let message = if context.owner_agent_id.is_some()
5777 || context.tool_authority.is_some()
5778 {
5779 readonly_refusal(&rejection, readonly_enforced_lane_available(context))
5780 } else {
5781 format!(
5782 "{rejection}. Use a read-only inspection command, or switch to Work mode (`/mode work`) for write-capable shell work."
5783 )
5784 };
5785 return Err(ToolError::permission_denied(message));
5786 }
5787 }
5788 ShellPolicy::Full => {}
5789 }
5790 enforce_readonly_network_reads(command, context)?;
5791 let requested_timeout_ms = if self.optional_timeout {
5792 input
5793 .get("timeout_ms")
5794 .map(|value| {
5795 value
5796 .as_u64()
5797 .ok_or_else(|| type_mismatch("timeout_ms", value, "a positive integer"))
5798 })
5799 .transpose()?
5800 } else {
5801 Some(optional_u64(&input, "timeout_ms", 120_000)?.min(600_000))
5802 };
5803 let background = optional_bool(&input, "background", false)?;
5804 let interactive = optional_bool(&input, "interactive", false)?;
5805 let combined_output = optional_bool(&input, "combined_output", false)?;
5806 let tty = optional_bool(&input, "tty", false)? || (combined_output && background);
5807 let timeout_ms = contract_bash_timeout_ms(
5808 self.optional_timeout,
5809 requested_timeout_ms,
5810 background,
5811 interactive,
5812 );
5813 let timeout_value_ms = timeout_ms.unwrap_or(BASH_MAX_TIMEOUT_MS);
5814 // Strict types (2026-08-04 review): a non-string here used to be
5815 // silently dropped — the command then ran with NO stdin and reported
5816 // success, the exact silent-drop failure the alias hardening closed
5817 // for misspelled names. A wrong type is an error, never a no-op.
5818 let stdin_data = match first_present_field(&input, &["stdin", "input", "data"]) {
5819 None => None,
5820 Some((name, value)) => Some(
5821 value
5822 .as_str()
5823 .ok_or_else(|| type_mismatch(name, value, "a string"))?
5824 .to_string(),
5825 ),
5826 };
5827
5828 if interactive && background {
5829 return Ok(ToolResult::error(
5830 "Interactive commands cannot run in background mode.",
5831 ));
5832 }
5833 if interactive && (tty || combined_output) {
5834 return Ok(ToolResult::error(
5835 "Interactive mode cannot be combined with TTY or combined_output sessions.",
5836 ));
5837 }
5838 if interactive && stdin_data.is_some() {
5839 return Ok(ToolResult::error(
5840 "Interactive mode cannot be combined with stdin data.",
5841 ));
5842 }
5843
5844 let persist = optional_bool(&input, "persist", false)?;
5845 if persist {
5846 if !background {
5847 return Err(ToolError::invalid_input(
5848 "persist:true requires background:true; a persisted service must be started as a background task.",
5849 ));
5850 }
5851 if interactive || tty {
5852 return Err(ToolError::invalid_input(
5853 "persist:true cannot be combined with interactive or TTY modes.",
5854 ));
5855 }
5856 if stdin_data.is_some() {
5857 return Err(ToolError::invalid_input(
5858 "persist:true spawns the service with null stdio; stdin data is not accepted.",
5859 ));
5860 }
5861 if !persistent_services_enabled_for(context) {
5862 return Err(ToolError::not_available(
5863 "persistent background services (persist:true) are only available on Unix in the real headless `codewhale exec` host under an explicit danger-full-access / full shell authority. They are rejected in interactive sessions, desktop/app-server hosts, Fleet/sub-agents, restricted or external sandboxes, and TTY/interactive/stdin modes.",
5864 ));
5865 }
5866 }
5867
5868 let background = background || tty;
5869
5870 let mut execpolicy_decision: Option<RuleDecision> = None;
5871 if context.features.enabled(Feature::ExecPolicy)
5872 && let Some(policy) = tokio::task::spawn_blocking(load_default_policy)
5873 .await
5874 .map_err(|e| {
5875 ToolError::execution_failed(format!("execpolicy load task failed: {e}"))
5876 })?
5877 .map_err(|e| ToolError::execution_failed(format!("execpolicy load failed: {e}")))?
5878 {
5879 let decision = policy.evaluate(command);
5880 execpolicy_decision = Some(decision.clone());
5881 if let RuleDecision::Deny(reason) = decision {
5882 return Ok(ToolResult {
5883 content: format!("BLOCKED: {reason}"),
5884 success: false,
5885 metadata: Some(json!({
5886 "execpolicy": {
5887 "decision": "deny",
5888 "reason": reason,
5889 }
5890 })),
5891 });
5892 }
5893 }
5894
5895 // Safety analysis (always run for metadata, but only block when not in YOLO mode)
5896 let safety = analyze_command(command);
5897 if !context.auto_approve {
5898 match safety.level {
5899 SafetyLevel::Dangerous => {
5900 let reasons = safety.reasons.join("; ");
5901 let suggestions = if safety.suggestions.is_empty() {
5902 String::new()
5903 } else {
5904 format!("\nSuggestions: {}", safety.suggestions.join("; "))
5905 };
5906 return Ok(ToolResult {
5907 content: format!(
5908 "BLOCKED: This command was blocked for safety reasons.\n\nReasons: {reasons}{suggestions}\n\nNote: allow_shell=true exposes shell tools, but it does not disable built-in shell safety validation."
5909 ),
5910 success: false,
5911 metadata: Some(json!({
5912 "safety_level": "dangerous",
5913 "blocked": true,
5914 "reasons": safety.reasons,
5915 "suggestions": safety.suggestions,
5916 })),
5917 });
5918 }
5919 SafetyLevel::RequiresApproval | SafetyLevel::Safe | SafetyLevel::WorkspaceSafe => {
5920 // Proceed normally
5921 }
5922 }
5923 }
5924
5925 // This explicit mode only narrows the caller's posture. No approval or
5926 // inherited full-access override can disable the required sandbox.
5927 let policy_override = if enforced_readonly {
5928 Some(ExecutionSandboxPolicy::ReadOnly)
5929 } else {
5930 context.elevated_sandbox_policy.clone()
5931 };
5932 // Strict types: a non-string cwd used to silently run the command in
5933 // the workspace default instead of erroring (2026-08-04 review).
5934 let working_dir = match first_present_field(&input, &["cwd", "working_dir"])
5935 .map(|(name, value)| {
5936 value
5937 .as_str()
5938 .ok_or_else(|| type_mismatch(name, value, "a string"))
5939 .and_then(|dir| require_no_nul(dir, name))
5940 })
5941 .transpose()?
5942 {
5943 Some(dir) => {
5944 // Validate cwd against workspace boundary (same as file tools)
5945 let resolved = context.resolve_path(dir)?;
5946 Some(resolved.to_string_lossy().to_string())
5947 }
5948 // Default to the tool context's workspace (which reflects the
5949 // child agent's worktree when `worktree: true` was used), not the
5950 // shared ShellManager's parent-workspace default_workspace.
5951 None => Some(context.workspace.display().to_string()),
5952 };
5953 if matches!(context.shell_policy, ShellPolicy::ReadOnly) && !enforced_readonly {
5954 let effective_cwd = working_dir
5955 .as_deref()
5956 .map(std::path::Path::new)
5957 .unwrap_or(&context.workspace);
5958 enforce_readonly_workspace_operands(command, &context.workspace, effective_cwd)?;
5959 }
5960
5961 // #456 — collect env from any configured `shell_env` hooks. Runs
5962 // synchronously, captures stdout, parses `KEY=VAL` lines, audit-logs
5963 // the keys (never the values). Empty / no-op when no hook is
5964 // configured.
5965 let read_only_shell =
5966 matches!(context.shell_policy, ShellPolicy::ReadOnly) || enforced_readonly;
5967 let mut extra_env = if read_only_shell {
5968 // shell_env hooks are arbitrary operator-configured processes.
5969 // They cannot run inside the evidence-only execution boundary.
5970 HashMap::new()
5971 } else if let Some(hook_executor) = &context.runtime.hook_executor {
5972 let hook_ctx = crate::hooks::HookContext::new()
5973 .with_tool_context(context)
5974 .with_tool_name("exec_shell")
5975 .with_tool_args(&input);
5976 let executor = Arc::clone(hook_executor);
5977 let policy = crate::plugins::activation::extension_host_policy_enabled();
5978 #[cfg(test)]
5979 let env_scope = crate::test_support::env_scope_ticket();
5980 tokio::task::spawn_blocking(move || {
5981 let _policy = crate::plugins::activation::PolicyScope::propagate(policy);
5982 #[cfg(test)]
5983 let _env_scope = crate::test_support::join_env_scope(env_scope);
5984 executor.collect_shell_env(&hook_ctx)
5985 })
5986 .await
5987 .unwrap_or_default()
5988 } else {
5989 std::collections::HashMap::new()
5990 };
5991 if read_only_shell {
5992 let null_device = if cfg!(windows) { "NUL" } else { "/dev/null" };
5993 let inert_git_helper = if cfg!(windows) {
5994 "cmd.exe /d /c exit 1"
5995 } else {
5996 "/usr/bin/false"
5997 };
5998 // Read-only Bash is intentionally a small inspection surface. Git
5999 // can otherwise invoke operator/repository configured helpers
6000 // while performing nominal reads (a pager or fsmonitor), and it
6001 // may opportunistically refresh the index. These environment
6002 // overrides make those reads non-interactive and suppress the
6003 // optional mutation/extension seams; the command classifier and
6004 // machine authority gate remain authoritative as well.
6005 extra_env.insert("GIT_OPTIONAL_LOCKS".to_string(), "0".to_string());
6006 extra_env.insert("GIT_NO_LAZY_FETCH".to_string(), "1".to_string());
6007 extra_env.insert("GIT_PAGER".to_string(), String::new());
6008 extra_env.insert("GH_PAGER".to_string(), String::new());
6009 extra_env.insert("GH_PROMPT_DISABLED".to_string(), "1".to_string());
6010 extra_env.insert("GH_NO_UPDATE_NOTIFIER".to_string(), "1".to_string());
6011 // The classifier rejects explicit GHES repo/URL targets. Pin the
6012 // implicit environment side too, so inherited GH_HOST/GH_REPO
6013 // cannot redirect the call after api.github.com was approved.
6014 extra_env.insert("GH_HOST".to_string(), "github.com".to_string());
6015 extra_env.insert("GH_REPO".to_string(), String::new());
6016 extra_env.insert("PAGER".to_string(), String::new());
6017 extra_env.insert("ENV".to_string(), String::new());
6018 extra_env.insert("BASH_ENV".to_string(), String::new());
6019 extra_env.insert("CDPATH".to_string(), String::new());
6020 extra_env.insert("RIPGREP_CONFIG_PATH".to_string(), String::new());
6021 // Ignore user/system Git configuration and replace any repository
6022 // external diff helper with a fixed inert executable. Repository
6023 // config and attributes are attacker-controlled evidence inputs;
6024 // a nominal `git diff/log/show` must not turn them into programs.
6025 extra_env.insert("GIT_CONFIG_NOSYSTEM".to_string(), "1".to_string());
6026 extra_env.insert("GIT_CONFIG_GLOBAL".to_string(), null_device.to_string());
6027 extra_env.insert("GIT_CONFIG_PARAMETERS".to_string(), String::new());
6028 extra_env.insert(
6029 "GIT_EXTERNAL_DIFF".to_string(),
6030 inert_git_helper.to_string(),
6031 );
6032 extra_env.insert("GIT_ATTR_NOSYSTEM".to_string(), "1".to_string());
6033 if let Some(path) = readonly_sanitized_path(&context.workspace) {
6034 extra_env.insert("PATH".to_string(), path);
6035 }
6036 extra_env.insert("GIT_CONFIG_KEY_0".to_string(), "core.fsmonitor".to_string());
6037 extra_env.insert("GIT_CONFIG_VALUE_0".to_string(), "false".to_string());
6038 extra_env.insert("GIT_CONFIG_KEY_1".to_string(), "core.hooksPath".to_string());
6039 extra_env.insert("GIT_CONFIG_VALUE_1".to_string(), null_device.to_string());
6040 extra_env.insert(
6041 "GIT_CONFIG_KEY_2".to_string(),
6042 "log.showSignature".to_string(),
6043 );
6044 extra_env.insert("GIT_CONFIG_VALUE_2".to_string(), "false".to_string());
6045 // A working-tree read (diff, status, blame) runs the repository's
6046 // clean filters, which no command-line flag disables.
6047 let git_dirs = readonly_git_dirs(
6048 command,
6049 working_dir
6050 .as_deref()
6051 .map_or(context.workspace.as_path(), std::path::Path::new),
6052 );
6053 let overrides = tokio::task::spawn_blocking(move || {
6054 let mut overrides = std::collections::BTreeSet::new();
6055 for dir in git_dirs {
6056 overrides.extend(crate::dependencies::Git::review_filter_overrides(&dir)?);
6057 }
6058 anyhow::Ok(overrides)
6059 })
6060 .await
6061 .map_err(|e| ToolError::execution_failed(format!("git task panicked: {e}")))?
6062 .map_err(|e| {
6063 ToolError::permission_denied(format!(
6064 "Read-only shell could not inspect the repository's Git filters: {e:#}"
6065 ))
6066 })?;
6067 let mut count = 3;
6068 for (key, value) in overrides {
6069 extra_env.insert(format!("GIT_CONFIG_KEY_{count}"), key);
6070 extra_env.insert(format!("GIT_CONFIG_VALUE_{count}"), value.to_string());
6071 count += 1;
6072 }
6073 extra_env.insert("GIT_CONFIG_COUNT".to_string(), count.to_string());
6074 extra_env.insert(READONLY_ENV_MARKER.to_string(), "1".to_string());
6075 }
6076
6077 let command_expense = infer_command_expense(command);
6078 let heavy_permit = acquire_heavy_command_permit(command, context.cancel_token.as_ref())
6079 .await
6080 .map_err(|error| ToolError::execution_failed(error.to_string()))?;
6081 let admission_wait_ms = heavy_permit
6082 .as_ref()
6083 .map(|permit| u64::try_from(permit.queued_for().as_millis()).unwrap_or(u64::MAX));
6084 let admission_limit = heavy_permit.as_ref().map(HeavyCommandPermit::limit);
6085 let admission_memory = heavy_permit
6086 .as_ref()
6087 .map(HeavyCommandPermit::memory_pressure);
6088
6089 // Route through external sandbox backend when configured.
6090 if let Some(backend) = &context.sandbox_backend {
6091 if self.optional_timeout {
6092 return Err(ToolError::not_available(
6093 "bash is unavailable with this external sandbox backend because it cannot preserve combined streaming output and timeout semantics. Use the native sandbox or search for the backend-specific shell tool.",
6094 ));
6095 }
6096 if matches!(context.shell_policy, ShellPolicy::ReadOnly) {
6097 return Err(ToolError::permission_denied(
6098 "Read-only Scout shell cannot use an external sandbox backend because that interface accepts a raw command string rather than the classifier-approved argv. Use File read/search, or run this Scout without the external backend.",
6099 ));
6100 }
6101 if interactive {
6102 return Ok(ToolResult::error(
6103 "Interactive mode is not supported with external sandbox backends.",
6104 ));
6105 }
6106 if background {
6107 return Ok(ToolResult::error(
6108 "Background mode is not supported with external sandbox backends.",
6109 ));
6110 }
6111 if tty {
6112 return Ok(ToolResult::error(
6113 "TTY mode is not supported with external sandbox backends.",
6114 ));
6115 }
6116
6117 let started = std::time::Instant::now();
6118 let backend_result = backend.exec(command, &extra_env).await;
6119
6120 let result = match backend_result {
6121 Ok(output) => {
6122 let (stdout, stdout_meta) = truncate_with_meta(&output.stdout);
6123 let (stderr, stderr_meta) = truncate_with_meta(&output.stderr);
6124 ShellResult {
6125 task_id: None,
6126 status: if output.exit_code == 0 {
6127 ShellStatus::Completed
6128 } else {
6129 ShellStatus::Failed
6130 },
6131 exit_code: Some(i64::from(output.exit_code)),
6132 stdout,
6133 stderr,
6134 duration_ms: u64::try_from(started.elapsed().as_millis())
6135 .unwrap_or(u64::MAX),
6136 stdout_len: stdout_meta.original_len,
6137 stderr_len: stderr_meta.original_len,
6138 stdout_omitted: stdout_meta.omitted,
6139 stderr_omitted: stderr_meta.omitted,
6140 stdout_truncated: stdout_meta.truncated,
6141 stderr_truncated: stderr_meta.truncated,
6142 sandboxed: true,
6143 sandbox_type: Some("opensandbox".to_string()),
6144 sandbox_denied: false,
6145 }
6146 }
6147 Err(e) => {
6148 return Ok(ToolResult::error(format!("Sandbox backend error: {e}")));
6149 }
6150 };
6151
6152 // Build result (reuse the existing output rendering below).
6153 let stdout_summary = summarize_output(&result.stdout);
6154 let stderr_summary = summarize_output(&result.stderr);
6155 let summary = if !stderr_summary.is_empty() {
6156 stderr_summary.clone()
6157 } else {
6158 stdout_summary.clone()
6159 };
6160 let python_dependency_hint = python_build_dependency_hint(command, &result);
6161 let mut output = if result.stdout.is_empty() && result.stderr.is_empty() {
6162 "(no output)".to_string()
6163 } else if result.stderr.is_empty() {
6164 result.stdout.clone()
6165 } else {
6166 format!("{}\n\nSTDERR:\n{}", result.stdout, result.stderr)
6167 };
6168 if let Some(hint) = python_dependency_hint {
6169 output = format!("{hint}\n\n{output}");
6170 }
6171
6172 let mut metadata = json!({
6173 "exit_code": result.exit_code,
6174 "exit_code_hex": exit_code_hex(result.exit_code),
6175 "status": format!("{:?}", result.status),
6176 "duration_ms": result.duration_ms,
6177 "sandboxed": true,
6178 "sandbox_type": "opensandbox",
6179 "sandbox_denied": false,
6180 "task_id": result.task_id,
6181 "stdout_len": result.stdout_len,
6182 "stderr_len": result.stderr_len,
6183 "stdout_truncated": result.stdout_truncated,
6184 "stderr_truncated": result.stderr_truncated,
6185 "stdout_omitted": result.stdout_omitted,
6186 "stderr_omitted": result.stderr_omitted,
6187 "summary": summary,
6188 "stdout_summary": stdout_summary,
6189 "stderr_summary": stderr_summary,
6190 "safety_level": format!("{:?}", safety.level),
6191 "interactive": false,
6192 "canceled": false,
6193 "sandbox_backend": backend.kind().as_str(),
6194 "expense_class": match command_expense {
6195 CommandExpense::Heavy => "heavy",
6196 CommandExpense::Normal => "normal",
6197 },
6198 "resource_admission_wait_ms": admission_wait_ms,
6199 "resource_admission_limit": admission_limit,
6200 });
6201 attach_shell_owner_metadata(&mut metadata, context);
6202 attach_cargo_failure_summary(&mut metadata, command, &result);
6203 attach_python_build_dependency_hint(&mut metadata, python_dependency_hint);
6204
6205 return Ok(ToolResult {
6206 content: output,
6207 success: result.status == ShellStatus::Completed,
6208 metadata: Some(metadata),
6209 });
6210 }
6211
6212 let mut lifecycle_warning = None;
6213 let mut receipt_identity = None;
6214 // #6689: the receipt exists for completion hooks to read. Without one
6215 // registered it would only ride along in tool metadata, which the
6216 // Runtime API persists and emits for every call.
6217 let wants_receipt = context.runtime.hook_executor.as_ref().is_some_and(|hooks| {
6218 hooks.has_hooks_for_event(crate::hooks::HookEvent::ToolCallAfter)
6219 || hooks.has_hooks_for_event(crate::hooks::HookEvent::OnError)
6220 });
6221 let result = if interactive {
6222 let mut manager = context
6223 .shell_manager
6224 .lock()
6225 .map_err(|_| ToolError::execution_failed("shell manager lock poisoned"))?;
6226 let work_lifecycle = shell_work_lifecycle_from_context(context);
6227 let task_id = format!("shell_{}", &Uuid::new_v4().to_string()[..8]);
6228 let mut spawn_guard = ShellSpawnIntentGuard::new(work_lifecycle, &task_id, command);
6229 // Only a registered operation is observed (#6435).
6230 let work_lifecycle = spawn_guard.lifecycle.clone();
6231 let result = manager.execute_interactive_with_policy_env(
6232 command,
6233 working_dir.as_deref(),
6234 timeout_value_ms,
6235 policy_override,
6236 extra_env,
6237 );
6238 match result {
6239 Ok(result) => {
6240 // The process result is authoritative once execution has
6241 // completed. Disarm before observing it so a graph-write
6242 // failure cannot relabel a successful command as Failed.
6243 spawn_guard.disarm();
6244 if let Some(lifecycle) = work_lifecycle.as_ref() {
6245 let raw_bytes = result.stdout_len.saturating_add(result.stderr_len);
6246 if let Err(err) = lifecycle.observe(&task_id, &result.status, 1, raw_bytes)
6247 {
6248 tracing::warn!(shell_id = %task_id, error = %err, "interactive shell completed but Work lifecycle reconciliation failed");
6249 lifecycle_warning = Some(err.to_string());
6250 }
6251 }
6252 Ok(result)
6253 }
6254 Err(err) => Err(err),
6255 }
6256 } else if background {
6257 let mut manager = context
6258 .shell_manager
6259 .lock()
6260 .map_err(|_| ToolError::execution_failed("shell manager lock poisoned"))?;
6261 let result = manager.execute_with_options_env_for_owner_and_work(
6262 command,
6263 working_dir.as_deref(),
6264 timeout_value_ms,
6265 true,
6266 stdin_data.as_deref(),
6267 tty,
6268 policy_override,
6269 extra_env,
6270 shell_job_owner_from_context(context),
6271 context.state_namespace.clone(),
6272 context.origin_tool_call_id.clone(),
6273 context.origin_turn_id.clone(),
6274 shell_work_lifecycle_from_context(context),
6275 None,
6276 persist,
6277 (1_000, 600_000),
6278 );
6279 if let (Ok(result), Some(permit)) = (&result, heavy_permit)
6280 && let Some(task_id) = result.task_id.as_deref()
6281 {
6282 manager
6283 .attach_heavy_permit(task_id, permit)
6284 .map_err(|error| ToolError::execution_failed(error.to_string()))?;
6285 }
6286 result
6287 } else {
6288 execute_foreground_via_background(
6289 context,
6290 command,
6291 heavy_permit,
6292 working_dir,
6293 timeout_ms,
6294 stdin_data.as_deref(),
6295 combined_output,
6296 policy_override,
6297 extra_env,
6298 matches!(context.shell_policy, ShellPolicy::ReadOnly) && !enforced_readonly,
6299 if self.optional_timeout {
6300 (1, BASH_MAX_TIMEOUT_MS)
6301 } else {
6302 (1_000, 600_000)
6303 },
6304 wants_receipt,
6305 &mut receipt_identity,
6306 )
6307 .await
6308 };
6309
6310 match result {
6311 Ok(result) => {
6312 let backgrounded_foreground =
6313 !background && !interactive && result.status == ShellStatus::Running;
6314 if (background || backgrounded_foreground)
6315 && let (Some(shell_id), Some(task_id)) = (
6316 result.task_id.as_deref(),
6317 context.runtime.active_task_id.clone(),
6318 )
6319 && let Ok(mut manager) = context.shell_manager.lock()
6320 {
6321 let _ = manager.tag_linked_task(shell_id, Some(task_id));
6322 }
6323
6324 let was_cancelled = context
6325 .cancel_token
6326 .as_ref()
6327 .is_some_and(|token| token.is_cancelled());
6328 // Lowercase `bash` joins stdout and stderr on one pipe, so
6329 // its preview is combined output, not stdout.
6330 let execution_receipt = receipt_identity.as_ref().and_then(|identity| {
6331 shell_execution_receipt(
6332 identity,
6333 &result,
6334 if self.optional_timeout {
6335 "combined"
6336 } else {
6337 "separate"
6338 },
6339 )
6340 });
6341 if self.optional_timeout {
6342 return finish_contract_bash_result(
6343 result,
6344 timeout_ms,
6345 context,
6346 execution_receipt,
6347 );
6348 }
6349 let task_id_str = result.task_id.clone().unwrap_or_default();
6350 let stdout_summary = summarize_output(&result.stdout);
6351 let stderr_summary = summarize_output(&result.stderr);
6352 let summary = if !stderr_summary.is_empty() {
6353 stderr_summary.clone()
6354 } else {
6355 stdout_summary.clone()
6356 };
6357 let network_restricted_hint =
6358 shell_network_restricted_hint(context, command, &result).map(str::to_string);
6359 let sandbox_denied_hint = if network_restricted_hint.is_none() {
6360 shell_sandbox_denied_hint(context, &result)
6361 } else {
6362 None
6363 };
6364 let provenance_hint = macos_provenance_hint(&result);
6365 let python_dependency_hint = python_build_dependency_hint(command, &result);
6366 let mut output = if interactive {
6367 format!(
6368 "Interactive command completed (exit code: {:?})",
6369 result.exit_code
6370 )
6371 } else if result.status == ShellStatus::Completed {
6372 if result.stdout.is_empty() && result.stderr.is_empty() {
6373 "(no output)".to_string()
6374 } else if result.stderr.is_empty() {
6375 result.stdout.clone()
6376 } else {
6377 format!("{}\n\nSTDERR:\n{}", result.stdout, result.stderr)
6378 }
6379 } else if persist && result.status == ShellStatus::Running {
6380 format!(
6381 "Persistent service staged: {task_id_str}. Probe readiness with a separate command. Codewhale will transfer ownership only if this exec finishes successfully."
6382 )
6383 } else if result.status == ShellStatus::Running {
6384 let completion_contract = if context.owner_agent_id.is_some() {
6385 "completion stays in task/status and is not injected into the parent model."
6386 } else {
6387 "completion is delivered to the model as an internal runtime event and shown in task/status state."
6388 };
6389 if backgrounded_foreground {
6390 format!(
6391 "Foreground shell wait moved to /jobs: {task_id_str}\n\nReturns immediately; {completion_contract} Keep working; call Bash action=\"wait\" task_id=\"{task_id_str}\" at a true dependency to block until completion or timeout."
6392 )
6393 } else {
6394 format!(
6395 "Background task started: {task_id_str}\n\nReturns immediately; {completion_contract} Codewhale terminates this task when the session exits. If a service must survive a successful headless exec, start it with background=true and persist=true. Keep working; call Bash action=\"wait\" task_id=\"{task_id_str}\" at a true dependency to block until completion or timeout."
6396 )
6397 }
6398 } else if result.status == ShellStatus::Killed && was_cancelled {
6399 format!(
6400 "Command canceled; process killed.\n\nSTDOUT:\n{}\n\nSTDERR:\n{}",
6401 result.stdout, result.stderr
6402 )
6403 } else if result.status == ShellStatus::TimedOut {
6404 format!(
6405 "Command timed out after {timeout_value_ms}ms; process killed.\n\n{FOREGROUND_TIMEOUT_RECOVERY_HINT}\n\nSTDOUT:\n{}\n\nSTDERR:\n{}",
6406 result.stdout, result.stderr
6407 )
6408 } else {
6409 format!(
6410 "Command failed ({})\n\nSTDOUT:\n{}\n\nSTDERR:\n{}",
6411 exit_code_label(result.exit_code),
6412 result.stdout,
6413 result.stderr
6414 )
6415 };
6416 if let Some(hint) = network_restricted_hint.as_deref() {
6417 output = format!("{hint}\n\n{output}");
6418 }
6419 if let Some(hint) = sandbox_denied_hint.as_deref() {
6420 output = format!("{hint}\n\n{output}");
6421 }
6422 if let Some(hint) = provenance_hint {
6423 output = format!("{hint}\n\n{output}");
6424 }
6425 if let Some(hint) = python_dependency_hint {
6426 output = format!("{hint}\n\n{output}");
6427 }
6428
6429 let mut metadata = json!({
6430 "exit_code": result.exit_code,
6431 "exit_code_hex": exit_code_hex(result.exit_code),
6432 "status": format!("{:?}", result.status),
6433 "duration_ms": result.duration_ms,
6434 "sandboxed": result.sandboxed,
6435 "sandbox_type": result.sandbox_type,
6436 "sandbox_denied": result.sandbox_denied,
6437 "task_id": result.task_id,
6438 "stdout_len": result.stdout_len,
6439 "stderr_len": result.stderr_len,
6440 "stdout_truncated": result.stdout_truncated,
6441 "stderr_truncated": result.stderr_truncated,
6442 "stdout_omitted": result.stdout_omitted,
6443 "stderr_omitted": result.stderr_omitted,
6444 "lifecycle_warning": lifecycle_warning,
6445 "expense_class": match command_expense {
6446 CommandExpense::Heavy => "heavy",
6447 CommandExpense::Normal => "normal",
6448 },
6449 "resource_admission_wait_ms": admission_wait_ms,
6450 "resource_admission_limit": admission_limit,
6451 "resource_admission_memory": match admission_memory {
6452 Some(MemoryPressure::Critical) => "critical",
6453 Some(MemoryPressure::Constrained) => "constrained",
6454 Some(MemoryPressure::Nominal) | None => "nominal",
6455 Some(MemoryPressure::Unknown) => "unknown",
6456 },
6457 "summary": summary,
6458 "stdout_summary": stdout_summary,
6459 "stderr_summary": stderr_summary,
6460 "safety_level": format!("{:?}", safety.level),
6461 "interactive": interactive,
6462 "combined_output": combined_output,
6463 "canceled": was_cancelled,
6464 "execpolicy": execpolicy_decision.as_ref().map(|decision| match decision {
6465 RuleDecision::Allow => json!({
6466 "decision": "allow",
6467 }),
6468 RuleDecision::Deny(reason) => json!({
6469 "decision": "deny",
6470 "reason": reason,
6471 }),
6472 RuleDecision::AskUser(reason) => json!({
6473 "decision": "ask_user",
6474 "reason": reason,
6475 }),
6476 }),
6477 });
6478 if let Some(receipt) = execution_receipt {
6479 metadata["execution_receipt"] = receipt;
6480 }
6481 metadata["backgrounded"] = json!(background || backgrounded_foreground);
6482 if persist {
6483 metadata["persist_requested"] = json!(true);
6484 metadata["ownership"] = json!("managed_pending_exec_success");
6485 metadata["background_policy"] = json!("pending_ownership_transfer");
6486 metadata["auto_resume_on_completion"] = json!(false);
6487 metadata["completion_surface"] = json!("headless_exec_release_receipt");
6488 } else if background || backgrounded_foreground {
6489 let child_owned = context.owner_agent_id.is_some();
6490 metadata["auto_resume_on_completion"] = json!(!child_owned);
6491 metadata["completion_surface"] = if child_owned {
6492 json!("task_status_and_explicit_wait")
6493 } else {
6494 json!("runtime_event_and_task_status")
6495 };
6496 metadata["background_policy"] = json!("nonblocking");
6497 }
6498 if result.status == ShellStatus::TimedOut && !background && !interactive {
6499 metadata["foreground_timeout_recovery"] = json!({
6500 "process_killed": true,
6501 "hint": FOREGROUND_TIMEOUT_RECOVERY_HINT,
6502 "recommended_tools": ["Bash", "task_shell_start", "task_shell_wait"],
6503 "rerun_as": {"tool": "Bash", "action": "run", "background": true},
6504 "poll_with": [
6505 {"tool": "Bash", "action": "wait"},
6506 {"tool": "task_shell_wait"}
6507 ]
6508 });
6509 }
6510 if let Some(hint) = network_restricted_hint {
6511 metadata["sandbox_network_restricted"] = json!(true);
6512 metadata["sandbox_network_denied_hint"] = json!(hint);
6513 }
6514 if let Some(hint) = sandbox_denied_hint {
6515 metadata["sandbox_denied_hint"] = json!(hint);
6516 }
6517 if provenance_hint.is_some() {
6518 metadata["macos_provenance_restricted"] = json!(true);
6519 }
6520 attach_shell_owner_metadata(&mut metadata, context);
6521 attach_cargo_failure_summary(&mut metadata, command, &result);
6522 attach_python_build_dependency_hint(&mut metadata, python_dependency_hint);
6523
6524 Ok(ToolResult {
6525 content: output,
6526 success: result.status == ShellStatus::Completed
6527 || result.status == ShellStatus::Running,
6528 metadata: Some(metadata),
6529 })
6530 }
6531 Err(e) => Ok(ToolResult::error(shell_execution_failed_message(&e))),
6532 }
6533 }
6534 }
6535
6536 /// Render a spawn/stream failure for the model and the user: the full cause
6537 /// chain (an anyhow context alone hides the `ENOSPC`/`EMFILE` underneath) plus,
6538 /// when the innermost error looks like host resource exhaustion, what to do
6539 /// about it. The shell tool keeps no state from a failed spawn, so retrying is
6540 /// always safe.
6541 fn shell_execution_failed_message(error: &anyhow::Error) -> String {
6542 let hint = error
6543 .chain()
6544 .filter_map(|cause| cause.downcast_ref::<io::Error>())
6545 .find_map(output::resource_exhaustion_hint);
6546 match hint {
6547 Some(hint) => format!(
6548 "Shell execution failed: {error:#}. Likely host resource exhaustion — {hint}. The shell tool itself is still usable; the next call starts fresh."
6549 ),
6550 None => format!("Shell execution failed: {error:#}"),
6551 }
6552 }
6553
6554 /// Maximum deliberate dependency-barrier wait accepted by `exec_shell_wait`.
6555 pub(crate) const EXEC_SHELL_WAIT_MAX_TIMEOUT_MS: u64 = 600_000;
6556
6557 impl BashTool {
6558 async fn execute_wait(
6559 &self,
6560 input: &serde_json::Value,
6561 context: &ToolContext,
6562 ) -> Result<ToolResult, ToolError> {
6563 // Multi-task wait (#5549): task_ids + until=any|all. The validated
6564 // single-task path stays unchanged below when task_ids is absent.
6565 if let Some(task_ids) = input.get("task_ids") {
6566 let ids = task_ids
6567 .as_array()
6568 .ok_or_else(|| type_mismatch("task_ids", task_ids, "an array of strings"))?;
6569 let ids = ids
6570 .iter()
6571 .map(|value| {
6572 value
6573 .as_str()
6574 .ok_or_else(|| type_mismatch("task_ids", value, "strings"))
6575 .map(str::to_string)
6576 })
6577 .collect::<Result<Vec<_>, _>>()?;
6578 return self.execute_wait_many(input, context, &ids).await;
6579 }
6580 let task_id = required_task_id(input)?;
6581 let wait = match first_present_field(input, &["wait", "block"]) {
6582 None => true,
6583 Some((name, value)) => value
6584 .as_bool()
6585 .ok_or_else(|| type_mismatch(name, value, "a boolean"))?,
6586 };
6587 let timeout_ms = wait_timeout_ms(input)?;
6588
6589 let (delta, wait_canceled) = if wait {
6590 wait_for_shell_delta_cancellable(context, task_id, timeout_ms).await?
6591 } else {
6592 let mut manager = context
6593 .shell_manager
6594 .lock()
6595 .map_err(|_| ToolError::execution_failed("shell manager lock poisoned"))?;
6596 let delta = manager
6597 .get_output_delta_for_session(&context.state_namespace, task_id, false, timeout_ms)
6598 .map_err(|err| ToolError::execution_failed(err.to_string()))?;
6599 (delta, false)
6600 };
6601
6602 let status = delta.result.status.clone();
6603 let mut result = build_shell_delta_tool_result(delta, context);
6604 if let Some(metadata) = result.metadata.as_mut()
6605 && let Some(object) = metadata.as_object_mut()
6606 {
6607 object.insert("wait_timeout_ms".to_string(), json!(timeout_ms));
6608 }
6609 if wait_canceled {
6610 if matches!(status, ShellStatus::Running) {
6611 result.content = format!(
6612 "Wait canceled; background shell task {task_id} is still running.\n\n{}",
6613 result.content
6614 );
6615 }
6616 if let Some(metadata) = result.metadata.as_mut()
6617 && let Some(object) = metadata.as_object_mut()
6618 {
6619 object.insert("wait_canceled".to_string(), json!(true));
6620 }
6621 }
6622
6623 Ok(result)
6624 }
6625
6626 /// Wait on several background tasks at once (#5549).
6627 ///
6628 /// `until` selects the completion condition: `any` returns as soon as one
6629 /// task is terminal, `all` waits for every one. Unknown task ids fail the
6630 /// call (same contract as the single-task path); a timeout reports the
6631 /// still-running ids so the caller can cancel or wait again.
6632 async fn execute_wait_many(
6633 &self,
6634 input: &serde_json::Value,
6635 context: &ToolContext,
6636 task_ids: &[String],
6637 ) -> Result<ToolResult, ToolError> {
6638 let until = match input.get("until").and_then(serde_json::Value::as_str) {
6639 None | Some("all") => "all",
6640 Some("any") => "any",
6641 Some(other) => {
6642 return Err(ToolError::invalid_input(format!(
6643 "until must be \"any\" or \"all\", got {other}"
6644 )));
6645 }
6646 };
6647 let wait = match first_present_field(input, &["wait", "block"]) {
6648 None => true,
6649 Some((name, value)) => value
6650 .as_bool()
6651 .ok_or_else(|| type_mismatch(name, value, "a boolean"))?,
6652 };
6653 let timeout_ms = wait_timeout_ms(input)?;
6654 let deadline = std::time::Instant::now() + Duration::from_millis(timeout_ms);
6655 let mut timed_out = false;
6656 let mut wait_canceled = false;
6657
6658 let snapshot =
6659 |manager: &mut ShellManager| -> Result<Vec<(String, ShellStatus)>, ToolError> {
6660 task_ids
6661 .iter()
6662 .map(|id| {
6663 let detail = manager
6664 .inspect_job(id)
6665 .map_err(|err| ToolError::execution_failed(err.to_string()))?;
6666 Ok((id.clone(), detail.snapshot.status))
6667 })
6668 .collect()
6669 };
6670
6671 let mut poll_tick_ms: u64 = FOREGROUND_POLL_INITIAL_MS;
6672 let statuses = loop {
6673 let current = {
6674 let mut manager = context
6675 .shell_manager
6676 .lock()
6677 .map_err(|_| ToolError::execution_failed("shell manager lock poisoned"))?;
6678 snapshot(&mut manager)?
6679 };
6680 if context
6681 .cancel_token
6682 .as_ref()
6683 .is_some_and(|token| token.is_cancelled())
6684 {
6685 wait_canceled = true;
6686 break current;
6687 }
6688 let running = current
6689 .iter()
6690 .filter(|(_, status)| *status == ShellStatus::Running)
6691 .count();
6692 let terminal = current.len().saturating_sub(running);
6693 let satisfied = if until == "any" {
6694 terminal >= 1
6695 } else {
6696 running == 0
6697 };
6698 if !wait || satisfied {
6699 break current;
6700 }
6701 if std::time::Instant::now() >= deadline {
6702 timed_out = true;
6703 break current;
6704 }
6705 tokio::time::sleep(Duration::from_millis(poll_tick_ms)).await;
6706 poll_tick_ms = (poll_tick_ms * 2).min(FOREGROUND_POLL_MAX_MS);
6707 };
6708
6709 let running_after = statuses
6710 .iter()
6711 .filter(|(_, status)| *status == ShellStatus::Running)
6712 .count();
6713 let settled = statuses.len().saturating_sub(running_after);
6714 let lines: Vec<String> = statuses
6715 .iter()
6716 .map(|(id, status)| format!("{id}: {status:?}"))
6717 .collect();
6718 let still_running: Vec<&str> = statuses
6719 .iter()
6720 .filter(|(_, status)| *status == ShellStatus::Running)
6721 .map(|(id, _)| id.as_str())
6722 .collect();
6723 let summary = if timed_out {
6724 format!(
6725 "wait timed out after {}ms; still running: {}",
6726 timeout_ms,
6727 if still_running.is_empty() {
6728 "none".to_string()
6729 } else {
6730 still_running.join(", ")
6731 }
6732 )
6733 } else if wait_canceled {
6734 format!(
6735 "wait canceled; still running: {}",
6736 if still_running.is_empty() {
6737 "none".to_string()
6738 } else {
6739 still_running.join(", ")
6740 }
6741 )
6742 } else {
6743 format!(
6744 "{settled} of {} background command{} settled; remaining running: {}",
6745 statuses.len(),
6746 if statuses.len() == 1 { "" } else { "s" },
6747 if still_running.is_empty() {
6748 "none".to_string()
6749 } else {
6750 still_running.join(", ")
6751 }
6752 )
6753 };
6754 let content = format!("{summary}\n{}\n", lines.join("\n"));
6755 let mut metadata = serde_json::Map::new();
6756 metadata.insert(
6757 "statuses".to_string(),
6758 serde_json::Value::Object(
6759 statuses
6760 .iter()
6761 .map(|(id, status)| (id.clone(), serde_json::json!(format!("{status:?}"))))
6762 .collect::<serde_json::Map<_, _>>(),
6763 ),
6764 );
6765 metadata.insert("wait_timeout_ms".to_string(), json!(timeout_ms));
6766 metadata.insert("until".to_string(), json!(until));
6767 metadata.insert("timed_out".to_string(), json!(timed_out));
6768 if wait_canceled {
6769 metadata.insert("wait_canceled".to_string(), json!(true));
6770 }
6771 Ok(ToolResult {
6772 content: content.trim().to_string(),
6773 success: true,
6774 metadata: Some(serde_json::Value::Object(metadata)),
6775 })
6776 }
6777
6778 async fn execute_interact(
6779 &self,
6780 input: &serde_json::Value,
6781 context: &ToolContext,
6782 ) -> Result<ToolResult, ToolError> {
6783 let task_id = required_task_id(input)?;
6784 let close_stdin = optional_bool(input, "close_stdin", false)?;
6785 let timeout_ms = optional_u64(input, "timeout_ms", 1_000)?;
6786 // Same strict-type contract as `run` (2026-08-04): a non-string here
6787 // was silently dropped, so an `interact` call reported success while
6788 // writing nothing to the child's stdin. Alias order also matches
6789 // `run` now — `stdin` first — so the same payload reaches the same
6790 // place whichever spelling the model uses.
6791 let interaction_input = match first_present_field(input, &["stdin", "input", "data"]) {
6792 None => "",
6793 Some((name, value)) => value
6794 .as_str()
6795 .ok_or_else(|| type_mismatch(name, value, "a string"))?,
6796 };
6797
6798 {
6799 let mut manager = context
6800 .shell_manager
6801 .lock()
6802 .map_err(|_| ToolError::execution_failed("shell manager lock poisoned"))?;
6803 if !interaction_input.is_empty() || close_stdin {
6804 manager
6805 .write_stdin_for_session(
6806 &context.state_namespace,
6807 task_id,
6808 interaction_input,
6809 close_stdin,
6810 )
6811 .map_err(|err| ToolError::execution_failed(err.to_string()))?;
6812 }
6813 }
6814
6815 let mut elapsed = 0u64;
6816 loop {
6817 if context
6818 .cancel_token
6819 .as_ref()
6820 .is_some_and(|token| token.is_cancelled())
6821 {
6822 let mut manager = context
6823 .shell_manager
6824 .lock()
6825 .map_err(|_| ToolError::execution_failed("shell manager lock poisoned"))?;
6826 let delta = manager
6827 .get_output_delta_for_session(&context.state_namespace, task_id, false, 0)
6828 .map_err(|err| ToolError::execution_failed(err.to_string()))?;
6829 let mut result = build_shell_delta_tool_result(delta, context);
6830 if let Some(metadata) = result.metadata.as_mut()
6831 && let Some(object) = metadata.as_object_mut()
6832 {
6833 object.insert("wait_canceled".to_string(), json!(true));
6834 }
6835 return Ok(result);
6836 }
6837
6838 let delta = {
6839 let mut manager = context
6840 .shell_manager
6841 .lock()
6842 .map_err(|_| ToolError::execution_failed("shell manager lock poisoned"))?;
6843 manager
6844 .get_output_delta_for_session(&context.state_namespace, task_id, false, 0)
6845 .map_err(|err| ToolError::execution_failed(err.to_string()))?
6846 };
6847
6848 if !delta.result.stdout.is_empty()
6849 || !delta.result.stderr.is_empty()
6850 || delta.result.status != ShellStatus::Running
6851 || elapsed >= timeout_ms
6852 {
6853 return Ok(build_shell_delta_tool_result(delta, context));
6854 }
6855
6856 tokio::time::sleep(Duration::from_millis(50)).await;
6857 elapsed = elapsed.saturating_add(50);
6858 }
6859 }
6860
6861 async fn execute_cancel(
6862 &self,
6863 input: &serde_json::Value,
6864 context: &ToolContext,
6865 ) -> Result<ToolResult, ToolError> {
6866 let cancel_all = optional_bool(input, "all", false)?;
6867 let mut manager = context
6868 .shell_manager
6869 .lock()
6870 .map_err(|_| ToolError::execution_failed("shell manager lock poisoned"))?;
6871
6872 if cancel_all {
6873 let results = manager
6874 .kill_running_for_session(&context.state_namespace)
6875 .map_err(|err| ToolError::execution_failed(err.to_string()))?;
6876 if results.is_empty() {
6877 return Ok(ToolResult {
6878 content: "No running background commands.".to_string(),
6879 success: true,
6880 metadata: Some(json!({
6881 "status": "Noop",
6882 "canceled": 0,
6883 "task_ids": [],
6884 })),
6885 });
6886 }
6887
6888 let task_ids = results
6889 .iter()
6890 .filter_map(|result| result.task_id.clone())
6891 .collect::<Vec<_>>();
6892 return Ok(ToolResult {
6893 content: format!(
6894 "Canceled {} background command{}: {}",
6895 task_ids.len(),
6896 if task_ids.len() == 1 { "" } else { "s" },
6897 task_ids.join(", ")
6898 ),
6899 success: true,
6900 metadata: Some(json!({
6901 "status": "Killed",
6902 "canceled": task_ids.len(),
6903 "task_ids": task_ids,
6904 })),
6905 });
6906 }
6907
6908 let task_id = required_task_id(input)?;
6909 let result = manager
6910 .kill_for_session(&context.state_namespace, task_id)
6911 .map_err(|err| ToolError::execution_failed(err.to_string()))?;
6912 let task_id = result
6913 .task_id
6914 .clone()
6915 .unwrap_or_else(|| task_id.to_string());
6916 Ok(ToolResult {
6917 content: format!("Canceled background command: {task_id}"),
6918 success: true,
6919 metadata: Some(json!({
6920 "status": format!("{:?}", result.status),
6921 "task_id": task_id,
6922 "exit_code": result.exit_code,
6923 "duration_ms": result.duration_ms,
6924 })),
6925 })
6926 }
6927 }
6928
6929 fn required_task_id(input: &serde_json::Value) -> Result<&str, ToolError> {
6930 // A present-but-non-string task_id is a type error, not a missing field:
6931 // "missing required field" sends the model's retry in the wrong
6932 // direction when it already supplied `task_id: 42` (2026-08-04 review).
6933 match first_present_field(input, &["task_id", "id"]) {
6934 None => Err(ToolError::missing_field("task_id")),
6935 Some((name, value)) => value
6936 .as_str()
6937 .ok_or_else(|| type_mismatch(name, value, "a string")),
6938 }
6939 }
6940
6941 /// First PRESENT value among aliased spellings of one field. `null` counts
6942 /// as absent, matching the `is_absent` rule the shared typed helpers use.
6943 fn first_present_field<'a>(
6944 input: &'a serde_json::Value,
6945 names: &[&'static str],
6946 ) -> Option<(&'static str, &'a serde_json::Value)> {
6947 names.iter().find_map(|name| match input.get(*name) {
6948 None | Some(serde_json::Value::Null) => None,
6949 Some(value) => Some((*name, value)),
6950 })
6951 }
6952
6953 /// Effective `action=wait` timeout in milliseconds. `timeout_ms` is
6954 /// canonical; `timeout_secs` (seconds) and bare `timeout` (milliseconds) are
6955 /// honored so a habit formed on other wait tools gets the duration it asked
6956 /// for instead of silently falling back to the 30 s default.
6957 fn wait_timeout_ms(input: &serde_json::Value) -> Result<u64, ToolError> {
6958 match first_present_field(input, &["timeout_ms", "timeout_secs", "timeout"]) {
6959 None => Ok(30_000),
6960 Some(("timeout_secs", value)) => {
6961 let secs = value
6962 .as_u64()
6963 .ok_or_else(|| type_mismatch("timeout_secs", value, "an integer"))?;
6964 Ok(secs.saturating_mul(1_000))
6965 }
6966 Some((name, value)) => value
6967 .as_u64()
6968 .ok_or_else(|| type_mismatch(name, value, "an integer")),
6969 }
6970 }
6971
6972 fn build_shell_delta_tool_result(delta: ShellDeltaResult, context: &ToolContext) -> ToolResult {
6973 let result = delta.result;
6974 let network_restricted_hint =
6975 shell_network_restricted_hint(context, &delta.command, &result).map(str::to_string);
6976 let sandbox_denied_hint = if network_restricted_hint.is_none() {
6977 shell_sandbox_denied_hint(context, &result)
6978 } else {
6979 None
6980 };
6981 let provenance_hint = macos_provenance_hint(&result);
6982 let python_dependency_hint = python_build_dependency_hint(&delta.command, &result);
6983 let stdout_summary = summarize_output(&result.stdout);
6984 let stderr_summary = summarize_output(&result.stderr);
6985 let summary = if !stderr_summary.is_empty() {
6986 stderr_summary.clone()
6987 } else {
6988 stdout_summary.clone()
6989 };
6990
6991 let mut output = if result.stdout.is_empty() && result.stderr.is_empty() {
6992 match result.status {
6993 ShellStatus::Running => "Background task running (no new output).".to_string(),
6994 ShellStatus::Completed => "(no new output)".to_string(),
6995 ShellStatus::Failed => {
6996 format!("Command failed ({})", exit_code_label(result.exit_code))
6997 }
6998 ShellStatus::TimedOut => "Command timed out (no new output).".to_string(),
6999 ShellStatus::Killed => "Command killed (no new output).".to_string(),
7000 }
7001 } else if result.stderr.is_empty() {
7002 result.stdout.clone()
7003 } else {
7004 format!("{}\n\nSTDERR:\n{}", result.stdout, result.stderr)
7005 };
7006 // The model cannot see metadata, so surface the real elapsed time in the
7007 // visible content. Without it every wait result looks identical whether
7008 // the task just started or has been running for minutes, which biases the
7009 // model into busy-polling short waits and misjudging long ones.
7010 output = format!("{}\n\n{output}", wait_timing_line(&result));
7011
7012 if let Some(hint) = network_restricted_hint.as_deref() {
7013 output = format!("{hint}\n\n{output}");
7014 }
7015 if let Some(hint) = sandbox_denied_hint.as_deref() {
7016 output = format!("{hint}\n\n{output}");
7017 }
7018 if let Some(hint) = provenance_hint {
7019 output = format!("{hint}\n\n{output}");
7020 }
7021 if let Some(hint) = python_dependency_hint {
7022 output = format!("{hint}\n\n{output}");
7023 }
7024
7025 let mut metadata = json!({
7026 "exit_code": result.exit_code,
7027 "exit_code_hex": exit_code_hex(result.exit_code),
7028 "status": format!("{:?}", result.status),
7029 "duration_ms": result.duration_ms,
7030 "sandboxed": result.sandboxed,
7031 "sandbox_type": result.sandbox_type,
7032 "sandbox_denied": result.sandbox_denied,
7033 "task_id": result.task_id,
7034 "stdout_len": result.stdout_len,
7035 "stderr_len": result.stderr_len,
7036 "stdout_truncated": result.stdout_truncated,
7037 "stderr_truncated": result.stderr_truncated,
7038 "stdout_omitted": result.stdout_omitted,
7039 "stderr_omitted": result.stderr_omitted,
7040 "stdout_total_len": delta.stdout_total_len,
7041 "stderr_total_len": delta.stderr_total_len,
7042 "summary": summary,
7043 "stdout_summary": stdout_summary,
7044 "stderr_summary": stderr_summary,
7045 "command": delta.command,
7046 "stream_delta": true,
7047 });
7048 attach_shell_owner_metadata(&mut metadata, context);
7049 attach_cargo_failure_summary(&mut metadata, &delta.command, &result);
7050 attach_python_build_dependency_hint(&mut metadata, python_dependency_hint);
7051
7052 let mut tool_result = ToolResult {
7053 content: output,
7054 success: matches!(result.status, ShellStatus::Completed | ShellStatus::Running),
7055 metadata: Some(metadata),
7056 };
7057 if let Some(hint) = network_restricted_hint
7058 && let Some(metadata) = tool_result.metadata.as_mut()
7059 && let Some(object) = metadata.as_object_mut()
7060 {
7061 object.insert("sandbox_network_restricted".to_string(), json!(true));
7062 object.insert("sandbox_network_denied_hint".to_string(), json!(hint));
7063 }
7064 if let Some(hint) = sandbox_denied_hint
7065 && let Some(metadata) = tool_result.metadata.as_mut()
7066 && let Some(object) = metadata.as_object_mut()
7067 {
7068 object.insert("sandbox_denied_hint".to_string(), json!(hint));
7069 }
7070 if provenance_hint.is_some()
7071 && let Some(metadata) = tool_result.metadata.as_mut()
7072 && let Some(object) = metadata.as_object_mut()
7073 {
7074 object.insert("macos_provenance_restricted".to_string(), json!(true));
7075 }
7076 tool_result
7077 }
7078
7079 /// Human-readable elapsed time for a shell task ("450 ms", "12.3 s", "2m5s").
7080 fn format_elapsed_ms(ms: u64) -> String {
7081 if ms < 1_000 {
7082 format!("{ms} ms")
7083 } else if ms < 60_000 {
7084 let secs = ms as f64 / 1_000.0;
7085 format!("{secs} s")
7086 } else {
7087 let total_secs = ms / 1_000;
7088 format!("{}m{}s", total_secs / 60, total_secs % 60)
7089 }
7090 }
7091
7092 /// One-line status + elapsed summary for wait/delta results, placed at the top
7093 /// of the visible content so the model can judge how long it actually waited.
7094 fn wait_timing_line(result: &ShellResult) -> String {
7095 let status_phrase = match result.status {
7096 ShellStatus::Running => "still running",
7097 ShellStatus::Completed => "completed",
7098 ShellStatus::Failed => "failed",
7099 ShellStatus::Killed => "killed",
7100 ShellStatus::TimedOut => "timed out",
7101 };
7102 let elapsed = format_elapsed_ms(result.duration_ms);
7103 match result.task_id.as_deref() {
7104 Some(task_id) => format!("Task {task_id} {status_phrase} after {elapsed}."),
7105 None => format!("Task {status_phrase} after {elapsed}."),
7106 }
7107 }
7108
7109 async fn wait_for_shell_delta_cancellable(
7110 context: &ToolContext,
7111 task_id: &str,
7112 timeout_ms: u64,
7113 ) -> Result<(ShellDeltaResult, bool), ToolError> {
7114 let timeout_ms = timeout_ms.clamp(1000, EXEC_SHELL_WAIT_MAX_TIMEOUT_MS);
7115 let deadline = Instant::now() + Duration::from_millis(timeout_ms);
7116 let mut stdout_accum = String::new();
7117 let mut stderr_accum = String::new();
7118
7119 let mut poll_tick_ms: u64 = FOREGROUND_POLL_INITIAL_MS;
7120 let (command, result, stdout_total_len, stderr_total_len) = loop {
7121 if context
7122 .cancel_token
7123 .as_ref()
7124 .is_some_and(|token| token.is_cancelled())
7125 {
7126 let mut manager = context
7127 .shell_manager
7128 .lock()
7129 .map_err(|_| ToolError::execution_failed("shell manager lock poisoned"))?;
7130 let delta = manager
7131 .get_output_delta_for_session(&context.state_namespace, task_id, false, 0)
7132 .map_err(|err| ToolError::execution_failed(err.to_string()))?;
7133 append_shell_delta_output(&mut stdout_accum, &mut stderr_accum, &delta.result);
7134 return Ok((
7135 shell_delta_with_accumulated_output(
7136 delta.command,
7137 delta.result,
7138 &stdout_accum,
7139 &stderr_accum,
7140 delta.stdout_total_len,
7141 delta.stderr_total_len,
7142 ),
7143 true,
7144 ));
7145 }
7146
7147 let delta = {
7148 let mut manager = context
7149 .shell_manager
7150 .lock()
7151 .map_err(|_| ToolError::execution_failed("shell manager lock poisoned"))?;
7152 manager
7153 .get_output_delta_for_session(&context.state_namespace, task_id, false, 0)
7154 .map_err(|err| ToolError::execution_failed(err.to_string()))?
7155 };
7156
7157 let stdout_total_len = delta.stdout_total_len;
7158 let stderr_total_len = delta.stderr_total_len;
7159 let command = delta.command.clone();
7160 append_shell_delta_output(&mut stdout_accum, &mut stderr_accum, &delta.result);
7161
7162 let status = delta.result.status.clone();
7163 if status != ShellStatus::Running || Instant::now() >= deadline {
7164 break (command, delta.result, stdout_total_len, stderr_total_len);
7165 }
7166
7167 tokio::time::sleep(Duration::from_millis(poll_tick_ms)).await;
7168 poll_tick_ms = (poll_tick_ms * 2).min(FOREGROUND_POLL_MAX_MS);
7169 };
7170
7171 Ok((
7172 shell_delta_with_accumulated_output(
7173 command,
7174 result,
7175 &stdout_accum,
7176 &stderr_accum,
7177 stdout_total_len,
7178 stderr_total_len,
7179 ),
7180 false,
7181 ))
7182 }
7183
7184 fn append_shell_delta_output(
7185 stdout_accum: &mut String,
7186 stderr_accum: &mut String,
7187 result: &ShellResult,
7188 ) {
7189 if !result.stdout.is_empty() {
7190 stdout_accum.push_str(&result.stdout);
7191 }
7192 if !result.stderr.is_empty() {
7193 stderr_accum.push_str(&result.stderr);
7194 }
7195 }
7196
7197 fn shell_delta_with_accumulated_output(
7198 command: String,
7199 mut result: ShellResult,
7200 stdout_accum: &str,
7201 stderr_accum: &str,
7202 stdout_total_len: usize,
7203 stderr_total_len: usize,
7204 ) -> ShellDeltaResult {
7205 let (stdout, stdout_meta) = truncate_with_meta(stdout_accum);
7206 let (stderr, stderr_meta) = truncate_with_meta(stderr_accum);
7207 result.stdout = stdout;
7208 result.stderr = stderr;
7209 result.stdout_len = stdout_meta.original_len;
7210 result.stderr_len = stderr_meta.original_len;
7211 result.stdout_omitted = stdout_meta.omitted;
7212 result.stderr_omitted = stderr_meta.omitted;
7213 result.stdout_truncated = stdout_meta.truncated;
7214 result.stderr_truncated = stderr_meta.truncated;
7215
7216 ShellDeltaResult {
7217 command,
7218 result,
7219 stdout_total_len,
7220 stderr_total_len,
7221 }
7222 }
7223
7224 /// Refuse a notes target this auto-approved tool must not append to.
7225 ///
7226 /// The file itself must not be a symlink (a committed `notes.md -> ~/.zshrc`
7227 /// would redirect the append), and a target placed inside the workspace must
7228 /// resolve inside it after symlinked parent directories are followed.
7229 async fn ensure_notes_target_is_safe(
7230 notes_path: &std::path::Path,
7231 workspace: &std::path::Path,
7232 ) -> Result<(), ToolError> {
7233 if let Ok(meta) = tokio::fs::symlink_metadata(notes_path).await
7234 && (meta.file_type().is_symlink() || !meta.is_file())
7235 {
7236 return Err(ToolError::permission_denied(format!(
7237 "Refusing to append a note to {}: the notes path is a symlink or not a regular file.",
7238 notes_path.display()
7239 )));
7240 }
7241 if notes_path.starts_with(workspace)
7242 && let (Some(parent), Ok(root)) = (
7243 notes_path.parent(),
7244 tokio::fs::canonicalize(workspace).await,
7245 )
7246 {
7247 // Walk up to the nearest existing ancestor: the rest is created
7248 // below as real directories, so only existing links can redirect.
7249 let mut existing = parent.to_path_buf();
7250 while !tokio::fs::try_exists(&existing).await.unwrap_or(false) {
7251 if !existing.pop() {
7252 break;
7253 }
7254 }
7255 if let Ok(resolved) = tokio::fs::canonicalize(&existing).await
7256 && !resolved.starts_with(&root)
7257 {
7258 return Err(ToolError::permission_denied(format!(
7259 "Refusing to append a note to {}: its directory resolves outside the workspace.",
7260 notes_path.display()
7261 )));
7262 }
7263 }
7264 Ok(())
7265 }
7266
7267 /// Tool for appending notes to a notes file.
7268 pub struct NoteTool;
7269
7270 #[async_trait]
7271 impl ToolSpec for NoteTool {
7272 fn name(&self) -> &'static str {
7273 "note"
7274 }
7275
7276 fn description(&self) -> &'static str {
7277 "Append a note to the agent notes file for persistent context across sessions."
7278 }
7279
7280 fn input_schema(&self) -> serde_json::Value {
7281 json!({
7282 "type": "object",
7283 "properties": {
7284 "content": {
7285 "type": "string",
7286 "description": "The note content to append"
7287 }
7288 },
7289 "required": ["content"]
7290 })
7291 }
7292
7293 fn capabilities(&self) -> Vec<ToolCapability> {
7294 vec![ToolCapability::WritesFiles]
7295 }
7296
7297 fn approval_requirement(&self) -> ApprovalRequirement {
7298 ApprovalRequirement::Auto // Notes are low-risk
7299 }
7300
7301 async fn execute(
7302 &self,
7303 input: serde_json::Value,
7304 context: &ToolContext,
7305 ) -> Result<ToolResult, ToolError> {
7306 let note_content = required_str(&input, "content")?;
7307 ensure_notes_target_is_safe(&context.notes_path, &context.workspace).await?;
7308
7309 // Ensure parent directory exists. Tool handlers run on the Tokio
7310 // runtime, so filesystem calls use tokio::fs (blocking-call
7311 // convention, #6149).
7312 if let Some(parent) = context.notes_path.parent() {
7313 tokio::fs::create_dir_all(parent).await.map_err(|e| {
7314 ToolError::execution_failed(format!("Failed to create notes directory: {e}"))
7315 })?;
7316 }
7317
7318 // Append to notes file
7319 let mut file = tokio::fs::OpenOptions::new()
7320 .create(true)
7321 .append(true)
7322 .open(&context.notes_path)
7323 .await
7324 .map_err(|e| ToolError::execution_failed(format!("Failed to open notes file: {e}")))?;
7325
7326 use tokio::io::AsyncWriteExt;
7327 file.write_all(format!("\n---\n{note_content}\n").as_bytes())
7328 .await
7329 .map_err(|e| ToolError::execution_failed(format!("Failed to write note: {e}")))?;
7330 // tokio's File finishes a write on a background thread; report
7331 // success only once the bytes reached the file.
7332 file.flush()
7333 .await
7334 .map_err(|e| ToolError::execution_failed(format!("Failed to write note: {e}")))?;
7335
7336 Ok(ToolResult::success(format!(
7337 "Note appended to {}",
7338 context.notes_path.display()
7339 )))
7340 }
7341 }
7342
7343 #[cfg(test)]
7344 #[path = "shell/tests/enforced_readonly.rs"]
7345 mod enforced_readonly_tests;
7346 #[cfg(test)]
7347 mod tests;
7348
7349 pub(crate) fn foreground_command_requests_detach(command: &str) -> bool {
7350 let mut single_quoted = false;
7351 let mut double_quoted = false;
7352 let mut escaped = false;
7353 let chars = command.chars().collect::<Vec<_>>();
7354 for (index, ch) in chars.iter().copied().enumerate() {
7355 if escaped {
7356 escaped = false;
7357 continue;
7358 }
7359 if ch == '\\' && !single_quoted {
7360 escaped = true;
7361 continue;
7362 }
7363 if ch == '\'' && !double_quoted {
7364 single_quoted = !single_quoted;
7365 continue;
7366 }
7367 if ch == '"' && !single_quoted {
7368 double_quoted = !double_quoted;
7369 continue;
7370 }
7371 if ch != '&' || single_quoted || double_quoted {
7372 continue;
7373 }
7374 let previous = index.checked_sub(1).and_then(|i| chars.get(i)).copied();
7375 let next = chars.get(index + 1).copied();
7376 // `&&`, `&>`/`&>>`, and `>&` are chaining/redirection rather than a
7377 // detached child. Any other unquoted ampersand is a background
7378 // control operator and is unavailable in ACP.
7379 if previous != Some('&') && next != Some('&') && next != Some('>') && previous != Some('>')
7380 {
7381 return true;
7382 }
7383 }
7384
7385 shell_words::split(command).is_ok_and(|words| {
7386 words.iter().any(|word| {
7387 matches!(
7388 word.to_ascii_lowercase().as_str(),
7389 "nohup" | "disown" | "setsid" | "daemonize"
7390 )
7391 })
7392 })
7393 }
7394
7394 lines RUST