Skip to main content

supercode_runtime/
background.rs

1//! P5-6 (COMPOSABLE-HARNESS-DESIGN.md §2 module 4 `tools.background`: "D1
2//! background exec + monitor/event feed; D10 bg-manager; D3 self-paced/
3//! scheduled loops"; §2.1 "tools.background → permissions.approvals
4//! (auto-policy) [C6 as dep]"; §2.2 C6): the pure, agent-independent data
5//! shapes and bounded-buffer arithmetic a runtime agent's
6//! `background_exec`/`background_status`/`background_list`/`background_kill`
7//! intrinsics build on — kept separate from `agent.rs` so the bounded-
8//! capture truncation logic and job-id shape are unit-testable without a
9//! full `Agent`/mock-`Provider`/real-subprocess harness, the same
10//! "pure config → set, testable without the loop" precedent
11//! subagent runtime documents for itself (P5-3).
12//!
13//! **Activation.** Everything here is inert until `Agent` actually consults
14//! it, which only happens when `Config::tools_background_enabled` is `true`
15//! (`capabilities.tools_background.enabled`, default `false`) — importing
16//! this module changes nothing for an agent that never turns the module on.
17//!
18//! **Process-kill reuse.** The actual OS-process spawn/kill machinery lives
19//! in the agent loop (it needs `tokio::process::Command`/`Child`, which this
20//! module deliberately does not depend on, keeping it synchronous and
21//! trivially unit-testable). The concurrency bound reuses
22//! the subagent concurrency guard verbatim — the same
23//! generic `Arc<AtomicUsize>` gauge machinery, just a second, independent
24//! gauge instance scoped to background JOBS rather than subagent SPAWNS
25//! (`Agent::background_concurrency_gauge`, distinct from
26//! `Agent::subagent_concurrency_gauge`).
27
28use std::sync::atomic::{AtomicU64, Ordering};
29use std::sync::Mutex;
30
31/// P5-6 (build brief "cap the buffer like P5-2's 16MiB caps"): the default
32/// per-job bounded-capture ceiling, matching `crate::mcp::MCP_MAX_RESPONSE_BYTES`'s
33/// hardening precedent — generous for real command output while bounding
34/// how much memory one background job (let alone `max_concurrent` of them
35/// at once) can force this process to hold.
36pub const DEFAULT_MAX_OUTPUT_BYTES: usize = 16 * 1024 * 1024;
37
38/// P5-6 (resource bound, mirroring `crate::subagents`'s "max concurrent...
39/// cap, fail-closed... configurable" precedent): the default maximum number
40/// of background jobs this agent may have in flight at once.
41pub const DEFAULT_MAX_CONCURRENT: usize = 4;
42
43/// A background job's run state, as observed by `background_status`/
44/// `background_list`.
45#[derive(Debug, Clone, Copy, PartialEq, Eq)]
46pub enum JobStatus {
47    /// Still running (no exit observed yet).
48    Running,
49    /// Exited on its own; `Some(code)` when the platform reported one
50    /// (`None` covers a signal-terminated exit with no portable code, same
51    /// convention `BashTool::execute` already uses via
52    /// `status.code().unwrap_or(-1)` — this type keeps the `Option`
53    /// instead of collapsing it, so callers can tell "exit code 0" from
54    /// "no code available" if they care to).
55    Exited(Option<i32>),
56    /// Killed via `background_kill` (or reclaimed on agent drop) before it
57    /// exited on its own.
58    Killed,
59}
60
61impl JobStatus {
62    /// The `"status"` string a tool result JSON reports.
63    pub fn as_str(self) -> &'static str {
64        match self {
65            JobStatus::Running => "running",
66            JobStatus::Exited(_) => "exited",
67            JobStatus::Killed => "killed",
68        }
69    }
70}
71
72/// Bounded, incrementally-appended output capture shared (via `Arc`)
73/// between a job's stdout/stderr reader tasks and whatever later polls it
74/// (`background_status`/`background_list`). Thread-safe; every method is
75/// fail-soft on a poisoned lock (treats it as "temporarily unavailable",
76/// the same posture `crate::permissions::approval::ApprovalCache` already
77/// documents for itself) rather than panicking a reader task or a tool
78/// call.
79#[derive(Debug, Default)]
80pub struct CapturedOutput {
81    inner: Mutex<CaptureState>,
82}
83
84#[derive(Debug, Default)]
85struct CaptureState {
86    buf: String,
87    truncated: bool,
88    /// Byte offset into `buf` already handed back by a previous
89    /// `drain_new` call — the event-feed cursor.
90    drained: usize,
91}
92
93impl CapturedOutput {
94    /// A fresh, empty capture.
95    pub fn new() -> Self {
96        CapturedOutput::default()
97    }
98
99    /// Append `chunk`, never growing the retained buffer past `cap` bytes —
100    /// bytes beyond the cap are DROPPED (fail-closed, never buffered) and
101    /// `truncated` latches `true` the first time that happens and stays
102    /// true thereafter. A caller must keep reading the underlying pipe past
103    /// this point regardless (to avoid blocking the child on a full,
104    /// undrained pipe) — this method only bounds what's RETAINED in
105    /// memory, not what's read off the pipe.
106    pub fn append(&self, chunk: &str, cap: usize) {
107        if chunk.is_empty() {
108            return;
109        }
110        let Ok(mut st) = self.inner.lock() else {
111            return;
112        };
113        if st.buf.len() >= cap {
114            st.truncated = true;
115            return;
116        }
117        let remaining = cap - st.buf.len();
118        if chunk.len() <= remaining {
119            st.buf.push_str(chunk);
120        } else {
121            // Largest char boundary <= remaining, same approach
122            // `Agent::cap_tool_output` already uses.
123            let mut end = remaining;
124            while end > 0 && !chunk.is_char_boundary(end) {
125                end -= 1;
126            }
127            st.buf.push_str(&chunk[..end]);
128            st.truncated = true;
129        }
130    }
131
132    /// The full captured text so far, and whether it was ever truncated.
133    pub fn snapshot(&self) -> (String, bool) {
134        self.inner
135            .lock()
136            .map(|st| (st.buf.clone(), st.truncated))
137            .unwrap_or_default()
138    }
139
140    /// Text appended since the last `drain_new` call (or since creation, on
141    /// the first call) — the event-feed's per-poll delta. Advances the
142    /// drain cursor even on an empty result, so polling twice in a row with
143    /// no new output between them returns `""` the second time, never a
144    /// repeat of the first drain.
145    pub fn drain_new(&self) -> String {
146        let Ok(mut st) = self.inner.lock() else {
147            return String::new();
148        };
149        let new = st.buf[st.drained..].to_string();
150        st.drained = st.buf.len();
151        new
152    }
153}
154
155/// P5-6: process-wide sequence number backing [`next_job_id`] —
156/// disambiguates two jobs spawned in the same millisecond, mirroring
157/// `crate::agent`'s own `SUBAGENT_ID_SEQ` precedent (kept as a second,
158/// independent counter rather than sharing that one, since a job id and a
159/// subagent id are never compared against each other).
160static JOB_ID_SEQ: AtomicU64 = AtomicU64::new(0);
161
162/// A fresh, process-unique background job id (`"bg-<hex-ts>-<hex-seq>"`),
163/// given the caller's own millisecond timestamp (kept as a parameter rather
164/// than reading the clock in here, so this stays a pure function for the
165/// unit tests below).
166pub fn next_job_id(now_ms: i64) -> String {
167    let seq = JOB_ID_SEQ.fetch_add(1, Ordering::Relaxed);
168    format!("bg-{now_ms:x}-{seq:x}")
169}