Skip to main content

kranz_engine/
backend_acp.rs

1//! ACP (Agent Client Protocol) worker backend: JSON-RPC 2.0 over NDJSON
2//! stdin/stdout. Protocol version 1 is required; optional capabilities are
3//! negotiated separately from that major version. Released adapter pins and
4//! live-proof limits are recorded in `docs/acp-compatibility.md`.
5//!
6//! `initialize` and `session/new` synthesize `Init`. Its session ID is the
7//! peer's ID; `AgentSession::session_id()` retains the engine ID. The raw
8//! event records both IDs, capabilities and the configured model label.
9//! Model attribution comes from the peer's model config option or legacy
10//! `models.currentModelId`, else `unreported`. Kranz does not apply the
11//! configured model/effort label as an ACP model-selection request.
12//!
13//! Existing role instructions and task text travel in the prompt, as with
14//! the Codex/Droid backends. There is no native JSON-schema enforcement.
15//! Message chunks become `Text`; tool announcements/results become
16//! `ToolUse`/`ToolResult`; remaining valid notifications become `Other`.
17//! A `session/prompt` response creates `Result` from the accumulated final
18//! message. Only `end_turn` is a successful result: cancellation, truncation,
19//! refusal and missing/unknown stop reasons fail honestly.
20//!
21//! Session updates and permissions must name the established peer session.
22//! Malformed JSON-RPC, duplicate JSON keys, invalid UTF-8, oversized frames
23//! and oversized retained message/tool state fail closed. A malformed frame
24//! that can be retained is a diagnostic `Other`, never a later valid report.
25//!
26//! Context-window usage is not a billable input/output token split. Tokens
27//! remain unavailable (the existing zero default). A finite nonnegative USD
28//! session cost becomes a turn delta only when adjacent totals are known;
29//! missing telemetry or a decreasing total yields no attributed turn cost.
30//!
31//! Permission decisions use current action identity, kind and raw arguments.
32//! Deny patterns run before the read-only posture. Unclassified operations
33//! and mode changes are refused. A display title cannot stand in for shell
34//! arguments. Only an offered, unambiguous one-time option may be selected;
35//! durable or malformed options cancel. The engine records each request and
36//! resolution before its separately callable responder can send an answer.
37//! Understood calls require one-time operator consent; prohibitions are denied
38//! by policy. Output continues while a request waits, with a five-minute limit.
39//! Delivery receipts are separate from the eventual tool outcome.
40//! This is a cooperative permission policy, not shell parsing or containment:
41//! an adapter can perform actions without asking. `allowed_tools` is not an
42//! enforced allowlist here, and read-only roles may still execute commands.
43//!
44//! Client fs/terminal services and same-feature resume remain unsupported.
45//! Unexpected client requests receive -32601. Released adapters may support
46//! load/resume methods; that does not imply Kranz negotiates or uses them.
47//! Configuration restricts ACP to opt-in worker use. Enforced missions require
48//! an explicit qualified `acpProfile`; arbitrary commands, validator roles and
49//! automatic backend promotion remain refused.
50//! The direct backend API has a Docker containment proof path on macOS/Linux:
51//! a pinned, trusted Linux image must supply /usr/local/bin/python3. It reuses
52//! the container mount policy and a private host lease checked by guest PID 1.
53//! Profile admission pins that image, startup policy and credential channel.
54//!
55//! Child environments are cleared through `agent_session_env`. ACP has no
56//! implicit credential selection: direct callers supply SessionSpec.env, while
57//! profiles seed only the operator-selected credential file in a private home.
58//! Ambient HOME and provider credentials are never restored by this backend.
59//!
60//! Native single-shot completion closes stdin, allows a short exit grace and then
61//! cleans up the owned process group. macOS/Linux observe exit with WNOWAIT
62//! so group identity stays owned until cleanup, before reaping. Windows
63//! requires Job Object assignment. Writes, cancellation, reap and diagnostic
64//! drain have deadlines. Same-group cleanup is not containment against a
65//! descendant that escapes its group; that requires the separate S6 proof.
66
67use crate::backend::{
68    AgentBackend, AgentEvent, AgentSession, PromptMode, SessionExit, SessionSpec,
69};
70#[cfg(windows)]
71use crate::backend_claude::win_job;
72use crate::error::{EngineError, Result};
73use crate::stream_bounds::{drain_to_tail, BoundedLines, STDERR_TAIL_CAP, STDOUT_LINE_CAP};
74use crate::types::TokenUsage;
75use serde_json::{json, Value};
76use std::collections::{HashMap, VecDeque};
77use std::path::PathBuf;
78use std::process::{ExitStatus, Stdio};
79use std::sync::{Arc, Mutex};
80use tokio::process::{Child, ChildStdin, ChildStdout};
81use tokio::task::JoinHandle;
82
83/// Max characters kept in tool-use / tool-result summaries.
84const SUMMARY_MAX_CHARS: usize = 200;
85/// Max characters of captured stderr included in failure messages.
86const STDERR_TAIL_CHARS: usize = 500;
87
88/// The only ACP protocol version this backend speaks (stable schema v1; see
89/// module docs). A peer negotiating anything else is refused at `initialize`.
90const ACP_PROTOCOL_VERSION: u64 = 1;
91
92/// Deadline for one handshake request (`initialize`, `session/new`). Both
93/// may involve adapter startup and authentication. The limit keeps
94/// `start()` from parking the run loop forever. The prompt turn itself is
95/// deliberately unbounded here: turn/stall budgets are the engine's call
96/// (runner-level), not the transport's.
97const HANDSHAKE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(30);
98const CLEANUP_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(1);
99const COMPLETION_GRACE: std::time::Duration = std::time::Duration::from_millis(200);
100const CONTAINER_COMPLETION_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
101const WRITE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(3);
102const CANCEL_WRITE_TIMEOUT: std::time::Duration = std::time::Duration::from_millis(100);
103
104/// JSON-RPC method names, stable schema v1 (`meta.json`).
105mod method {
106    pub(crate) const INITIALIZE: &str = "initialize";
107    pub(crate) const SESSION_NEW: &str = "session/new";
108    pub(crate) const SESSION_PROMPT: &str = "session/prompt";
109    pub(crate) const SESSION_CANCEL: &str = "session/cancel";
110    pub(crate) const SESSION_UPDATE: &str = "session/update";
111    pub(crate) const REQUEST_PERMISSION: &str = "session/request_permission";
112}
113
114/// Claude-style tool name → ACP `kind`s it matches, for interpreting
115/// `SessionSpec::disallowed_tools` patterns against ACP tool calls. ACP
116/// kinds are the only tool identity the wire carries (`title` is
117/// free-form display text), so the mapping is necessarily coarse; it is
118/// used ONLY for deny decisions, never to auto-approve.
119const TOOL_NAME_KINDS: &[(&str, &[&str])] = &[
120    ("Bash", &["execute"]),
121    ("Read", &["read"]),
122    ("Write", &["edit"]),
123    ("Edit", &["edit"]),
124    ("NotebookEdit", &["edit"]),
125    ("Glob", &["search"]),
126    ("Grep", &["search"]),
127    ("WebSearch", &["search", "fetch"]),
128    ("WebFetch", &["fetch"]),
129];
130
131/// ACP `kind`s that mutate the filesystem; refused outright in read-only
132/// (`writable: false`) sessions. `execute` is deliberately not in this set —
133/// see the module-docs permission section.
134const MUTATING_KINDS: &[&str] = &["edit", "delete", "move"];
135
136// ---------------------------------------------------------------------------
137// JSON-RPC framing
138// ---------------------------------------------------------------------------
139
140/// One stdout line classified by JSON-RPC shape. The wire is symmetric
141/// (both sides issue requests), so a line is one of: a response to a client
142/// request, a peer request we must answer, or a notification.
143#[derive(Debug)]
144enum Frame {
145    /// In-process broker messages; never constructible from peer JSON.
146    PermissionAnswer(crate::live_permission::Answer),
147    PermissionExpired,
148    /// `result`/`error` for a client-issued request id.
149    Response {
150        id: u64,
151        outcome: RpcOutcome,
152    },
153    /// Peer→client request (`session/request_permission`, or an unsupported
154    /// client method). The `id` is echoed back verbatim — it may be a string
155    /// or a number per JSON-RPC, so it is kept as a raw [`Value`].
156    Request {
157        id: Value,
158        method: String,
159        params: Value,
160        raw: Value,
161    },
162    /// Peer→client notification (`session/update`, or anything else).
163    Notification {
164        method: String,
165        params: Value,
166        raw: Value,
167    },
168    /// Malformed protocol input: retained diagnostic, then session failure.
169    Unrecognized(Value),
170}
171
172/// The payload of a JSON-RPC response: peer ids are always numbers in our
173/// exchanges with the peer's client side, but the error path keeps the raw
174/// object for the transcript.
175#[derive(Debug)]
176enum RpcOutcome {
177    Result(Value),
178    Error(Value),
179}
180
181/// Classify one stdout line. Unparseable lines become
182/// [`Frame::Unrecognized`] with `raw = {"unparsed": <line>}` so nothing is
183/// ever dropped from transcripts (mirrors `backend_codex::parse_codex_line`).
184fn classify_line(line: &str) -> Frame {
185    let value = match crate::strict_json::parse(line.as_bytes()) {
186        Ok(value) => value,
187        Err(_) => return Frame::Unrecognized(json!({ "unparsed": line })),
188    };
189    classify_value(value)
190}
191
192fn classify_value(value: Value) -> Frame {
193    let obj = match value.as_object() {
194        Some(obj) => obj,
195        None => return Frame::Unrecognized(value),
196    };
197    if obj.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
198        return Frame::Unrecognized(value);
199    }
200    let has_id = obj.contains_key("id");
201    if obj.contains_key("method") && (obj.contains_key("result") || obj.contains_key("error")) {
202        return Frame::Unrecognized(value);
203    }
204    if has_id && !(obj["id"].is_string() || obj["id"].is_i64() || obj["id"].is_u64()) {
205        return Frame::Unrecognized(value);
206    }
207    let method = obj.get("method").and_then(Value::as_str);
208    match (method, has_id) {
209        // Request: method + id.
210        (Some(method), true) => Frame::Request {
211            id: obj.get("id").cloned().unwrap_or(Value::Null),
212            method: method.to_string(),
213            params: obj.get("params").cloned().unwrap_or(Value::Null),
214            raw: value,
215        },
216        // Notification: method, no id.
217        (Some(method), false) => Frame::Notification {
218            method: method.to_string(),
219            params: obj.get("params").cloned().unwrap_or(Value::Null),
220            raw: value,
221        },
222        // Response: no method, carries result or error. Our request ids are
223        // numbers; anything else is not a response to us.
224        (None, true) => {
225            let id = obj.get("id").and_then(Value::as_u64);
226            match (id, obj.get("result"), obj.get("error")) {
227                (Some(id), Some(result), None) => Frame::Response {
228                    id,
229                    outcome: RpcOutcome::Result(result.clone()),
230                },
231                (Some(id), None, Some(error)) => Frame::Response {
232                    id,
233                    outcome: RpcOutcome::Error(error.clone()),
234                },
235                _ => Frame::Unrecognized(value),
236            }
237        }
238        (None, false) => Frame::Unrecognized(value),
239    }
240}
241
242// ---------------------------------------------------------------------------
243// Permission policy (pure; the seam the ticket pins)
244// ---------------------------------------------------------------------------
245
246/// What the seam decided about one `session/request_permission`.
247#[derive(Debug, Clone, PartialEq, Eq)]
248enum PermissionDecision {
249    Allow,
250    /// Human-readable reason; lands on the synthesized denial event's summary.
251    Deny(String),
252}
253
254/// Everything known about the tool call a permission request covers: the
255/// `tool_call`/`tool_call_update` fields tracked by [`AcpSession`], merged
256/// with whatever the request itself carries (the request's embedded
257/// `ToolCallUpdate` may be the first place a title/kind appears).
258#[derive(Debug, Clone, Default)]
259struct ToolCallInfo {
260    kind: String,
261    title: String,
262    /// Command (`execute`) or path (file kinds) extracted from
263    /// `rawInput`/`locations`; the subject glob patterns match against.
264    subject: String,
265}
266
267/// Extract the matchable subject from a tool-call-shaped value
268/// (`rawInput.command`/`cmd` for execute, `locations[0].path` or
269/// `rawInput.path` for file kinds). Explicit raw input invalidates display
270/// fallback; execute never trusts a title as its command.
271fn tool_call_subject(kind: &str, title: &str, call: &Value) -> String {
272    let raw_input = call.get("rawInput").cloned().unwrap_or(Value::Null);
273    let str_at = |value: &Value, keys: &[&str]| -> Option<String> {
274        keys.iter()
275            .find_map(|k| value.get(*k).and_then(Value::as_str).map(str::to_string))
276    };
277    if kind == "execute" {
278        // Display titles and file locations are not executable arguments.
279        return str_at(&raw_input, &["command", "cmd"]).unwrap_or_default();
280    }
281    if let Some(path) = call
282        .get("locations")
283        .and_then(Value::as_array)
284        .and_then(|locs| locs.first())
285        .and_then(|loc| loc.get("path"))
286        .and_then(Value::as_str)
287    {
288        return path.to_string();
289    }
290    if let Some(path) = str_at(&raw_input, &["path", "filePath", "file_path"]) {
291        return path;
292    }
293    if call.get("rawInput").is_some() || call.get("locations").is_some() {
294        return String::new();
295    }
296    title.to_string()
297}
298
299/// `*`-wildcard match (a bare `*` spans any text including empty; every
300/// other character matches literally, case-sensitively — shell commands and
301/// paths are case-sensitive on the platforms this guards).
302fn wildcard_match(pattern: &str, text: &str) -> bool {
303    let parts = pattern.split('*');
304    let anchored_start = !pattern.starts_with('*');
305    let anchored_end = !pattern.ends_with('*');
306    let mut rest = text;
307    let mut first = true;
308    for part in parts {
309        if part.is_empty() {
310            first = false;
311            continue;
312        }
313        match rest.find(part) {
314            Some(idx) if !first || !anchored_start || idx == 0 => {
315                rest = &rest[idx + part.len()..];
316            }
317            _ => return false,
318        }
319        first = false;
320    }
321    // After the last literal, a pattern not ending in `*` must end exactly
322    // there (nothing left over).
323    !anchored_end || rest.is_empty()
324}
325
326/// Split one claude-shaped permission pattern (`Bash(git push*)`, `Write`,
327/// …) into its tool name and subject glob. A bare name carries the glob `*`.
328fn split_pattern(pattern: &str) -> (&str, &str) {
329    match pattern.split_once('(') {
330        Some((name, rest)) => (name.trim(), rest.strip_suffix(')').unwrap_or(rest)),
331        None => (pattern.trim(), "*"),
332    }
333}
334
335/// Whether one permission pattern's tool name maps to this ACP `kind` at
336/// all, regardless of subject. Separate from [`pattern_matches`] because a
337/// pattern that COVERS a call but cannot be evaluated against it is a
338/// refusal, not a pass.
339fn pattern_covers_kind(pattern: &str, kind: &str) -> bool {
340    let (name, _) = split_pattern(pattern);
341    TOOL_NAME_KINDS
342        .iter()
343        .find(|(n, _)| n.eq_ignore_ascii_case(name))
344        .is_some_and(|(_, kinds)| kinds.contains(&kind))
345}
346
347/// Whether one claude-shaped permission pattern (`Bash(git push*)`,
348/// `Write`, …) covers an ACP tool call of `kind` with `subject`.
349fn pattern_matches(pattern: &str, kind: &str, subject: &str) -> bool {
350    let (_, glob) = split_pattern(pattern);
351    if pattern_covers_kind(pattern, kind) {
352        wildcard_match(glob, subject)
353    } else {
354        false
355    }
356}
357
358/// The seam: decide one permission request from the [`SessionSpec`]. Deny
359/// rules are evaluated before the read-only posture so the recorded reason
360/// names the most specific rule that fired.
361///
362/// Fails CLOSED on missing wire fields. ACP v1 leaves `title` optional and
363/// `rawInput` free-form, so a peer can send a tool call with no subject at
364/// all; every glob then matches nothing and the deny list silently vacates.
365/// The same holds for the read-only posture when the peer omits `kind`
366/// (it arrives as `"other"`, which `MUTATING_KINDS` cannot classify). In
367/// both cases the guard cannot be evaluated, so the call is refused and the
368/// reason names the field the peer left out.
369fn decide_permission(spec: &SessionSpec, call: &ToolCallInfo) -> PermissionDecision {
370    if call.kind == "switch_mode" {
371        return PermissionDecision::Deny(
372            "agent mode changes require an explicit human decision".into(),
373        );
374    }
375    if ![
376        "read",
377        "edit",
378        "delete",
379        "move",
380        "search",
381        "execute",
382        "think",
383        "fetch",
384        "switch_mode",
385    ]
386    .contains(&call.kind.as_str())
387    {
388        return PermissionDecision::Deny("tool call carries an unknown or missing ACP kind".into());
389    }
390    if call.subject.trim().is_empty() {
391        if let Some(pattern) = spec
392            .disallowed_tools
393            .iter()
394            .find(|pattern| pattern_covers_kind(pattern, &call.kind))
395        {
396            return PermissionDecision::Deny(format!(
397                "tool call carries no subject (no rawInput command or path, no locations[0].path, \
398                 no title), so deny pattern {pattern:?} for ACP kind {:?} cannot be evaluated",
399                call.kind
400            ));
401        }
402    }
403    for pattern in &spec.disallowed_tools {
404        if pattern_matches(pattern, &call.kind, &call.subject) {
405            return PermissionDecision::Deny(format!(
406                "matches SessionSpec.disallowed_tools pattern {pattern:?}"
407            ));
408        }
409    }
410    if !spec.writable {
411        let kind = call.kind.trim();
412        if kind.is_empty() || kind == "other" {
413            return PermissionDecision::Deny(format!(
414                "read-only session (writable: false): tool call carries no usable ACP kind \
415                 ({kind:?}), so whether it mutates the filesystem cannot be decided"
416            ));
417        }
418        if MUTATING_KINDS.contains(&kind) {
419            return PermissionDecision::Deny(format!(
420                "read-only session (writable: false): ACP kind {:?} mutates the filesystem",
421                call.kind
422            ));
423        }
424    }
425    PermissionDecision::Allow
426}
427
428/// Select only a well-formed, unambiguous one-time option. Adapter option
429/// IDs are opaque: the kind supplies semantics, never an ID spelling.
430fn permission_response(decision: &PermissionDecision, options: &[Value]) -> Value {
431    let mut ids = std::collections::HashSet::new();
432    let valid = !options.is_empty()
433        && options.len() <= 32
434        && options.iter().all(|option| {
435            option
436                .get("optionId")
437                .and_then(Value::as_str)
438                .filter(|id| !id.trim().is_empty() && id.len() <= 256)
439                .is_some_and(|id| ids.insert(id))
440        });
441    let kind = match decision {
442        PermissionDecision::Allow => "allow_once",
443        PermissionDecision::Deny(_) => "reject_once",
444    };
445    let selected = valid
446        .then(|| {
447            let mut matching = options
448                .iter()
449                .filter(|option| option.get("kind").and_then(Value::as_str) == Some(kind));
450            let first = matching.next()?;
451            matching.next().is_none().then_some(first)
452        })
453        .flatten();
454    match selected {
455        Some(option) => {
456            json!({ "outcome": { "outcome": "selected", "optionId": option["optionId"] } })
457        }
458        None => json!({ "outcome": { "outcome": "cancelled" } }),
459    }
460}
461
462// ---------------------------------------------------------------------------
463// Event mapping
464// ---------------------------------------------------------------------------
465
466/// Keep at most `max` characters (not bytes — never splits a code point).
467fn truncate_chars(text: &str, max: usize) -> String {
468    if text.chars().count() <= max {
469        text.to_string()
470    } else {
471        text.chars().take(max).collect()
472    }
473}
474
475/// Last `max` characters of `text` (for stderr tails in failure messages).
476fn last_chars(text: &str, max: usize) -> String {
477    let chars: Vec<char> = text.chars().collect();
478    let start = chars.len().saturating_sub(max);
479    chars[start..].iter().collect()
480}
481
482/// Summarize a tool call's content for a [`AgentEvent::ToolResult`]: the
483/// first text content block, else the updated title, else the status.
484fn tool_result_summary(update: &Value, tracked: &ToolCallInfo, status: &str) -> String {
485    if let Some(content) = update.get("content").and_then(Value::as_array) {
486        for item in content {
487            if item.get("type").and_then(Value::as_str) == Some("content") {
488                if let Some(text) = item
489                    .get("content")
490                    .and_then(|c| c.get("text"))
491                    .and_then(Value::as_str)
492                {
493                    return truncate_chars(text, SUMMARY_MAX_CHARS);
494                }
495            }
496        }
497    }
498    if !tracked.title.is_empty() {
499        return truncate_chars(&tracked.title, SUMMARY_MAX_CHARS);
500    }
501    status.to_string()
502}
503
504// ---------------------------------------------------------------------------
505// Backend
506// ---------------------------------------------------------------------------
507
508/// The [`AgentBackend`] for an ACP-speaking agent executable (KRZ-301).
509///
510/// There is no canonical ACP binary name and no `--version` convention, so
511/// there is deliberately no discovery/probe: the configured command is
512/// spawned directly and the `initialize` handshake IS the probe — a
513/// non-ACP executable fails there, loudly, before any prompt turn.
514#[derive(Debug, Clone)]
515pub struct AcpBackend {
516    program: PathBuf,
517    args: Vec<String>,
518    profile: Option<crate::acp_worker::AcpWorkerProfile>,
519}
520
521impl AcpBackend {
522    /// Spawn `program args` as the ACP agent (no validation performed; the
523    /// handshake at session start is the validation).
524    pub fn new(program: impl Into<PathBuf>, args: Vec<String>) -> Self {
525        AcpBackend {
526            program: program.into(),
527            args,
528            profile: None,
529        }
530    }
531
532    /// Ordinary-worker construction. Profile argv and startup policy are fixed;
533    /// the credential is read only at session start, after boundary checks.
534    pub fn for_worker(cfg: &crate::types::RoleConfig) -> Result<Self> {
535        if let Some(profile) = &cfg.acp_profile {
536            profile.validate_config(
537                crate::types::Role::Worker,
538                cfg,
539                crate::types::WorkerIsolation::Worktree,
540            )?;
541            let definition = profile.definition()?;
542            Ok(Self {
543                program: definition.program.into(),
544                args: definition.args.iter().map(|s| (*s).to_owned()).collect(),
545                profile: Some(profile.clone()),
546            })
547        } else {
548            let command = cfg.acp_command.as_ref().ok_or_else(|| {
549                EngineError::Config("ACP worker requires acpCommand or acpProfile".into())
550            })?;
551            Ok(Self::new(command, cfg.acp_args.clone()))
552        }
553    }
554
555    /// The program this backend spawns.
556    pub fn program(&self) -> &std::path::Path {
557        &self.program
558    }
559}
560
561fn effective_prompt(spec: &SessionSpec) -> String {
562    let text = match &spec.prompt {
563        PromptMode::SingleShot(text) | PromptMode::Streaming(text) => text,
564    };
565    match &spec.append_system_prompt {
566        Some(system) if !system.is_empty() => format!("{system}\n\n{text}"),
567        _ => text.clone(),
568    }
569}
570
571fn peer_reported_model(session: &Value) -> Option<&str> {
572    session
573        .get("configOptions")
574        .and_then(Value::as_array)
575        .and_then(|options| {
576            options.iter().find_map(|option| {
577                (option.get("category").and_then(Value::as_str) == Some("model"))
578                    .then(|| option.get("currentValue").and_then(Value::as_str))
579                    .flatten()
580            })
581        })
582        .or_else(|| session.get("models")?.get("currentModelId")?.as_str())
583        .filter(|model| !model.trim().is_empty() && model.len() <= 256)
584}
585
586#[async_trait::async_trait]
587impl AgentBackend for AcpBackend {
588    async fn start(&self, mut spec: SessionSpec) -> Result<Box<dyn AgentSession>> {
589        // Configuration admission is not the only caller of this public
590        // backend. Never silently discard an embedding caller's boundary.
591        let container_requested = spec.sandbox.as_ref().is_some_and(|sandbox| {
592            cfg!(any(target_os = "macos", target_os = "linux"))
593                && sandbox.backend == crate::sandbox::SandboxBackend::Container
594                && sandbox.container.is_some()
595        });
596        if spec.sandbox.is_some() && !container_requested {
597            return Err(EngineError::Backend(
598                "acp backend: enforced containment is not certified; refusing the supplied sandbox before spawn".into(),
599            ));
600        }
601        if spec.resume.is_some() {
602            return Err(EngineError::Backend(
603                "acp backend: resume is unsupported (session/load is an optional v1 \
604                 capability this backend does not negotiate)"
605                    .to_string(),
606            ));
607        }
608        let profile_home = self
609            .profile
610            .as_ref()
611            .map(|profile| profile.prepare(&mut spec))
612            .transpose()?;
613        let model = spec.model.clone();
614
615        #[cfg(any(target_os = "macos", target_os = "linux"))]
616        let (container, mut command) = if container_requested {
617            let (container, command) =
618                crate::acp_container::OwnedContainer::prepare(&spec, &self.program, &self.args)
619                    .await?;
620            (Some(container), command)
621        } else {
622            (None, self.native_command(&spec))
623        };
624        #[cfg(not(any(target_os = "macos", target_os = "linux")))]
625        let mut command = self.native_command(&spec);
626        command
627            .current_dir(&spec.cwd)
628            // ACP is bidirectional: stdin carries the client's requests.
629            .stdin(Stdio::piped())
630            .stdout(Stdio::piped())
631            .stderr(Stdio::piped())
632            .kill_on_drop(true);
633        // Unix: make the child the leader of a fresh process group so aborts
634        // can kill the whole tree, mirroring `backend_claude::ClaudeBackend`.
635        #[cfg(unix)]
636        command.process_group(0);
637
638        let mut child = command.spawn().map_err(|e| {
639            EngineError::Backend(format!(
640                "failed to spawn acp agent {}: {e}",
641                self.program.display()
642            ))
643        })?;
644
645        // Windows: kill-on-close Job Object, mirroring `backend_claude`.
646        #[cfg(windows)]
647        let job = match child.raw_handle() {
648            Some(handle) => match win_job::JobHandle::create_and_assign(handle) {
649                Ok(job) => Some(job),
650                Err(e) => {
651                    let _ = child.kill().await;
652                    return Err(EngineError::Backend(format!(
653                        "acp Job Object assignment failed: {e}"
654                    )));
655                }
656            },
657            None => {
658                let _ = child.kill().await;
659                return Err(EngineError::Backend(
660                    "acp child has no process handle".into(),
661                ));
662            }
663        };
664
665        let stdin = child
666            .stdin
667            .take()
668            .ok_or_else(|| EngineError::Backend("acp child has no stdin pipe".to_string()))?;
669        let stdout = child
670            .stdout
671            .take()
672            .ok_or_else(|| EngineError::Backend("acp child has no stdout pipe".to_string()))?;
673        let stderr = child
674            .stderr
675            .take()
676            .ok_or_else(|| EngineError::Backend("acp child has no stderr pipe".to_string()))?;
677
678        // Capture stderr concurrently so a chatty child never blocks on a
679        // full pipe and failure messages can include the tail. The stream is
680        // drained to EOF but only a bounded tail is retained — a noisy or
681        // malicious peer must not exhaust host memory (stream_bounds).
682        let stderr_buf = Arc::new(Mutex::new(String::new()));
683        let stderr_task = {
684            let buf = Arc::clone(&stderr_buf);
685            tokio::spawn(async move {
686                let tail = drain_to_tail(stderr, STDERR_TAIL_CAP).await;
687                *buf.lock().expect("stderr buffer lock") = tail;
688            })
689        };
690
691        let (permission_responder, permission_answers) =
692            crate::live_permission::PermissionResponder::channel();
693        let mut session = AcpSession {
694            session_id: spec.session_id.clone(),
695            acp_session_id: None,
696            model,
697            spec,
698            child,
699            #[cfg(any(target_os = "macos", target_os = "linux"))]
700            container,
701            profile_home,
702            #[cfg(windows)]
703            job,
704            stdin: Some(stdin),
705            child_status: None,
706            cleanup_failure: None,
707            drain_deadline: None,
708            lines: BoundedLines::new_strict(stdout),
709            stderr_buf,
710            stderr_task: Some(stderr_task),
711            queue: VecDeque::new(),
712            next_request_id: 1,
713            prompt_request_id: None,
714            tool_calls: HashMap::new(),
715            permission_responder,
716            permission_answers,
717            deferred_permission_answer: None,
718            pending_permissions: HashMap::new(),
719            seen_permission_ids: std::collections::HashSet::new(),
720            message_text: String::new(),
721            message_id: None,
722            last_usage: None,
723            previous_cost_total: None,
724            handshake_bytes: 0,
725            handshake_complete: false,
726            saw_result: false,
727            exit: None,
728        };
729
730        // The handshake is eager: a non-ACP executable (or a hung peer)
731        // fails start() loudly rather than mid-run. A handshake failure
732        // kills the child before the error crosses back.
733        if let Err(e) = session.handshake().await {
734            session.kill_child().await;
735            let message = match e {
736                EngineError::Backend(message) => message,
737                other => other.to_string(),
738            };
739            return Err(EngineError::Backend(format!(
740                "{message}; startup stderr after cleanup: {}{}",
741                session.stderr_tail(),
742                session
743                    .cleanup_failure
744                    .map(|cause| format!("; cleanup unconfirmed: {cause}"))
745                    .unwrap_or_default(),
746            )));
747        }
748        Ok(Box::new(session))
749    }
750}
751
752impl AcpBackend {
753    fn native_command(&self, spec: &SessionSpec) -> tokio::process::Command {
754        let mut command = tokio::process::Command::new(&self.program);
755        command
756            .args(&self.args)
757            // agent-env-clear: ACP has no implicit credential channel.
758            // Only explicit session variables cross the native boundary.
759            .env_clear()
760            .envs(crate::agent_env::agent_session_env(
761                &spec.env,
762                &spec.session_id,
763                None,
764            ));
765        command
766    }
767}
768
769// ---------------------------------------------------------------------------
770// Session
771// ---------------------------------------------------------------------------
772
773/// A live ACP session (the [`AgentSession`] impl).
774///
775/// Reading is inline in `next_event`. Every write is bounded: a peer that
776/// stops reading stdin cannot block cancellation indefinitely.
777struct PendingPermission {
778    proposal: crate::live_permission::Proposal,
779    expires_at: tokio::time::Instant,
780}
781
782pub struct AcpSession {
783    session_id: String,
784    /// The peer-issued session id (`session/new` response).
785    acp_session_id: Option<String>,
786    model: String,
787    /// Kept for permission decisions (`disallowed_tools`, `writable`).
788    spec: SessionSpec,
789    child: Child,
790    #[cfg(any(target_os = "macos", target_os = "linux"))]
791    container: Option<crate::acp_container::OwnedContainer>,
792    profile_home: Option<crate::acp_worker::PreparedProfile>,
793    #[cfg(windows)]
794    job: Option<win_job::JobHandle>,
795    stdin: Option<ChildStdin>,
796    child_status: Option<ExitStatus>,
797    cleanup_failure: Option<&'static str>,
798    drain_deadline: Option<tokio::time::Instant>,
799    lines: BoundedLines<ChildStdout>,
800    stderr_buf: Arc<Mutex<String>>,
801    stderr_task: Option<JoinHandle<()>>,
802    /// Converted events not yet surfaced; popped one per `next_event`.
803    queue: VecDeque<AgentEvent>,
804    next_request_id: u64,
805    /// The in-flight `session/prompt` request id; its response synthesizes
806    /// the terminal `Result`. `None` outside a prompt turn.
807    prompt_request_id: Option<u64>,
808    /// Tracked tool calls by `toolCallId` (kind/title/subject), so a
809    /// `tool_call_update` or permission request resolves to what is known
810    /// about the call.
811    tool_calls: HashMap<String, ToolCallInfo>,
812    permission_responder: crate::live_permission::PermissionResponder,
813    permission_answers: tokio::sync::mpsc::Receiver<crate::live_permission::Answer>,
814    deferred_permission_answer: Option<crate::live_permission::Answer>,
815    pending_permissions: HashMap<String, PendingPermission>,
816    seen_permission_ids: std::collections::HashSet<String>,
817    /// Accumulated text of the CURRENT assistant message (chunks with the
818    /// same `messageId` concatenate; a changed/missing-`messageId` boundary
819    /// starts a new message, and the last message wins the terminal Result,
820    /// mirroring `backend_codex`'s last-agent_message rule).
821    message_text: String,
822    message_id: Option<String>,
823    /// Latest `usage_update` (context state + optional cumulative USD cost);
824    /// its cost lands on the terminal `Result` (see module docs).
825    last_usage: Option<Value>,
826    previous_cost_total: Option<f64>,
827    handshake_bytes: usize,
828    handshake_complete: bool,
829    saw_result: bool,
830    exit: Option<SessionExit>,
831}
832
833#[cfg(unix)]
834impl Drop for AcpSession {
835    fn drop(&mut self) {
836        crate::backend_claude::kill_unreaped_group(&self.child);
837    }
838}
839
840async fn wait_for_peer_exit(child: &mut Child) -> std::io::Result<()> {
841    #[cfg(any(target_os = "macos", target_os = "linux"))]
842    {
843        let pid = child
844            .id()
845            .ok_or_else(|| std::io::Error::other("acp child already reaped"))?;
846        crate::command_exec::control_leader_exited(pid).await
847    }
848    #[cfg(not(any(target_os = "macos", target_os = "linux")))]
849    {
850        child.wait().await.map(|_| ())
851    }
852}
853
854fn kill_owned_peer(child: &mut Child, #[cfg(windows)] job: Option<&win_job::JobHandle>) {
855    #[cfg(unix)]
856    crate::backend_claude::kill_unreaped_group(child);
857    #[cfg(windows)]
858    if let Some(job) = job {
859        job.kill();
860    }
861    // Also target the still-owned leader if it changed its process group.
862    if child.id().is_some() {
863        let _ = child.start_kill();
864    }
865}
866
867impl AcpSession {
868    // -- wire helpers -------------------------------------------------------
869
870    /// Serialize one JSON-RPC message as a single NDJSON line on the peer's
871    /// stdin. Keys are inserted in sorted order by serde_json's default map,
872    /// which also puts `"id"` first — mock peers and debugging tools rely on
873    /// nothing more than JSON semantics, but stable key order keeps captured
874    /// transcripts diffable.
875    async fn write_message(&mut self, message: Value) -> Result<()> {
876        use tokio::io::AsyncWriteExt;
877        let mut line = serde_json::to_string(&message)
878            .map_err(|e| EngineError::Backend(format!("failed to encode acp message: {e}")))?;
879        line.push('\n');
880        if line.len() > STDOUT_LINE_CAP {
881            return Err(EngineError::Backend(
882                "acp request exceeds the frame byte limit".into(),
883            ));
884        }
885        let stdin = self
886            .stdin
887            .as_mut()
888            .ok_or_else(|| EngineError::Backend("acp stdin is closed".into()))?;
889        tokio::time::timeout(WRITE_TIMEOUT, async {
890            stdin.write_all(line.as_bytes()).await?;
891            stdin.flush().await
892        })
893        .await
894        .map_err(|_| EngineError::Backend("acp stdin write timed out".into()))?
895        .map_err(|e| EngineError::Backend(format!("failed to write to acp agent stdin: {e}")))?;
896        Ok(())
897    }
898
899    /// Send a client request and return its id (the caller either awaits the
900    /// response via [`AcpSession::pump_until_response`] or, for
901    /// `session/prompt`, leaves it to the `next_event` loop).
902    async fn send_request(&mut self, method: &str, params: Value) -> Result<u64> {
903        let id = self.next_request_id;
904        self.next_request_id += 1;
905        self.write_message(json!({
906            "jsonrpc": "2.0",
907            "id": id,
908            "method": method,
909            "params": params,
910        }))
911        .await?;
912        Ok(id)
913    }
914
915    /// `initialize` + `session/new`, then the first `session/prompt` (its
916    /// response streams in through `next_event` like any later turn).
917    async fn handshake(&mut self) -> Result<()> {
918        #[cfg(any(target_os = "macos", target_os = "linux"))]
919        if let Some(container) = &mut self.container {
920            container
921                .write_launch(self.stdin.as_mut().expect("startup stdin"))
922                .await?;
923        }
924        let init_id = self
925            .send_request(
926                method::INITIALIZE,
927                json!({
928                    "protocolVersion": ACP_PROTOCOL_VERSION,
929                    "clientCapabilities": {
930                        // fs/terminal unsupported: a conformant peer never
931                        // calls fs/* or terminal/*; one that does is answered
932                        // with -32601 rather than silently served.
933                        "fs": { "readTextFile": false, "writeTextFile": false },
934                        "terminal": false,
935                    },
936                    "clientInfo": {
937                        "name": "kranz",
938                        "title": "kranz mission engine",
939                        "version": env!("CARGO_PKG_VERSION"),
940                    },
941                }),
942            )
943            .await?;
944        let init_result = self.pump_until_response(init_id, HANDSHAKE_TIMEOUT).await?;
945        let peer_version = init_result
946            .get("protocolVersion")
947            .and_then(Value::as_u64)
948            .unwrap_or(0);
949        if peer_version != ACP_PROTOCOL_VERSION {
950            return Err(EngineError::Backend(format!(
951                "acp agent negotiated protocol version {peer_version}, but this backend speaks \
952                 only stable version {ACP_PROTOCOL_VERSION} (schema v1)"
953            )));
954        }
955
956        let new_id = self
957            .send_request(
958                method::SESSION_NEW,
959                json!({
960                    "cwd": self.spec.cwd.display().to_string(),
961                    "mcpServers": [],
962                }),
963            )
964            .await?;
965        let new_result = self.pump_until_response(new_id, HANDSHAKE_TIMEOUT).await?;
966        let acp_session_id = new_result
967            .get("sessionId")
968            .and_then(Value::as_str)
969            .filter(|id| !id.trim().is_empty() && id.len() <= 256)
970            .ok_or_else(|| {
971                EngineError::Backend("acp session/new response carried no sessionId".to_string())
972            })?
973            .to_string();
974        self.acp_session_id = Some(acp_session_id.clone());
975        self.handshake_complete = true;
976        let reported_model = peer_reported_model(&new_result).map(str::to_owned);
977        #[cfg(any(target_os = "macos", target_os = "linux"))]
978        let containment = self.container.as_ref().map(|container| container.receipt());
979        #[cfg(not(any(target_os = "macos", target_os = "linux")))]
980        let containment: Option<Value> = None;
981        // Configured attribution and peer-reported state are different facts.
982        self.queue.push_back(AgentEvent::Init {
983            session_id: acp_session_id,
984            model: reported_model
985                .clone()
986                .unwrap_or_else(|| "unreported".into()),
987            raw: json!({
988                "initialize": init_result,
989                "sessionNew": new_result,
990                "engineSessionId": self.session_id,
991                "configuredModel": self.model,
992                "modelSource": if reported_model.is_some() { "peer" } else { "unreported" },
993                "configuredModelSelectionApplied": false,
994                "containment": containment,
995                "workerProfile": self.profile_home.as_ref().map(|p| &p.receipt),
996                "synthesizedBy": "backend_acp",
997            }),
998        });
999
1000        let prompt_text = effective_prompt(&self.spec);
1001        self.send_prompt(&prompt_text).await
1002    }
1003
1004    /// Send one `session/prompt` request and mark its id as the turn whose
1005    /// response synthesizes the terminal `Result`.
1006    async fn send_prompt(&mut self, text: &str) -> Result<()> {
1007        let acp_session_id = self
1008            .acp_session_id
1009            .clone()
1010            .ok_or_else(|| EngineError::Backend("acp session not established yet".to_string()))?;
1011        let id = self
1012            .send_request(
1013                method::SESSION_PROMPT,
1014                json!({
1015                    "sessionId": acp_session_id,
1016                    "prompt": [ { "type": "text", "text": text } ],
1017                }),
1018            )
1019            .await?;
1020        self.prompt_request_id = Some(id);
1021        // A new turn begins: the terminal Result stitches only this turn's
1022        // last message.
1023        self.message_text.clear();
1024        self.message_id = None;
1025        self.last_usage = None;
1026        Ok(())
1027    }
1028
1029    /// Read frames until the response to `id` arrives, converting everything
1030    /// else through the normal frame path (notifications become queued
1031    /// events; permission requests are answered). Used by the handshake —
1032    /// the only place a response is awaited synchronously.
1033    async fn pump_until_response(
1034        &mut self,
1035        id: u64,
1036        timeout: std::time::Duration,
1037    ) -> Result<Value> {
1038        let pump = async {
1039            loop {
1040                let frame = match self.read_frame().await? {
1041                    Some(frame) => frame,
1042                    None => {
1043                        return Err(EngineError::Backend(format!(
1044                            "acp agent closed stdout before answering request id {id}; \
1045                             stderr tail: {}",
1046                            self.stderr_tail()
1047                        )))
1048                    }
1049                };
1050                match frame {
1051                    Frame::Response {
1052                        id: response_id,
1053                        outcome,
1054                    } if response_id == id => {
1055                        return match outcome {
1056                            RpcOutcome::Result(result) => Ok(result),
1057                            RpcOutcome::Error(error) => Err(EngineError::Backend(format!(
1058                                "acp request id {id} failed: {error}"
1059                            ))),
1060                        };
1061                    }
1062                    other => self.handle_frame(other).await?,
1063                }
1064            }
1065        };
1066        match tokio::time::timeout(timeout, pump).await {
1067            Ok(result) => result,
1068            Err(_) => Err(EngineError::Backend(format!(
1069                "acp agent did not answer request id {id} within {}s (handshake timeout)",
1070                timeout.as_secs()
1071            ))),
1072        }
1073    }
1074
1075    /// Read and classify the next stdout line; `None` at EOF. Blank lines
1076    /// are skipped (NDJSON tolerates them; a peer's pretty-printing or
1077    /// keepalive must not fabricate events).
1078    async fn read_frame(&mut self) -> Result<Option<Frame>> {
1079        loop {
1080            let permission_wait = self
1081                .pending_permissions
1082                .values()
1083                .map(|p| {
1084                    p.expires_at
1085                        .saturating_duration_since(tokio::time::Instant::now())
1086                        .min(
1087                            (p.proposal.deadline - chrono::Utc::now())
1088                                .to_std()
1089                                .unwrap_or_default(),
1090                        )
1091                })
1092                .min()
1093                .unwrap_or(std::time::Duration::from_secs(86400));
1094            if permission_wait.is_zero() {
1095                return Ok(Some(Frame::PermissionExpired));
1096            }
1097            let can_answer = !self.lines.has_partial_line();
1098            let read = {
1099                let line = self.lines.next_line();
1100                tokio::pin!(line);
1101                let read = if let Some(deadline) = self.drain_deadline {
1102                    tokio::time::timeout_at(deadline, &mut line)
1103                        .await
1104                        .map_err(|_| {
1105                            EngineError::Backend(
1106                                "acp stdout remained open after peer cleanup".into(),
1107                            )
1108                        })?
1109                } else {
1110                    tokio::select! {
1111                        biased;
1112                        // Drain already-buffered action changes before applying
1113                        // a queued answer to the older invocation description.
1114                        read = &mut line => read,
1115                        answer = async {
1116                            if let Some(answer) = self.deferred_permission_answer.take() {
1117                                Some(answer)
1118                            } else {
1119                                self.permission_answers.recv().await
1120                            }
1121                        }, if can_answer => {
1122                            return Ok(answer.map(Frame::PermissionAnswer));
1123                        }
1124                        _ = tokio::time::sleep(permission_wait), if !self.pending_permissions.is_empty() => {
1125                            return Ok(Some(Frame::PermissionExpired));
1126                        }
1127                        exited = wait_for_peer_exit(&mut self.child) => {
1128                            exited.map_err(|e| EngineError::Backend(format!("acp process observation failed: {e}")))?;
1129                            kill_owned_peer(&mut self.child, #[cfg(windows)] self.job.as_ref());
1130                            self.child_status = Some(self.child.wait().await.map_err(|e| {
1131                                EngineError::Backend(format!("acp process reap failed: {e}"))
1132                            })?);
1133                            let deadline = tokio::time::Instant::now() + CLEANUP_TIMEOUT;
1134                            self.drain_deadline = Some(deadline);
1135                            tokio::time::timeout_at(deadline, &mut line).await
1136                                .map_err(|_| EngineError::Backend("acp stdout remained open after peer cleanup".into()))?
1137                        }
1138                    }
1139                };
1140                read
1141            };
1142            match read {
1143                Ok(Some(line)) if line.trim().is_empty() => continue,
1144                Ok(Some(line)) => {
1145                    if self.profile_home.is_some()
1146                        && crate::strict_json::parse(line.as_bytes()).is_err()
1147                    {
1148                        // Do not retain malformed credential-bearing peer input,
1149                        // including duplicate keys hiding an escaped token.
1150                        return Err(EngineError::Backend(
1151                            "acp profile peer emitted invalid unique-key JSON; frame refused"
1152                                .into(),
1153                        ));
1154                    }
1155                    if self
1156                        .profile_home
1157                        .as_ref()
1158                        .is_some_and(|p| p.contains_secret(&line))
1159                    {
1160                        return Err(EngineError::Backend(
1161                            "acp peer exposed a configured credential; frame refused".into(),
1162                        ));
1163                    }
1164                    if !self.handshake_complete {
1165                        self.handshake_bytes = self.handshake_bytes.saturating_add(line.len());
1166                        if self.handshake_bytes > STDOUT_LINE_CAP {
1167                            return Err(EngineError::Backend(
1168                                "acp handshake output exceeded its byte limit".into(),
1169                            ));
1170                        }
1171                    }
1172                    return Ok(Some(classify_line(&line)));
1173                }
1174                Ok(None) => return Ok(None),
1175                Err(e) => {
1176                    return Err(EngineError::Backend(format!(
1177                        "error reading acp agent stdout: {e}; stderr tail: {}",
1178                        self.stderr_tail()
1179                    )))
1180                }
1181            }
1182        }
1183    }
1184
1185    /// Convert one frame into queued events and side effects (permission
1186    /// answers, tool tracking, usage capture, terminal-Result synthesis).
1187    async fn handle_frame(&mut self, frame: Frame) -> Result<()> {
1188        match frame {
1189            Frame::PermissionAnswer(answer) => self.answer_permission(answer).await?,
1190            Frame::PermissionExpired => {
1191                // Do not manufacture or replay a denial resolution. Terminating
1192                // the peer closes its requests without authorizing an effect.
1193                return Err(EngineError::Backend(
1194                    "live permission deadline expired".into(),
1195                ));
1196            }
1197            Frame::Notification {
1198                method,
1199                params,
1200                raw,
1201            } => {
1202                if method == method::SESSION_UPDATE {
1203                    self.handle_session_update(&params, raw)?;
1204                } else {
1205                    self.queue.push_back(AgentEvent::Other { raw });
1206                }
1207            }
1208            Frame::Request {
1209                id,
1210                method,
1211                params,
1212                raw,
1213            } => {
1214                if method == method::REQUEST_PERMISSION {
1215                    self.handle_permission_request(id, &params, raw).await?;
1216                } else {
1217                    // A client capability we did not advertise (fs/*,
1218                    // terminal/*, elicitation/*): refuse with JSON-RPC
1219                    // -32601 rather than silently serving or hanging.
1220                    self.write_message(json!({
1221                        "jsonrpc": "2.0",
1222                        "id": id,
1223                        "error": {
1224                            "code": -32601,
1225                            "message": format!("kranz acp backend does not support {method:?}"),
1226                        },
1227                    }))
1228                    .await?;
1229                    self.queue.push_back(AgentEvent::Other { raw });
1230                }
1231            }
1232            Frame::Response { id, outcome } => {
1233                if Some(id) == self.prompt_request_id {
1234                    self.prompt_request_id = None;
1235                    self.synthesize_result(outcome, id);
1236                } else {
1237                    // A response to nothing outstanding (a straggler from a
1238                    // cancelled turn, or a peer bug): transcript only.
1239                    self.queue.push_back(AgentEvent::Other {
1240                        raw: match outcome {
1241                            RpcOutcome::Result(result) => {
1242                                json!({ "unmatchedResponse": { "id": id, "result": result } })
1243                            }
1244                            RpcOutcome::Error(error) => {
1245                                json!({ "unmatchedResponse": { "id": id, "error": error } })
1246                            }
1247                        },
1248                    });
1249                }
1250            }
1251            Frame::Unrecognized(raw) => {
1252                self.queue.push_back(AgentEvent::Other { raw });
1253                return Err(EngineError::Backend(
1254                    "acp peer emitted malformed JSON-RPC".into(),
1255                ));
1256            }
1257        }
1258        Ok(())
1259    }
1260
1261    fn track_tool_call(&mut self, id: String, info: ToolCallInfo) -> Result<()> {
1262        if id.trim().is_empty() || id.len() > 256 {
1263            return Err(EngineError::Backend(
1264                "acp tool update has no valid toolCallId".into(),
1265            ));
1266        }
1267        let bytes = info.kind.len() + info.title.len() + info.subject.len();
1268        let retained = self
1269            .tool_calls
1270            .iter()
1271            .filter(|(key, _)| *key != &id)
1272            .map(|(key, value)| {
1273                key.len() + value.kind.len() + value.title.len() + value.subject.len()
1274            })
1275            .sum::<usize>();
1276        if self.tool_calls.len() >= 1024
1277            || retained.saturating_add(bytes).saturating_add(id.len()) > STDOUT_LINE_CAP
1278        {
1279            return Err(EngineError::Backend(
1280                "acp tool-call tracking exceeded its limit".into(),
1281            ));
1282        }
1283        self.tool_calls.insert(id, info);
1284        Ok(())
1285    }
1286
1287    /// Map one `session/update` notification onto events (see module docs).
1288    fn handle_session_update(&mut self, params: &Value, raw: Value) -> Result<()> {
1289        let Some(expected) = self.acp_session_id.as_deref() else {
1290            // No prompt has been sent: pre-session notices are diagnostic only.
1291            self.queue.push_back(AgentEvent::Other { raw });
1292            return Ok(());
1293        };
1294        if params.get("sessionId").and_then(Value::as_str) != Some(expected) {
1295            return Err(EngineError::Backend(
1296                "acp update has a foreign or missing sessionId".into(),
1297            ));
1298        }
1299        let update = params.get("update").cloned().unwrap_or(Value::Null);
1300        if let Some(call_id) = update.get("toolCallId").and_then(Value::as_str) {
1301            if let Some(pending) = self
1302                .pending_permissions
1303                .values()
1304                .map(|p| &p.proposal)
1305                .find(|p| p.tool_call_id == call_id)
1306            {
1307                // Compare complete effect-bearing fields, including content.
1308                // A same-title/same-path edit can still contain different bytes.
1309                let changed = ["kind", "rawInput", "locations", "content"]
1310                    .iter()
1311                    .any(|key| {
1312                        update
1313                            .get(*key)
1314                            .is_some_and(|value| pending.action.get(*key) != Some(value))
1315                    });
1316                let terminal = update
1317                    .get("status")
1318                    .and_then(Value::as_str)
1319                    .is_some_and(|s| s == "completed" || s == "failed" || s == "in_progress");
1320                if changed || terminal {
1321                    return Err(EngineError::Backend(
1322                        "ACP invocation changed or started while consent was pending".into(),
1323                    ));
1324                }
1325            }
1326        }
1327        match update.get("sessionUpdate").and_then(Value::as_str) {
1328            Some("agent_message_chunk") => {
1329                let text = update
1330                    .get("content")
1331                    .and_then(|c| c.get("text"))
1332                    .and_then(Value::as_str)
1333                    .unwrap_or("");
1334                if text.is_empty() {
1335                    self.queue.push_back(AgentEvent::Other { raw });
1336                    return Ok(());
1337                }
1338                // Message boundaries: a changed messageId starts a new
1339                // message; the LAST message wins the terminal Result.
1340                let chunk_id = update
1341                    .get("messageId")
1342                    .and_then(Value::as_str)
1343                    .map(str::to_string);
1344                if chunk_id.is_some() && chunk_id != self.message_id {
1345                    self.message_text.clear();
1346                    self.message_id = chunk_id;
1347                }
1348                if self.message_text.len().saturating_add(text.len()) > STDOUT_LINE_CAP {
1349                    return Err(EngineError::Backend(
1350                        "acp assistant message exceeded its byte limit".into(),
1351                    ));
1352                }
1353                self.message_text.push_str(text);
1354                self.queue.push_back(AgentEvent::Text {
1355                    text: text.to_string(),
1356                    raw,
1357                });
1358            }
1359            Some("tool_call") => {
1360                let id = update
1361                    .get("toolCallId")
1362                    .and_then(Value::as_str)
1363                    .unwrap_or_default()
1364                    .to_string();
1365                let kind = update
1366                    .get("kind")
1367                    .and_then(Value::as_str)
1368                    .unwrap_or("other")
1369                    .to_string();
1370                let title = update
1371                    .get("title")
1372                    .and_then(Value::as_str)
1373                    .unwrap_or_default()
1374                    .to_string();
1375                let subject = tool_call_subject(&kind, &title, &update);
1376                self.track_tool_call(
1377                    id,
1378                    ToolCallInfo {
1379                        kind: kind.clone(),
1380                        title: title.clone(),
1381                        subject,
1382                    },
1383                )?;
1384                self.queue.push_back(AgentEvent::ToolUse {
1385                    tool: kind,
1386                    summary: truncate_chars(&title, SUMMARY_MAX_CHARS),
1387                    raw,
1388                });
1389            }
1390            Some("tool_call_update") => {
1391                let id = update
1392                    .get("toolCallId")
1393                    .and_then(Value::as_str)
1394                    .unwrap_or_default()
1395                    .to_string();
1396                let status = update
1397                    .get("status")
1398                    .and_then(Value::as_str)
1399                    .unwrap_or("")
1400                    .to_string();
1401                let mut tracked = self.tool_calls.get(&id).cloned().unwrap_or_default();
1402                if let Some(kind) = update.get("kind").and_then(Value::as_str) {
1403                    tracked.kind = kind.to_string();
1404                }
1405                if let Some(title) = update.get("title").and_then(Value::as_str) {
1406                    tracked.title = title.to_string();
1407                }
1408                if update.get("kind").is_some()
1409                    || update.get("rawInput").is_some()
1410                    || update.get("locations").is_some()
1411                {
1412                    tracked.subject = tool_call_subject(&tracked.kind, &tracked.title, &update);
1413                }
1414                self.track_tool_call(id.clone(), tracked.clone())?;
1415                match status.as_str() {
1416                    // Terminal statuses surface as first-class ToolResult
1417                    // events; progress updates (pending/in_progress) are
1418                    // transcript-only. `failed` is a normal failure, NOT a
1419                    // denial (mirrors backend_codex) — denials are
1420                    // synthesized at the permission seam.
1421                    "completed" | "failed" => {
1422                        self.tool_calls.remove(&id);
1423                        self.queue.push_back(AgentEvent::ToolResult {
1424                            tool: Some(tracked.kind.clone()),
1425                            denied: false,
1426                            summary: tool_result_summary(&update, &tracked, &status),
1427                            raw,
1428                        });
1429                    }
1430                    _ => self.queue.push_back(AgentEvent::Other { raw }),
1431                }
1432            }
1433            Some("usage_update") => {
1434                // Remembered for the terminal Result's cost; the update
1435                // itself is transcript-only (no AgentEvent kind for
1436                // mid-stream usage, and `used`/`size` are context-window
1437                // state, not a billable token split).
1438                self.last_usage = Some(update.clone());
1439                self.queue.push_back(AgentEvent::Other { raw });
1440            }
1441            _ => self.queue.push_back(AgentEvent::Other { raw }),
1442        }
1443        Ok(())
1444    }
1445
1446    /// Answer one `session/request_permission` at the seam; a refusal is
1447    /// synthesized into a `denied` ToolResult so the guardrail firing is a
1448    /// first-class event (the peer may report nothing itself).
1449    async fn handle_permission_request(
1450        &mut self,
1451        id: Value,
1452        params: &Value,
1453        raw: Value,
1454    ) -> Result<()> {
1455        if self.acp_session_id.is_none()
1456            || params.get("sessionId").and_then(Value::as_str) != self.acp_session_id.as_deref()
1457        {
1458            self.write_message(json!({
1459                "jsonrpc": "2.0", "id": id,
1460                "result": { "outcome": { "outcome": "cancelled" } },
1461            }))
1462            .await?;
1463            return Err(EngineError::Backend(
1464                "acp permission has a foreign or missing sessionId".into(),
1465            ));
1466        }
1467        let call_update = params.get("toolCall").cloned().unwrap_or(Value::Null);
1468        let call_id = call_update
1469            .get("toolCallId")
1470            .and_then(Value::as_str)
1471            .filter(|id| !id.trim().is_empty() && id.len() <= 256);
1472        let missing_id = call_id.is_none();
1473        let call_id = call_id.unwrap_or_default().to_string();
1474        // Merge what the request carries with what the tool_call/update
1475        // stream already told us about this call.
1476        let mut info = self.tool_calls.get(&call_id).cloned().unwrap_or_default();
1477        if let Some(kind) = call_update.get("kind").and_then(Value::as_str) {
1478            info.kind = kind.to_string();
1479        }
1480        if info.kind.is_empty() {
1481            info.kind = "other".to_string();
1482        }
1483        if let Some(title) = call_update.get("title").and_then(Value::as_str) {
1484            info.title = title.to_string();
1485        }
1486        if info.subject.is_empty()
1487            || call_update.get("kind").is_some()
1488            || call_update.get("rawInput").is_some()
1489            || call_update.get("locations").is_some()
1490        {
1491            info.subject = tool_call_subject(&info.kind, &info.title, &call_update);
1492        }
1493        if !missing_id {
1494            self.track_tool_call(call_id.clone(), info.clone())?;
1495        }
1496
1497        let decision = if missing_id {
1498            PermissionDecision::Deny("tool call has no valid action identity (toolCallId)".into())
1499        } else {
1500            match decide_permission(&self.spec, &info) {
1501                PermissionDecision::Allow
1502                    if !call_update.get("rawInput").is_some_and(Value::is_object)
1503                        || info.subject.trim().is_empty() =>
1504                {
1505                    PermissionDecision::Deny(
1506                        "the complete action is unavailable for one-call consent".into(),
1507                    )
1508                }
1509                decision => decision,
1510            }
1511        };
1512        let peer_id = serde_json::to_string(&id)?;
1513        if self.pending_permissions.len() >= crate::live_permission::MAX_PENDING
1514            || self.seen_permission_ids.len() >= 1024
1515            || !self.seen_permission_ids.insert(peer_id)
1516            || self
1517                .pending_permissions
1518                .values()
1519                .any(|p| p.proposal.tool_call_id == call_id)
1520        {
1521            return Err(EngineError::Backend(
1522                "duplicate or excessive ACP permission requests".into(),
1523            ));
1524        }
1525        let options = params
1526            .get("options")
1527            .and_then(Value::as_array)
1528            .cloned()
1529            .unwrap_or_default();
1530        let mut action = call_update;
1531        if let Some(fields) = action.as_object_mut() {
1532            fields.insert("kind".into(), Value::String(info.kind));
1533        }
1534        let now = chrono::Utc::now();
1535        let mut proposal = crate::live_permission::Proposal {
1536            id: format!("permission-{}", uuid::Uuid::new_v4()),
1537            engine_session_id: self.session_id.clone(),
1538            peer_session_id: self.acp_session_id.clone().expect("validated above"),
1539            peer_request_id: id,
1540            tool_call_id: call_id,
1541            action_digest: crate::live_permission::digest(&action)?,
1542            options_digest: crate::live_permission::digest(&options)?,
1543            action,
1544            options,
1545            observed_at: now,
1546            deadline: now + chrono::Duration::seconds(crate::live_permission::REQUEST_TTL_SECS),
1547            prohibition: match decision {
1548                PermissionDecision::Deny(reason) => Some(reason),
1549                PermissionDecision::Allow => None,
1550            },
1551        };
1552        if proposal.option(true).is_none() && proposal.prohibition.is_none() {
1553            proposal.prohibition = Some("no unique certified allow_once option was offered".into());
1554        }
1555        proposal.validate()?;
1556        self.pending_permissions.insert(
1557            proposal.id.clone(),
1558            PendingPermission {
1559                proposal: proposal.clone(),
1560                expires_at: tokio::time::Instant::now()
1561                    + std::time::Duration::from_secs(
1562                        crate::live_permission::REQUEST_TTL_SECS as u64,
1563                    ),
1564            },
1565        );
1566        self.queue.push_back(AgentEvent::PermissionRequested {
1567            proposal: Box::new(proposal),
1568            raw,
1569        });
1570        Ok(())
1571    }
1572
1573    async fn answer_permission(&mut self, answer: crate::live_permission::Answer) -> Result<()> {
1574        // The biased read may have consumed a fragment before yielding. Do
1575        // not authorize against an earlier description until that frame has
1576        // been classified. The existing permission deadline bounds this wait.
1577        if self.lines.has_partial_line() {
1578            self.deferred_permission_answer = Some(answer);
1579            return Ok(());
1580        }
1581        let Some(pending) = self.pending_permissions.get(&answer.proposal.id) else {
1582            return Err(EngineError::Backend(
1583                "permission response names no live request".into(),
1584            ));
1585        };
1586        if pending.proposal != answer.proposal
1587            || chrono::Utc::now() >= pending.proposal.deadline
1588            || tokio::time::Instant::now() >= pending.expires_at
1589            || (answer.allow
1590                && (pending.proposal.prohibition.is_some()
1591                    || pending.proposal.option(true).is_none()))
1592        {
1593            return Err(EngineError::Backend(
1594                "stale or prohibited permission response".into(),
1595            ));
1596        }
1597        let proposal = self
1598            .pending_permissions
1599            .remove(&answer.proposal.id)
1600            .expect("checked above")
1601            .proposal;
1602        let decision = if answer.allow {
1603            PermissionDecision::Allow
1604        } else {
1605            PermissionDecision::Deny("one-call consent refused".into())
1606        };
1607        let result = permission_response(&decision, &proposal.options);
1608        let sent = self
1609            .write_message(json!({
1610                "jsonrpc": "2.0", "id": proposal.peer_request_id, "result": result,
1611            }))
1612            .await;
1613        let delivery = if sent.is_ok() {
1614            crate::live_permission::Delivery::Sent
1615        } else {
1616            crate::live_permission::Delivery::Uncertain
1617        };
1618        self.queue.push_back(AgentEvent::PermissionResponded {
1619            request_id: proposal.id.clone(),
1620            delivery: delivery.clone(),
1621            raw: json!({"permissionResponse":proposal.id,"delivery":delivery}),
1622        });
1623        sent
1624    }
1625
1626    /// Synthesize the terminal `Result` from a `session/prompt` response
1627    /// (see module docs): text stitched from the turn's last assistant
1628    /// message, `is_error` ⇔ `stopReason != "end_turn"`, cost only when the
1629    /// peer reported a USD amount — absent data stays absent.
1630    ///
1631    /// WHY non-`end_turn` is an error, not just `refusal` (12th-pass review):
1632    /// ACP v1's stop reasons are `end_turn` (natural completion), `refusal`,
1633    /// `max_tokens`, `max_turn_requests`, and `cancelled`. Mapping only
1634    /// `refusal` to `is_error` let a turn cut short by `max_tokens`/
1635    /// `max_turn_requests` — or answered `cancelled`, or carrying a missing
1636    /// or unrecognized reason — surface as a SUCCESSFUL result, so a
1637    /// truncated validator report could pass validation. Fail-closed is the
1638    /// only honest mapping: anything but a natural completion is an error,
1639    /// and the reason string rides the raw payload so the failure is
1640    /// diagnosable (`null` when the peer omitted the field entirely).
1641    fn synthesize_result(&mut self, outcome: RpcOutcome, request_id: u64) {
1642        let cumulative_cost_usd = self
1643            .last_usage
1644            .as_ref()
1645            .and_then(|u| u.get("cost"))
1646            .filter(|cost| {
1647                cost.get("currency").and_then(Value::as_str) == Some("USD")
1648                    && cost.get("amount").and_then(Value::as_f64).is_some()
1649            })
1650            .and_then(|cost| cost.get("amount").and_then(Value::as_f64))
1651            .filter(|amount| amount.is_finite() && *amount >= 0.0);
1652        let last_cost_usd = match (
1653            self.saw_result,
1654            self.previous_cost_total,
1655            cumulative_cost_usd,
1656        ) {
1657            (false, _, total) => total,
1658            (true, Some(previous), Some(total)) if total >= previous => Some(total - previous),
1659            _ => None,
1660        };
1661        self.previous_cost_total = cumulative_cost_usd;
1662        let (text, is_error, raw) = match outcome {
1663            RpcOutcome::Result(result) => {
1664                let stop_reason = result.get("stopReason").and_then(Value::as_str);
1665                (
1666                    std::mem::take(&mut self.message_text),
1667                    stop_reason != Some("end_turn"),
1668                    json!({
1669                        "promptResponse": result,
1670                        "usageUpdate": self.last_usage,
1671                        "costScope": "turn_delta_from_reported_session_total",
1672                        "synthesizedBy": "backend_acp",
1673                        // The classification input, verbatim: exactly what the
1674                        // peer sent, `null` when it sent nothing — the raw
1675                        // payload always explains WHY a non-end_turn failed.
1676                        "stopReason": stop_reason,
1677                    }),
1678                )
1679            }
1680            RpcOutcome::Error(error) => (
1681                format!("acp session/prompt failed: {error}"),
1682                true,
1683                json!({
1684                    "promptError": error,
1685                    "requestId": request_id,
1686                    "synthesizedBy": "backend_acp",
1687                }),
1688            ),
1689        };
1690        let event = AgentEvent::Result {
1691            text,
1692            is_error,
1693            usage: TokenUsage::default(),
1694            cost_usd: last_cost_usd,
1695            num_turns: Some(1),
1696            raw,
1697        };
1698        self.observe(&event);
1699        self.queue.push_back(event);
1700    }
1701
1702    fn observe(&mut self, event: &AgentEvent) {
1703        if let AgentEvent::Result { .. } = event {
1704            self.saw_result = true;
1705        }
1706    }
1707
1708    /// Kill only while the leader is still owned. Never signal a cached PID
1709    /// after reaping; the same-group cleanup precedes Child::wait.
1710    async fn kill_child(&mut self) {
1711        self.stdin.take();
1712        if self.child_status.is_none() {
1713            kill_owned_peer(
1714                &mut self.child,
1715                #[cfg(windows)]
1716                self.job.as_ref(),
1717            );
1718            if let Ok(Ok(status)) = tokio::time::timeout(CLEANUP_TIMEOUT, self.child.wait()).await {
1719                self.child_status = Some(status);
1720            }
1721        }
1722        #[cfg(any(target_os = "macos", target_os = "linux"))]
1723        if let Some(container) = self.container.as_mut() {
1724            if let Err(error) = container.remove().await {
1725                self.cleanup_failure
1726                    .get_or_insert("container removal failed");
1727                tracing::error!(%error, "ACP cleanup requires recovery");
1728            }
1729        }
1730        if let Some(profile) = self.profile_home.as_mut() {
1731            if let Err(error) = profile.close() {
1732                self.cleanup_failure
1733                    .get_or_insert("private credential home cleanup failed");
1734                tracing::error!(%error, "ACP private home requires recovery");
1735            }
1736        }
1737        self.finish_stderr().await;
1738    }
1739
1740    async fn finish_stderr(&mut self) {
1741        if let Some(mut task) = self.stderr_task.take() {
1742            match tokio::time::timeout(CLEANUP_TIMEOUT, &mut task).await {
1743                Ok(Ok(())) => {}
1744                Ok(Err(_)) => {
1745                    self.cleanup_failure
1746                        .get_or_insert("stderr drain task failed");
1747                }
1748                Err(_) => {
1749                    self.cleanup_failure.get_or_insert("stderr drain timed out");
1750                    task.abort();
1751                    let _ = task.await;
1752                }
1753            }
1754        }
1755    }
1756
1757    /// A single-shot ACP turn ends with its response, even if the adapter is
1758    /// a long-lived server. Give it a short stdin-EOF grace, then terminate
1759    /// our owned process group. Contained peers terminate under the trusted
1760    /// supervisor; the Docker client's exit status must arrive before success.
1761    async fn finish_session(&mut self) {
1762        self.stdin.take();
1763        #[cfg(any(target_os = "macos", target_os = "linux"))]
1764        let contained = self.container.is_some();
1765        #[cfg(not(any(target_os = "macos", target_os = "linux")))]
1766        let contained = false;
1767        let completion_timeout = if contained {
1768            CONTAINER_COMPLETION_TIMEOUT
1769        } else {
1770            COMPLETION_GRACE
1771        };
1772        let mut forced = false;
1773        if self.child_status.is_none() {
1774            match tokio::time::timeout(completion_timeout, wait_for_peer_exit(&mut self.child))
1775                .await
1776            {
1777                Ok(Ok(())) => {}
1778                Ok(Err(error)) => {
1779                    self.kill_child().await;
1780                    self.exit = Some(SessionExit::Failed(format!(
1781                        "acp process observation failed: {error}"
1782                    )));
1783                    return;
1784                }
1785                Err(_) if contained => {
1786                    self.cleanup_failure
1787                        .get_or_insert("container completion deadline expired");
1788                }
1789                Err(_) => forced = true,
1790            }
1791            self.kill_child().await;
1792        } else {
1793            self.kill_child().await;
1794        }
1795        if let Some(cause) = self.cleanup_failure {
1796            self.exit = Some(SessionExit::Failed(format!(
1797                "acp cleanup could not be confirmed: {cause}"
1798            )));
1799            return;
1800        }
1801        self.exit = Some(match self.child_status {
1802            Some(status) if self.saw_result && (status.success() || forced) => {
1803                SessionExit::Completed
1804            }
1805            Some(status) => SessionExit::Failed(format!(
1806                "acp agent exited with {status}{}; stderr tail: {}",
1807                if self.saw_result {
1808                    ""
1809                } else {
1810                    " without answering session/prompt"
1811                },
1812                self.stderr_tail(),
1813            )),
1814            None => SessionExit::Failed("acp process did not reap within cleanup deadline".into()),
1815        });
1816    }
1817
1818    fn stderr_tail(&self) -> String {
1819        let captured = self
1820            .stderr_buf
1821            .lock()
1822            .map(|guard| guard.clone())
1823            .unwrap_or_default();
1824        let captured = self
1825            .profile_home
1826            .as_ref()
1827            .map_or_else(|| captured.clone(), |p| p.scrub(captured.clone()));
1828        last_chars(captured.trim_end(), STDERR_TAIL_CHARS)
1829    }
1830}
1831
1832#[async_trait::async_trait]
1833impl AgentSession for AcpSession {
1834    fn permission_responder(&self) -> Option<crate::live_permission::PermissionResponder> {
1835        Some(self.permission_responder.clone())
1836    }
1837
1838    fn session_id(&self) -> String {
1839        self.session_id.clone()
1840    }
1841
1842    async fn next_event(&mut self) -> Result<Option<AgentEvent>> {
1843        loop {
1844            if let Some(event) = self.queue.pop_front() {
1845                return Ok(Some(event));
1846            }
1847            if self.exit.is_some() {
1848                return Ok(None);
1849            }
1850            if self.saw_result && matches!(self.spec.prompt, PromptMode::SingleShot(_)) {
1851                self.finish_session().await;
1852                return Ok(None);
1853            }
1854            let frame = match self.read_frame().await {
1855                Ok(Some(frame)) => frame,
1856                Ok(None) => {
1857                    self.finish_session().await;
1858                    return Ok(None);
1859                }
1860                Err(e) => {
1861                    self.kill_child().await;
1862                    self.exit = Some(SessionExit::Failed(e.to_string()));
1863                    return Ok(None);
1864                }
1865            };
1866            if let Err(e) = self.handle_frame(frame).await {
1867                // A transport error mid-session (stdin write failed — the
1868                // peer is gone): fail the session honestly rather than hang.
1869                self.kill_child().await;
1870                self.exit = Some(SessionExit::Failed(e.to_string()));
1871                // Preserve the rejected frame's diagnostic event.
1872                continue;
1873            }
1874        }
1875    }
1876
1877    async fn send_user_message(&mut self, text: &str) -> Result<()> {
1878        if self.exit.is_some() {
1879            return Err(EngineError::Backend(
1880                "acp session is closed; cannot send further messages".to_string(),
1881            ));
1882        }
1883        if self.acp_session_id.is_none() {
1884            return Err(EngineError::Backend(
1885                "acp session not established yet; cannot send a message".to_string(),
1886            ));
1887        }
1888        if matches!(self.spec.prompt, PromptMode::SingleShot(_)) || self.prompt_request_id.is_some()
1889        {
1890            return Err(EngineError::Backend(
1891                "acp session cannot accept an overlapping or single-shot follow-up".into(),
1892            ));
1893        }
1894        self.send_prompt(text).await
1895    }
1896
1897    async fn abort(&mut self) -> Result<()> {
1898        if self.exit.is_some() {
1899            return Ok(());
1900        }
1901        // Best-effort graceful cancel first (the peer MAY stop its turn
1902        // cleanly and answer the prompt with stopReason "cancelled"), then
1903        // the house tree-kill regardless — abort must never depend on the
1904        // peer honoring the notification.
1905        if let Some(acp_session_id) = self.acp_session_id.clone() {
1906            let _ = tokio::time::timeout(
1907                CANCEL_WRITE_TIMEOUT,
1908                self.write_message(json!({
1909                    "jsonrpc": "2.0",
1910                    "method": method::SESSION_CANCEL,
1911                    "params": { "sessionId": acp_session_id },
1912                })),
1913            )
1914            .await;
1915        }
1916        self.kill_child().await;
1917        #[cfg(any(target_os = "macos", target_os = "linux"))]
1918        let container_cleanup_failed = self.container.is_some()
1919            && (self.cleanup_failure.is_some() || self.child_status.is_none());
1920        #[cfg(not(any(target_os = "macos", target_os = "linux")))]
1921        let container_cleanup_failed = false;
1922        if container_cleanup_failed {
1923            let cause = self
1924                .cleanup_failure
1925                .unwrap_or("host process did not reap within cleanup deadline");
1926            let message = format!("acp abort cleanup could not be confirmed: {cause}");
1927            self.exit = Some(SessionExit::Failed(message.clone()));
1928            return Err(EngineError::Backend(message));
1929        }
1930        self.exit = Some(SessionExit::Aborted);
1931        Ok(())
1932    }
1933
1934    fn exit_status(&self) -> Option<SessionExit> {
1935        self.exit.clone()
1936    }
1937}
1938
1939// ---------------------------------------------------------------------------
1940
1941#[cfg(test)]
1942mod tests {
1943    use super::*;
1944
1945    #[test]
1946    fn backend_acp_classify_distinguishes_response_request_notification() {
1947        let response =
1948            classify_line(r#"{"jsonrpc":"2.0","id":3,"result":{"stopReason":"end_turn"}}"#);
1949        assert!(matches!(
1950            response,
1951            Frame::Response {
1952                id: 3,
1953                outcome: RpcOutcome::Result(_)
1954            }
1955        ));
1956        let error =
1957            classify_line(r#"{"jsonrpc":"2.0","id":4,"error":{"code":-32603,"message":"boom"}}"#);
1958        assert!(matches!(
1959            error,
1960            Frame::Response {
1961                id: 4,
1962                outcome: RpcOutcome::Error(_)
1963            }
1964        ));
1965        let request = classify_line(
1966            r#"{"jsonrpc":"2.0","id":100,"method":"session/request_permission","params":{}}"#,
1967        );
1968        assert!(matches!(
1969            request,
1970            Frame::Request { ref method, .. } if method == "session/request_permission"
1971        ));
1972        let notification = classify_line(
1973            r#"{"jsonrpc":"2.0","method":"session/update","params":{"sessionId":"s","update":{"sessionUpdate":"plan"}}}"#,
1974        );
1975        assert!(matches!(
1976            notification,
1977            Frame::Notification { ref method, .. } if method == "session/update"
1978        ));
1979        // Classification retains the diagnostic; the session fails on it.
1980        let torn = classify_line(r#"{"jsonrpc":"2.0","method":"session/upda"#);
1981        assert!(matches!(torn, Frame::Unrecognized(_)));
1982    }
1983
1984    #[test]
1985    fn backend_acp_wildcard_match_anchors_like_a_shell_glob() {
1986        assert!(wildcard_match("git push*", "git push origin main"));
1987        assert!(wildcard_match("git push*", "git push"));
1988        assert!(!wildcard_match("git push*", "git pull"));
1989        assert!(wildcard_match("*", "anything"));
1990        assert!(wildcard_match(
1991            "cargo * --workspace",
1992            "cargo test --workspace"
1993        ));
1994        assert!(!wildcard_match(
1995            "cargo * --workspace",
1996            "cargo test --package x"
1997        ));
1998        assert!(wildcard_match("*/etc/passwd", "/etc/passwd"));
1999        assert!(!wildcard_match("*/etc/passwd", "/etc/passwd.bak"));
2000    }
2001
2002    #[test]
2003    fn backend_acp_pattern_matches_maps_claude_names_to_acp_kinds() {
2004        assert!(pattern_matches(
2005            "Bash(git push*)",
2006            "execute",
2007            "git push origin main"
2008        ));
2009        assert!(!pattern_matches("Bash(git push*)", "execute", "git pull"));
2010        assert!(!pattern_matches(
2011            "Bash(git push*)",
2012            "edit",
2013            "git push origin main"
2014        ));
2015        assert!(pattern_matches("Write", "edit", "/repo/src/main.rs"));
2016        assert!(!pattern_matches("Write", "read", "/repo/src/main.rs"));
2017        assert!(pattern_matches("Edit(/etc/*)", "edit", "/etc/hosts"));
2018        assert!(pattern_matches("Read", "read", "/anywhere"));
2019        // Unknown tool names match nothing (deny-only mapping; never widens).
2020        assert!(!pattern_matches("NotAClaudeTool(*)", "execute", "x"));
2021    }
2022
2023    fn spec_with(writable: bool, disallowed: &[&str]) -> SessionSpec {
2024        SessionSpec {
2025            cwd: PathBuf::from("."),
2026            prompt: PromptMode::SingleShot("do the thing".to_string()),
2027            append_system_prompt: None,
2028            model: "acp-model".to_string(),
2029            effort: "high".to_string(),
2030            session_id: "sess-1".to_string(),
2031            resume: None,
2032            permission_mode: None,
2033            allowed_tools: vec![],
2034            disallowed_tools: disallowed.iter().map(|s| s.to_string()).collect(),
2035            tools: vec![],
2036            writable,
2037            settings_json: None,
2038            json_schema: None,
2039            max_budget_usd: None,
2040            max_turns: None,
2041            env: Default::default(),
2042            sandbox: None,
2043            hook_status: None,
2044        }
2045    }
2046
2047    #[test]
2048    fn backend_acp_permission_denies_disallowed_and_mutating_kinds() {
2049        let spec = spec_with(true, &["Bash(git push*)"]);
2050        let push = ToolCallInfo {
2051            kind: "execute".to_string(),
2052            title: "git push origin main".to_string(),
2053            subject: "git push origin main".to_string(),
2054        };
2055        assert!(matches!(
2056            decide_permission(&spec, &push),
2057            PermissionDecision::Deny(reason) if reason.contains("Bash(git push*)")
2058        ));
2059        let test = ToolCallInfo {
2060            kind: "execute".to_string(),
2061            title: "cargo test".to_string(),
2062            subject: "cargo test".to_string(),
2063        };
2064        assert_eq!(decide_permission(&spec, &test), PermissionDecision::Allow);
2065
2066        // Read-only posture: mutating kinds refused, execute/read allowed.
2067        let ro = spec_with(false, &[]);
2068        let edit = ToolCallInfo {
2069            kind: "edit".to_string(),
2070            title: "write src/main.rs".to_string(),
2071            subject: "/repo/src/main.rs".to_string(),
2072        };
2073        assert!(matches!(
2074            decide_permission(&ro, &edit),
2075            PermissionDecision::Deny(reason) if reason.contains("writable: false")
2076        ));
2077        assert_eq!(decide_permission(&ro, &test), PermissionDecision::Allow);
2078        let read = ToolCallInfo {
2079            kind: "read".to_string(),
2080            title: "read src/main.rs".to_string(),
2081            subject: "/repo/src/main.rs".to_string(),
2082        };
2083        assert_eq!(decide_permission(&ro, &read), PermissionDecision::Allow);
2084    }
2085
2086    /// The permission seam used to fail OPEN when the peer omitted the
2087    /// subject: `wildcard_match("git push*", "")` is false, so no deny fired
2088    /// and the decision was Allow. Every glob-carrying deny rule was
2089    /// bypassable that way, by a peer that need not even be hostile.
2090    #[test]
2091    fn backend_acp_permission_denies_when_the_subject_is_missing() {
2092        let spec = spec_with(true, &["Bash(git push*)"]);
2093        let no_subject = ToolCallInfo {
2094            kind: "execute".to_string(),
2095            title: String::new(),
2096            subject: String::new(),
2097        };
2098        assert!(
2099            matches!(
2100                decide_permission(&spec, &no_subject),
2101                PermissionDecision::Deny(ref reason)
2102                    if reason.contains("no subject") && reason.contains("Bash(git push*)")
2103            ),
2104            "got {:?}",
2105            decide_permission(&spec, &no_subject)
2106        );
2107
2108        // Whitespace is no subject either.
2109        let blank_subject = ToolCallInfo {
2110            subject: "   ".to_string(),
2111            ..no_subject.clone()
2112        };
2113        assert!(matches!(
2114            decide_permission(&spec, &blank_subject),
2115            PermissionDecision::Deny(_)
2116        ));
2117
2118        // Precise: a deny list that does not cover this call's kind is not
2119        // made to fire by a missing subject.
2120        let read_no_subject = ToolCallInfo {
2121            kind: "read".to_string(),
2122            ..no_subject.clone()
2123        };
2124        assert_eq!(
2125            decide_permission(&spec_with(true, &["Bash(git push*)"]), &read_no_subject),
2126            PermissionDecision::Allow
2127        );
2128    }
2129
2130    /// Read-only posture: `MUTATING_KINDS` cannot classify a call whose kind
2131    /// the peer omitted (it arrives as `"other"`), so the containment claim
2132    /// cannot be checked and the call is refused.
2133    #[test]
2134    fn backend_acp_read_only_denies_an_unclassifiable_kind() {
2135        let ro = spec_with(false, &[]);
2136        for kind in ["", "other"] {
2137            let call = ToolCallInfo {
2138                kind: kind.to_string(),
2139                title: "do something".to_string(),
2140                subject: "/repo/src/main.rs".to_string(),
2141            };
2142            assert!(
2143                matches!(
2144                    decide_permission(&ro, &call),
2145                    PermissionDecision::Deny(ref reason) if reason.contains("kind")
2146                ),
2147                "kind {kind:?} got {:?}",
2148                decide_permission(&ro, &call)
2149            );
2150        }
2151        // Unknown kinds cannot bypass deny rules in writable sessions either.
2152        let writable = spec_with(true, &[]);
2153        let other = ToolCallInfo {
2154            kind: "other".to_string(),
2155            title: "think".to_string(),
2156            subject: "think".to_string(),
2157        };
2158        assert!(matches!(
2159            decide_permission(&writable, &other),
2160            PermissionDecision::Deny(_)
2161        ));
2162    }
2163
2164    #[test]
2165    fn backend_acp_permission_response_picks_options_or_cancels() {
2166        let options = vec![
2167            json!({ "optionId": "allow-1", "name": "Allow", "kind": "allow_once" }),
2168            json!({ "optionId": "reject-1", "name": "Reject", "kind": "reject_once" }),
2169        ];
2170        let allow = permission_response(&PermissionDecision::Allow, &options);
2171        assert_eq!(allow["outcome"]["optionId"], json!("allow-1"));
2172        let deny = permission_response(&PermissionDecision::Deny("nope".to_string()), &options);
2173        assert_eq!(deny["outcome"]["optionId"], json!("reject-1"));
2174        // No reject option offered: refusal degrades to the cancelled
2175        // outcome — never an invented selection.
2176        let deny_no_reject = permission_response(
2177            &PermissionDecision::Deny("nope".to_string()),
2178            &[options[0].clone()],
2179        );
2180        assert_eq!(deny_no_reject["outcome"]["outcome"], json!("cancelled"));
2181        for (kind, decision) in [
2182            ("allow_once", PermissionDecision::Allow),
2183            ("reject_once", PermissionDecision::Deny("refused".into())),
2184        ] {
2185            let ambiguous = vec![
2186                json!({"optionId":"a","kind":kind}),
2187                json!({"optionId":"b","kind":kind}),
2188            ];
2189            assert_eq!(
2190                permission_response(&decision, &ambiguous)["outcome"]["outcome"],
2191                "cancelled"
2192            );
2193        }
2194    }
2195}