Skip to main content

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}