返回 CodeWhale
terminal_session.rs
根目录 / crates / tui / src / tools / terminal_session.rs
1 //! Stateful, PTY-backed terminal sessions.
2 //!
3 //! Live PTY processes remain deliberately process-local, while a non-secret
4 //! durable summary records identity, last-known cwd, lifecycle state, and
5 //! replacement history. A later process reports the shell as stale/lost and
6 //! starts a new identity; it never claims to reattach from a reused PID.
7
8 #[cfg(unix)]
9 use std::collections::{HashMap, VecDeque};
10 #[cfg(unix)]
11 use std::io::Write;
12 #[cfg(unix)]
13 use std::path::{Path, PathBuf};
14 #[cfg(unix)]
15 use std::sync::{Arc, Mutex, OnceLock};
16 #[cfg(unix)]
17 use std::time::{Duration, Instant};
18
19 use async_trait::async_trait;
20 #[cfg(unix)]
21 use serde::{Deserialize, Serialize};
22 use serde_json::json;
23 #[cfg(unix)]
24 use sha2::{Digest, Sha256};
25 #[cfg(unix)]
26 use uuid::Uuid;
27
28 use super::spec::{
29 ApprovalRequirement, ToolCapability, ToolContext, ToolError, ToolResult, ToolSpec,
30 };
31 #[cfg(unix)]
32 use super::spec::{optional_u64, required_str};
33
34 #[cfg(unix)]
35 const BUFFER_LIMIT: usize = 512 * 1024;
36 #[cfg(unix)]
37 const OUTPUT_LIMIT: usize = 12 * 1024;
38 /// Ceiling for one [`OutputBuffer::read_since`] response. The ring is already
39 /// bounded at [`BUFFER_LIMIT`]; this bounds a single frame so a client cannot
40 /// ask for the whole window at once.
41 #[cfg(unix)]
42 pub(crate) const READ_LIMIT: usize = 64 * 1024;
43 #[cfg(unix)]
44 const DEFAULT_TIMEOUT_SECS: u64 = 120;
45 #[cfg(unix)]
46 const MAX_TIMEOUT_SECS: u64 = 600;
47 #[cfg(unix)]
48 const CANCEL_CONFIRM_TIMEOUT: Duration = Duration::from_secs(2);
49 #[cfg(unix)]
50 const CANCEL_SENTINEL_RETRY_INTERVAL: Duration = Duration::from_millis(50);
51
52 #[cfg(unix)]
53 pub(crate) struct TerminalSession {
54 writer: Arc<Mutex<Box<dyn Write + Send>>>,
55 /// The pty master is retained after the reader and writer clones are taken
56 /// so [`TerminalSession::resize`] can reach the kernel's window size. The
57 /// clones keep the pty alive; nothing else clones the master.
58 master: Box<dyn portable_pty::MasterPty + Send>,
59 child: Box<dyn portable_pty::Child + Send>,
60 output: Arc<Mutex<OutputBuffer>>,
61 read_cursor: u64,
62 command: Option<CommandState>,
63 durable: DurableTerminalRecord,
64 durable_path: PathBuf,
65 /// True when the live PTY was started through a real sandbox backend.
66 /// A later narrowed posture must not reuse an unsandboxed shell.
67 sandbox_confined: bool,
68 }
69
70 #[cfg(unix)]
71 #[derive(Clone, Debug, Serialize, Deserialize)]
72 struct DurableTerminalRecord {
73 schema_version: u32,
74 session_id: String,
75 runtime_nonce: String,
76 process_id: u32,
77 workspace: PathBuf,
78 name: String,
79 shell: String,
80 last_known_cwd: PathBuf,
81 environment_summary: EnvironmentSummary,
82 state: DurableTerminalState,
83 updated_at: String,
84 #[serde(default, skip_serializing_if = "Option::is_none")]
85 previous: Option<Box<DurableTerminalRecord>>,
86 }
87
88 #[cfg(unix)]
89 #[derive(Clone, Debug, Serialize, Deserialize)]
90 struct EnvironmentSummary {
91 term: Option<String>,
92 virtual_env_active: bool,
93 conda_env_active: bool,
94 nix_shell_active: bool,
95 note: String,
96 }
97
98 #[cfg(unix)]
99 #[derive(Clone, Copy, Debug, Serialize, Deserialize, PartialEq, Eq)]
100 #[serde(rename_all = "snake_case")]
101 enum DurableTerminalState {
102 Running,
103 Idle,
104 Canceled,
105 Failed,
106 StaleLost,
107 Reset,
108 }
109
110 #[cfg(unix)]
111 struct CommandState {
112 marker: String,
113 }
114
115 #[cfg(unix)]
116 #[derive(Default)]
117 struct OutputBuffer {
118 bytes: VecDeque<u8>,
119 total: u64,
120 }
121
122 /// One absolute-offset slice of a session's output.
123 ///
124 /// Offsets are absolute for the life of the live PTY: they count every byte
125 /// the reader thread ever appended, not the bytes still retained. That is what
126 /// lets a client resume from a cursor it stored earlier, and what lets this
127 /// type tell it the truth when the retained window has moved past it.
128 ///
129 /// Field names match the `/v1/terminal/{name}/output` wire so the payload is a
130 /// projection, not a translation (the same vocabulary the jobs byte stream
131 /// uses).
132 #[cfg(unix)]
133 #[derive(Clone, Debug, PartialEq, Eq)]
134 pub(crate) struct OutputChunk {
135 pub(crate) bytes: Vec<u8>,
136 /// Absolute offset of `bytes[0]` in the session's lifetime output.
137 pub(crate) offset: u64,
138 /// Offset to pass as the next `cursor`.
139 pub(crate) next_cursor: u64,
140 /// Every byte the session has produced, retained or not.
141 pub(crate) total: u64,
142 /// Leading bytes the ring has permanently discarded.
143 pub(crate) dropped: u64,
144 /// True when the requested cursor predates the retained window. The bytes
145 /// in between cannot be recovered from this process — report, not repair.
146 pub(crate) gap: bool,
147 }
148
149 #[cfg(unix)]
150 impl OutputBuffer {
151 fn append(&mut self, data: &[u8]) {
152 self.total = self.total.saturating_add(data.len() as u64);
153 self.bytes.extend(data);
154 while self.bytes.len() > BUFFER_LIMIT {
155 let _ = self.bytes.pop_front();
156 }
157 }
158
159 fn text(&self) -> String {
160 String::from_utf8_lossy(&self.bytes.iter().copied().collect::<Vec<_>>()).into_owned()
161 }
162
163 /// Absolute offset of the oldest retained byte.
164 fn oldest(&self) -> u64 {
165 self.total.saturating_sub(self.bytes.len() as u64)
166 }
167
168 /// Read from an absolute cursor without consuming: the caller owns its own
169 /// position, so two readers can replay the same bytes and a repeated read
170 /// is idempotent. Does not advance `read_cursor` — the consuming
171 /// tool-result path is untouched. This is the cursor arithmetic only;
172 /// response-size policy belongs to the caller.
173 ///
174 /// A cursor past the head clamps to `total`: nothing is there yet, and
175 /// echoing the future cursor back as `next_cursor` would make every byte
176 /// produced before the stream reached it invisible to that client.
177 fn read_since(&self, cursor: u64, max_bytes: usize) -> OutputChunk {
178 let dropped = self.oldest();
179 let gap = cursor < dropped;
180 let start = cursor.max(dropped).min(self.total);
181 let skip = usize::try_from(start - dropped).unwrap_or(usize::MAX);
182 let take = max_bytes.min(self.bytes.len().saturating_sub(skip));
183 let bytes = self.bytes.iter().skip(skip).take(take).copied().collect();
184 OutputChunk {
185 bytes,
186 offset: start,
187 next_cursor: start + take as u64,
188 total: self.total,
189 dropped,
190 gap,
191 }
192 }
193 }
194
195 #[cfg(unix)]
196 pub(crate) type SharedSession = Arc<Mutex<TerminalSession>>;
197
198 #[cfg(unix)]
199 #[derive(Clone, Debug, Hash, PartialEq, Eq)]
200 struct SessionKey {
201 workspace: PathBuf,
202 name: String,
203 }
204
205 #[cfg(unix)]
206 static SESSIONS: OnceLock<Mutex<HashMap<SessionKey, SharedSession>>> = OnceLock::new();
207
208 #[cfg(unix)]
209 static RUNTIME_NONCE: OnceLock<String> = OnceLock::new();
210
211 #[cfg(unix)]
212 fn runtime_nonce() -> &'static str {
213 RUNTIME_NONCE.get_or_init(|| Uuid::new_v4().to_string())
214 }
215
216 #[cfg(unix)]
217 fn sessions() -> &'static Mutex<HashMap<SessionKey, SharedSession>> {
218 SESSIONS.get_or_init(|| Mutex::new(HashMap::new()))
219 }
220
221 #[cfg(unix)]
222 fn session_key(name: &str, workspace: &Path) -> SessionKey {
223 SessionKey {
224 workspace: workspace
225 .canonicalize()
226 .unwrap_or_else(|_| workspace.to_path_buf()),
227 name: name.to_string(),
228 }
229 }
230
231 #[cfg(unix)]
232 fn durable_path(name: &str, workspace: &Path) -> Result<PathBuf, String> {
233 let workspace = workspace
234 .canonicalize()
235 .unwrap_or_else(|_| workspace.to_path_buf());
236 let mut hasher = Sha256::new();
237 hasher.update(workspace.to_string_lossy().as_bytes());
238 hasher.update(b"\0");
239 hasher.update(name.as_bytes());
240 let digest = hasher
241 .finalize()
242 .iter()
243 .map(|byte| format!("{byte:02x}"))
244 .collect::<String>();
245 #[cfg(test)]
246 let state_dir = workspace.join(".codewhale-test-terminal-sessions");
247 #[cfg(not(test))]
248 let state_dir = codewhale_config::ensure_state_dir("terminal-sessions")
249 .map_err(|error| format!("failed to resolve terminal session state directory: {error}"))?;
250 std::fs::create_dir_all(&state_dir)
251 .map_err(|error| format!("failed to create terminal session state directory: {error}"))?;
252 Ok(state_dir.join(format!("{}.json", &digest[..32])))
253 }
254
255 #[cfg(unix)]
256 fn load_durable(path: &Path) -> Option<DurableTerminalRecord> {
257 let bytes = std::fs::read(path).ok()?;
258 serde_json::from_slice(&bytes).ok()
259 }
260
261 #[cfg(unix)]
262 fn persist_durable(path: &Path, record: &DurableTerminalRecord) -> Result<(), String> {
263 let payload = serde_json::to_vec_pretty(record)
264 .map_err(|error| format!("failed to encode terminal session state: {error}"))?;
265 crate::utils::write_atomic(path, &payload)
266 .map_err(|error| format!("failed to persist terminal session state: {error}"))
267 }
268
269 #[cfg(unix)]
270 fn pty_lacks_required_sandbox(
271 policy: &crate::sandbox::SandboxPolicy,
272 applied: crate::sandbox::SandboxType,
273 ) -> bool {
274 policy.should_sandbox() && matches!(applied, crate::sandbox::SandboxType::None)
275 }
276
277 #[cfg(unix)]
278 fn pty_sandbox_policy(context: &ToolContext) -> crate::sandbox::SandboxPolicy {
279 context.elevated_sandbox_policy.clone().unwrap_or_default()
280 }
281
282 #[cfg(unix)]
283 fn prepare_pty_shell(
284 shell: &str,
285 workspace: &std::path::Path,
286 policy: crate::sandbox::SandboxPolicy,
287 ) -> Result<(crate::sandbox::ExecEnv, bool), String> {
288 let spec = crate::sandbox::CommandSpec {
289 program: shell.to_string(),
290 args: vec!["-i".to_string()],
291 cwd: workspace.to_path_buf(),
292 env: std::collections::HashMap::new(),
293 timeout: Duration::from_secs(24 * 60 * 60),
294 sandbox_policy: policy.clone(),
295 justification: Some("persistent PTY shell".to_string()),
296 requested_command: None,
297 };
298 let prepared = crate::sandbox::SandboxManager::new().prepare(&spec);
299 if pty_lacks_required_sandbox(&policy, prepared.sandbox_type) {
300 return Err(
301 "terminal PTY tools cannot run unsandboxed under a narrowed filesystem posture; use Full Access / ExternalSandbox or enable seatbelt/bwrap"
302 .to_string(),
303 );
304 }
305 if prepared.command.is_empty() {
306 return Err("sandbox prepare returned an empty PTY command".to_string());
307 }
308 let confined = !matches!(prepared.sandbox_type, crate::sandbox::SandboxType::None);
309 Ok((prepared, confined))
310 }
311
312 #[cfg(unix)]
313 fn create_session(
314 name: &str,
315 workspace: &std::path::Path,
316 policy: crate::sandbox::SandboxPolicy,
317 ) -> Result<SharedSession, String> {
318 // Unit tests exercise the persistent-PTY contract, not a developer's
319 // interactive shell startup files. Keep their deadlines deterministic and
320 // avoid racing process-global HOME/SHELL overrides from parallel tests.
321 #[cfg(test)]
322 let shell = "/bin/sh".to_string();
323 #[cfg(not(test))]
324 let shell = std::env::var("SHELL").unwrap_or_else(|_| "/bin/sh".to_string());
325 let workspace = workspace
326 .canonicalize()
327 .unwrap_or_else(|_| workspace.to_path_buf());
328 let durable_path = durable_path(name, &workspace)?;
329 let previous = load_durable(&durable_path).map(|mut record| {
330 // A persisted shell is historical evidence only. A random per-process
331 // nonce makes PID reuse irrelevant and deliberately forbids reattach.
332 record.state = DurableTerminalState::StaleLost;
333 record.updated_at = chrono::Utc::now().to_rfc3339();
334 Box::new(record)
335 });
336 let pty = portable_pty::native_pty_system();
337 let pair = pty
338 .openpty(portable_pty::PtySize {
339 rows: 24,
340 cols: 120,
341 pixel_width: 0,
342 pixel_height: 0,
343 })
344 .map_err(|e| format!("failed to open PTY: {e}"))?;
345
346 let (prepared, sandbox_confined) = prepare_pty_shell(&shell, &workspace, policy)?;
347 let mut command = portable_pty::CommandBuilder::new(&prepared.command[0]);
348 for arg in &prepared.command[1..] {
349 command.arg(arg);
350 }
351 command.cwd(&prepared.cwd);
352 // Same sanitized environment as the `exec_shell` PTY path; sandbox
353 // markers from `prepared.env` are applied as explicit overrides.
354 crate::child_env::apply_to_pty_command(
355 &mut command,
356 crate::child_env::string_map_env(&prepared.env),
357 );
358 let child = pair
359 .slave
360 .spawn_command(command)
361 .map_err(|e| format!("failed to start shell {shell}: {e}"))?;
362 drop(pair.slave);
363
364 let reader = pair
365 .master
366 .try_clone_reader()
367 .map_err(|e| format!("failed to read PTY: {e}"))?;
368 let writer = pair
369 .master
370 .take_writer()
371 .map_err(|e| format!("failed to write PTY: {e}"))?;
372 let output = Arc::new(Mutex::new(OutputBuffer::default()));
373 let reader_output = Arc::clone(&output);
374 std::thread::spawn(move || {
375 let mut reader = reader;
376 let mut buf = [0u8; 8192];
377 loop {
378 match std::io::Read::read(&mut reader, &mut buf) {
379 Ok(0) | Err(_) => break,
380 Ok(n) => {
381 if let Ok(mut output) = reader_output.lock() {
382 output.append(&buf[..n]);
383 }
384 }
385 }
386 }
387 });
388
389 let durable = DurableTerminalRecord {
390 schema_version: 1,
391 session_id: Uuid::new_v4().to_string(),
392 runtime_nonce: runtime_nonce().to_string(),
393 process_id: std::process::id(),
394 workspace: workspace.clone(),
395 name: name.to_string(),
396 shell,
397 last_known_cwd: workspace,
398 environment_summary: EnvironmentSummary {
399 term: std::env::var("TERM").ok().filter(|value| !value.trim().is_empty()),
400 virtual_env_active: std::env::var_os("VIRTUAL_ENV").is_some(),
401 conda_env_active: std::env::var_os("CONDA_PREFIX").is_some(),
402 nix_shell_active: std::env::var_os("IN_NIX_SHELL").is_some(),
403 note: "Values and secrets are not persisted; in-shell environment changes are process-local."
404 .to_string(),
405 },
406 state: DurableTerminalState::Idle,
407 updated_at: chrono::Utc::now().to_rfc3339(),
408 previous,
409 };
410 persist_durable(&durable_path, &durable)?;
411
412 Ok(Arc::new(Mutex::new(TerminalSession {
413 writer: Arc::new(Mutex::new(writer)),
414 master: pair.master,
415 child,
416 output,
417 read_cursor: 0,
418 command: None,
419 durable,
420 durable_path,
421 sandbox_confined,
422 })))
423 }
424
425 #[cfg(unix)]
426 pub(crate) fn get_or_create(
427 name: &str,
428 workspace: &std::path::Path,
429 policy: crate::sandbox::SandboxPolicy,
430 ) -> Result<SharedSession, String> {
431 let key = session_key(name, workspace);
432 let mut registry = sessions()
433 .lock()
434 .map_err(|_| "terminal session registry lock poisoned".to_string())?;
435 if let Some(session) = registry.get(&key) {
436 let confined = session
437 .lock()
438 .ok()
439 .map(|guard| guard.sandbox_confined)
440 .unwrap_or(false);
441 if policy.should_sandbox() && !confined {
442 return Err(
443 "existing PTY session was started without a sandbox; start a new session name under the current posture or use Full Access"
444 .to_string(),
445 );
446 }
447 return Ok(Arc::clone(session));
448 }
449 let session = create_session(name, workspace, policy)?;
450 registry.insert(key, Arc::clone(&session));
451 Ok(session)
452 }
453
454 #[cfg(unix)]
455 fn find(name: &str, workspace: &Path) -> Result<SharedSession, String> {
456 let live = sessions()
457 .lock()
458 .map_err(|_| "terminal session registry lock poisoned".to_string())?
459 .get(&session_key(name, workspace))
460 .cloned();
461 if let Some(live) = live {
462 return Ok(live);
463 }
464 if let Ok(path) = durable_path(name, workspace)
465 && let Some(mut record) = load_durable(&path)
466 {
467 record.state = DurableTerminalState::StaleLost;
468 record.updated_at = chrono::Utc::now().to_rfc3339();
469 let _ = persist_durable(&path, &record);
470 return Err(format!(
471 "terminal session '{name}' is stale/lost after restart (last cwd: {}); run terminal/run with this name to start a replacement while preserving the historical summary",
472 record.last_known_cwd.display()
473 ));
474 }
475 Err(format!(
476 "terminal session '{name}' does not exist in workspace {}",
477 workspace.display()
478 ))
479 }
480
481 #[cfg(unix)]
482 pub(crate) fn write_bytes(session: &TerminalSession, bytes: &[u8]) -> Result<(), String> {
483 let mut writer = session
484 .writer
485 .lock()
486 .map_err(|_| "terminal PTY writer lock poisoned".to_string())?;
487 writer
488 .write_all(bytes)
489 .map_err(|e| format!("PTY write failed: {e}"))?;
490 writer.flush().map_err(|e| format!("PTY flush failed: {e}"))
491 }
492
493 #[cfg(unix)]
494 fn output_snapshot(session: &TerminalSession) -> String {
495 session
496 .output
497 .lock()
498 .map(|output| output.text())
499 .unwrap_or_default()
500 }
501
502 #[cfg(unix)]
503 fn take_output(session: &mut TerminalSession) -> String {
504 let Ok(output) = session.output.lock() else {
505 return String::new();
506 };
507 let chunk = output.read_since(session.read_cursor, usize::MAX);
508 session.read_cursor = output.total;
509 String::from_utf8_lossy(&chunk.bytes).into_owned()
510 }
511
512 /// Resize the live pty's window. The kernel updates its `winsize` and signals
513 /// the child, which is what makes an interactive app redraw at the new size.
514 #[cfg(unix)]
515 pub(crate) fn resize_session(
516 session: &TerminalSession,
517 rows: u16,
518 cols: u16,
519 ) -> Result<(), String> {
520 session
521 .master
522 .resize(portable_pty::PtySize {
523 rows,
524 cols,
525 pixel_width: 0,
526 pixel_height: 0,
527 })
528 .map_err(|e| format!("PTY resize failed: {e}"))
529 }
530
531 /// Absolute-offset, non-consuming read. `max_bytes` is clamped to
532 /// [`READ_LIMIT`] so one caller cannot ask for the whole retained window.
533 ///
534 /// Known limitation: this reports a [`OutputChunk::gap`]; it does not repair
535 /// one. Bytes dropped by the ring are gone with the process, and nothing here
536 /// re-reads them from disk — the durable record is identity and lifecycle,
537 /// never output.
538 #[cfg(unix)]
539 pub(crate) fn read_session_since(
540 session: &TerminalSession,
541 cursor: u64,
542 max_bytes: usize,
543 ) -> Result<OutputChunk, String> {
544 let output = session
545 .output
546 .lock()
547 .map_err(|_| "terminal output lock poisoned".to_string())?;
548 Ok(output.read_since(cursor, max_bytes.min(READ_LIMIT)))
549 }
550
551 /// Poll the shell without blocking; `None` means it is still running.
552 #[cfg(unix)]
553 pub(crate) fn session_exit_status(
554 session: &mut TerminalSession,
555 ) -> Result<Option<portable_pty::ExitStatus>, String> {
556 session
557 .child
558 .try_wait()
559 .map_err(|e| format!("PTY wait failed: {e}"))
560 }
561
562 /// Terminate the shell. The caller still observes the exit through
563 /// [`session_exit_status`].
564 #[cfg(unix)]
565 pub(crate) fn kill_session(session: &mut TerminalSession) -> Result<(), String> {
566 session
567 .child
568 .kill()
569 .map_err(|e| format!("PTY kill failed: {e}"))
570 }
571
572 /// Resolve a live session without creating one.
573 ///
574 /// `/v1/terminal` attaches to shells the Engine already owns (the agent's
575 /// terminal tools create them). A request for a name that has no live session
576 /// is a 404, never a new process: an HTTP client must not be able to conjure a
577 /// shell the Engine does not know about.
578 #[cfg(unix)]
579 pub(crate) fn lookup(name: &str, workspace: &Path) -> Option<SharedSession> {
580 let key = session_key(name, workspace);
581 sessions()
582 .lock()
583 .ok()
584 .and_then(|registry| registry.get(&key).map(Arc::clone))
585 }
586
587 #[cfg(unix)]
588 fn prune_output(input: &str) -> String {
589 if input.len() <= OUTPUT_LIMIT {
590 return input.to_string();
591 }
592 let head = OUTPUT_LIMIT / 3;
593 let tail = OUTPUT_LIMIT - head;
594 let head_end = input
595 .char_indices()
596 .find(|(index, _)| *index >= head)
597 .map_or(input.len(), |(index, _)| index);
598 // The earliest boundary that still fits: the tail is where a build's
599 // error and the command's completion marker are.
600 let tail_start = input
601 .char_indices()
602 .map(|(index, _)| index)
603 .find(|index| *index >= head_end && input.len() - *index <= tail)
604 .unwrap_or(input.len());
605 format!(
606 "{}\n… [output truncated: {} bytes omitted] …\n{}",
607 &input[..head_end],
608 tail_start - head_end,
609 &input[tail_start..]
610 )
611 }
612
613 #[cfg(unix)]
614 fn completion(session: &TerminalSession) -> Option<(i32, String)> {
615 let state = session.command.as_ref()?;
616 let output = output_snapshot(session);
617 let marker = format!("\n{}:", state.marker);
618 let line = output
619 .rsplit(&marker)
620 .next()
621 .and_then(|tail| tail.lines().next())?;
622 let (status, cwd) = line.split_once(':')?;
623 Some((status.parse().ok()?, cwd.to_string()))
624 }
625
626 #[cfg(unix)]
627 fn start_command(session: &mut TerminalSession, command: &str) -> Result<(), String> {
628 if session.command.is_some() && completion(session).is_none() {
629 return Err("terminal session already has a running foreground command".to_string());
630 }
631 let marker = format!("__CODEWHALE_TERM_{}__", Uuid::new_v4().simple());
632 // The command must run in the CURRENT shell — a subshell would discard
633 // exactly the state (cd, exports, functions, activated envs) this tool
634 // exists to preserve (EXEC-001). The sentinel line is typed after the
635 // command; the tty line discipline holds it until the foreground command
636 // finishes reading input.
637 let wrapped = format!(
638 "{command}\n__cw_status=$?; printf '\\n{marker}:%s:%s\\n' \"$__cw_status\" \"$PWD\"\n"
639 );
640 write_bytes(session, wrapped.as_bytes())?;
641 session.command = Some(CommandState { marker });
642 session.durable.state = DurableTerminalState::Running;
643 session.durable.updated_at = chrono::Utc::now().to_rfc3339();
644 persist_durable(&session.durable_path, &session.durable)?;
645 Ok(())
646 }
647
648 #[cfg(unix)]
649 fn write_completion_sentinel(
650 session: &TerminalSession,
651 marker: &str,
652 status: i32,
653 ) -> Result<(), String> {
654 let sentinel =
655 format!("__cw_status={status}; printf '\\n{marker}:%s:%s\\n' \"$__cw_status\" \"$PWD\"\n");
656 write_bytes(session, sentinel.as_bytes())
657 }
658
659 #[cfg(unix)]
660 #[cfg(test)]
661 fn wait_session(session: &mut TerminalSession, timeout: Duration) -> (Option<(i32, String)>, bool) {
662 let deadline = Instant::now() + timeout;
663 loop {
664 if let Some(done) = completion(session) {
665 return (Some(done), false);
666 }
667 if Instant::now() >= deadline {
668 return (None, true);
669 }
670 std::thread::sleep(Duration::from_millis(25));
671 }
672 }
673
674 /// Wait without monopolizing the per-session lock. `terminal/send` and
675 /// `terminal/cancel` must be able to acquire the lock while a foreground
676 /// command is active; otherwise interactive input and cancellation deadlock
677 /// behind the waiter.
678 #[cfg(unix)]
679 fn wait_shared_session(
680 session: &SharedSession,
681 timeout: Duration,
682 ) -> Result<(Option<(i32, String)>, bool), ToolError> {
683 let deadline = Instant::now() + timeout;
684 loop {
685 let done = {
686 let session = session
687 .lock()
688 .map_err(|_| ToolError::execution_failed("terminal session lock poisoned"))?;
689 completion(&session)
690 };
691 if let Some(done) = done {
692 return Ok((Some(done), false));
693 }
694 if Instant::now() >= deadline {
695 return Ok((None, true));
696 }
697 std::thread::sleep(Duration::from_millis(25));
698 }
699 }
700
701 /// Interrupt a foreground command and wait until the persistent shell confirms
702 /// that it is ready again.
703 ///
704 /// Canonical PTYs may flush queued input when they process ETX. Resubmit the
705 /// completion sentinel until the shell acknowledges it so a busy host cannot
706 /// strand the session merely because the one post-interrupt write raced that
707 /// flush.
708 #[cfg(unix)]
709 fn cancel_shared_session(session: &SharedSession) -> Result<(i32, String), ToolError> {
710 let marker = {
711 let session = session
712 .lock()
713 .map_err(|_| ToolError::execution_failed("terminal session lock poisoned"))?;
714 if session.command.is_none() || completion(&session).is_some() {
715 return Err(ToolError::execution_failed(
716 "terminal session has no running foreground command",
717 ));
718 }
719 write_bytes(&session, &[3]).map_err(ToolError::execution_failed)?;
720 session
721 .command
722 .as_ref()
723 .expect("running command checked above")
724 .marker
725 .clone()
726 };
727 let deadline = Instant::now() + CANCEL_CONFIRM_TIMEOUT;
728
729 loop {
730 std::thread::sleep(CANCEL_SENTINEL_RETRY_INTERVAL);
731 let session = session
732 .lock()
733 .map_err(|_| ToolError::execution_failed("terminal session lock poisoned"))?;
734 if let Some(done) = completion(&session) {
735 return Ok(done);
736 }
737 if Instant::now() >= deadline {
738 return Err(ToolError::execution_failed(
739 "terminal interrupt was sent, but cancellation was not confirmed within 2 seconds",
740 ));
741 }
742 write_completion_sentinel(&session, &marker, 130).map_err(ToolError::execution_failed)?;
743 }
744 }
745
746 #[cfg(unix)]
747 fn session_result(
748 session: &mut TerminalSession,
749 done: Option<(i32, String)>,
750 timed_out: bool,
751 ) -> ToolResult {
752 let output = prune_output(&take_output(session));
753 let finished = done.is_some();
754 let (exit_code, cwd) = done.map_or((None, String::new()), |(code, cwd)| (Some(code), cwd));
755 let status = if timed_out {
756 "timed_out"
757 } else if finished {
758 "completed"
759 } else {
760 "running"
761 };
762 if !cwd.is_empty() {
763 session.durable.last_known_cwd = PathBuf::from(&cwd);
764 }
765 session.durable.state = if timed_out || !finished {
766 DurableTerminalState::Running
767 } else if exit_code == Some(0) {
768 DurableTerminalState::Idle
769 } else if exit_code == Some(130) {
770 DurableTerminalState::Canceled
771 } else {
772 DurableTerminalState::Failed
773 };
774 session.durable.updated_at = chrono::Utc::now().to_rfc3339();
775 let persistence_error = persist_durable(&session.durable_path, &session.durable).err();
776 let previous = session.durable.previous.as_deref().map(|record| {
777 json!({
778 "session_id": record.session_id,
779 "state": record.state,
780 "last_known_cwd": record.last_known_cwd,
781 "updated_at": record.updated_at,
782 })
783 });
784 ToolResult {
785 content: output,
786 // A successful send while the command is still running is itself a
787 // successful tool operation. Completed commands still report their
788 // real exit status, and timeouts remain unsuccessful.
789 success: !timed_out && (!finished || exit_code == Some(0)),
790 metadata: Some(json!({
791 "status": status,
792 "exit_code": exit_code,
793 "cwd": cwd,
794 "session_persistent": true,
795 "durability": "live shell is process-local; identity and last-known summary persist",
796 "terminal_session_id": session.durable.session_id,
797 "terminal_state": session.durable.state,
798 "state_path": session.durable_path,
799 "previous_session": previous,
800 "persistence_error": persistence_error,
801 })),
802 }
803 }
804
805 #[cfg(unix)]
806 fn session_name(input: &serde_json::Value, required: bool) -> Result<&str, ToolError> {
807 match input.get("session").and_then(serde_json::Value::as_str) {
808 Some(name) if !name.is_empty() => Ok(name),
809 Some(_) => Err(ToolError::execution_failed("session must not be empty")),
810 None if required => Err(ToolError::missing_field("session")),
811 None => Ok("term-1"),
812 }
813 }
814
815 #[cfg(unix)]
816 fn timeout_secs(input: &serde_json::Value, key: &str) -> Result<Duration, ToolError> {
817 Ok(Duration::from_secs(
818 optional_u64(input, key, DEFAULT_TIMEOUT_SECS)?.clamp(1, MAX_TIMEOUT_SECS),
819 ))
820 }
821
822 fn shell_allowed(context: &ToolContext, name: &str) -> Result<(), ToolError> {
823 crate::core::engine::tool_catalog::enforce_tool_denial(context, name, &json!({}))?;
824 if matches!(name, "terminal/run" | "terminal/send" | "terminal/reset")
825 && context.shell_policy != crate::worker_profile::ShellPolicy::Full
826 {
827 return Err(ToolError::permission_denied(
828 "Persistent terminal execution and input require full shell permission.",
829 ));
830 }
831 if context.shell_policy.allows_shell() {
832 Ok(())
833 } else {
834 Err(ToolError::execution_failed(
835 "Shell tools are disabled by the active permission profile.",
836 ))
837 }
838 }
839
840 #[cfg(not(unix))]
841 fn unsupported() -> ToolResult {
842 ToolResult::error("Stateful terminal sessions are currently supported on Unix only.")
843 }
844
845 macro_rules! terminal_tool_common {
846 ($name:literal, $description:literal) => {
847 fn name(&self) -> &'static str {
848 $name
849 }
850 fn description(&self) -> &'static str {
851 $description
852 }
853 fn capabilities(&self) -> Vec<ToolCapability> {
854 vec![
855 ToolCapability::ExecutesCode,
856 ToolCapability::RequiresApproval,
857 ]
858 }
859 fn approval_requirement(&self) -> ApprovalRequirement {
860 ApprovalRequirement::Required
861 }
862 };
863 }
864
865 pub struct TerminalRunTool;
866 #[async_trait]
867 impl ToolSpec for TerminalRunTool {
868 terminal_tool_common!(
869 "terminal/run",
870 "Run a command in a persistent PTY shell session. cd, exports, shell functions, and activated environments persist across calls in this process. Identity and a non-secret last-known summary persist across restarts; prior shells are surfaced as stale/lost and are never reattached. On timeout the wait is abandoned but the command keeps running in the session (use terminal/wait or cancel). (Unix only)"
871 );
872 fn input_schema(&self) -> serde_json::Value {
873 json!({"type":"object","properties":{"command":{"type":"string"},"session":{"type":"string","default":"term-1"},"timeout_secs":{"type":"integer","default":120}},"required":["command"]})
874 }
875 async fn execute(
876 &self,
877 input: serde_json::Value,
878 context: &ToolContext,
879 ) -> Result<ToolResult, ToolError> {
880 shell_allowed(context, self.name())?;
881 #[cfg(unix)]
882 {
883 let command = required_str(&input, "command")?.to_string();
884 let name = session_name(&input, false)?.to_string();
885 let session = get_or_create(&name, &context.workspace, pty_sandbox_policy(context))
886 .map_err(ToolError::execution_failed)?;
887 let timeout = timeout_secs(&input, "timeout_secs")?;
888 return tokio::task::spawn_blocking(move || {
889 {
890 let mut session = session.lock().map_err(|_| {
891 ToolError::execution_failed("terminal session lock poisoned")
892 })?;
893 start_command(&mut session, &command).map_err(ToolError::execution_failed)?;
894 }
895 let (done, timed_out) = wait_shared_session(&session, timeout)?;
896 let mut session = session
897 .lock()
898 .map_err(|_| ToolError::execution_failed("terminal session lock poisoned"))?;
899 Ok(session_result(&mut session, done, timed_out))
900 })
901 .await
902 .map_err(|e| ToolError::execution_failed(e.to_string()))?;
903 }
904 #[cfg(not(unix))]
905 {
906 let _ = input;
907 Ok(unsupported())
908 }
909 }
910 }
911
912 pub struct TerminalSendTool;
913 #[async_trait]
914 impl ToolSpec for TerminalSendTool {
915 terminal_tool_common!(
916 "terminal/send",
917 "Send raw input to a live persistent terminal session. Use a literal ETX control byte to interrupt an interactive process. A prior-process shell is reported as stale/lost rather than reattached. (Unix only)"
918 );
919 fn input_schema(&self) -> serde_json::Value {
920 json!({"type":"object","properties":{"session":{"type":"string"},"text":{"type":"string"},"wait_ms":{"type":"integer","default":250}},"required":["session","text"]})
921 }
922 async fn execute(
923 &self,
924 input: serde_json::Value,
925 context: &ToolContext,
926 ) -> Result<ToolResult, ToolError> {
927 shell_allowed(context, self.name())?;
928 #[cfg(unix)]
929 {
930 let name = session_name(&input, true)?.to_string();
931 let text = required_str(&input, "text")?.as_bytes().to_vec();
932 let session = find(&name, &context.workspace).map_err(ToolError::execution_failed)?;
933 let wait = Duration::from_millis(optional_u64(&input, "wait_ms", 250)?.min(60_000));
934 return tokio::task::spawn_blocking(move || {
935 let mut session = session
936 .lock()
937 .map_err(|_| ToolError::execution_failed("terminal session lock poisoned"))?;
938 write_bytes(&session, &text).map_err(ToolError::execution_failed)?;
939 std::thread::sleep(wait);
940 let done = completion(&session);
941 Ok(session_result(&mut session, done, false))
942 })
943 .await
944 .map_err(|e| ToolError::execution_failed(e.to_string()))?;
945 }
946 #[cfg(not(unix))]
947 {
948 let _ = input;
949 Ok(unsupported())
950 }
951 }
952 }
953
954 pub struct TerminalWaitTool;
955 #[async_trait]
956 impl ToolSpec for TerminalWaitTool {
957 terminal_tool_common!(
958 "terminal/wait",
959 "Wait for the current foreground command in a live persistent terminal session and return buffered output. A prior-process shell is reported as stale/lost rather than reattached. Output buffer holds at most 512KiB; older bytes are dropped silently. (Unix only)"
960 );
961 fn input_schema(&self) -> serde_json::Value {
962 json!({"type":"object","properties":{"session":{"type":"string"},"timeout_secs":{"type":"integer","default":120}},"required":["session"]})
963 }
964 async fn execute(
965 &self,
966 input: serde_json::Value,
967 context: &ToolContext,
968 ) -> Result<ToolResult, ToolError> {
969 shell_allowed(context, self.name())?;
970 #[cfg(unix)]
971 {
972 let name = session_name(&input, true)?.to_string();
973 let session = find(&name, &context.workspace).map_err(ToolError::execution_failed)?;
974 let timeout = timeout_secs(&input, "timeout_secs")?;
975 return tokio::task::spawn_blocking(move || {
976 let (done, timed_out) = wait_shared_session(&session, timeout)?;
977 let mut session = session
978 .lock()
979 .map_err(|_| ToolError::execution_failed("terminal session lock poisoned"))?;
980 Ok(session_result(&mut session, done, timed_out))
981 })
982 .await
983 .map_err(|e| ToolError::execution_failed(e.to_string()))?;
984 }
985 #[cfg(not(unix))]
986 {
987 let _ = input;
988 Ok(unsupported())
989 }
990 }
991 }
992
993 pub struct TerminalCancelTool;
994 #[async_trait]
995 impl ToolSpec for TerminalCancelTool {
996 terminal_tool_common!(
997 "terminal/cancel",
998 "Interrupt the running foreground command with ETX. The live terminal session survives and can be reused; its non-secret summary persists. (Unix only)"
999 );
1000 fn input_schema(&self) -> serde_json::Value {
1001 json!({"type":"object","properties":{"session":{"type":"string"}},"required":["session"]})
1002 }
1003 async fn execute(
1004 &self,
1005 input: serde_json::Value,
1006 context: &ToolContext,
1007 ) -> Result<ToolResult, ToolError> {
1008 shell_allowed(context, self.name())?;
1009 #[cfg(unix)]
1010 {
1011 let name = session_name(&input, true)?.to_string();
1012 let session = find(&name, &context.workspace).map_err(ToolError::execution_failed)?;
1013 return tokio::task::spawn_blocking(move || {
1014 let done = cancel_shared_session(&session)?;
1015 let mut session = session
1016 .lock()
1017 .map_err(|_| ToolError::execution_failed("terminal session lock poisoned"))?;
1018 let mut result = session_result(&mut session, Some(done), false);
1019 result.success = true;
1020 if let Some(metadata) = result.metadata.as_mut() {
1021 metadata["status"] = json!("canceled");
1022 metadata["canceled"] = json!(true);
1023 }
1024 Ok(result)
1025 })
1026 .await
1027 .map_err(|e| ToolError::execution_failed(e.to_string()))?;
1028 }
1029 #[cfg(not(unix))]
1030 {
1031 let _ = input;
1032 Ok(unsupported())
1033 }
1034 }
1035 }
1036
1037 pub struct TerminalResetTool;
1038 #[async_trait]
1039 impl ToolSpec for TerminalResetTool {
1040 terminal_tool_common!(
1041 "terminal/reset",
1042 "Kill and recreate a persistent terminal session with a fresh environment. This loses live cd, exports, functions, activated environments, and running work while retaining the prior historical summary. (Unix only)"
1043 );
1044 fn input_schema(&self) -> serde_json::Value {
1045 json!({"type":"object","properties":{"session":{"type":"string"}},"required":["session"]})
1046 }
1047 async fn execute(
1048 &self,
1049 input: serde_json::Value,
1050 context: &ToolContext,
1051 ) -> Result<ToolResult, ToolError> {
1052 shell_allowed(context, self.name())?;
1053 #[cfg(unix)]
1054 {
1055 let name = session_name(&input, true)?.to_string();
1056 let old = find(&name, &context.workspace).map_err(ToolError::execution_failed)?;
1057 let workspace = context.workspace.clone();
1058 let policy = pty_sandbox_policy(context);
1059 return tokio::task::spawn_blocking(move || {
1060 if let Ok(mut old) = old.lock() { let _ = old.child.kill(); }
1061 if let Ok(mut old) = old.lock() {
1062 old.durable.state = DurableTerminalState::Reset;
1063 old.durable.updated_at = chrono::Utc::now().to_rfc3339();
1064 let _ = persist_durable(&old.durable_path, &old.durable);
1065 }
1066 let fresh = create_session(&name, &workspace, policy)
1067 .map_err(ToolError::execution_failed)?;
1068 sessions().lock().map_err(|_| ToolError::execution_failed("terminal session registry lock poisoned"))?.insert(session_key(&name, &workspace), fresh);
1069 Ok(ToolResult { content: format!("Reset terminal session '{name}'. Lost shell state and any running command."), success: true, metadata: Some(json!({"session":name,"reset":true,"lost_state":["cwd","environment","functions","activated environments","running command"]})) })
1070 }).await.map_err(|e| ToolError::execution_failed(e.to_string()))?;
1071 }
1072 #[cfg(not(unix))]
1073 {
1074 let _ = input;
1075 Ok(unsupported())
1076 }
1077 }
1078 }
1079
1080 #[cfg(all(test, unix))]
1081 mod tests {
1082 use super::*;
1083
1084 fn fresh(name: &str) -> SharedSession {
1085 let session = get_or_create(
1086 name,
1087 std::path::Path::new("/tmp"),
1088 crate::sandbox::SandboxPolicy::DangerFullAccess,
1089 )
1090 .unwrap();
1091 let mut session_guard = session.lock().unwrap();
1092 let _ = session_guard.child.kill();
1093 drop(session_guard);
1094 let replacement = create_session(
1095 name,
1096 std::path::Path::new("/tmp"),
1097 crate::sandbox::SandboxPolicy::DangerFullAccess,
1098 )
1099 .unwrap();
1100 sessions().lock().unwrap().insert(
1101 session_key(name, Path::new("/tmp")),
1102 Arc::clone(&replacement),
1103 );
1104 replacement
1105 }
1106
1107 fn run(session: &SharedSession, command: &str, timeout: Duration) -> ToolResult {
1108 let mut session = session.lock().unwrap();
1109 start_command(&mut session, command).unwrap();
1110 let (done, timed_out) = wait_session(&mut session, timeout);
1111 session_result(&mut session, done, timed_out)
1112 }
1113
1114 #[test]
1115 fn prune_output_keeps_the_tail_and_counts_what_it_dropped() {
1116 let input = format!(
1117 "first-line\n{}\nerror: the last line é\n",
1118 "x".repeat(OUTPUT_LIMIT * 2)
1119 );
1120 let pruned = prune_output(&input);
1121 let (head, rest) = pruned
1122 .split_once("\n… [output truncated: ")
1123 .expect("truncation marker");
1124 let (omitted, tail) = rest.split_once(" bytes omitted] …\n").expect("count");
1125 let omitted: usize = omitted.parse().expect("omitted count");
1126 assert!(head.starts_with("first-line"));
1127 assert!(tail.ends_with("error: the last line é\n"), "{tail:?}");
1128 assert!(tail.len() > OUTPUT_LIMIT / 2, "tail kept only {tail:?}");
1129 assert_eq!(head.len() + omitted + tail.len(), input.len());
1130 }
1131
1132 // The whole family is compiled and registered on Unix only, a run's
1133 // timeout abandons the wait without stopping the command, and the output
1134 // ring silently drops the oldest bytes past 512KiB — the descriptions
1135 // must say all three so the model can plan around them.
1136 #[test]
1137 fn terminal_descriptions_disclose_platform_and_wait_semantics() {
1138 let descriptions = [
1139 (TerminalRunTool.name(), TerminalRunTool.description()),
1140 (TerminalSendTool.name(), TerminalSendTool.description()),
1141 (TerminalWaitTool.name(), TerminalWaitTool.description()),
1142 (TerminalCancelTool.name(), TerminalCancelTool.description()),
1143 (TerminalResetTool.name(), TerminalResetTool.description()),
1144 ];
1145 for (name, description) in descriptions {
1146 assert!(
1147 description.contains("(Unix only)"),
1148 "{name} must disclose its platform restriction: {description}"
1149 );
1150 }
1151 assert!(
1152 TerminalRunTool
1153 .description()
1154 .contains("the command keeps running in the session"),
1155 "terminal/run must disclose that a timeout abandons the wait, not the command"
1156 );
1157 assert!(
1158 TerminalWaitTool.description().contains("512KiB"),
1159 "terminal/wait must disclose the retained-output bound"
1160 );
1161 }
1162
1163 #[test]
1164 #[cfg(unix)]
1165 fn cd_persists_between_runs() {
1166 let session = fresh("test-cd");
1167 let _ = run(
1168 &session,
1169 "mkdir -p /tmp/cw-term-cd-proof && cd /tmp/cw-term-cd-proof",
1170 Duration::from_secs(3),
1171 );
1172 // A separate run must still be inside the directory — this is the
1173 // whole point of the stateful session (EXEC-001).
1174 let result = run(&session, "pwd", Duration::from_secs(3));
1175 assert!(result.content.contains("cw-term-cd-proof"), "{}", {
1176 &result.content
1177 });
1178 }
1179
1180 #[test]
1181 #[cfg(unix)]
1182 fn export_persists_and_sessions_are_isolated() {
1183 let one = fresh("test-env-one");
1184 let two = fresh("test-env-two");
1185 let _ = run(&one, "export CW_TERM_TEST=present", Duration::from_secs(3));
1186 assert!(
1187 run(&one, "printf %s $CW_TERM_TEST", Duration::from_secs(3))
1188 .content
1189 .contains("present")
1190 );
1191 assert!(
1192 !run(
1193 &two,
1194 "printf %s ${CW_TERM_TEST-unset}",
1195 Duration::from_secs(3)
1196 )
1197 .content
1198 .contains("present")
1199 );
1200 }
1201
1202 #[test]
1203 #[cfg(unix)]
1204 fn reset_replaces_shell_environment() {
1205 let session = fresh("test-reset");
1206 let _ = run(
1207 &session,
1208 "export CW_TERM_RESET=present",
1209 Duration::from_secs(3),
1210 );
1211 assert!(
1212 run(&session, "printf %s $CW_TERM_RESET", Duration::from_secs(3))
1213 .content
1214 .contains("present")
1215 );
1216 let _ = session.lock().unwrap().child.kill();
1217 let replacement = create_session(
1218 "test-reset",
1219 std::path::Path::new("/tmp"),
1220 crate::sandbox::SandboxPolicy::DangerFullAccess,
1221 )
1222 .unwrap();
1223 sessions().lock().unwrap().insert(
1224 session_key("test-reset", Path::new("/tmp")),
1225 Arc::clone(&replacement),
1226 );
1227 assert!(
1228 !run(
1229 &replacement,
1230 "printf %s ${CW_TERM_RESET-unset}",
1231 Duration::from_secs(3)
1232 )
1233 .content
1234 .contains("present")
1235 );
1236 }
1237
1238 #[test]
1239 #[cfg(unix)]
1240 fn timeout_leaves_session_alive() {
1241 let session = fresh("test-timeout");
1242 let result = run(
1243 &session,
1244 "printf before; sleep 2",
1245 Duration::from_millis(100),
1246 );
1247 assert_eq!(result.metadata.as_ref().unwrap()["status"], "timed_out");
1248 let result = {
1249 let mut session_guard = session.lock().unwrap();
1250 let (done, timed_out) = wait_session(&mut session_guard, Duration::from_secs(3));
1251 session_result(&mut session_guard, done, timed_out)
1252 };
1253 assert_eq!(result.metadata.as_ref().unwrap()["status"], "completed");
1254 let result = run(&session, "printf after", Duration::from_secs(3));
1255 assert!(result.content.contains("after"));
1256 }
1257
1258 #[test]
1259 #[cfg(unix)]
1260 fn cancel_interrupts_sleep_and_session_survives() {
1261 let session = fresh("test-cancel");
1262 // Prove the interactive shell is initialized before interrupting it.
1263 // A SIGINT delivered while `sh -i` is still starting can kill the
1264 // shell itself, which later surfaces as EIO on the next PTY write
1265 // (hosted macOS run 34716492759 under load).
1266 assert!(
1267 run(&session, "printf ready", Duration::from_secs(10))
1268 .content
1269 .contains("ready")
1270 );
1271 let worker = Arc::clone(&session);
1272 {
1273 let mut guard = worker.lock().unwrap();
1274 // The quoted split keeps the marker out of the echoed command
1275 // line, so seeing it proves the shell reached this command.
1276 start_command(&mut guard, "printf 'st''arted'; sleep 10").unwrap();
1277 }
1278 let deadline = Instant::now() + Duration::from_secs(10);
1279 loop {
1280 let started = output_snapshot(&session.lock().unwrap()).contains("started");
1281 if started {
1282 break;
1283 }
1284 assert!(
1285 Instant::now() < deadline,
1286 "shell never reached the sleep command"
1287 );
1288 std::thread::sleep(Duration::from_millis(25));
1289 }
1290 let handle = std::thread::spawn(move || {
1291 let (done, timed_out) = wait_shared_session(&worker, Duration::from_secs(30)).unwrap();
1292 let mut guard = worker.lock().unwrap();
1293 session_result(&mut guard, done, timed_out)
1294 });
1295 std::thread::sleep(Duration::from_millis(150));
1296 let done = cancel_shared_session(&session).unwrap();
1297 let result = handle.join().unwrap();
1298 assert_eq!(done.0, 130);
1299 assert_eq!(result.metadata.as_ref().unwrap()["exit_code"], 130);
1300 assert!(
1301 run(&session, "printf alive", Duration::from_secs(3))
1302 .content
1303 .contains("alive")
1304 );
1305 }
1306
1307 #[test]
1308 fn narrowed_posture_without_a_backend_fails_closed() {
1309 assert!(pty_lacks_required_sandbox(
1310 &crate::sandbox::SandboxPolicy::ReadOnly,
1311 crate::sandbox::SandboxType::None,
1312 ));
1313 assert!(!pty_lacks_required_sandbox(
1314 &crate::sandbox::SandboxPolicy::DangerFullAccess,
1315 crate::sandbox::SandboxType::None,
1316 ));
1317 }
1318
1319 #[test]
1320 fn full_access_prepares_an_unsandboxed_pty() {
1321 let (prepared, confined) = prepare_pty_shell(
1322 "/bin/sh",
1323 Path::new("/tmp"),
1324 crate::sandbox::SandboxPolicy::DangerFullAccess,
1325 )
1326 .expect("full access may start a PTY");
1327 assert!(!confined);
1328 assert_eq!(prepared.command[0], "/bin/sh");
1329 assert_eq!(prepared.command[1], "-i");
1330 }
1331
1332 #[test]
1333 #[cfg(unix)]
1334 fn session_names_are_scoped_to_workspace() {
1335 let first_workspace = tempfile::tempdir().unwrap();
1336 let second_workspace = tempfile::tempdir().unwrap();
1337 let first = get_or_create(
1338 "shared-name",
1339 first_workspace.path(),
1340 crate::sandbox::SandboxPolicy::DangerFullAccess,
1341 )
1342 .unwrap();
1343 let second = get_or_create(
1344 "shared-name",
1345 second_workspace.path(),
1346 crate::sandbox::SandboxPolicy::DangerFullAccess,
1347 )
1348 .unwrap();
1349 assert!(!Arc::ptr_eq(&first, &second));
1350 assert!(find("shared-name", first_workspace.path()).is_ok());
1351 assert!(find("shared-name", second_workspace.path()).is_ok());
1352 assert!(find("shared-name", Path::new("/tmp")).is_err());
1353 }
1354
1355 #[test]
1356 #[cfg(unix)]
1357 fn durable_summary_marks_prior_process_shell_stale_and_preserves_history() {
1358 let workspace = tempfile::tempdir().unwrap();
1359 let canonical_workspace = workspace.path().canonicalize().unwrap();
1360 let name = format!("restart-proof-{}", Uuid::new_v4());
1361 let first = create_session(
1362 &name,
1363 workspace.path(),
1364 crate::sandbox::SandboxPolicy::DangerFullAccess,
1365 )
1366 .unwrap();
1367 let first_record = first.lock().unwrap().durable.clone();
1368 assert_eq!(first_record.state, DurableTerminalState::Idle);
1369 assert_eq!(first_record.last_known_cwd, canonical_workspace);
1370 assert!(!first_record.session_id.is_empty());
1371
1372 // Removing only the process-local registry entry models a restart.
1373 // The durable record remains, but cannot be used to reattach.
1374 sessions()
1375 .lock()
1376 .unwrap()
1377 .remove(&session_key(&name, workspace.path()));
1378 let stale = match find(&name, workspace.path()) {
1379 Ok(_) => panic!("persisted session must not be reattached"),
1380 Err(error) => error,
1381 };
1382 assert!(stale.contains("stale/lost"), "{stale}");
1383 assert!(stale.contains("start a replacement"), "{stale}");
1384
1385 let replacement = create_session(
1386 &name,
1387 workspace.path(),
1388 crate::sandbox::SandboxPolicy::DangerFullAccess,
1389 )
1390 .unwrap();
1391 let replacement = replacement.lock().unwrap();
1392 assert_ne!(replacement.durable.session_id, first_record.session_id);
1393 let previous = replacement.durable.previous.as_deref().unwrap();
1394 assert_eq!(previous.session_id, first_record.session_id);
1395 assert_eq!(previous.state, DurableTerminalState::StaleLost);
1396 assert_eq!(previous.last_known_cwd, canonical_workspace);
1397 assert_ne!(replacement.durable.runtime_nonce, "");
1398 }
1399
1400 #[test]
1401 #[cfg(unix)]
1402 fn output_is_capped_with_notice() {
1403 let session = fresh("test-output-cap");
1404 let result = run(&session, "yes x | head -n 100000", Duration::from_secs(3));
1405 assert!(result.content.len() <= OUTPUT_LIMIT + 100);
1406 assert!(result.content.contains("output truncated"));
1407 }
1408
1409 /// The cursor contract the Engine byte stream rides on: offsets are
1410 /// absolute for the life of the PTY, reads never consume, and a cursor the
1411 /// retained window has moved past is reported rather than silently
1412 /// answered from the middle of the stream.
1413 #[test]
1414 #[cfg(unix)]
1415 fn bounded_replay_is_absolute_non_consuming_and_reports_its_gap() {
1416 let mut buffer = OutputBuffer::default();
1417 buffer.append(b"hello");
1418 let first = buffer.read_since(0, 4);
1419 assert_eq!(first.bytes, b"hell");
1420 assert_eq!(first.offset, 0);
1421 assert_eq!(first.next_cursor, 4);
1422 assert_eq!(first.total, 5);
1423 assert_eq!(first.dropped, 0);
1424 assert!(!first.gap);
1425 // A second read at the same cursor returns the same bytes: the caller
1426 // owns the position, so replay is idempotent.
1427 assert_eq!(buffer.read_since(0, 4), first);
1428 let rest = buffer.read_since(first.next_cursor, 4);
1429 assert_eq!(rest.bytes, b"o");
1430 assert_eq!(rest.offset, 4);
1431 assert_eq!(rest.next_cursor, 5);
1432
1433 // Past the ring: the first 32 bytes are gone and the chunk says so.
1434 let mut wrapped = OutputBuffer::default();
1435 wrapped.append(&vec![b'x'; BUFFER_LIMIT + 32]);
1436 let lost = wrapped.read_since(0, 8);
1437 assert!(lost.gap, "a cursor below the retained window is a gap");
1438 assert_eq!(lost.dropped, 32);
1439 assert_eq!(
1440 lost.offset, 32,
1441 "the chunk starts at the oldest retained byte"
1442 );
1443 assert_eq!(lost.total, BUFFER_LIMIT as u64 + 32);
1444 assert_eq!(lost.bytes, vec![b'x'; 8]);
1445 assert_eq!(lost.next_cursor, 40);
1446
1447 // A cursor at the head is not a gap and returns nothing new.
1448 let current = wrapped.read_since(BUFFER_LIMIT as u64 + 32, 8);
1449 assert!(!current.gap);
1450 assert!(current.bytes.is_empty());
1451 assert_eq!(current.next_cursor, BUFFER_LIMIT as u64 + 32);
1452
1453 // A cursor past the head clamps to it instead of being echoed back:
1454 // a client that continues from `next_cursor` must not skip the bytes
1455 // the stream produces before it reaches the bogus position.
1456 let beyond = wrapped.read_since(BUFFER_LIMIT as u64 + 4096, 8);
1457 assert!(!beyond.gap);
1458 assert!(beyond.bytes.is_empty());
1459 assert_eq!(beyond.offset, BUFFER_LIMIT as u64 + 32);
1460 assert_eq!(beyond.next_cursor, BUFFER_LIMIT as u64 + 32);
1461 }
1462
1463 /// The session-level entry point an Engine byte stream will call: absolute
1464 /// cursor, non-consuming, clamped, and honest about a cursor ahead of the
1465 /// stream.
1466 #[test]
1467 #[cfg(unix)]
1468 fn session_read_since_is_non_consuming_and_clamped() {
1469 let session = fresh("test-read-since");
1470 let _ = run(
1471 &session,
1472 "printf 'cw-replay-proof\\n'",
1473 Duration::from_secs(3),
1474 );
1475 // The tool-result path already consumed this output through its own
1476 // cursor; an absolute cursor still reads it, which is the point.
1477 let printed = read_session_since(&session.lock().unwrap(), 0, usize::MAX).unwrap();
1478 assert!(
1479 String::from_utf8_lossy(&printed.bytes).contains("cw-replay-proof"),
1480 "{}",
1481 String::from_utf8_lossy(&printed.bytes)
1482 );
1483 let again = read_session_since(&session.lock().unwrap(), 0, usize::MAX).unwrap();
1484 assert_eq!(again.bytes, printed.bytes, "a replay read must not consume");
1485
1486 // A cursor ahead of the stream is not a gap: nothing was lost, there
1487 // is simply nothing there yet — and the answer clamps to the head so
1488 // continuing from it cannot skip what arrives next.
1489 let ahead =
1490 read_session_since(&session.lock().unwrap(), printed.next_cursor + 4096, 16).unwrap();
1491 assert!(!ahead.gap);
1492 assert!(ahead.bytes.is_empty());
1493 assert_eq!(ahead.next_cursor, ahead.total);
1494
1495 // More than one response's worth of output proves the clamp.
1496 let large = fresh("test-read-since-clamp");
1497 let _ = run(&large, "yes x | head -n 100000", Duration::from_secs(3));
1498 let chunk = read_session_since(&large.lock().unwrap(), 0, usize::MAX).unwrap();
1499 assert_eq!(chunk.bytes.len(), READ_LIMIT);
1500 assert!(!chunk.gap);
1501 assert_eq!(chunk.next_cursor, READ_LIMIT as u64);
1502 }
1503
1504 /// Resize has to reach the kernel, not just a field: `get_size` reads the
1505 /// window back from the pty, and the shell reports it through `stty`.
1506 #[test]
1507 #[cfg(unix)]
1508 fn resize_reaches_the_kernel_and_the_live_shell() {
1509 let session = fresh("test-resize");
1510 {
1511 let guard = session.lock().unwrap();
1512 assert_eq!(guard.master.get_size().unwrap().rows, 24);
1513 }
1514 resize_session(&session.lock().unwrap(), 40, 100).unwrap();
1515 let size = session.lock().unwrap().master.get_size().unwrap();
1516 assert_eq!((size.rows, size.cols), (40, 100));
1517 let result = run(&session, "stty size", Duration::from_secs(3));
1518 assert!(
1519 result.content.contains("40 100"),
1520 "the shell should see the new window: {}",
1521 result.content
1522 );
1523 }
1524
1525 /// Exit is observable: a killed shell reports a status instead of looking
1526 /// alive forever (the "pretending to reattach" failure the durable record
1527 /// is written to avoid).
1528 #[test]
1529 #[cfg(unix)]
1530 fn killed_shell_reports_an_exit_status() {
1531 let session = fresh("test-exit-status");
1532 {
1533 let mut guard = session.lock().unwrap();
1534 assert!(
1535 session_exit_status(&mut guard).unwrap().is_none(),
1536 "fresh shell is alive"
1537 );
1538 kill_session(&mut guard).unwrap();
1539 }
1540 let deadline = Instant::now() + Duration::from_secs(5);
1541 loop {
1542 let exited = session_exit_status(&mut session.lock().unwrap())
1543 .unwrap()
1544 .is_some();
1545 if exited {
1546 break;
1547 }
1548 assert!(
1549 Instant::now() < deadline,
1550 "a killed shell must report its exit"
1551 );
1552 std::thread::sleep(Duration::from_millis(20));
1553 }
1554 }
1555
1556 #[test]
1557 #[cfg(unix)]
1558 fn new_session_does_not_inherit_parent_secret_env() {
1559 use crate::test_support::{EnvVarGuard, lock_test_env};
1560 let _env_lock = lock_test_env();
1561 let _secret = EnvVarGuard::set("CODEWHALE_TEST_PTY_SECRET", "pty-secret-value");
1562 let session = fresh("test-env-scrub");
1563 let result = run(
1564 &session,
1565 "printf 'se''cret=%s ho''me=%s' \"${CODEWHALE_TEST_PTY_SECRET-unset}\" \"${HOME-none}\"",
1566 Duration::from_secs(10),
1567 );
1568 assert!(
1569 !result.content.contains("pty-secret-value"),
1570 "{}",
1571 result.content
1572 );
1573 assert!(
1574 result.content.contains("secret=unset"),
1575 "{}",
1576 result.content
1577 );
1578 let home = std::env::var("HOME").unwrap_or_default();
1579 if !home.is_empty() {
1580 assert!(
1581 result.content.contains(&format!("home={home}")),
1582 "allowlisted variables still reach the shell: {}",
1583 result.content
1584 );
1585 }
1586 }
1587 }
1588
1588 lines RUST