cosh_tools/subagent/mod.rs
1//! Tool for calling sub-agents.
2//!
3//! The model sees a single `subagent_call` tool, but the harness dispatches
4//! it to two implementations (see [`SubAgent`]): an external agent harness —
5//! driven through the [Agent Client Protocol (ACP)](https://agentclientprotocol.com/)
6//! as a full client turn (initialize → `session/new` or `session/resume` →
7//! `session/prompt`), with agent message chunks streamed in real time —
8//! only the LAST message (post-tool-calls) becomes the returned report (see
9//! [`closure`]) — or, when `agent` is omitted/empty, an
10//! internal agent (a nested harness that reports only its final answer).
11//!
12//! Two turn behaviors apply to the external path:
13//!
14//! - **Session resume by default** (Phase 4): the agent's most recent ACP
15//! session is reused on the next call (`continue_session`, omitted =
16//! resume), keeping the sub-agent's context across calls; `false` starts
17//! a brand-new session. A failed or timed-out turn stores nothing.
18//! - **Cancellation** (Phase 5): the agent loop's shared stop flag is raced
19//! against the prompt await — a trigger sends `session/cancel` and the
20//! turn ends with the agent's own `StopReason::Cancelled`, preserving the
21//! output streamed so far.
22//!
23//! # Supported agents
24//!
25//! Each agent must speak ACP over stdio — the protocol contract replaces the
26//! fragile per-CLI flag scraping. See [`acp::ACP_AGENTS`] for the registry.
27//! Agents without ACP support are not registered.
28//!
29//! | Name | ACP invocation |
30//! |------|----------------|
31//! | `gemini` | `gemini --experimental-acp` |
32//! | `goose` | `goose acp` |
33//! | `opencode` | `opencode acp` |
34//! | `kilo` | `kilo acp` |
35//! | `cline` | `cline --acp` |
36//! | `devin` | `devin acp` |
37//! | `claude` | `npx -y @agentclientprotocol/claude-agent-acp@latest` (official adapter) |
38//! | `codex` | `npx -y @agentclientprotocol/codex-acp@latest` (official adapter) |
39
40pub mod acp;
41/// The turn's final-message accumulator: which streamed text is the report
42/// the caller receives (the ACP analogue of the internal harness's
43/// `ContextItem::Closure`). Market-standard contract: only the LAST message
44/// returns to the parent; earlier narration stays in the live timeline.
45pub mod closure;
46/// Typed event stream mapped from the ACP `session/update` notifications:
47/// the backend foundation for the TUI sub-agent box (Phase 3).
48pub mod events;
49/// Synthesis of typed events for the INTERNAL sub-agent (a nested harness):
50/// maps the nested loop's plain event stream to the same
51/// [`events::SubagentEvent`] stream the external ACP path emits, so both
52/// sub-agent flavors render identically in the TUI box.
53pub mod internal;
54/// Kernel-pinned filesystem sandbox for the ACP `fs/*` handlers (Unix only;
55/// other platforms use the validate-then-serve fallback in [`acp`]).
56#[cfg(unix)]
57pub mod sandbox;
58/// Severity DSL for code-review reports: the `code_review` prompt contract
59/// (`<!-- severity: ... -->` header) and its extraction. The header tints
60/// the sub-agent box green/orange/red (the wire's middle value stays named
61/// `yellow`; only the rendered color is orange) and is never rendered.
62pub mod severity;
63pub mod types;
64
65use std::collections::HashMap;
66use std::sync::{Arc, Mutex};
67
68use crate::ToolDescription;
69pub use types::{SubAgentCallInput, SubAgentCallOutput};
70
71/// Tool for calling sub-agents.
72///
73/// A SINGLE visible tool dispatches to two implementations — the model
74/// sees only one: an external ACP agent harness when `agent` is provided,
75/// or an internal agent (a nested harness: fresh empty context,
76/// auto-approve, no persistence, final report only) when `agent` is omitted
77/// or empty. The harness routes the call; this struct owns the visible
78/// schema, the optional description `note` (see [`set_note`](Self::set_note)),
79/// and the cross-call memory: the last input message (retry support) and
80/// the last ACP session id per agent (Phase 4 session resume).
81///
82/// Each instance keeps the last input message sent to a sub-agent, so a
83/// retry after a failed call does not require re-writing the whole prompt,
84/// plus the last ACP session id per agent name, so the next call with
85/// `continue_session` (the default) resumes that session. The harness
86/// creates one instance per agent loop, so neither ever leaks across
87/// sessions and no explicit `clean()` is needed.
88pub struct SubAgent {
89 /// MCP Tool description for `call`.
90 pub description_call: ToolDescription,
91 /// Optional informational chunk interpolated naturally into the tool
92 /// description (see [`set_note`](Self::set_note)).
93 note: String,
94 /// Last input message sent to a sub-agent in this session, reused when
95 /// a call omits `input`.
96 last_input: Mutex<Option<String>>,
97 /// Last ACP session id per agent name (Phase 4): the session each
98 /// harness returned from `session/new`, resumed by the next call with
99 /// `continue_session` (the default) so the sub-agent keeps its context
100 /// across calls. Failed turns store nothing — a session whose
101 /// turn errored is not trusted.
102 ///
103 /// Behind an `Arc` so a BACKGROUND turn (spawned detached on the
104 /// process-wide runtime, long after this dispatcher's borrow ends) can
105 /// still record the session id its harness returned. The inner `Mutex`
106 /// keeps every other access pattern unchanged.
107 ///
108 /// Concurrency invariant: the read-then-store-after-await pattern in
109 /// the dispatch layer is safe because tool dispatch is sequential per
110 /// agent loop and each harness owns its own `SubAgent` — a parallel
111 /// dispatch change would need its own synchronization here. Background
112 /// turns are the one deliberate exception: exactly ONE detached task
113 /// stores per spawned turn (the store happens once, at turn end), and
114 /// the dispatch path's busy guard (`background::running_task_for`)
115 /// rejects a second background spawn for the same agent while one is
116 /// running, so two turns never share (and interleave in) one remote
117 /// session.
118 last_session: Arc<Mutex<HashMap<String, String>>>,
119}
120
121impl Default for SubAgent {
122 fn default() -> Self {
123 Self::new()
124 }
125}
126
127impl SubAgent {
128 /// Create a new `SubAgent` with a tool description tailored to
129 /// only the ACP agent harnesses that are actually installed in PATH.
130 /// Detection runs once per process (cached by `detect_installed()`).
131 #[must_use]
132 pub fn new() -> Self {
133 Self {
134 description_call: Self::build_tool_description(""),
135 note: String::new(),
136 last_input: Mutex::new(None),
137 last_session: Arc::new(Mutex::new(HashMap::new())),
138 }
139 }
140
141 /// Set an informational chunk that is interpolated NATURALLY into the
142 /// tool description (inside the description prose — not prepended as a
143 /// notice). The harness uses this to tell the model that omitting
144 /// `agent` routes the call to an internal agent instead of an external
145 /// ACP harness.
146 ///
147 /// Rebuilds [`Self::description_call`] with the new chunk; an empty
148 /// chunk keeps the original description byte-for-byte.
149 pub fn set_note(&mut self, note: impl Into<String>) {
150 self.note = note.into().trim().to_string();
151 self.description_call = Self::build_tool_description(&self.note);
152 }
153
154 /// Build the full tool description (prose + input schema), tailoring the
155 /// agent list to the ACP harnesses actually installed in PATH, and
156 /// interpolating the given `note` into the description's natural flow.
157 ///
158 /// # Panics
159 ///
160 /// Panics if an agent returned by `detect_installed()` is not present
161 /// in [`ACP_AGENTS`](acp::ACP_AGENTS). This is a logic invariant —
162 /// detection only returns names that exist in the table.
163 #[allow(clippy::expect_used, clippy::format_collect)]
164 fn build_tool_description(note: &str) -> ToolDescription {
165 let installed = acp::detect_installed();
166
167 // The note lands as the tail of the FIRST paragraph, so it reads as
168 // part of the instructions ("...return its output. When the `agent`
169 // argument is omitted or empty, ...") instead of a prepended notice.
170 // An empty note reproduces the original text exactly.
171 let note = if note.is_empty() {
172 String::new()
173 } else {
174 format!(" {note}")
175 };
176 let common = format!(
177 "Call a supported agent ACP harness with the given input message and \
178 return its output.{note}\n\
179 `input` is optional: if omitted, the last message sent to a \
180 sub-agent in this session is reused automatically, so a failed \
181 call can be retried without re-writing the prompt. If no \
182 sub-agent has been called yet, an error is returned."
183 );
184 let usage = "## When to use\n\
185 - Use `subagent_call` with `agent` (and optionally `input`) \
186 to delegate a task to another agent ACP harness.\n\
187 - Omit `input` to reuse the last message sent to a sub-agent.\n\
188 - Consecutive calls to the same agent resume that agent's \
189 most recent session by default, keeping its context (ideal \
190 for review/iteration follow-ups). Pass `continue_session: \
191 false` when the new task is unrelated and needs clean context.\n\
192 - Omit `agent` (or pass an empty string) to call the internal \
193 agent instead of an external ACP harness.\n\
194 - Set `run_in_background: true` to spawn the sub-agent without \
195 blocking: the call returns a `task_id` immediately and the final \
196 report is delivered later as an automated completion notification. \
197 Query progress with `subagent_status`. Keep background spawns \
198 focused: prefer 3–5 parallel sub-agents, and use the synchronous \
199 mode when you need the result before continuing.\n\
200 - Use `bash_run` for regular shell commands. \
201 These are separate tools with different purposes.";
202
203 let description = if installed.is_empty() {
204 // No agents installed — the LLM will see this and likely
205 // avoid calling the tool, but the error message is helpful.
206 format!(
207 "{common}\n\
208 ## Supported agents\n\
209 (None detected — install one of the ACP-capable agents \
210 (gemini, goose, opencode, kilo, claude, codex) and restart cosh.)\n\
211 {usage}"
212 )
213 } else {
214 let agents_desc = installed
215 .iter()
216 .map(|name| {
217 // SAFETY: `name` comes from detect_installed() which only
218 // returns entries present in ACP_AGENTS.
219 let entry = acp::ACP_AGENTS
220 .iter()
221 .find(|a| a.name == *name)
222 .expect("installed agent must be in ACP_AGENTS");
223 let invocation = entry.invocation();
224 format!("- `{name}` → `{invocation}`\n")
225 })
226 .collect::<String>();
227
228 format!(
229 "{common}\n\
230 ## Supported agents\n\
231 {agents_desc}\n\
232 {usage}"
233 )
234 };
235
236 let enum_values: Vec<serde_json::Value> = if installed.is_empty() {
237 // Even with no agents detected, keep the full enum so the
238 // LLM can still attempt the tool if we missed one.
239 acp::ACP_AGENTS
240 .iter()
241 .map(|a| serde_json::Value::String(a.name.to_string()))
242 .collect()
243 } else {
244 installed
245 .iter()
246 .map(|n| serde_json::Value::String(n.to_string()))
247 .collect()
248 };
249
250 serde_json::json!({
251 "name": "subagent_call",
252 "description": description,
253 "inputSchema": {
254 "type": "object",
255 "properties": {
256 "agent": {
257 "type": "string",
258 "description": "The agent ACP harness to call. Optional: if omitted (or empty), an internal agent with an empty context runs the task instead and returns only its final report.",
259 "enum": enum_values,
260 },
261 "input": {
262 "type": "string",
263 "description": "The message to send to the sub-agent as input. Optional: if omitted (or empty), the last message sent to a sub-agent in this session is reused automatically, so a failed call can be retried without re-writing the prompt. If no sub-agent has been called yet, an error is returned.",
264 },
265 "code_review": {
266 "type": "boolean",
267 "description": "Set to true when the task is a CODE REVIEW. The sub-agent's final report then starts with a `<!-- severity: green|yellow|red -->` header consumed by the client: it tints the sub-agent box (green = at most cosmetic details, yellow = minor issues / bad practice, red = something critical found). Ignored for non-review tasks. Optional, defaults to false.",
268 },
269 "continue_session": {
270 "type": "boolean",
271 "description": "Resume the agent's most recent session, keeping its context (default). Set false to start a brand-new session when the new task is unrelated and needs clean context.",
272 },
273 "run_in_background": {
274 "type": "boolean",
275 "description": "Set true to spawn the sub-agent WITHOUT blocking: the call returns a task_id immediately and the final report is delivered later as an automated completion notification. Query progress with the subagent_status tool. Omitted or false (default): the call blocks until the sub-agent finishes and returns its report directly.",
276 },
277 },
278 "required": [],
279 },
280 })
281 }
282
283 /// Resolve the effective input message for a sub-agent call.
284 ///
285 /// When `input` is provided (and non-empty), it is stored as the last
286 /// message sent to a sub-agent in this session and returned. When
287 /// omitted, the stored message is reused so the calling agent does not
288 /// have to re-write a long prompt after a failed call.
289 ///
290 /// # Errors
291 ///
292 /// Returns an error if `input` is omitted and no message is stored yet
293 /// (i.e. no sub-agent has been called in this session).
294 ///
295 /// # Panics
296 ///
297 /// Panics if the internal mutex is poisoned (only possible if another
298 /// thread panicked while holding the lock).
299 #[allow(clippy::unwrap_used)]
300 pub fn resolve_input(&self, input: Option<String>) -> Result<String, String> {
301 match input {
302 Some(text) => {
303 let text = text.trim().to_string();
304 if text.is_empty() {
305 self.stored_input()
306 } else {
307 *self.last_input.lock().unwrap() = Some(text.clone());
308 Ok(text)
309 }
310 }
311 None => self.stored_input(),
312 }
313 }
314
315 /// Return the stored last input message, or an error explaining that
316 /// there is no message in sub-agent storage yet.
317 fn stored_input(&self) -> Result<String, String> {
318 self.last_input.lock().unwrap().clone().ok_or_else(|| {
319 "subagent_call was called without an 'input' argument, but there is \
320 no stored sub-agent message in this session yet. Provide an 'input' \
321 argument to send the first message."
322 .to_string()
323 })
324 }
325
326 /// Store the ACP session id a harness returned for `agent` (Phase 4).
327 ///
328 /// Only successful turns are stored: a failed or timed-out turn leaves
329 /// the previous id untouched (so `None` from [`stored_session`] makes
330 /// the next call open a fresh session, and an older id that a failed
331 /// turn ran on is retried on the next default call — a stale one
332 /// self-heals via the `session/new` fallback in
333 /// [`open_or_resume_session`]).
334 #[allow(clippy::unwrap_used)]
335 pub fn store_session(&self, agent: &str, session_id: String) {
336 self.last_session
337 .lock()
338 .unwrap()
339 .insert(agent.to_string(), session_id);
340 }
341
342 /// Cloneable recorder of the ACP session map, for BACKGROUND turns.
343 ///
344 /// A background `subagent_call` spawns a detached task and returns
345 /// immediately, so it cannot call [`Self::store_session`] after the
346 /// (much later) turn ends — the dispatcher borrow is long gone. The
347 /// detached task keeps this handle and records the session id at turn
348 /// end, exactly where the synchronous path would. Sharing the same
349 /// `Arc<Mutex<..>>` keeps background and synchronous bookkeeping in ONE
350 /// map, so `resume_id` sees both.
351 #[must_use]
352 pub fn session_recorder(&self) -> Arc<Mutex<HashMap<String, String>>> {
353 Arc::clone(&self.last_session)
354 }
355
356 /// The session id to resume for `agent`, when `continue_session` is
357 /// requested and one is stored. `None` → the caller opens a fresh
358 /// session (first call for this agent, or no id has ever been stored —
359 /// failed turns do not clear a previously stored id).
360 #[allow(clippy::unwrap_used)]
361 pub fn stored_session(&self, agent: &str) -> Option<String> {
362 self.last_session.lock().unwrap().get(agent).cloned()
363 }
364
365 /// The `resume` argument for [`acp::call`](crate::subagent::acp::call):
366 /// the stored session id when `continue_session` is requested
367 /// (resume-by-default), `None` for the opt-out (`continue_session:
368 /// false`) or when nothing is stored yet. Keeping the flag→id decision
369 /// here makes the opt-out mapping unit-testable; the dispatch layer
370 /// just forwards the result.
371 #[allow(clippy::unwrap_used)]
372 pub fn resume_id(&self, agent: &str, continue_session: bool) -> Option<String> {
373 if continue_session {
374 self.stored_session(agent)
375 } else {
376 None
377 }
378 }
379}
380
381#[cfg(test)]
382mod test;