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::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    /// Concurrency invariant: the read-then-store-after-await pattern in
104    /// the dispatch layer is safe because tool dispatch is sequential per
105    /// agent loop and each harness owns its own `SubAgent` — a parallel
106    /// dispatch change would need its own synchronization here.
107    last_session: Mutex<HashMap<String, String>>,
108}
109
110impl Default for SubAgent {
111    fn default() -> Self {
112        Self::new()
113    }
114}
115
116impl SubAgent {
117    /// Create a new `SubAgent` with a tool description tailored to
118    /// only the ACP agent harnesses that are actually installed in PATH.
119    /// Detection runs once per process (cached by `detect_installed()`).
120    #[must_use]
121    pub fn new() -> Self {
122        Self {
123            description_call: Self::build_tool_description(""),
124            note: String::new(),
125            last_input: Mutex::new(None),
126            last_session: Mutex::new(HashMap::new()),
127        }
128    }
129
130    /// Set an informational chunk that is interpolated NATURALLY into the
131    /// tool description (inside the description prose — not prepended as a
132    /// notice). The harness uses this to tell the model that omitting
133    /// `agent` routes the call to an internal agent instead of an external
134    /// ACP harness.
135    ///
136    /// Rebuilds [`Self::description_call`] with the new chunk; an empty
137    /// chunk keeps the original description byte-for-byte.
138    pub fn set_note(&mut self, note: impl Into<String>) {
139        self.note = note.into().trim().to_string();
140        self.description_call = Self::build_tool_description(&self.note);
141    }
142
143    /// Build the full tool description (prose + input schema), tailoring the
144    /// agent list to the ACP harnesses actually installed in PATH, and
145    /// interpolating the given `note` into the description's natural flow.
146    ///
147    /// # Panics
148    ///
149    /// Panics if an agent returned by `detect_installed()` is not present
150    /// in [`ACP_AGENTS`](acp::ACP_AGENTS). This is a logic invariant —
151    /// detection only returns names that exist in the table.
152    #[allow(clippy::expect_used, clippy::format_collect)]
153    fn build_tool_description(note: &str) -> ToolDescription {
154        let installed = acp::detect_installed();
155
156        // The note lands as the tail of the FIRST paragraph, so it reads as
157        // part of the instructions ("...return its output. When the `agent`
158        // argument is omitted or empty, ...") instead of a prepended notice.
159        // An empty note reproduces the original text exactly.
160        let note = if note.is_empty() {
161            String::new()
162        } else {
163            format!(" {note}")
164        };
165        let common = format!(
166            "Call a supported agent ACP harness with the given input message and \
167             return its output.{note}\n\
168             `input` is optional: if omitted, the last message sent to a \
169             sub-agent in this session is reused automatically, so a failed \
170             call can be retried without re-writing the prompt. If no \
171             sub-agent has been called yet, an error is returned."
172        );
173        let usage = "## When to use\n\
174                     - Use `subagent_call` with `agent` (and optionally `input`) \
175                     to delegate a task to another agent ACP harness.\n\
176                     - Omit `input` to reuse the last message sent to a sub-agent.\n\
177                     - Consecutive calls to the same agent resume that agent's \
178                     most recent session by default, keeping its context (ideal \
179                     for review/iteration follow-ups). Pass `continue_session: \
180                     false` when the new task is unrelated and needs clean context.\n\
181                     - Omit `agent` (or pass an empty string) to call the internal \
182                     agent instead of an external ACP harness.\n\
183                     - Use `bash_run` for regular shell commands. \
184                     These are separate tools with different purposes.";
185
186        let description = if installed.is_empty() {
187            // No agents installed — the LLM will see this and likely
188            // avoid calling the tool, but the error message is helpful.
189            format!(
190                "{common}\n\
191                 ## Supported agents\n\
192                 (None detected — install one of the ACP-capable agents \
193                  (gemini, goose, opencode, kilo, claude, codex) and restart cosh.)\n\
194                 {usage}"
195            )
196        } else {
197            let agents_desc = installed
198                .iter()
199                .map(|name| {
200                    // SAFETY: `name` comes from detect_installed() which only
201                    // returns entries present in ACP_AGENTS.
202                    let entry = acp::ACP_AGENTS
203                        .iter()
204                        .find(|a| a.name == *name)
205                        .expect("installed agent must be in ACP_AGENTS");
206                    let invocation = entry.invocation();
207                    format!("- `{name}` → `{invocation}`\n")
208                })
209                .collect::<String>();
210
211            format!(
212                "{common}\n\
213                 ## Supported agents\n\
214                 {agents_desc}\n\
215                 {usage}"
216            )
217        };
218
219        let enum_values: Vec<serde_json::Value> = if installed.is_empty() {
220            // Even with no agents detected, keep the full enum so the
221            // LLM can still attempt the tool if we missed one.
222            acp::ACP_AGENTS
223                .iter()
224                .map(|a| serde_json::Value::String(a.name.to_string()))
225                .collect()
226        } else {
227            installed
228                .iter()
229                .map(|n| serde_json::Value::String(n.to_string()))
230                .collect()
231        };
232
233        serde_json::json!({
234            "name": "subagent_call",
235            "description": description,
236            "inputSchema": {
237                "type": "object",
238                "properties": {
239                    "agent": {
240                        "type": "string",
241                        "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.",
242                        "enum": enum_values,
243                    },
244                    "input": {
245                        "type": "string",
246                        "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.",
247                    },
248                    "code_review": {
249                        "type": "boolean",
250                        "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.",
251                    },
252                    "continue_session": {
253                        "type": "boolean",
254                        "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.",
255                    },
256                },
257                "required": [],
258            },
259        })
260    }
261
262    /// Resolve the effective input message for a sub-agent call.
263    ///
264    /// When `input` is provided (and non-empty), it is stored as the last
265    /// message sent to a sub-agent in this session and returned. When
266    /// omitted, the stored message is reused so the calling agent does not
267    /// have to re-write a long prompt after a failed call.
268    ///
269    /// # Errors
270    ///
271    /// Returns an error if `input` is omitted and no message is stored yet
272    /// (i.e. no sub-agent has been called in this session).
273    ///
274    /// # Panics
275    ///
276    /// Panics if the internal mutex is poisoned (only possible if another
277    /// thread panicked while holding the lock).
278    #[allow(clippy::unwrap_used)]
279    pub fn resolve_input(&self, input: Option<String>) -> Result<String, String> {
280        match input {
281            Some(text) => {
282                let text = text.trim().to_string();
283                if text.is_empty() {
284                    self.stored_input()
285                } else {
286                    *self.last_input.lock().unwrap() = Some(text.clone());
287                    Ok(text)
288                }
289            }
290            None => self.stored_input(),
291        }
292    }
293
294    /// Return the stored last input message, or an error explaining that
295    /// there is no message in sub-agent storage yet.
296    fn stored_input(&self) -> Result<String, String> {
297        self.last_input.lock().unwrap().clone().ok_or_else(|| {
298            "subagent_call was called without an 'input' argument, but there is \
299                 no stored sub-agent message in this session yet. Provide an 'input' \
300                 argument to send the first message."
301                .to_string()
302        })
303    }
304
305    /// Store the ACP session id a harness returned for `agent` (Phase 4).
306    ///
307    /// Only successful turns are stored: a failed or timed-out turn leaves
308    /// the previous id untouched (so `None` from [`stored_session`] makes
309    /// the next call open a fresh session, and an older id that a failed
310    /// turn ran on is retried on the next default call — a stale one
311    /// self-heals via the `session/new` fallback in
312    /// [`open_or_resume_session`]).
313    #[allow(clippy::unwrap_used)]
314    pub fn store_session(&self, agent: &str, session_id: String) {
315        self.last_session
316            .lock()
317            .unwrap()
318            .insert(agent.to_string(), session_id);
319    }
320
321    /// The session id to resume for `agent`, when `continue_session` is
322    /// requested and one is stored. `None` → the caller opens a fresh
323    /// session (first call for this agent, or no id has ever been stored —
324    /// failed turns do not clear a previously stored id).
325    #[allow(clippy::unwrap_used)]
326    pub fn stored_session(&self, agent: &str) -> Option<String> {
327        self.last_session.lock().unwrap().get(agent).cloned()
328    }
329
330    /// The `resume` argument for [`acp::call`](crate::subagent::acp::call):
331    /// the stored session id when `continue_session` is requested
332    /// (resume-by-default), `None` for the opt-out (`continue_session:
333    /// false`) or when nothing is stored yet. Keeping the flag→id decision
334    /// here makes the opt-out mapping unit-testable; the dispatch layer
335    /// just forwards the result.
336    #[allow(clippy::unwrap_used)]
337    pub fn resume_id(&self, agent: &str, continue_session: bool) -> Option<String> {
338        if continue_session {
339            self.stored_session(agent)
340        } else {
341            None
342        }
343    }
344}
345
346#[cfg(test)]
347mod test;