1use 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
97const SUMMARY_MAX_CHARS: usize = 200;
99const STDERR_TAIL_CHARS: usize = 500;
101
102const 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
114const 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
131const MUTATING_KINDS: &[&str] = &["edit", "delete", "move"];
135
136#[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#[derive(Debug, Clone, PartialEq, Eq)]
157enum PermissionDecision {
158 Allow,
159 Deny(String),
161}
162
163#[derive(Debug, Clone, Default)]
168struct ToolCallInfo {
169 kind: String,
170 title: String,
171 subject: Vec<String>,
174}
175
176fn 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
225fn 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
246fn wildcard_match(pattern: &str, text: &str) -> bool {
250 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
273fn 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
282fn 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
294fn 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
305fn 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
392fn 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
426fn 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
439fn 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
446fn 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#[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 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 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 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 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 .stdin(Stdio::piped())
618 .stdout(Stdio::piped())
619 .stderr(Stdio::piped())
620 .kill_on_drop(true);
621 #[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 #[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 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 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 .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
766struct 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 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 queue: VecDeque<AgentEvent>,
804 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 message_text: String,
818 message_id: Option<String>,
819 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 if child.id().is_some() {
859 let _ = child.start_kill();
860 }
861}
862
863impl AcpSession {
864 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 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 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 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 self.protocol.complete_prompt(request.id);
969 return Err(error);
970 }
971 self.message_text.clear();
974 self.message_id = None;
975 self.last_usage = None;
976 Ok(())
977 }
978
979 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 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 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 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 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 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(¶ms, 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, ¶ms, raw).await?;
1187 } else {
1188 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 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 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 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 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 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 "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 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 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 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 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 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 "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 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 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 self.kill_child().await;
1925 self.exit = Some(SessionExit::Failed(e.to_string()));
1926 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 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#[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 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 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 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 #[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 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 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 #[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 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 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}