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