cosh_tools/subagent/acp.rs
1//! ACP (Agent Client Protocol) integration for external sub-agents.
2//!
3//! External sub-agents are driven through the [Agent Client Protocol
4//! (ACP)](https://agentclientprotocol.com/) — the JSON-RPC 2.0 protocol
5//! created by the Zed team for client↔agent communication — using the
6//! official [`agent_client_protocol`] crate.
7//!
8//! Instead of spawning a one-shot CLI process and scraping its stdout (a
9//! per-agent contract of flags that breaks on every CLI update), cosh acts
10//! as an ACP **client** and follows the protocol contract:
11//!
12//! 1. Spawns the agent harness as a subprocess speaking ACP over stdio.
13//! 2. `initialize` — negotiates the protocol version and capabilities.
14//! 3. `authenticate` — when the agent advertises auth methods.
15//! 4. `session/new` — opens a session rooted at the workspace directory.
16//! 5. `session/prompt` — sends the task as the user prompt.
17//! 6. `session/update` notifications — mapped events stream live to the
18//! TUI's sub-agent box, and message chunks feed the turn's closure
19//! ([`TurnClosure`](crate::subagent::closure)): only the message written
20//! after the LAST tool call is the report returned to the caller.
21//! 7. The turn ends when the `session/prompt` response arrives with a
22//! [`StopReason`](agent_client_protocol::schema::v1::StopReason).
23//!
24//! Because the interaction follows a protocol contract, ANY harness that
25//! speaks ACP can act as a sub-agent — including third-party harnesses —
26//! with zero per-agent maintenance. Client-side capabilities are honored:
27//! permission requests are auto-approved (first option, YOLO style) so a
28//! headless call never blocks, and the client advertises the `fs` capability
29//! so agents may delegate `fs/read_text_file` / `fs/write_text_file` to the
30//! real workspace files (served below).
31//!
32//! # Supported agents
33//!
34//! Only agents with ACP support are registered. Agents without ACP support
35//! were removed (they cannot satisfy the protocol contract):
36//!
37//! | Name | ACP invocation |
38//! |------|----------------|
39//! | `gemini` | `gemini --experimental-acp` |
40//! | `goose` | `goose acp` |
41//! | `opencode` | `opencode acp` |
42//! | `kilo` | `kilo acp` |
43//! | `cline` | `cline --acp` |
44//! | `devin` | `devin acp` |
45//! | `claude` | `npx -y @agentclientprotocol/claude-agent-acp@latest` (official adapter) |
46//! | `codex` | `npx -y @agentclientprotocol/codex-acp@latest` (official adapter) |
47//!
48//! # Adding new agents
49//!
50//! Add an entry to [`ACP_AGENTS`] — an [`Agent`] of
51//! `(name, command, [args], [required binaries], install_hint)`. Anything
52//! registered in the official ACP registry works with no code changes
53//! beyond the table entry.
54
55use std::iter;
56use std::path::{Path, PathBuf};
57use std::sync::atomic::{AtomicBool, Ordering};
58use std::sync::{Arc, Mutex, OnceLock};
59use std::time::Duration;
60
61use agent_client_protocol::schema::ProtocolVersion;
62use agent_client_protocol::schema::v1::{
63 AgentCapabilities, AuthenticateRequest, CancelNotification, ClientCapabilities, ContentBlock,
64 FileSystemCapabilities, InitializeRequest, LoadSessionRequest, LoadSessionResponse,
65 NewSessionRequest, PermissionOptionKind, PromptRequest, ReadTextFileRequest,
66 ReadTextFileResponse, RequestPermissionOutcome, RequestPermissionRequest,
67 RequestPermissionResponse, ResumeSessionRequest, ResumeSessionResponse,
68 SelectedPermissionOutcome, SessionConfigId, SessionConfigKind, SessionConfigOption,
69 SessionConfigSelectOptions, SessionConfigValueId, SessionId, SessionNotification,
70 SetSessionConfigOptionRequest, StopReason, TextContent, WriteTextFileRequest,
71 WriteTextFileResponse,
72};
73use agent_client_protocol::{
74 AcpAgent, Agent as AcpRole, Client, ConnectionTo, Error as AcpError, LineDirection,
75 on_receive_notification, on_receive_request,
76};
77
78use super::events::SubagentEvent;
79
80/// A single registered ACP agent harness.
81///
82/// Maps a public `name` (used by the LLM in the tool call) to the command
83/// that launches the harness in ACP mode. Unlike the previous one-shot
84/// integration there are no per-agent prompt flags: the prompt travels as
85/// the ACP `session/prompt` payload, so adding an agent is purely a table
86/// entry.
87///
88/// `requires` lists the binaries that must be present in PATH for the agent
89/// to be considered installed — the launch command itself plus whatever
90/// engine backs it. The official `claude`/`codex` entries spawn the ACP
91/// adapter via `npx`, so they additionally require `npx`.
92#[derive(Debug, Clone, Copy)]
93pub struct Agent {
94 /// Public name used by the LLM in the `subagent_call` tool.
95 pub name: &'static str,
96 /// Executable that speaks ACP over stdio.
97 pub command: &'static str,
98 /// Static arguments that put the executable into ACP server mode.
99 pub args: &'static [&'static str],
100 /// Binaries that must be in PATH for this agent to be offered.
101 pub requires: &'static [&'static str],
102 /// Installation guidance surfaced when the spawn fails.
103 pub install_hint: &'static str,
104 /// Preferred session model, selected through the standard ACP
105 /// `session/set_config_option` request when the harness advertises a
106 /// `model` config option in `session/new`.
107 ///
108 /// Harness defaults are frequently a *paid-tier* model, while the same
109 /// CLI's own `run` command uses a working free/default model — the
110 /// harness's ACP default and its CLI default need not agree. An empty
111 /// value leaves the harness default untouched.
112 pub model: &'static str,
113}
114
115impl Agent {
116 /// Human-readable ACP launch command (e.g. `gemini --experimental-acp`)
117 /// used in the tool description and docs. The prompt travels as the ACP
118 /// `session/prompt` payload, so only the launch command is shown.
119 #[must_use]
120 pub fn invocation(&self) -> String {
121 let mut parts: Vec<&str> = Vec::with_capacity(self.args.len() + 1);
122 parts.push(self.command);
123 parts.extend_from_slice(self.args);
124 parts.join(" ")
125 }
126}
127
128/// Registered ACP agent harnesses. This is the whole integration surface
129/// for adding a new sub-agent: any entry of the official ACP registry works.
130pub const ACP_AGENTS: &[Agent] = &[
131 Agent {
132 name: "gemini",
133 command: "gemini",
134 args: &["--experimental-acp"],
135 requires: &["gemini"],
136 install_hint: "Install: npm install -g @google/gemini-cli. Then: gemini auth login. More: https://github.com/google-gemini/gemini-cli",
137 model: "",
138 },
139 Agent {
140 name: "goose",
141 command: "goose",
142 args: &["acp"],
143 requires: &["goose"],
144 install_hint: "Install: curl -fsSL https://github.com/block/goose/releases/download/stable/download_cli.sh | bash. More: https://block.github.io/goose",
145 model: "",
146 },
147 Agent {
148 name: "opencode",
149 command: "opencode",
150 args: &["acp"],
151 requires: &["opencode"],
152 install_hint: "Install: curl -fsSL https://opencode.ai/install | bash (or: npm install -g opencode-ai). More: https://opencode.ai/docs",
153 model: "",
154 },
155 Agent {
156 name: "kilo",
157 command: "kilo",
158 args: &["acp"],
159 requires: &["kilo"],
160 install_hint: "Install: curl -fsSL https://kilo.ai/install.sh | sh. More: https://kilo.ai/docs/code-with-ai/platforms/cli",
161 // The harness's ACP default (`kilo/google/gemini-3-pro-image`) is a
162 // Kilo paid-tier model; the CLI's own default is this free Nvidia
163 // model, which works with the provider keys the user already has.
164 model: "kilo/nvidia/nemotron-3-ultra-550b-a55b:free",
165 },
166 Agent {
167 name: "cline",
168 command: "cline",
169 args: &["--acp"],
170 requires: &["cline"],
171 install_hint: "Install: npm install -g cline. Then: cline auth (or sign in from the first session). More: https://docs.cline.bot/usage/acp",
172 model: "",
173 },
174 Agent {
175 name: "devin",
176 command: "devin",
177 args: &["acp"],
178 requires: &["devin"],
179 install_hint: "Install: curl -fsSL https://cli.devin.ai/install.sh | bash. Then: devin auth login (or set WINDSURF_API_KEY). More: https://docs.devin.ai/cli",
180 model: "",
181 },
182 Agent {
183 name: "claude",
184 // Claude Code is not ACP-native yet: the official route is Zed's
185 // adapter, the same one `AcpAgent::claude_agent()` in the SDK uses.
186 command: "npx",
187 args: &["-y", "@agentclientprotocol/claude-agent-acp@latest"],
188 requires: &["claude", "npx"],
189 install_hint: "Install: npm install -g @anthropic-ai/claude-code (the ACP adapter runs via npx). More: https://code.claude.com/docs",
190 model: "",
191 },
192 Agent {
193 name: "codex",
194 // Codex CLI is not ACP-native yet: the official route is Zed's
195 // adapter, the same one `AcpAgent::codex()` in the SDK uses.
196 command: "npx",
197 args: &["-y", "@agentclientprotocol/codex-acp@latest"],
198 requires: &["codex", "npx"],
199 install_hint: "Install: npm install -g @openai/codex (the ACP adapter runs via npx). More: https://developers.openai.com/codex",
200 model: "",
201 },
202];
203
204/// Build the [`AcpAgent`] launcher for a registered agent.
205///
206/// On Windows, npm-style launch commands (`.cmd`/`.bat` shims, e.g. `npx`)
207/// are routed through `cmd /d /s /c …` (see [`windows_script_launcher`]).
208pub(crate) fn agent_launcher(entry: &Agent) -> Result<AcpAgent, String> {
209 let args: Vec<String> = iter::once(entry.command)
210 .chain(entry.args.iter().copied())
211 .map(String::from)
212 .collect();
213 let args = windows_script_launcher(&args).unwrap_or(args);
214 AcpAgent::from_args(args)
215 .map_err(|e| format!("invalid ACP launch command for '{}': {e}", entry.name))
216}
217
218/// Build an ACP JSON-RPC error carrying a human-readable `message`.
219pub(crate) fn acp_error(message: impl Into<String>) -> AcpError {
220 let mut error = AcpError::internal_error();
221 error.message = message.into();
222 error
223}
224
225/// Check whether `binary` exists in PATH (with the `.exe` variant on Windows).
226fn binary_in_path(binary: &str) -> bool {
227 std::env::var_os("PATH").is_some_and(|path| {
228 std::env::split_paths(&path).any(|dir| {
229 dir.join(binary).is_file()
230 // Windows PATHEXT resolution: only relevant on Windows. npm
231 // installs CLIs as `.cmd`/`.bat` batch shims (plus an
232 // extensionless sh script), all of which are launchable.
233 || (cfg!(windows)
234 && ["exe", "cmd", "bat"]
235 .into_iter()
236 .any(|ext| dir.join(format!("{binary}.{ext}")).is_file()))
237 })
238 })
239}
240
241/// Whether `binary` resolves to a native `.exe` image in PATH (Windows only).
242///
243/// `CreateProcessW` (used by `AcpAgent::spawn_process`) can only launch
244/// `.exe` images directly; batch-file shims need the `cmd` detour below.
245#[cfg(windows)]
246fn exe_in_path(binary: &str) -> bool {
247 std::env::var_os("PATH").is_some_and(|path| {
248 std::env::split_paths(&path).any(|dir| dir.join(format!("{binary}.exe")).is_file())
249 })
250}
251
252/// Whether `binary` resolves to a `.cmd`/`.bat` batch shim in PATH (Windows
253/// only) — the shape npm takes when installing global CLIs.
254#[cfg(windows)]
255fn script_shim_in_path(binary: &str) -> bool {
256 std::env::var_os("PATH").is_some_and(|path| {
257 std::env::split_paths(&path).any(|dir| {
258 ["cmd", "bat"]
259 .into_iter()
260 .any(|ext| dir.join(format!("{binary}.{ext}")).is_file())
261 })
262 })
263}
264
265/// Windows workaround for npm-style launch commands (`npx …`, or any CLI
266/// installed through `npm install -g`): those are `.cmd`/`.bat` batch shims,
267/// and the direct `CreateProcessW` spawn inside `AcpAgent::spawn_process`
268/// fails on them with `program not found`. When the launch command is a
269/// script shim, the whole invocation is routed through
270/// `cmd /d /s /c "<full command line>"` — `/s` makes `cmd` strip the outer
271/// quotes that the process-spawning layer adds around the joined line.
272///
273/// The joined line is interpreted by `cmd.exe`, so any argument containing
274/// cmd metacharacters (`& | ^ % < > "`) would be executed/expanded rather
275/// than passed through. Those cannot appear in the static [`ACP_AGENTS`]
276/// registry, and the guard below refuses to wrap if one ever does — the
277/// direct spawn then fails loudly with `program not found` instead of
278/// running something unintended.
279///
280/// Returns `Some(wrapped_args)` when a `cmd` detour is needed, `None`
281/// otherwise (native `.exe`, non-Windows, or unsafe-to-wrap argument).
282#[cfg(windows)]
283fn windows_script_launcher(args: &[String]) -> Option<Vec<String>> {
284 const CMD_METACHARACTERS: [char; 7] = ['&', '|', '^', '%', '<', '>', '"'];
285 if args
286 .iter()
287 .any(|arg| arg.chars().any(|c| CMD_METACHARACTERS.contains(&c)))
288 {
289 return None;
290 }
291 let command = args.first()?;
292 if exe_in_path(command) || !script_shim_in_path(command) {
293 return None;
294 }
295 Some(vec![
296 "cmd".to_string(),
297 "/d".to_string(),
298 "/s".to_string(),
299 "/c".to_string(),
300 args.join(" "),
301 ])
302}
303
304/// Non-Windows no-op: Unix spawn resolves scripts through the shebang line.
305#[cfg(not(windows))]
306fn windows_script_launcher(_args: &[String]) -> Option<Vec<String>> {
307 None
308}
309
310/// Return installation / configuration guidance for a given agent.
311#[must_use]
312pub fn install_hint(agent: &str) -> &'static str {
313 ACP_AGENTS
314 .iter()
315 .find(|entry| entry.name == agent)
316 .map_or("", |entry| entry.install_hint)
317}
318
319/// Check which ACP agent harnesses the user has installed (all required
320/// binaries found in PATH).
321///
322/// Results are cached in a `OnceLock` so detection runs exactly once
323/// per process lifetime. Subsequent calls return the cached list.
324pub fn detect_installed() -> &'static Vec<&'static str> {
325 static INSTALLED: OnceLock<Vec<&'static str>> = OnceLock::new();
326 INSTALLED.get_or_init(|| {
327 ACP_AGENTS
328 .iter()
329 .filter(|agent| agent.requires.iter().all(|bin| binary_in_path(bin)))
330 .map(|agent| agent.name)
331 .collect()
332 })
333}
334
335/// Validate-then-serve path sandbox (non-Unix fallback, and the unit-test
336/// subject for sandbox semantics).
337///
338/// On Unix the fs handlers use the kernel-pinned component walk in
339/// [`super::sandbox`] instead — the read/write descriptor is obtained during
340/// the walk, so there is no validate-then-use window — while this helper
341/// remains the non-Unix fallback and the direct unit-test subject.
342///
343/// ACP agents send absolute paths; a request that resolves OUTSIDE the
344/// workspace — a different absolute prefix, or the same prefix escaping it
345/// through `..` segments or symlinks — is rejected instead of served.
346///
347/// Three gates run in sequence:
348///
349/// 1. **Canonical root** — the workspace root is canonicalized first, so
350/// containment runs in fully-resolved space (on macOS `/tmp` is a symlink
351/// to `/private/tmp`).
352/// 2. **Lexical** — the request must be rooted at the workspace (either its
353/// given or canonical spelling) and must not climb out through `..`
354/// components.
355/// 3. **Component walk** — each relative component is appended to the
356/// resolved base and checked with `symlink_metadata` (which does NOT
357/// follow the final entry): an encountered symlink is canonicalized
358/// immediately and its target must stay inside the workspace; a DANGLING
359/// symlink therefore fails closed (its target cannot be verified), while
360/// a plain MISSING component is appended as-is — a write legitimately
361/// creates new files and directories. Because the walk never descends
362/// into a non-existent component, no later component can hide a symlink.
363///
364/// NOTE: on Unix this check alone is NOT sufficient for serving (a TOCTOU
365/// window opens between this validation and the subsequent `std::fs`
366/// operation); it is the non-Unix fallback where that window is accepted,
367/// and a fast lexical pre-filter in tests.
368///
369/// The returned path is the resolved one the handlers operate on (symlinks
370/// encountered on the way are already resolved).
371///
372/// # Errors
373///
374/// Returns an ACP error when the path escapes the workspace root, a symlink
375/// inside it dangles or points outside, or the workspace itself is not
376/// accessible.
377#[cfg(any(not(unix), test))]
378pub(crate) fn ensure_path_within(root: &Path, requested: &Path) -> Result<PathBuf, AcpError> {
379 let canonical_root = root.canonicalize().map_err(|e| {
380 acp_error(format!(
381 "session workspace '{}' is not accessible: {e}",
382 root.display()
383 ))
384 })?;
385
386 // Lexical gate (see above).
387 let inside = requested
388 .strip_prefix(root)
389 .ok()
390 .or_else(|| requested.strip_prefix(&canonical_root).ok());
391 let Some(relative) = inside else {
392 return Err(outside_workspace_error(requested));
393 };
394 if relative
395 .components()
396 .any(|c| c == std::path::Component::ParentDir)
397 {
398 return Err(outside_workspace_error(requested));
399 }
400
401 // Component walk (see above).
402 let mut resolved = canonical_root.clone();
403 for component in relative.components() {
404 let candidate = resolved.join(component);
405 match candidate.symlink_metadata() {
406 Ok(meta) if meta.file_type().is_symlink() => {
407 // The entry itself exists and is a symlink: resolve it NOW.
408 // A dangling symlink (unresolvable target) fails closed.
409 let target = candidate
410 .canonicalize()
411 .map_err(|_| outside_workspace_error(requested))?;
412 if !target.starts_with(&canonical_root) {
413 return Err(outside_workspace_error(requested));
414 }
415 resolved = target;
416 }
417 Ok(_) => resolved = candidate,
418 // Missing entry: plain append. Everything below a missing
419 // component is necessarily missing too (no hidden symlinks),
420 // and the fs handlers create the intermediate directories.
421 Err(_) => resolved = candidate,
422 }
423 }
424 if !resolved.starts_with(&canonical_root) {
425 return Err(outside_workspace_error(requested));
426 }
427 Ok(resolved)
428}
429
430/// The rejection error for a path that escapes the session workspace.
431pub(crate) fn outside_workspace_error(requested: &Path) -> AcpError {
432 acp_error(format!(
433 "path '{}' resolves outside the session workspace; refusing to serve it",
434 requested.display()
435 ))
436}
437
438/// Apply the protocol's 1-based `line`/`limit` window to file content
439/// (shared by the sandboxed and fallback read paths).
440#[must_use]
441pub(crate) fn slice_lines(content: &str, line: Option<u32>, limit: Option<u32>) -> String {
442 let start = line.map_or(0, |line| line.saturating_sub(1) as usize);
443 match limit {
444 None if start == 0 => content.to_string(),
445 limit => content
446 .lines()
447 .skip(start)
448 .take(limit.map_or(usize::MAX, |l| l as usize))
449 .collect::<Vec<_>>()
450 .join("\n"),
451 }
452}
453
454/// Serve `fs/read_text_file`, sandboxed to the workspace root.
455///
456/// Unix: kernel-pinned component walk (see [`sandbox`](super::sandbox)) —
457/// the read descriptor is obtained DURING the walk, so a concurrent swap of
458/// a validated directory for an outside symlink is caught at open time and
459/// there is no validate-then-use window. Other platforms fall back to
460/// validate-then-serve via [`ensure_path_within`].
461#[cfg(unix)]
462fn serve_read(
463 root: &super::sandbox::PinnedRoot,
464 requested: &Path,
465 line: Option<u32>,
466 limit: Option<u32>,
467) -> Result<String, AcpError> {
468 super::sandbox::read(root, requested, line, limit)
469}
470
471/// See [`serve_read`]; non-Unix platforms use validate-then-serve.
472#[cfg(not(unix))]
473fn serve_read(
474 root: &Path,
475 requested: &Path,
476 line: Option<u32>,
477 limit: Option<u32>,
478) -> Result<String, AcpError> {
479 let path = ensure_path_within(root, requested)?;
480 let content = std::fs::read_to_string(&path)
481 .map_err(|e| acp_error(format!("failed to read {}: {e}", requested.display())))?;
482 Ok(slice_lines(&content, line, limit))
483}
484
485/// Serve `fs/write_text_file`, sandboxed to the workspace root (see
486/// [`serve_read`] for the platform split). Missing intermediate directories
487/// are created inside the workspace.
488#[cfg(unix)]
489fn serve_write(
490 root: &super::sandbox::PinnedRoot,
491 requested: &Path,
492 content: &str,
493) -> Result<(), AcpError> {
494 super::sandbox::write(root, requested, content)
495}
496
497/// See [`serve_write`]; non-Unix platforms use validate-then-serve.
498#[cfg(not(unix))]
499fn serve_write(root: &Path, requested: &Path, content: &str) -> Result<(), AcpError> {
500 let path = ensure_path_within(root, requested)?;
501 if let Some(parent) = path.parent() {
502 let _ = std::fs::create_dir_all(parent);
503 }
504 std::fs::write(&path, content)
505 .map_err(|e| acp_error(format!("failed to write {}: {e}", requested.display())))
506}
507
508/// Stable, serde-shaped string for an ACP [`StopReason`].
509///
510/// The persisted `SubAgentCallOutput.stop_reason` must not depend on `Debug`
511/// formatting (which can change with crate versions); these snake_case names
512/// mirror the protocol's wire spelling. Unknown variants (the enum is
513/// `#[non_exhaustive]`) degrade to their debug form instead of panicking.
514#[must_use]
515pub(crate) fn stop_reason_str(reason: StopReason) -> String {
516 match reason {
517 StopReason::EndTurn => "end_turn".to_string(),
518 StopReason::MaxTokens => "max_tokens".to_string(),
519 StopReason::MaxTurnRequests => "max_turn_requests".to_string(),
520 StopReason::Refusal => "refusal".to_string(),
521 StopReason::Cancelled => "cancelled".to_string(),
522 other => format!("{other:?}"),
523 }
524}
525
526/// Validate that `agent` is registered in [`ACP_AGENTS`].
527///
528/// # Errors
529///
530/// Returns an error if the agent name is not found.
531pub fn validate_agent(agent: &str) -> Result<(), String> {
532 if ACP_AGENTS.iter().any(|entry| entry.name == agent) {
533 Ok(())
534 } else {
535 let supported: Vec<&str> = ACP_AGENTS.iter().map(|a| a.name).collect();
536 Err(format!(
537 "Unsupported agent '{agent}'. Supported agents (ACP): {}. Use bash_run for shell commands.",
538 supported.join(", "),
539 ))
540 }
541}
542
543/// Call a sub-agent harness over ACP and return its final output.
544///
545/// 1. Looks up `agent` in [`ACP_AGENTS`] to get the ACP launch command.
546/// 2. Runs a full ACP client turn (see the [module docs](self)):
547/// `initialize` → `authenticate` (when advertised) → `session/new` or,
548/// when resuming, `session/resume`/`session/load` (falling back to
549/// `session/new` on any resume failure) → `session/prompt`,
550/// auto-approving permission requests and serving `fs/*` requests.
551/// The shared stop flag is raced against the prompt await; a trigger
552/// sends `session/cancel` and the agent ends the turn itself.
553/// 3. Streams typed [`SubagentEvent`]s through `chunk_tx` (for the TUI)
554/// while feeding the message text to the turn's closure
555/// ([`TurnClosure`](crate::subagent::closure::TurnClosure)).
556/// 4. Returns `(final_report, stop_reason, Option<session_id>)` when
557/// the turn ends — the report is the message written after the LAST
558/// tool call, with a fallback to the last completed message — and the
559/// session id only when it ended protocol-clean
560/// (`Some` is withheld on error arms) — or an error when
561/// the harness fails before producing any output.
562///
563/// The ACP session runs on a dedicated current-thread runtime inside a
564/// blocking task, isolating the SDK's connection machinery from the ambient
565/// runtime flavor and keeping the caller's future `Send`.
566///
567/// # Errors
568///
569/// Returns an error if the agent is unsupported, the harness cannot be
570/// launched, or the ACP turn fails before any output. There is NO time
571/// limit: the turn runs until the agent ends it or the user stops it.
572///
573/// # Panics
574///
575/// Panics if the internal mutex protecting the output buffer is poisoned
576/// (only possible if another thread panicked while holding the lock).
577#[allow(clippy::unwrap_used)]
578pub async fn call(
579 agent: &str,
580 input: &str,
581 resume: Option<String>,
582 cwd: PathBuf,
583 stop_signal: Arc<AtomicBool>,
584 chunk_tx: tokio::sync::mpsc::UnboundedSender<SubagentEvent>,
585) -> Result<(String, String, Option<String>), String> {
586 validate_agent(agent)?;
587 let entry = ACP_AGENTS
588 .iter()
589 .find(|a| a.name == agent)
590 .ok_or_else(|| format!("unknown agent '{agent}'"))?;
591 let launcher = agent_launcher(entry)?
592 // Surface harness stderr in the logs to ease debugging of
593 // unauthenticated/failed launches.
594 .with_debug(|line, direction| {
595 if matches!(direction, LineDirection::Stderr) {
596 log::debug!("subagent ACP stderr: {line}");
597 }
598 });
599
600 let accumulated: Arc<Mutex<super::closure::TurnClosure>> =
601 Arc::new(Mutex::new(super::closure::TurnClosure::default()));
602 let input = input.to_string();
603
604 // The turn runs to completion with NO time limit: a sub-agent may work
605 // for hours and the only legitimate way to end it is the user's stop
606 // signal (raced against the prompt await inside `run_session`), which
607 // makes the agent end the turn itself with `StopReason::Cancelled`.
608 let turn = {
609 let accumulated = accumulated.clone();
610 tokio::task::spawn_blocking(move || {
611 tokio::runtime::Builder::new_current_thread()
612 .enable_all()
613 .build()
614 .map_err(|e| format!("failed to build ACP runtime: {e}"))?
615 .block_on(run_session(
616 launcher,
617 input,
618 resume,
619 cwd,
620 entry.model,
621 accumulated,
622 stop_signal,
623 chunk_tx,
624 ))
625 })
626 .await
627 .map_err(|e| format!("sub-agent task failed: {e}"))?
628 };
629
630 let final_output = accumulated.lock().unwrap().output();
631
632 match turn {
633 // Turn completed: the closure already retained the final message
634 // (post-tool-calls) — see the `TurnClosure` contract on `output`.
635 Ok((stop_reason, session_id)) => Ok((final_output, stop_reason, Some(session_id))),
636 // Turn errored: partial output (if any) is still worth returning.
637 // The turn's own session id is NOT propagated — an errored turn's
638 // session state is unreliable, so its id is not (re)stored here.
639 // An OLDER stored id for this agent stays in place and is resumed
640 // again by the next default call; a stale one self-heals via the
641 // `session/new` fallback in `open_or_resume_session`.
642 Err(session_error) => {
643 if final_output.is_empty() {
644 Err(session_error)
645 } else {
646 log::warn!("sub-agent '{agent}' ACP turn failed: {session_error}");
647 Ok((final_output, "error".to_string(), None))
648 }
649 }
650 }
651}
652
653/// Whether the agent harness supports resuming sessions, read from the
654/// capabilities advertised at `initialize`. Both ACP mechanisms count:
655/// `sessionCapabilities.resume` (`session/resume`) and the older top-level
656/// `session/load` capability (`session/load`).
657fn session_resume_support(capabilities: &AgentCapabilities) -> SessionResumeSupport {
658 if capabilities.session_capabilities.resume.is_some() {
659 SessionResumeSupport::Resume
660 } else if capabilities.load_session {
661 SessionResumeSupport::Load
662 } else {
663 SessionResumeSupport::None
664 }
665}
666
667/// The ACP resume mechanism a harness advertises (Phase 4).
668#[derive(Debug, Clone, Copy, PartialEq, Eq)]
669enum SessionResumeSupport {
670 /// `session/resume` (`sessionCapabilities.resume`).
671 Resume,
672 /// The older top-level `session/load` capability.
673 Load,
674 /// No resume support — every call opens a fresh session.
675 None,
676}
677
678/// Open (or resume) the session for one ACP prompt turn.
679///
680/// With `resume: Some(id)` and advertised support, the stored session id is
681/// resumed (`session/resume`, falling back to the older `session/load`).
682/// ANY failure on the resume path — no advertised capability, a stale id
683/// after a harness restart, a transport error — silently falls back to a
684/// fresh `session/new`: resume is a default, never a hard dependency.
685///
686/// Returns the session id to prompt against plus the config options the
687/// session came back with (used by [`select_session_model`]; `None` for a
688/// plain `session/new` response is carried through as-is).
689///
690/// # Errors
691///
692/// Only a failed `session/new` is fatal — the turn cannot proceed without
693/// a session.
694async fn open_or_resume_session(
695 connection: &ConnectionTo<AcpRole>,
696 resume: Option<String>,
697 cwd: &Path,
698 support: SessionResumeSupport,
699) -> Result<(SessionId, Option<Vec<SessionConfigOption>>), AcpError> {
700 // Exhaustive over (resume, support): no `unreachable!()` arm — the
701 // `None`-support and no-id cases fall through to `session/new` below.
702 match (resume, support) {
703 (Some(id), SessionResumeSupport::Resume) => {
704 let response = connection
705 .send_request(ResumeSessionRequest::new(
706 SessionId::new(id.as_str()),
707 cwd.to_path_buf(),
708 ))
709 .block_task()
710 .await;
711 match response.map(|r: ResumeSessionResponse| (r.config_options,)) {
712 Ok((config_options,)) => {
713 log::debug!("sub-agent ACP session resumed: {id}");
714 return Ok((SessionId::new(id), config_options));
715 }
716 Err(e) => {
717 // Stale id (harness restart) or a resume quirk: fall
718 // through to a fresh session instead of failing the
719 // turn — the sub-agent just loses its previous context.
720 log::warn!(
721 "sub-agent ACP resume of session '{id}' failed ({e}); \
722 falling back to session/new"
723 );
724 }
725 }
726 }
727 (Some(id), SessionResumeSupport::Load) => {
728 let response = connection
729 .send_request(LoadSessionRequest::new(
730 SessionId::new(id.as_str()),
731 cwd.to_path_buf(),
732 ))
733 .block_task()
734 .await;
735 match response.map(|r: LoadSessionResponse| (r.config_options,)) {
736 Ok((config_options,)) => {
737 log::debug!("sub-agent ACP session loaded: {id}");
738 return Ok((SessionId::new(id), config_options));
739 }
740 Err(e) => {
741 log::warn!(
742 "sub-agent ACP load of session '{id}' failed ({e}); \
743 falling back to session/new"
744 );
745 }
746 }
747 }
748 (Some(_), SessionResumeSupport::None) => {
749 log::debug!(
750 "sub-agent harness advertises no session resume support; \
751 starting a fresh session"
752 );
753 }
754 (None, _) => {}
755 }
756
757 let session = connection
758 .send_request(NewSessionRequest::new(cwd.to_path_buf()))
759 .block_task()
760 .await
761 .map_err(|e| acp_error(format!("ACP session/new failed: {e}")))?;
762 Ok((session.session_id, session.config_options))
763}
764
765/// Select the agent's preferred session model, when one is registered and
766/// the harness advertises a `model` config option in `session/new`.
767///
768/// A mismatch here is fatal on some harnesses (e.g. kilo's ACP default is a
769/// paid-tier model that fails the prompt with "You need to sign in", while
770/// its own CLI default is a working free model), so a selection that the
771/// harness rejects aborts the turn. An unregistered model (empty string) or
772/// a harness without a `model` option is a no-op.
773///
774/// # Errors
775///
776/// Returns an error when the model is not among the harness's offered
777/// values, or when the `session/set_config_option` request fails.
778async fn select_session_model(
779 connection: &ConnectionTo<AcpRole>,
780 session_id: &SessionId,
781 config_options: Option<&[SessionConfigOption]>,
782 preferred_model: &str,
783) -> Result<(), String> {
784 if preferred_model.is_empty() {
785 return Ok(());
786 }
787 let Some(option) = config_options
788 .unwrap_or_default()
789 .iter()
790 .find(|option| option.id.0.as_ref() == "model")
791 else {
792 return Ok(());
793 };
794 let SessionConfigKind::Select(select) = &option.kind else {
795 return Ok(());
796 };
797 if select.current_value.0.as_ref() == preferred_model {
798 return Ok(());
799 }
800 let offered: Vec<&str> = match &select.options {
801 SessionConfigSelectOptions::Ungrouped(options) => {
802 options.iter().map(|o| o.value.0.as_ref()).collect()
803 }
804 SessionConfigSelectOptions::Grouped(groups) => groups
805 .iter()
806 .flat_map(|g| g.options.iter().map(|o| o.value.0.as_ref()))
807 .collect(),
808 _ => Vec::new(),
809 };
810 if !offered.contains(&preferred_model) {
811 return Err(format!(
812 "ACP session model '{preferred_model}' is not offered by this harness; \
813 update the agent's `model` entry in the registry",
814 ));
815 }
816 connection
817 .send_request(SetSessionConfigOptionRequest::new(
818 session_id.clone(),
819 SessionConfigId::new("model"),
820 SessionConfigValueId::new(preferred_model),
821 ))
822 .block_task()
823 .await
824 .map_err(|e| format!("ACP set model failed: {e}"))?;
825 Ok(())
826}
827
828/// Run one full ACP prompt turn against the given agent-side transport.
829///
830/// `transport` is anything implementing [`ConnectTo`] for the agent role: in
831/// production it is the [`AcpAgent`] subprocess launcher; tests use an
832/// in-memory [`Channel`](agent_client_protocol::Channel) instead. Returns
833/// the prompt's stop reason on success. `select_session_model` runs before
834/// the prompt, and the shared stop flag is raced against the prompt await
835/// (Phase 5 cancellation: a trigger sends `session/cancel` and the turn
836/// still ends with the agent's own answer). Message text feeds the
837/// [`TurnClosure`](super::closure::TurnClosure) — only the message written
838/// after the LAST tool call is the returned report (the market-standard
839/// "final message" contract), with a defensive fallback to the last
840/// completed message when the turn ends right after a tool call; every
841/// mapped session update is streamed through `chunk_tx` as a typed
842/// [`SubagentEvent`].
843///
844/// # Errors
845///
846/// Returns an error if any step of the ACP handshake or the prompt turn
847/// fails.
848#[allow(clippy::unwrap_used)]
849#[allow(clippy::too_many_arguments)]
850pub(crate) async fn run_session<T>(
851 transport: T,
852 input: String,
853 resume: Option<String>,
854 cwd: PathBuf,
855 preferred_model: &'static str,
856 accumulated: Arc<Mutex<super::closure::TurnClosure>>,
857 stop_signal: Arc<AtomicBool>,
858 chunk_tx: tokio::sync::mpsc::UnboundedSender<SubagentEvent>,
859) -> Result<(String, String), String>
860where
861 T: agent_client_protocol::ConnectTo<agent_client_protocol::Client> + 'static,
862{
863 // `fs/*` requests are served against the real filesystem, sandboxed to
864 // the session workspace root. The root is pinned to a VERIFIED directory
865 // descriptor once per turn (see `sandbox::PinnedRoot`): the handlers
866 // below clone that descriptor's handle and every `fs/*` resolution
867 // walks `openat`-relative from it — no path is ever re-resolved after
868 // verification, closing the validate-then-use window (including at the
869 // root itself).
870 // Non-Unix has no descriptor pinning; the handlers fall back to the
871 // lexical validate-then-serve path (`ensure_path_within`), so hand them
872 // the plain root path there.
873 #[cfg(unix)]
874 let read_root = super::sandbox::PinnedRoot::acquire(&cwd).map_err(|e| e.to_string())?;
875 #[cfg(unix)]
876 let write_root = super::sandbox::PinnedRoot::acquire(&cwd).map_err(|e| e.to_string())?;
877 #[cfg(not(unix))]
878 let read_root = cwd.clone();
879 #[cfg(not(unix))]
880 let write_root = cwd.clone();
881
882 // Route the transport through the generic connection machinery. The
883 // subprocess launcher gets stderr line logging; in-memory test channels
884 // pass through untouched.
885 Client
886 .builder()
887 .name("cosh")
888 .on_receive_notification(
889 async move |notification: SessionNotification, _cx| {
890 // Map every session update with a TUI representation into a
891 // typed event; message chunks additionally feed the turn's
892 // closure (the last-message accumulator that produces the
893 // returned report).
894 match SubagentEvent::from_session_update(notification.update) {
895 Some(event) => {
896 accumulated.lock().unwrap().observe(&event);
897 let _ = chunk_tx.send(event);
898 }
899 // Updates with no display representation (user chunks,
900 // available-commands refreshes, unstable variants):
901 // log instead of silently dropping.
902 None => {
903 log::debug!("sub-agent session update without a display mapping");
904 }
905 }
906 Ok(())
907 },
908 on_receive_notification!(),
909 )
910 // Headless automation: auto-approve (YOLO style), so the sub-agent
911 // never blocks on a dialog nobody can answer. Option ORDER is
912 // harness-defined — some harnesses list reject-like choices first —
913 // so an explicit allow-kind option is preferred over "first option"
914 // whenever one exists. With no options, cancel the request per the
915 // spec.
916 .on_receive_request(
917 async move |request: RequestPermissionRequest, responder, _connection| {
918 // Explicit preference order: AllowAlways first, then
919 // AllowOnce — both auto-approve, but Always avoids the same
920 // prompt coming back for every later call. Without any
921 // allow-kind option, the first listed option is picked as a
922 // last resort; its kind is logged so a reject-like auto-
923 // approval is visible in debug output.
924 let selected = request
925 .options
926 .iter()
927 .find(|option| option.kind == PermissionOptionKind::AllowAlways)
928 .or_else(|| {
929 request
930 .options
931 .iter()
932 .find(|option| option.kind == PermissionOptionKind::AllowOnce)
933 })
934 .or_else(|| request.options.first());
935 match selected {
936 Some(option) => {
937 log::debug!(
938 "sub-agent permission auto-approved: option '{}' (kind {:?}) for tool call {}",
939 option.option_id.0,
940 option.kind,
941 request.tool_call.tool_call_id.0
942 );
943 responder.respond(RequestPermissionResponse::new(
944 RequestPermissionOutcome::Selected(SelectedPermissionOutcome::new(
945 option.option_id.clone(),
946 )),
947 ))
948 }
949 None => responder.respond(RequestPermissionResponse::new(
950 RequestPermissionOutcome::Cancelled,
951 )),
952 }
953 },
954 on_receive_request!(),
955 )
956 // Serve the agent's filesystem capability against the real workspace
957 // (absolute paths, 1-based lines per the protocol contract), SANDBOXED
958 // to the session workspace root: a path that escapes `cwd` (through
959 // `..` segments or symlinks) is rejected with an ACP error instead of
960 // being served.
961 .on_receive_request(
962 async move |request: ReadTextFileRequest, responder, _connection| {
963 let content = serve_read(
964 &read_root,
965 &request.path,
966 request.line,
967 request.limit,
968 )?;
969 responder.respond(ReadTextFileResponse::new(content))
970 },
971 on_receive_request!(),
972 )
973 .on_receive_request(
974 async move |request: WriteTextFileRequest, responder, _connection| {
975 serve_write(&write_root, &request.path, &request.content)?;
976 responder.respond(WriteTextFileResponse::new())
977 },
978 on_receive_request!(),
979 )
980 .connect_with(transport, async move |connection: ConnectionTo<AcpRole>| {
981 // Advertise exactly what we serve below: fs read/write. Without
982 // this, spec-conformant agents never send `fs/*` requests.
983 let capabilities = ClientCapabilities::new().fs(FileSystemCapabilities::new()
984 .read_text_file(true)
985 .write_text_file(true));
986 let init = connection
987 .send_request(
988 InitializeRequest::new(ProtocolVersion::V1).client_capabilities(capabilities),
989 )
990 .block_task()
991 .await
992 .map_err(|e| acp_error(format!("ACP initialize failed: {e}")))?;
993
994 // Authenticate when the agent advertises methods (e.g. `goose
995 // acp`). The user's existing login is reused: the first method
996 // that succeeds wins, and a failing method is skipped rather
997 // than aborting — an adapter can advertise several (codex:
998 // `api-key` requires env keys the user may not have, while
999 // `chat-gpt` works with the stored `codex login` state). A
1000 // harness that manages auth itself advertises no methods and is
1001 // skipped.
1002 let mut auth_error = None;
1003 for method in &init.auth_methods {
1004 match connection
1005 .send_request(AuthenticateRequest::new(method.id().clone()))
1006 .block_task()
1007 .await
1008 {
1009 Ok(_) => {
1010 auth_error = None;
1011 break;
1012 }
1013 Err(e) => {
1014 log::debug!(
1015 "sub-agent ACP authenticate with method '{}' failed: {e}",
1016 method.id().0
1017 );
1018 auth_error = Some(format!("ACP authenticate failed: {e}"));
1019 }
1020 }
1021 }
1022 if let Some(e) = auth_error {
1023 return Err(acp_error(e));
1024 }
1025
1026 let (session_id, config_options) =
1027 open_or_resume_session(
1028 &connection,
1029 resume,
1030 &cwd,
1031 session_resume_support(&init.agent_capabilities),
1032 )
1033 .await?;
1034
1035 // Select the preferred session model when the agent registered
1036 // one and the harness advertises a `model` config option. A
1037 // mismatch here is fatal on some harnesses (e.g. kilo's ACP
1038 // default is a paid-tier model that fails the prompt with "You
1039 // need to sign in"), so a failed selection aborts the turn.
1040 // Resumed sessions carry their config options in the resume
1041 // response, so the selection applies to them too.
1042 select_session_model(
1043 &connection,
1044 &session_id,
1045 config_options.as_deref(),
1046 preferred_model,
1047 )
1048 .await
1049 .map_err(acp_error)?;
1050
1051 let prompt_request = connection
1052 .send_request(PromptRequest::new(
1053 session_id.clone(),
1054 vec![ContentBlock::Text(TextContent::new(input))],
1055 ))
1056 .block_task();
1057
1058 // Cancellation (Phase 5): the shared stop flag is raced against
1059 // the prompt await. On trigger, `session/cancel` is sent as a
1060 // fire-and-forget notification and the prompt is STILL awaited —
1061 // the ACP spec requires the agent to end the turn itself,
1062 // answering with `StopReason::Cancelled`. This is the ONLY way
1063 // the turn ends early: there is no time limit.
1064 tokio::pin!(prompt_request);
1065 let stop_wait = async {
1066 while !stop_signal.load(Ordering::Relaxed) {
1067 tokio::time::sleep(Duration::from_millis(50)).await;
1068 }
1069 };
1070 let prompt = tokio::select! {
1071 prompt = &mut prompt_request => prompt,
1072 () = stop_wait => {
1073 log::debug!("sub-agent ACP turn cancelled by stop signal");
1074 if let Err(e) =
1075 connection.send_notification(CancelNotification::new(session_id.clone()))
1076 {
1077 log::warn!("sub-agent ACP cancel notification failed: {e}");
1078 }
1079 (&mut prompt_request).await
1080 }
1081 };
1082 let prompt = prompt
1083 .map_err(|e| acp_error(format!("ACP session/prompt failed: {e}")))?;
1084
1085 // The session id flows back to the caller: it is stored per
1086 // agent name and reused by the next `continue_session` call.
1087 Ok((
1088 stop_reason_str(prompt.stop_reason),
1089 session_id.0.to_string(),
1090 ))
1091 })
1092 .await
1093 .map_err(|e| format!("ACP connection failed: {e}"))
1094}