Skip to main content

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;