1use crate::backend::{
68 AgentBackend, AgentEvent, AgentSession, PromptMode, SessionExit, SessionSpec,
69};
70#[cfg(windows)]
71use crate::backend_claude::win_job;
72use crate::error::{EngineError, Result};
73use crate::stream_bounds::{drain_to_tail, BoundedLines, STDERR_TAIL_CAP, STDOUT_LINE_CAP};
74use crate::types::TokenUsage;
75use serde_json::{json, Value};
76use std::collections::{HashMap, VecDeque};
77use std::path::PathBuf;
78use std::process::{ExitStatus, Stdio};
79use std::sync::{Arc, Mutex};
80use tokio::process::{Child, ChildStdin, ChildStdout};
81use tokio::task::JoinHandle;
82
83const SUMMARY_MAX_CHARS: usize = 200;
85const STDERR_TAIL_CHARS: usize = 500;
87
88const ACP_PROTOCOL_VERSION: u64 = 1;
91
92const HANDSHAKE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(30);
98const CLEANUP_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(1);
99const COMPLETION_GRACE: std::time::Duration = std::time::Duration::from_millis(200);
100const CONTAINER_COMPLETION_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
101const WRITE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(3);
102const CANCEL_WRITE_TIMEOUT: std::time::Duration = std::time::Duration::from_millis(100);
103
104mod method {
106 pub(crate) const INITIALIZE: &str = "initialize";
107 pub(crate) const SESSION_NEW: &str = "session/new";
108 pub(crate) const SESSION_PROMPT: &str = "session/prompt";
109 pub(crate) const SESSION_CANCEL: &str = "session/cancel";
110 pub(crate) const SESSION_UPDATE: &str = "session/update";
111 pub(crate) const REQUEST_PERMISSION: &str = "session/request_permission";
112}
113
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)]
144enum Frame {
145 PermissionAnswer(crate::live_permission::Answer),
147 PermissionExpired,
148 Response {
150 id: u64,
151 outcome: RpcOutcome,
152 },
153 Request {
157 id: Value,
158 method: String,
159 params: Value,
160 raw: Value,
161 },
162 Notification {
164 method: String,
165 params: Value,
166 raw: Value,
167 },
168 Unrecognized(Value),
170}
171
172#[derive(Debug)]
176enum RpcOutcome {
177 Result(Value),
178 Error(Value),
179}
180
181fn classify_line(line: &str) -> Frame {
185 let value = match crate::strict_json::parse(line.as_bytes()) {
186 Ok(value) => value,
187 Err(_) => return Frame::Unrecognized(json!({ "unparsed": line })),
188 };
189 classify_value(value)
190}
191
192fn classify_value(value: Value) -> Frame {
193 let obj = match value.as_object() {
194 Some(obj) => obj,
195 None => return Frame::Unrecognized(value),
196 };
197 if obj.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
198 return Frame::Unrecognized(value);
199 }
200 let has_id = obj.contains_key("id");
201 if obj.contains_key("method") && (obj.contains_key("result") || obj.contains_key("error")) {
202 return Frame::Unrecognized(value);
203 }
204 if has_id && !(obj["id"].is_string() || obj["id"].is_i64() || obj["id"].is_u64()) {
205 return Frame::Unrecognized(value);
206 }
207 let method = obj.get("method").and_then(Value::as_str);
208 match (method, has_id) {
209 (Some(method), true) => Frame::Request {
211 id: obj.get("id").cloned().unwrap_or(Value::Null),
212 method: method.to_string(),
213 params: obj.get("params").cloned().unwrap_or(Value::Null),
214 raw: value,
215 },
216 (Some(method), false) => Frame::Notification {
218 method: method.to_string(),
219 params: obj.get("params").cloned().unwrap_or(Value::Null),
220 raw: value,
221 },
222 (None, true) => {
225 let id = obj.get("id").and_then(Value::as_u64);
226 match (id, obj.get("result"), obj.get("error")) {
227 (Some(id), Some(result), None) => Frame::Response {
228 id,
229 outcome: RpcOutcome::Result(result.clone()),
230 },
231 (Some(id), None, Some(error)) => Frame::Response {
232 id,
233 outcome: RpcOutcome::Error(error.clone()),
234 },
235 _ => Frame::Unrecognized(value),
236 }
237 }
238 (None, false) => Frame::Unrecognized(value),
239 }
240}
241
242#[derive(Debug, Clone, PartialEq, Eq)]
248enum PermissionDecision {
249 Allow,
250 Deny(String),
252}
253
254#[derive(Debug, Clone, Default)]
259struct ToolCallInfo {
260 kind: String,
261 title: String,
262 subject: String,
265}
266
267fn tool_call_subject(kind: &str, title: &str, call: &Value) -> String {
272 let raw_input = call.get("rawInput").cloned().unwrap_or(Value::Null);
273 let str_at = |value: &Value, keys: &[&str]| -> Option<String> {
274 keys.iter()
275 .find_map(|k| value.get(*k).and_then(Value::as_str).map(str::to_string))
276 };
277 if kind == "execute" {
278 return str_at(&raw_input, &["command", "cmd"]).unwrap_or_default();
280 }
281 if let Some(path) = call
282 .get("locations")
283 .and_then(Value::as_array)
284 .and_then(|locs| locs.first())
285 .and_then(|loc| loc.get("path"))
286 .and_then(Value::as_str)
287 {
288 return path.to_string();
289 }
290 if let Some(path) = str_at(&raw_input, &["path", "filePath", "file_path"]) {
291 return path;
292 }
293 if call.get("rawInput").is_some() || call.get("locations").is_some() {
294 return String::new();
295 }
296 title.to_string()
297}
298
299fn wildcard_match(pattern: &str, text: &str) -> bool {
303 let parts = pattern.split('*');
304 let anchored_start = !pattern.starts_with('*');
305 let anchored_end = !pattern.ends_with('*');
306 let mut rest = text;
307 let mut first = true;
308 for part in parts {
309 if part.is_empty() {
310 first = false;
311 continue;
312 }
313 match rest.find(part) {
314 Some(idx) if !first || !anchored_start || idx == 0 => {
315 rest = &rest[idx + part.len()..];
316 }
317 _ => return false,
318 }
319 first = false;
320 }
321 !anchored_end || rest.is_empty()
324}
325
326fn split_pattern(pattern: &str) -> (&str, &str) {
329 match pattern.split_once('(') {
330 Some((name, rest)) => (name.trim(), rest.strip_suffix(')').unwrap_or(rest)),
331 None => (pattern.trim(), "*"),
332 }
333}
334
335fn pattern_covers_kind(pattern: &str, kind: &str) -> bool {
340 let (name, _) = split_pattern(pattern);
341 TOOL_NAME_KINDS
342 .iter()
343 .find(|(n, _)| n.eq_ignore_ascii_case(name))
344 .is_some_and(|(_, kinds)| kinds.contains(&kind))
345}
346
347fn pattern_matches(pattern: &str, kind: &str, subject: &str) -> bool {
350 let (_, glob) = split_pattern(pattern);
351 if pattern_covers_kind(pattern, kind) {
352 wildcard_match(glob, subject)
353 } else {
354 false
355 }
356}
357
358fn decide_permission(spec: &SessionSpec, call: &ToolCallInfo) -> PermissionDecision {
370 if call.kind == "switch_mode" {
371 return PermissionDecision::Deny(
372 "agent mode changes require an explicit human decision".into(),
373 );
374 }
375 if ![
376 "read",
377 "edit",
378 "delete",
379 "move",
380 "search",
381 "execute",
382 "think",
383 "fetch",
384 "switch_mode",
385 ]
386 .contains(&call.kind.as_str())
387 {
388 return PermissionDecision::Deny("tool call carries an unknown or missing ACP kind".into());
389 }
390 if call.subject.trim().is_empty() {
391 if let Some(pattern) = spec
392 .disallowed_tools
393 .iter()
394 .find(|pattern| pattern_covers_kind(pattern, &call.kind))
395 {
396 return PermissionDecision::Deny(format!(
397 "tool call carries no subject (no rawInput command or path, no locations[0].path, \
398 no title), so deny pattern {pattern:?} for ACP kind {:?} cannot be evaluated",
399 call.kind
400 ));
401 }
402 }
403 for pattern in &spec.disallowed_tools {
404 if pattern_matches(pattern, &call.kind, &call.subject) {
405 return PermissionDecision::Deny(format!(
406 "matches SessionSpec.disallowed_tools pattern {pattern:?}"
407 ));
408 }
409 }
410 if !spec.writable {
411 let kind = call.kind.trim();
412 if kind.is_empty() || kind == "other" {
413 return PermissionDecision::Deny(format!(
414 "read-only session (writable: false): tool call carries no usable ACP kind \
415 ({kind:?}), so whether it mutates the filesystem cannot be decided"
416 ));
417 }
418 if MUTATING_KINDS.contains(&kind) {
419 return PermissionDecision::Deny(format!(
420 "read-only session (writable: false): ACP kind {:?} mutates the filesystem",
421 call.kind
422 ));
423 }
424 }
425 PermissionDecision::Allow
426}
427
428fn permission_response(decision: &PermissionDecision, options: &[Value]) -> Value {
431 let mut ids = std::collections::HashSet::new();
432 let valid = !options.is_empty()
433 && options.len() <= 32
434 && options.iter().all(|option| {
435 option
436 .get("optionId")
437 .and_then(Value::as_str)
438 .filter(|id| !id.trim().is_empty() && id.len() <= 256)
439 .is_some_and(|id| ids.insert(id))
440 });
441 let kind = match decision {
442 PermissionDecision::Allow => "allow_once",
443 PermissionDecision::Deny(_) => "reject_once",
444 };
445 let selected = valid
446 .then(|| {
447 let mut matching = options
448 .iter()
449 .filter(|option| option.get("kind").and_then(Value::as_str) == Some(kind));
450 let first = matching.next()?;
451 matching.next().is_none().then_some(first)
452 })
453 .flatten();
454 match selected {
455 Some(option) => {
456 json!({ "outcome": { "outcome": "selected", "optionId": option["optionId"] } })
457 }
458 None => json!({ "outcome": { "outcome": "cancelled" } }),
459 }
460}
461
462fn truncate_chars(text: &str, max: usize) -> String {
468 if text.chars().count() <= max {
469 text.to_string()
470 } else {
471 text.chars().take(max).collect()
472 }
473}
474
475fn last_chars(text: &str, max: usize) -> String {
477 let chars: Vec<char> = text.chars().collect();
478 let start = chars.len().saturating_sub(max);
479 chars[start..].iter().collect()
480}
481
482fn tool_result_summary(update: &Value, tracked: &ToolCallInfo, status: &str) -> String {
485 if let Some(content) = update.get("content").and_then(Value::as_array) {
486 for item in content {
487 if item.get("type").and_then(Value::as_str) == Some("content") {
488 if let Some(text) = item
489 .get("content")
490 .and_then(|c| c.get("text"))
491 .and_then(Value::as_str)
492 {
493 return truncate_chars(text, SUMMARY_MAX_CHARS);
494 }
495 }
496 }
497 }
498 if !tracked.title.is_empty() {
499 return truncate_chars(&tracked.title, SUMMARY_MAX_CHARS);
500 }
501 status.to_string()
502}
503
504#[derive(Debug, Clone)]
515pub struct AcpBackend {
516 program: PathBuf,
517 args: Vec<String>,
518 profile: Option<crate::acp_worker::AcpWorkerProfile>,
519}
520
521impl AcpBackend {
522 pub fn new(program: impl Into<PathBuf>, args: Vec<String>) -> Self {
525 AcpBackend {
526 program: program.into(),
527 args,
528 profile: None,
529 }
530 }
531
532 pub fn for_worker(cfg: &crate::types::RoleConfig) -> Result<Self> {
535 if let Some(profile) = &cfg.acp_profile {
536 profile.validate_config(
537 crate::types::Role::Worker,
538 cfg,
539 crate::types::WorkerIsolation::Worktree,
540 )?;
541 let definition = profile.definition()?;
542 Ok(Self {
543 program: definition.program.into(),
544 args: definition.args.iter().map(|s| (*s).to_owned()).collect(),
545 profile: Some(profile.clone()),
546 })
547 } else {
548 let command = cfg.acp_command.as_ref().ok_or_else(|| {
549 EngineError::Config("ACP worker requires acpCommand or acpProfile".into())
550 })?;
551 Ok(Self::new(command, cfg.acp_args.clone()))
552 }
553 }
554
555 pub fn program(&self) -> &std::path::Path {
557 &self.program
558 }
559}
560
561fn effective_prompt(spec: &SessionSpec) -> String {
562 let text = match &spec.prompt {
563 PromptMode::SingleShot(text) | PromptMode::Streaming(text) => text,
564 };
565 match &spec.append_system_prompt {
566 Some(system) if !system.is_empty() => format!("{system}\n\n{text}"),
567 _ => text.clone(),
568 }
569}
570
571fn peer_reported_model(session: &Value) -> Option<&str> {
572 session
573 .get("configOptions")
574 .and_then(Value::as_array)
575 .and_then(|options| {
576 options.iter().find_map(|option| {
577 (option.get("category").and_then(Value::as_str) == Some("model"))
578 .then(|| option.get("currentValue").and_then(Value::as_str))
579 .flatten()
580 })
581 })
582 .or_else(|| session.get("models")?.get("currentModelId")?.as_str())
583 .filter(|model| !model.trim().is_empty() && model.len() <= 256)
584}
585
586#[async_trait::async_trait]
587impl AgentBackend for AcpBackend {
588 async fn start(&self, mut spec: SessionSpec) -> Result<Box<dyn AgentSession>> {
589 let container_requested = spec.sandbox.as_ref().is_some_and(|sandbox| {
592 cfg!(any(target_os = "macos", target_os = "linux"))
593 && sandbox.backend == crate::sandbox::SandboxBackend::Container
594 && sandbox.container.is_some()
595 });
596 if spec.sandbox.is_some() && !container_requested {
597 return Err(EngineError::Backend(
598 "acp backend: enforced containment is not certified; refusing the supplied sandbox before spawn".into(),
599 ));
600 }
601 if spec.resume.is_some() {
602 return Err(EngineError::Backend(
603 "acp backend: resume is unsupported (session/load is an optional v1 \
604 capability this backend does not negotiate)"
605 .to_string(),
606 ));
607 }
608 let profile_home = self
609 .profile
610 .as_ref()
611 .map(|profile| profile.prepare(&mut spec))
612 .transpose()?;
613 let model = spec.model.clone();
614
615 #[cfg(any(target_os = "macos", target_os = "linux"))]
616 let (container, mut command) = if container_requested {
617 let (container, command) =
618 crate::acp_container::OwnedContainer::prepare(&spec, &self.program, &self.args)
619 .await?;
620 (Some(container), command)
621 } else {
622 (None, self.native_command(&spec))
623 };
624 #[cfg(not(any(target_os = "macos", target_os = "linux")))]
625 let mut command = self.native_command(&spec);
626 command
627 .current_dir(&spec.cwd)
628 .stdin(Stdio::piped())
630 .stdout(Stdio::piped())
631 .stderr(Stdio::piped())
632 .kill_on_drop(true);
633 #[cfg(unix)]
636 command.process_group(0);
637
638 let mut child = command.spawn().map_err(|e| {
639 EngineError::Backend(format!(
640 "failed to spawn acp agent {}: {e}",
641 self.program.display()
642 ))
643 })?;
644
645 #[cfg(windows)]
647 let job = match child.raw_handle() {
648 Some(handle) => match win_job::JobHandle::create_and_assign(handle) {
649 Ok(job) => Some(job),
650 Err(e) => {
651 let _ = child.kill().await;
652 return Err(EngineError::Backend(format!(
653 "acp Job Object assignment failed: {e}"
654 )));
655 }
656 },
657 None => {
658 let _ = child.kill().await;
659 return Err(EngineError::Backend(
660 "acp child has no process handle".into(),
661 ));
662 }
663 };
664
665 let stdin = child
666 .stdin
667 .take()
668 .ok_or_else(|| EngineError::Backend("acp child has no stdin pipe".to_string()))?;
669 let stdout = child
670 .stdout
671 .take()
672 .ok_or_else(|| EngineError::Backend("acp child has no stdout pipe".to_string()))?;
673 let stderr = child
674 .stderr
675 .take()
676 .ok_or_else(|| EngineError::Backend("acp child has no stderr pipe".to_string()))?;
677
678 let stderr_buf = Arc::new(Mutex::new(String::new()));
683 let stderr_task = {
684 let buf = Arc::clone(&stderr_buf);
685 tokio::spawn(async move {
686 let tail = drain_to_tail(stderr, STDERR_TAIL_CAP).await;
687 *buf.lock().expect("stderr buffer lock") = tail;
688 })
689 };
690
691 let (permission_responder, permission_answers) =
692 crate::live_permission::PermissionResponder::channel();
693 let mut session = AcpSession {
694 session_id: spec.session_id.clone(),
695 acp_session_id: None,
696 model,
697 spec,
698 child,
699 #[cfg(any(target_os = "macos", target_os = "linux"))]
700 container,
701 profile_home,
702 #[cfg(windows)]
703 job,
704 stdin: Some(stdin),
705 child_status: None,
706 cleanup_failure: None,
707 drain_deadline: None,
708 lines: BoundedLines::new_strict(stdout),
709 stderr_buf,
710 stderr_task: Some(stderr_task),
711 queue: VecDeque::new(),
712 next_request_id: 1,
713 prompt_request_id: None,
714 tool_calls: HashMap::new(),
715 permission_responder,
716 permission_answers,
717 deferred_permission_answer: None,
718 pending_permissions: HashMap::new(),
719 seen_permission_ids: std::collections::HashSet::new(),
720 message_text: String::new(),
721 message_id: None,
722 last_usage: None,
723 previous_cost_total: None,
724 handshake_bytes: 0,
725 handshake_complete: false,
726 saw_result: false,
727 exit: None,
728 };
729
730 if let Err(e) = session.handshake().await {
734 session.kill_child().await;
735 let message = match e {
736 EngineError::Backend(message) => message,
737 other => other.to_string(),
738 };
739 return Err(EngineError::Backend(format!(
740 "{message}; startup stderr after cleanup: {}{}",
741 session.stderr_tail(),
742 session
743 .cleanup_failure
744 .map(|cause| format!("; cleanup unconfirmed: {cause}"))
745 .unwrap_or_default(),
746 )));
747 }
748 Ok(Box::new(session))
749 }
750}
751
752impl AcpBackend {
753 fn native_command(&self, spec: &SessionSpec) -> tokio::process::Command {
754 let mut command = tokio::process::Command::new(&self.program);
755 command
756 .args(&self.args)
757 .env_clear()
760 .envs(crate::agent_env::agent_session_env(
761 &spec.env,
762 &spec.session_id,
763 None,
764 ));
765 command
766 }
767}
768
769struct PendingPermission {
778 proposal: crate::live_permission::Proposal,
779 expires_at: tokio::time::Instant,
780}
781
782pub struct AcpSession {
783 session_id: String,
784 acp_session_id: Option<String>,
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 next_request_id: u64,
805 prompt_request_id: Option<u64>,
808 tool_calls: HashMap<String, ToolCallInfo>,
812 permission_responder: crate::live_permission::PermissionResponder,
813 permission_answers: tokio::sync::mpsc::Receiver<crate::live_permission::Answer>,
814 deferred_permission_answer: Option<crate::live_permission::Answer>,
815 pending_permissions: HashMap<String, PendingPermission>,
816 seen_permission_ids: std::collections::HashSet<String>,
817 message_text: String,
822 message_id: Option<String>,
823 last_usage: Option<Value>,
826 previous_cost_total: Option<f64>,
827 handshake_bytes: usize,
828 handshake_complete: bool,
829 saw_result: bool,
830 exit: Option<SessionExit>,
831}
832
833#[cfg(unix)]
834impl Drop for AcpSession {
835 fn drop(&mut self) {
836 crate::backend_claude::kill_unreaped_group(&self.child);
837 }
838}
839
840async fn wait_for_peer_exit(child: &mut Child) -> std::io::Result<()> {
841 #[cfg(any(target_os = "macos", target_os = "linux"))]
842 {
843 let pid = child
844 .id()
845 .ok_or_else(|| std::io::Error::other("acp child already reaped"))?;
846 crate::command_exec::control_leader_exited(pid).await
847 }
848 #[cfg(not(any(target_os = "macos", target_os = "linux")))]
849 {
850 child.wait().await.map(|_| ())
851 }
852}
853
854fn kill_owned_peer(child: &mut Child, #[cfg(windows)] job: Option<&win_job::JobHandle>) {
855 #[cfg(unix)]
856 crate::backend_claude::kill_unreaped_group(child);
857 #[cfg(windows)]
858 if let Some(job) = job {
859 job.kill();
860 }
861 if child.id().is_some() {
863 let _ = child.start_kill();
864 }
865}
866
867impl AcpSession {
868 async fn write_message(&mut self, message: Value) -> Result<()> {
876 use tokio::io::AsyncWriteExt;
877 let mut line = serde_json::to_string(&message)
878 .map_err(|e| EngineError::Backend(format!("failed to encode acp message: {e}")))?;
879 line.push('\n');
880 if line.len() > STDOUT_LINE_CAP {
881 return Err(EngineError::Backend(
882 "acp request exceeds the frame byte limit".into(),
883 ));
884 }
885 let stdin = self
886 .stdin
887 .as_mut()
888 .ok_or_else(|| EngineError::Backend("acp stdin is closed".into()))?;
889 tokio::time::timeout(WRITE_TIMEOUT, async {
890 stdin.write_all(line.as_bytes()).await?;
891 stdin.flush().await
892 })
893 .await
894 .map_err(|_| EngineError::Backend("acp stdin write timed out".into()))?
895 .map_err(|e| EngineError::Backend(format!("failed to write to acp agent stdin: {e}")))?;
896 Ok(())
897 }
898
899 async fn send_request(&mut self, method: &str, params: Value) -> Result<u64> {
903 let id = self.next_request_id;
904 self.next_request_id += 1;
905 self.write_message(json!({
906 "jsonrpc": "2.0",
907 "id": id,
908 "method": method,
909 "params": params,
910 }))
911 .await?;
912 Ok(id)
913 }
914
915 async fn handshake(&mut self) -> Result<()> {
918 #[cfg(any(target_os = "macos", target_os = "linux"))]
919 if let Some(container) = &mut self.container {
920 container
921 .write_launch(self.stdin.as_mut().expect("startup stdin"))
922 .await?;
923 }
924 let init_id = self
925 .send_request(
926 method::INITIALIZE,
927 json!({
928 "protocolVersion": ACP_PROTOCOL_VERSION,
929 "clientCapabilities": {
930 "fs": { "readTextFile": false, "writeTextFile": false },
934 "terminal": false,
935 },
936 "clientInfo": {
937 "name": "kranz",
938 "title": "kranz mission engine",
939 "version": env!("CARGO_PKG_VERSION"),
940 },
941 }),
942 )
943 .await?;
944 let init_result = self.pump_until_response(init_id, HANDSHAKE_TIMEOUT).await?;
945 let peer_version = init_result
946 .get("protocolVersion")
947 .and_then(Value::as_u64)
948 .unwrap_or(0);
949 if peer_version != ACP_PROTOCOL_VERSION {
950 return Err(EngineError::Backend(format!(
951 "acp agent negotiated protocol version {peer_version}, but this backend speaks \
952 only stable version {ACP_PROTOCOL_VERSION} (schema v1)"
953 )));
954 }
955
956 let new_id = self
957 .send_request(
958 method::SESSION_NEW,
959 json!({
960 "cwd": self.spec.cwd.display().to_string(),
961 "mcpServers": [],
962 }),
963 )
964 .await?;
965 let new_result = self.pump_until_response(new_id, HANDSHAKE_TIMEOUT).await?;
966 let acp_session_id = new_result
967 .get("sessionId")
968 .and_then(Value::as_str)
969 .filter(|id| !id.trim().is_empty() && id.len() <= 256)
970 .ok_or_else(|| {
971 EngineError::Backend("acp session/new response carried no sessionId".to_string())
972 })?
973 .to_string();
974 self.acp_session_id = Some(acp_session_id.clone());
975 self.handshake_complete = true;
976 let reported_model = peer_reported_model(&new_result).map(str::to_owned);
977 #[cfg(any(target_os = "macos", target_os = "linux"))]
978 let containment = self.container.as_ref().map(|container| container.receipt());
979 #[cfg(not(any(target_os = "macos", target_os = "linux")))]
980 let containment: Option<Value> = None;
981 self.queue.push_back(AgentEvent::Init {
983 session_id: acp_session_id,
984 model: reported_model
985 .clone()
986 .unwrap_or_else(|| "unreported".into()),
987 raw: json!({
988 "initialize": init_result,
989 "sessionNew": new_result,
990 "engineSessionId": self.session_id,
991 "configuredModel": self.model,
992 "modelSource": if reported_model.is_some() { "peer" } else { "unreported" },
993 "configuredModelSelectionApplied": false,
994 "containment": containment,
995 "workerProfile": self.profile_home.as_ref().map(|p| &p.receipt),
996 "synthesizedBy": "backend_acp",
997 }),
998 });
999
1000 let prompt_text = effective_prompt(&self.spec);
1001 self.send_prompt(&prompt_text).await
1002 }
1003
1004 async fn send_prompt(&mut self, text: &str) -> Result<()> {
1007 let acp_session_id = self
1008 .acp_session_id
1009 .clone()
1010 .ok_or_else(|| EngineError::Backend("acp session not established yet".to_string()))?;
1011 let id = self
1012 .send_request(
1013 method::SESSION_PROMPT,
1014 json!({
1015 "sessionId": acp_session_id,
1016 "prompt": [ { "type": "text", "text": text } ],
1017 }),
1018 )
1019 .await?;
1020 self.prompt_request_id = Some(id);
1021 self.message_text.clear();
1024 self.message_id = None;
1025 self.last_usage = None;
1026 Ok(())
1027 }
1028
1029 async fn pump_until_response(
1034 &mut self,
1035 id: u64,
1036 timeout: std::time::Duration,
1037 ) -> Result<Value> {
1038 let pump = async {
1039 loop {
1040 let frame = match self.read_frame().await? {
1041 Some(frame) => frame,
1042 None => {
1043 return Err(EngineError::Backend(format!(
1044 "acp agent closed stdout before answering request id {id}; \
1045 stderr tail: {}",
1046 self.stderr_tail()
1047 )))
1048 }
1049 };
1050 match frame {
1051 Frame::Response {
1052 id: response_id,
1053 outcome,
1054 } if response_id == id => {
1055 return match outcome {
1056 RpcOutcome::Result(result) => Ok(result),
1057 RpcOutcome::Error(error) => Err(EngineError::Backend(format!(
1058 "acp request id {id} failed: {error}"
1059 ))),
1060 };
1061 }
1062 other => self.handle_frame(other).await?,
1063 }
1064 }
1065 };
1066 match tokio::time::timeout(timeout, pump).await {
1067 Ok(result) => result,
1068 Err(_) => Err(EngineError::Backend(format!(
1069 "acp agent did not answer request id {id} within {}s (handshake timeout)",
1070 timeout.as_secs()
1071 ))),
1072 }
1073 }
1074
1075 async fn read_frame(&mut self) -> Result<Option<Frame>> {
1079 loop {
1080 let permission_wait = self
1081 .pending_permissions
1082 .values()
1083 .map(|p| {
1084 p.expires_at
1085 .saturating_duration_since(tokio::time::Instant::now())
1086 .min(
1087 (p.proposal.deadline - chrono::Utc::now())
1088 .to_std()
1089 .unwrap_or_default(),
1090 )
1091 })
1092 .min()
1093 .unwrap_or(std::time::Duration::from_secs(86400));
1094 if permission_wait.is_zero() {
1095 return Ok(Some(Frame::PermissionExpired));
1096 }
1097 let can_answer = !self.lines.has_partial_line();
1098 let read = {
1099 let line = self.lines.next_line();
1100 tokio::pin!(line);
1101 let read = if let Some(deadline) = self.drain_deadline {
1102 tokio::time::timeout_at(deadline, &mut line)
1103 .await
1104 .map_err(|_| {
1105 EngineError::Backend(
1106 "acp stdout remained open after peer cleanup".into(),
1107 )
1108 })?
1109 } else {
1110 tokio::select! {
1111 biased;
1112 read = &mut line => read,
1115 answer = async {
1116 if let Some(answer) = self.deferred_permission_answer.take() {
1117 Some(answer)
1118 } else {
1119 self.permission_answers.recv().await
1120 }
1121 }, if can_answer => {
1122 return Ok(answer.map(Frame::PermissionAnswer));
1123 }
1124 _ = tokio::time::sleep(permission_wait), if !self.pending_permissions.is_empty() => {
1125 return Ok(Some(Frame::PermissionExpired));
1126 }
1127 exited = wait_for_peer_exit(&mut self.child) => {
1128 exited.map_err(|e| EngineError::Backend(format!("acp process observation failed: {e}")))?;
1129 kill_owned_peer(&mut self.child, #[cfg(windows)] self.job.as_ref());
1130 self.child_status = Some(self.child.wait().await.map_err(|e| {
1131 EngineError::Backend(format!("acp process reap failed: {e}"))
1132 })?);
1133 let deadline = tokio::time::Instant::now() + CLEANUP_TIMEOUT;
1134 self.drain_deadline = Some(deadline);
1135 tokio::time::timeout_at(deadline, &mut line).await
1136 .map_err(|_| EngineError::Backend("acp stdout remained open after peer cleanup".into()))?
1137 }
1138 }
1139 };
1140 read
1141 };
1142 match read {
1143 Ok(Some(line)) if line.trim().is_empty() => continue,
1144 Ok(Some(line)) => {
1145 if self.profile_home.is_some()
1146 && crate::strict_json::parse(line.as_bytes()).is_err()
1147 {
1148 return Err(EngineError::Backend(
1151 "acp profile peer emitted invalid unique-key JSON; frame refused"
1152 .into(),
1153 ));
1154 }
1155 if self
1156 .profile_home
1157 .as_ref()
1158 .is_some_and(|p| p.contains_secret(&line))
1159 {
1160 return Err(EngineError::Backend(
1161 "acp peer exposed a configured credential; frame refused".into(),
1162 ));
1163 }
1164 if !self.handshake_complete {
1165 self.handshake_bytes = self.handshake_bytes.saturating_add(line.len());
1166 if self.handshake_bytes > STDOUT_LINE_CAP {
1167 return Err(EngineError::Backend(
1168 "acp handshake output exceeded its byte limit".into(),
1169 ));
1170 }
1171 }
1172 return Ok(Some(classify_line(&line)));
1173 }
1174 Ok(None) => return Ok(None),
1175 Err(e) => {
1176 return Err(EngineError::Backend(format!(
1177 "error reading acp agent stdout: {e}; stderr tail: {}",
1178 self.stderr_tail()
1179 )))
1180 }
1181 }
1182 }
1183 }
1184
1185 async fn handle_frame(&mut self, frame: Frame) -> Result<()> {
1188 match frame {
1189 Frame::PermissionAnswer(answer) => self.answer_permission(answer).await?,
1190 Frame::PermissionExpired => {
1191 return Err(EngineError::Backend(
1194 "live permission deadline expired".into(),
1195 ));
1196 }
1197 Frame::Notification {
1198 method,
1199 params,
1200 raw,
1201 } => {
1202 if method == method::SESSION_UPDATE {
1203 self.handle_session_update(¶ms, raw)?;
1204 } else {
1205 self.queue.push_back(AgentEvent::Other { raw });
1206 }
1207 }
1208 Frame::Request {
1209 id,
1210 method,
1211 params,
1212 raw,
1213 } => {
1214 if method == method::REQUEST_PERMISSION {
1215 self.handle_permission_request(id, ¶ms, raw).await?;
1216 } else {
1217 self.write_message(json!({
1221 "jsonrpc": "2.0",
1222 "id": id,
1223 "error": {
1224 "code": -32601,
1225 "message": format!("kranz acp backend does not support {method:?}"),
1226 },
1227 }))
1228 .await?;
1229 self.queue.push_back(AgentEvent::Other { raw });
1230 }
1231 }
1232 Frame::Response { id, outcome } => {
1233 if Some(id) == self.prompt_request_id {
1234 self.prompt_request_id = None;
1235 self.synthesize_result(outcome, id);
1236 } else {
1237 self.queue.push_back(AgentEvent::Other {
1240 raw: match outcome {
1241 RpcOutcome::Result(result) => {
1242 json!({ "unmatchedResponse": { "id": id, "result": result } })
1243 }
1244 RpcOutcome::Error(error) => {
1245 json!({ "unmatchedResponse": { "id": id, "error": error } })
1246 }
1247 },
1248 });
1249 }
1250 }
1251 Frame::Unrecognized(raw) => {
1252 self.queue.push_back(AgentEvent::Other { raw });
1253 return Err(EngineError::Backend(
1254 "acp peer emitted malformed JSON-RPC".into(),
1255 ));
1256 }
1257 }
1258 Ok(())
1259 }
1260
1261 fn track_tool_call(&mut self, id: String, info: ToolCallInfo) -> Result<()> {
1262 if id.trim().is_empty() || id.len() > 256 {
1263 return Err(EngineError::Backend(
1264 "acp tool update has no valid toolCallId".into(),
1265 ));
1266 }
1267 let bytes = info.kind.len() + info.title.len() + info.subject.len();
1268 let retained = self
1269 .tool_calls
1270 .iter()
1271 .filter(|(key, _)| *key != &id)
1272 .map(|(key, value)| {
1273 key.len() + value.kind.len() + value.title.len() + value.subject.len()
1274 })
1275 .sum::<usize>();
1276 if self.tool_calls.len() >= 1024
1277 || retained.saturating_add(bytes).saturating_add(id.len()) > STDOUT_LINE_CAP
1278 {
1279 return Err(EngineError::Backend(
1280 "acp tool-call tracking exceeded its limit".into(),
1281 ));
1282 }
1283 self.tool_calls.insert(id, info);
1284 Ok(())
1285 }
1286
1287 fn handle_session_update(&mut self, params: &Value, raw: Value) -> Result<()> {
1289 let Some(expected) = self.acp_session_id.as_deref() else {
1290 self.queue.push_back(AgentEvent::Other { raw });
1292 return Ok(());
1293 };
1294 if params.get("sessionId").and_then(Value::as_str) != Some(expected) {
1295 return Err(EngineError::Backend(
1296 "acp update has a foreign or missing sessionId".into(),
1297 ));
1298 }
1299 let update = params.get("update").cloned().unwrap_or(Value::Null);
1300 if let Some(call_id) = update.get("toolCallId").and_then(Value::as_str) {
1301 if let Some(pending) = self
1302 .pending_permissions
1303 .values()
1304 .map(|p| &p.proposal)
1305 .find(|p| p.tool_call_id == call_id)
1306 {
1307 let changed = ["kind", "rawInput", "locations", "content"]
1310 .iter()
1311 .any(|key| {
1312 update
1313 .get(*key)
1314 .is_some_and(|value| pending.action.get(*key) != Some(value))
1315 });
1316 let terminal = update
1317 .get("status")
1318 .and_then(Value::as_str)
1319 .is_some_and(|s| s == "completed" || s == "failed" || s == "in_progress");
1320 if changed || terminal {
1321 return Err(EngineError::Backend(
1322 "ACP invocation changed or started while consent was pending".into(),
1323 ));
1324 }
1325 }
1326 }
1327 match update.get("sessionUpdate").and_then(Value::as_str) {
1328 Some("agent_message_chunk") => {
1329 let text = update
1330 .get("content")
1331 .and_then(|c| c.get("text"))
1332 .and_then(Value::as_str)
1333 .unwrap_or("");
1334 if text.is_empty() {
1335 self.queue.push_back(AgentEvent::Other { raw });
1336 return Ok(());
1337 }
1338 let chunk_id = update
1341 .get("messageId")
1342 .and_then(Value::as_str)
1343 .map(str::to_string);
1344 if chunk_id.is_some() && chunk_id != self.message_id {
1345 self.message_text.clear();
1346 self.message_id = chunk_id;
1347 }
1348 if self.message_text.len().saturating_add(text.len()) > STDOUT_LINE_CAP {
1349 return Err(EngineError::Backend(
1350 "acp assistant message exceeded its byte limit".into(),
1351 ));
1352 }
1353 self.message_text.push_str(text);
1354 self.queue.push_back(AgentEvent::Text {
1355 text: text.to_string(),
1356 raw,
1357 });
1358 }
1359 Some("tool_call") => {
1360 let id = update
1361 .get("toolCallId")
1362 .and_then(Value::as_str)
1363 .unwrap_or_default()
1364 .to_string();
1365 let kind = update
1366 .get("kind")
1367 .and_then(Value::as_str)
1368 .unwrap_or("other")
1369 .to_string();
1370 let title = update
1371 .get("title")
1372 .and_then(Value::as_str)
1373 .unwrap_or_default()
1374 .to_string();
1375 let subject = tool_call_subject(&kind, &title, &update);
1376 self.track_tool_call(
1377 id,
1378 ToolCallInfo {
1379 kind: kind.clone(),
1380 title: title.clone(),
1381 subject,
1382 },
1383 )?;
1384 self.queue.push_back(AgentEvent::ToolUse {
1385 tool: kind,
1386 summary: truncate_chars(&title, SUMMARY_MAX_CHARS),
1387 raw,
1388 });
1389 }
1390 Some("tool_call_update") => {
1391 let id = update
1392 .get("toolCallId")
1393 .and_then(Value::as_str)
1394 .unwrap_or_default()
1395 .to_string();
1396 let status = update
1397 .get("status")
1398 .and_then(Value::as_str)
1399 .unwrap_or("")
1400 .to_string();
1401 let mut tracked = self.tool_calls.get(&id).cloned().unwrap_or_default();
1402 if let Some(kind) = update.get("kind").and_then(Value::as_str) {
1403 tracked.kind = kind.to_string();
1404 }
1405 if let Some(title) = update.get("title").and_then(Value::as_str) {
1406 tracked.title = title.to_string();
1407 }
1408 if update.get("kind").is_some()
1409 || update.get("rawInput").is_some()
1410 || update.get("locations").is_some()
1411 {
1412 tracked.subject = tool_call_subject(&tracked.kind, &tracked.title, &update);
1413 }
1414 self.track_tool_call(id.clone(), tracked.clone())?;
1415 match status.as_str() {
1416 "completed" | "failed" => {
1422 self.tool_calls.remove(&id);
1423 self.queue.push_back(AgentEvent::ToolResult {
1424 tool: Some(tracked.kind.clone()),
1425 denied: false,
1426 summary: tool_result_summary(&update, &tracked, &status),
1427 raw,
1428 });
1429 }
1430 _ => self.queue.push_back(AgentEvent::Other { raw }),
1431 }
1432 }
1433 Some("usage_update") => {
1434 self.last_usage = Some(update.clone());
1439 self.queue.push_back(AgentEvent::Other { raw });
1440 }
1441 _ => self.queue.push_back(AgentEvent::Other { raw }),
1442 }
1443 Ok(())
1444 }
1445
1446 async fn handle_permission_request(
1450 &mut self,
1451 id: Value,
1452 params: &Value,
1453 raw: Value,
1454 ) -> Result<()> {
1455 if self.acp_session_id.is_none()
1456 || params.get("sessionId").and_then(Value::as_str) != self.acp_session_id.as_deref()
1457 {
1458 self.write_message(json!({
1459 "jsonrpc": "2.0", "id": id,
1460 "result": { "outcome": { "outcome": "cancelled" } },
1461 }))
1462 .await?;
1463 return Err(EngineError::Backend(
1464 "acp permission has a foreign or missing sessionId".into(),
1465 ));
1466 }
1467 let call_update = params.get("toolCall").cloned().unwrap_or(Value::Null);
1468 let call_id = call_update
1469 .get("toolCallId")
1470 .and_then(Value::as_str)
1471 .filter(|id| !id.trim().is_empty() && id.len() <= 256);
1472 let missing_id = call_id.is_none();
1473 let call_id = call_id.unwrap_or_default().to_string();
1474 let mut info = self.tool_calls.get(&call_id).cloned().unwrap_or_default();
1477 if let Some(kind) = call_update.get("kind").and_then(Value::as_str) {
1478 info.kind = kind.to_string();
1479 }
1480 if info.kind.is_empty() {
1481 info.kind = "other".to_string();
1482 }
1483 if let Some(title) = call_update.get("title").and_then(Value::as_str) {
1484 info.title = title.to_string();
1485 }
1486 if info.subject.is_empty()
1487 || call_update.get("kind").is_some()
1488 || call_update.get("rawInput").is_some()
1489 || call_update.get("locations").is_some()
1490 {
1491 info.subject = tool_call_subject(&info.kind, &info.title, &call_update);
1492 }
1493 if !missing_id {
1494 self.track_tool_call(call_id.clone(), info.clone())?;
1495 }
1496
1497 let decision = if missing_id {
1498 PermissionDecision::Deny("tool call has no valid action identity (toolCallId)".into())
1499 } else {
1500 match decide_permission(&self.spec, &info) {
1501 PermissionDecision::Allow
1502 if !call_update.get("rawInput").is_some_and(Value::is_object)
1503 || info.subject.trim().is_empty() =>
1504 {
1505 PermissionDecision::Deny(
1506 "the complete action is unavailable for one-call consent".into(),
1507 )
1508 }
1509 decision => decision,
1510 }
1511 };
1512 let peer_id = serde_json::to_string(&id)?;
1513 if self.pending_permissions.len() >= crate::live_permission::MAX_PENDING
1514 || self.seen_permission_ids.len() >= 1024
1515 || !self.seen_permission_ids.insert(peer_id)
1516 || self
1517 .pending_permissions
1518 .values()
1519 .any(|p| p.proposal.tool_call_id == call_id)
1520 {
1521 return Err(EngineError::Backend(
1522 "duplicate or excessive ACP permission requests".into(),
1523 ));
1524 }
1525 let options = params
1526 .get("options")
1527 .and_then(Value::as_array)
1528 .cloned()
1529 .unwrap_or_default();
1530 let mut action = call_update;
1531 if let Some(fields) = action.as_object_mut() {
1532 fields.insert("kind".into(), Value::String(info.kind));
1533 }
1534 let now = chrono::Utc::now();
1535 let mut proposal = crate::live_permission::Proposal {
1536 id: format!("permission-{}", uuid::Uuid::new_v4()),
1537 engine_session_id: self.session_id.clone(),
1538 peer_session_id: self.acp_session_id.clone().expect("validated above"),
1539 peer_request_id: id,
1540 tool_call_id: call_id,
1541 action_digest: crate::live_permission::digest(&action)?,
1542 options_digest: crate::live_permission::digest(&options)?,
1543 action,
1544 options,
1545 observed_at: now,
1546 deadline: now + chrono::Duration::seconds(crate::live_permission::REQUEST_TTL_SECS),
1547 prohibition: match decision {
1548 PermissionDecision::Deny(reason) => Some(reason),
1549 PermissionDecision::Allow => None,
1550 },
1551 };
1552 if proposal.option(true).is_none() && proposal.prohibition.is_none() {
1553 proposal.prohibition = Some("no unique certified allow_once option was offered".into());
1554 }
1555 proposal.validate()?;
1556 self.pending_permissions.insert(
1557 proposal.id.clone(),
1558 PendingPermission {
1559 proposal: proposal.clone(),
1560 expires_at: tokio::time::Instant::now()
1561 + std::time::Duration::from_secs(
1562 crate::live_permission::REQUEST_TTL_SECS as u64,
1563 ),
1564 },
1565 );
1566 self.queue.push_back(AgentEvent::PermissionRequested {
1567 proposal: Box::new(proposal),
1568 raw,
1569 });
1570 Ok(())
1571 }
1572
1573 async fn answer_permission(&mut self, answer: crate::live_permission::Answer) -> Result<()> {
1574 if self.lines.has_partial_line() {
1578 self.deferred_permission_answer = Some(answer);
1579 return Ok(());
1580 }
1581 let Some(pending) = self.pending_permissions.get(&answer.proposal.id) else {
1582 return Err(EngineError::Backend(
1583 "permission response names no live request".into(),
1584 ));
1585 };
1586 if pending.proposal != answer.proposal
1587 || chrono::Utc::now() >= pending.proposal.deadline
1588 || tokio::time::Instant::now() >= pending.expires_at
1589 || (answer.allow
1590 && (pending.proposal.prohibition.is_some()
1591 || pending.proposal.option(true).is_none()))
1592 {
1593 return Err(EngineError::Backend(
1594 "stale or prohibited permission response".into(),
1595 ));
1596 }
1597 let proposal = self
1598 .pending_permissions
1599 .remove(&answer.proposal.id)
1600 .expect("checked above")
1601 .proposal;
1602 let decision = if answer.allow {
1603 PermissionDecision::Allow
1604 } else {
1605 PermissionDecision::Deny("one-call consent refused".into())
1606 };
1607 let result = permission_response(&decision, &proposal.options);
1608 let sent = self
1609 .write_message(json!({
1610 "jsonrpc": "2.0", "id": proposal.peer_request_id, "result": result,
1611 }))
1612 .await;
1613 let delivery = if sent.is_ok() {
1614 crate::live_permission::Delivery::Sent
1615 } else {
1616 crate::live_permission::Delivery::Uncertain
1617 };
1618 self.queue.push_back(AgentEvent::PermissionResponded {
1619 request_id: proposal.id.clone(),
1620 delivery: delivery.clone(),
1621 raw: json!({"permissionResponse":proposal.id,"delivery":delivery}),
1622 });
1623 sent
1624 }
1625
1626 fn synthesize_result(&mut self, outcome: RpcOutcome, request_id: u64) {
1642 let cumulative_cost_usd = self
1643 .last_usage
1644 .as_ref()
1645 .and_then(|u| u.get("cost"))
1646 .filter(|cost| {
1647 cost.get("currency").and_then(Value::as_str) == Some("USD")
1648 && cost.get("amount").and_then(Value::as_f64).is_some()
1649 })
1650 .and_then(|cost| cost.get("amount").and_then(Value::as_f64))
1651 .filter(|amount| amount.is_finite() && *amount >= 0.0);
1652 let last_cost_usd = match (
1653 self.saw_result,
1654 self.previous_cost_total,
1655 cumulative_cost_usd,
1656 ) {
1657 (false, _, total) => total,
1658 (true, Some(previous), Some(total)) if total >= previous => Some(total - previous),
1659 _ => None,
1660 };
1661 self.previous_cost_total = cumulative_cost_usd;
1662 let (text, is_error, raw) = match outcome {
1663 RpcOutcome::Result(result) => {
1664 let stop_reason = result.get("stopReason").and_then(Value::as_str);
1665 (
1666 std::mem::take(&mut self.message_text),
1667 stop_reason != Some("end_turn"),
1668 json!({
1669 "promptResponse": result,
1670 "usageUpdate": self.last_usage,
1671 "costScope": "turn_delta_from_reported_session_total",
1672 "synthesizedBy": "backend_acp",
1673 "stopReason": stop_reason,
1677 }),
1678 )
1679 }
1680 RpcOutcome::Error(error) => (
1681 format!("acp session/prompt failed: {error}"),
1682 true,
1683 json!({
1684 "promptError": error,
1685 "requestId": request_id,
1686 "synthesizedBy": "backend_acp",
1687 }),
1688 ),
1689 };
1690 let event = AgentEvent::Result {
1691 text,
1692 is_error,
1693 usage: TokenUsage::default(),
1694 cost_usd: last_cost_usd,
1695 num_turns: Some(1),
1696 raw,
1697 };
1698 self.observe(&event);
1699 self.queue.push_back(event);
1700 }
1701
1702 fn observe(&mut self, event: &AgentEvent) {
1703 if let AgentEvent::Result { .. } = event {
1704 self.saw_result = true;
1705 }
1706 }
1707
1708 async fn kill_child(&mut self) {
1711 self.stdin.take();
1712 if self.child_status.is_none() {
1713 kill_owned_peer(
1714 &mut self.child,
1715 #[cfg(windows)]
1716 self.job.as_ref(),
1717 );
1718 if let Ok(Ok(status)) = tokio::time::timeout(CLEANUP_TIMEOUT, self.child.wait()).await {
1719 self.child_status = Some(status);
1720 }
1721 }
1722 #[cfg(any(target_os = "macos", target_os = "linux"))]
1723 if let Some(container) = self.container.as_mut() {
1724 if let Err(error) = container.remove().await {
1725 self.cleanup_failure
1726 .get_or_insert("container removal failed");
1727 tracing::error!(%error, "ACP cleanup requires recovery");
1728 }
1729 }
1730 if let Some(profile) = self.profile_home.as_mut() {
1731 if let Err(error) = profile.close() {
1732 self.cleanup_failure
1733 .get_or_insert("private credential home cleanup failed");
1734 tracing::error!(%error, "ACP private home requires recovery");
1735 }
1736 }
1737 self.finish_stderr().await;
1738 }
1739
1740 async fn finish_stderr(&mut self) {
1741 if let Some(mut task) = self.stderr_task.take() {
1742 match tokio::time::timeout(CLEANUP_TIMEOUT, &mut task).await {
1743 Ok(Ok(())) => {}
1744 Ok(Err(_)) => {
1745 self.cleanup_failure
1746 .get_or_insert("stderr drain task failed");
1747 }
1748 Err(_) => {
1749 self.cleanup_failure.get_or_insert("stderr drain timed out");
1750 task.abort();
1751 let _ = task.await;
1752 }
1753 }
1754 }
1755 }
1756
1757 async fn finish_session(&mut self) {
1762 self.stdin.take();
1763 #[cfg(any(target_os = "macos", target_os = "linux"))]
1764 let contained = self.container.is_some();
1765 #[cfg(not(any(target_os = "macos", target_os = "linux")))]
1766 let contained = false;
1767 let completion_timeout = if contained {
1768 CONTAINER_COMPLETION_TIMEOUT
1769 } else {
1770 COMPLETION_GRACE
1771 };
1772 let mut forced = false;
1773 if self.child_status.is_none() {
1774 match tokio::time::timeout(completion_timeout, wait_for_peer_exit(&mut self.child))
1775 .await
1776 {
1777 Ok(Ok(())) => {}
1778 Ok(Err(error)) => {
1779 self.kill_child().await;
1780 self.exit = Some(SessionExit::Failed(format!(
1781 "acp process observation failed: {error}"
1782 )));
1783 return;
1784 }
1785 Err(_) if contained => {
1786 self.cleanup_failure
1787 .get_or_insert("container completion deadline expired");
1788 }
1789 Err(_) => forced = true,
1790 }
1791 self.kill_child().await;
1792 } else {
1793 self.kill_child().await;
1794 }
1795 if let Some(cause) = self.cleanup_failure {
1796 self.exit = Some(SessionExit::Failed(format!(
1797 "acp cleanup could not be confirmed: {cause}"
1798 )));
1799 return;
1800 }
1801 self.exit = Some(match self.child_status {
1802 Some(status) if self.saw_result && (status.success() || forced) => {
1803 SessionExit::Completed
1804 }
1805 Some(status) => SessionExit::Failed(format!(
1806 "acp agent exited with {status}{}; stderr tail: {}",
1807 if self.saw_result {
1808 ""
1809 } else {
1810 " without answering session/prompt"
1811 },
1812 self.stderr_tail(),
1813 )),
1814 None => SessionExit::Failed("acp process did not reap within cleanup deadline".into()),
1815 });
1816 }
1817
1818 fn stderr_tail(&self) -> String {
1819 let captured = self
1820 .stderr_buf
1821 .lock()
1822 .map(|guard| guard.clone())
1823 .unwrap_or_default();
1824 let captured = self
1825 .profile_home
1826 .as_ref()
1827 .map_or_else(|| captured.clone(), |p| p.scrub(captured.clone()));
1828 last_chars(captured.trim_end(), STDERR_TAIL_CHARS)
1829 }
1830}
1831
1832#[async_trait::async_trait]
1833impl AgentSession for AcpSession {
1834 fn permission_responder(&self) -> Option<crate::live_permission::PermissionResponder> {
1835 Some(self.permission_responder.clone())
1836 }
1837
1838 fn session_id(&self) -> String {
1839 self.session_id.clone()
1840 }
1841
1842 async fn next_event(&mut self) -> Result<Option<AgentEvent>> {
1843 loop {
1844 if let Some(event) = self.queue.pop_front() {
1845 return Ok(Some(event));
1846 }
1847 if self.exit.is_some() {
1848 return Ok(None);
1849 }
1850 if self.saw_result && matches!(self.spec.prompt, PromptMode::SingleShot(_)) {
1851 self.finish_session().await;
1852 return Ok(None);
1853 }
1854 let frame = match self.read_frame().await {
1855 Ok(Some(frame)) => frame,
1856 Ok(None) => {
1857 self.finish_session().await;
1858 return Ok(None);
1859 }
1860 Err(e) => {
1861 self.kill_child().await;
1862 self.exit = Some(SessionExit::Failed(e.to_string()));
1863 return Ok(None);
1864 }
1865 };
1866 if let Err(e) = self.handle_frame(frame).await {
1867 self.kill_child().await;
1870 self.exit = Some(SessionExit::Failed(e.to_string()));
1871 continue;
1873 }
1874 }
1875 }
1876
1877 async fn send_user_message(&mut self, text: &str) -> Result<()> {
1878 if self.exit.is_some() {
1879 return Err(EngineError::Backend(
1880 "acp session is closed; cannot send further messages".to_string(),
1881 ));
1882 }
1883 if self.acp_session_id.is_none() {
1884 return Err(EngineError::Backend(
1885 "acp session not established yet; cannot send a message".to_string(),
1886 ));
1887 }
1888 if matches!(self.spec.prompt, PromptMode::SingleShot(_)) || self.prompt_request_id.is_some()
1889 {
1890 return Err(EngineError::Backend(
1891 "acp session cannot accept an overlapping or single-shot follow-up".into(),
1892 ));
1893 }
1894 self.send_prompt(text).await
1895 }
1896
1897 async fn abort(&mut self) -> Result<()> {
1898 if self.exit.is_some() {
1899 return Ok(());
1900 }
1901 if let Some(acp_session_id) = self.acp_session_id.clone() {
1906 let _ = tokio::time::timeout(
1907 CANCEL_WRITE_TIMEOUT,
1908 self.write_message(json!({
1909 "jsonrpc": "2.0",
1910 "method": method::SESSION_CANCEL,
1911 "params": { "sessionId": acp_session_id },
1912 })),
1913 )
1914 .await;
1915 }
1916 self.kill_child().await;
1917 #[cfg(any(target_os = "macos", target_os = "linux"))]
1918 let container_cleanup_failed = self.container.is_some()
1919 && (self.cleanup_failure.is_some() || self.child_status.is_none());
1920 #[cfg(not(any(target_os = "macos", target_os = "linux")))]
1921 let container_cleanup_failed = false;
1922 if container_cleanup_failed {
1923 let cause = self
1924 .cleanup_failure
1925 .unwrap_or("host process did not reap within cleanup deadline");
1926 let message = format!("acp abort cleanup could not be confirmed: {cause}");
1927 self.exit = Some(SessionExit::Failed(message.clone()));
1928 return Err(EngineError::Backend(message));
1929 }
1930 self.exit = Some(SessionExit::Aborted);
1931 Ok(())
1932 }
1933
1934 fn exit_status(&self) -> Option<SessionExit> {
1935 self.exit.clone()
1936 }
1937}
1938
1939#[cfg(test)]
1942mod tests {
1943 use super::*;
1944
1945 #[test]
1946 fn backend_acp_classify_distinguishes_response_request_notification() {
1947 let response =
1948 classify_line(r#"{"jsonrpc":"2.0","id":3,"result":{"stopReason":"end_turn"}}"#);
1949 assert!(matches!(
1950 response,
1951 Frame::Response {
1952 id: 3,
1953 outcome: RpcOutcome::Result(_)
1954 }
1955 ));
1956 let error =
1957 classify_line(r#"{"jsonrpc":"2.0","id":4,"error":{"code":-32603,"message":"boom"}}"#);
1958 assert!(matches!(
1959 error,
1960 Frame::Response {
1961 id: 4,
1962 outcome: RpcOutcome::Error(_)
1963 }
1964 ));
1965 let request = classify_line(
1966 r#"{"jsonrpc":"2.0","id":100,"method":"session/request_permission","params":{}}"#,
1967 );
1968 assert!(matches!(
1969 request,
1970 Frame::Request { ref method, .. } if method == "session/request_permission"
1971 ));
1972 let notification = classify_line(
1973 r#"{"jsonrpc":"2.0","method":"session/update","params":{"sessionId":"s","update":{"sessionUpdate":"plan"}}}"#,
1974 );
1975 assert!(matches!(
1976 notification,
1977 Frame::Notification { ref method, .. } if method == "session/update"
1978 ));
1979 let torn = classify_line(r#"{"jsonrpc":"2.0","method":"session/upda"#);
1981 assert!(matches!(torn, Frame::Unrecognized(_)));
1982 }
1983
1984 #[test]
1985 fn backend_acp_wildcard_match_anchors_like_a_shell_glob() {
1986 assert!(wildcard_match("git push*", "git push origin main"));
1987 assert!(wildcard_match("git push*", "git push"));
1988 assert!(!wildcard_match("git push*", "git pull"));
1989 assert!(wildcard_match("*", "anything"));
1990 assert!(wildcard_match(
1991 "cargo * --workspace",
1992 "cargo test --workspace"
1993 ));
1994 assert!(!wildcard_match(
1995 "cargo * --workspace",
1996 "cargo test --package x"
1997 ));
1998 assert!(wildcard_match("*/etc/passwd", "/etc/passwd"));
1999 assert!(!wildcard_match("*/etc/passwd", "/etc/passwd.bak"));
2000 }
2001
2002 #[test]
2003 fn backend_acp_pattern_matches_maps_claude_names_to_acp_kinds() {
2004 assert!(pattern_matches(
2005 "Bash(git push*)",
2006 "execute",
2007 "git push origin main"
2008 ));
2009 assert!(!pattern_matches("Bash(git push*)", "execute", "git pull"));
2010 assert!(!pattern_matches(
2011 "Bash(git push*)",
2012 "edit",
2013 "git push origin main"
2014 ));
2015 assert!(pattern_matches("Write", "edit", "/repo/src/main.rs"));
2016 assert!(!pattern_matches("Write", "read", "/repo/src/main.rs"));
2017 assert!(pattern_matches("Edit(/etc/*)", "edit", "/etc/hosts"));
2018 assert!(pattern_matches("Read", "read", "/anywhere"));
2019 assert!(!pattern_matches("NotAClaudeTool(*)", "execute", "x"));
2021 }
2022
2023 fn spec_with(writable: bool, disallowed: &[&str]) -> SessionSpec {
2024 SessionSpec {
2025 cwd: PathBuf::from("."),
2026 prompt: PromptMode::SingleShot("do the thing".to_string()),
2027 append_system_prompt: None,
2028 model: "acp-model".to_string(),
2029 effort: "high".to_string(),
2030 session_id: "sess-1".to_string(),
2031 resume: None,
2032 permission_mode: None,
2033 allowed_tools: vec![],
2034 disallowed_tools: disallowed.iter().map(|s| s.to_string()).collect(),
2035 tools: vec![],
2036 writable,
2037 settings_json: None,
2038 json_schema: None,
2039 max_budget_usd: None,
2040 max_turns: None,
2041 env: Default::default(),
2042 sandbox: None,
2043 hook_status: None,
2044 }
2045 }
2046
2047 #[test]
2048 fn backend_acp_permission_denies_disallowed_and_mutating_kinds() {
2049 let spec = spec_with(true, &["Bash(git push*)"]);
2050 let push = ToolCallInfo {
2051 kind: "execute".to_string(),
2052 title: "git push origin main".to_string(),
2053 subject: "git push origin main".to_string(),
2054 };
2055 assert!(matches!(
2056 decide_permission(&spec, &push),
2057 PermissionDecision::Deny(reason) if reason.contains("Bash(git push*)")
2058 ));
2059 let test = ToolCallInfo {
2060 kind: "execute".to_string(),
2061 title: "cargo test".to_string(),
2062 subject: "cargo test".to_string(),
2063 };
2064 assert_eq!(decide_permission(&spec, &test), PermissionDecision::Allow);
2065
2066 let ro = spec_with(false, &[]);
2068 let edit = ToolCallInfo {
2069 kind: "edit".to_string(),
2070 title: "write src/main.rs".to_string(),
2071 subject: "/repo/src/main.rs".to_string(),
2072 };
2073 assert!(matches!(
2074 decide_permission(&ro, &edit),
2075 PermissionDecision::Deny(reason) if reason.contains("writable: false")
2076 ));
2077 assert_eq!(decide_permission(&ro, &test), PermissionDecision::Allow);
2078 let read = ToolCallInfo {
2079 kind: "read".to_string(),
2080 title: "read src/main.rs".to_string(),
2081 subject: "/repo/src/main.rs".to_string(),
2082 };
2083 assert_eq!(decide_permission(&ro, &read), PermissionDecision::Allow);
2084 }
2085
2086 #[test]
2091 fn backend_acp_permission_denies_when_the_subject_is_missing() {
2092 let spec = spec_with(true, &["Bash(git push*)"]);
2093 let no_subject = ToolCallInfo {
2094 kind: "execute".to_string(),
2095 title: String::new(),
2096 subject: String::new(),
2097 };
2098 assert!(
2099 matches!(
2100 decide_permission(&spec, &no_subject),
2101 PermissionDecision::Deny(ref reason)
2102 if reason.contains("no subject") && reason.contains("Bash(git push*)")
2103 ),
2104 "got {:?}",
2105 decide_permission(&spec, &no_subject)
2106 );
2107
2108 let blank_subject = ToolCallInfo {
2110 subject: " ".to_string(),
2111 ..no_subject.clone()
2112 };
2113 assert!(matches!(
2114 decide_permission(&spec, &blank_subject),
2115 PermissionDecision::Deny(_)
2116 ));
2117
2118 let read_no_subject = ToolCallInfo {
2121 kind: "read".to_string(),
2122 ..no_subject.clone()
2123 };
2124 assert_eq!(
2125 decide_permission(&spec_with(true, &["Bash(git push*)"]), &read_no_subject),
2126 PermissionDecision::Allow
2127 );
2128 }
2129
2130 #[test]
2134 fn backend_acp_read_only_denies_an_unclassifiable_kind() {
2135 let ro = spec_with(false, &[]);
2136 for kind in ["", "other"] {
2137 let call = ToolCallInfo {
2138 kind: kind.to_string(),
2139 title: "do something".to_string(),
2140 subject: "/repo/src/main.rs".to_string(),
2141 };
2142 assert!(
2143 matches!(
2144 decide_permission(&ro, &call),
2145 PermissionDecision::Deny(ref reason) if reason.contains("kind")
2146 ),
2147 "kind {kind:?} got {:?}",
2148 decide_permission(&ro, &call)
2149 );
2150 }
2151 let writable = spec_with(true, &[]);
2153 let other = ToolCallInfo {
2154 kind: "other".to_string(),
2155 title: "think".to_string(),
2156 subject: "think".to_string(),
2157 };
2158 assert!(matches!(
2159 decide_permission(&writable, &other),
2160 PermissionDecision::Deny(_)
2161 ));
2162 }
2163
2164 #[test]
2165 fn backend_acp_permission_response_picks_options_or_cancels() {
2166 let options = vec![
2167 json!({ "optionId": "allow-1", "name": "Allow", "kind": "allow_once" }),
2168 json!({ "optionId": "reject-1", "name": "Reject", "kind": "reject_once" }),
2169 ];
2170 let allow = permission_response(&PermissionDecision::Allow, &options);
2171 assert_eq!(allow["outcome"]["optionId"], json!("allow-1"));
2172 let deny = permission_response(&PermissionDecision::Deny("nope".to_string()), &options);
2173 assert_eq!(deny["outcome"]["optionId"], json!("reject-1"));
2174 let deny_no_reject = permission_response(
2177 &PermissionDecision::Deny("nope".to_string()),
2178 &[options[0].clone()],
2179 );
2180 assert_eq!(deny_no_reject["outcome"]["outcome"], json!("cancelled"));
2181 for (kind, decision) in [
2182 ("allow_once", PermissionDecision::Allow),
2183 ("reject_once", PermissionDecision::Deny("refused".into())),
2184 ] {
2185 let ambiguous = vec![
2186 json!({"optionId":"a","kind":kind}),
2187 json!({"optionId":"b","kind":kind}),
2188 ];
2189 assert_eq!(
2190 permission_response(&decision, &ambiguous)["outcome"]["outcome"],
2191 "cancelled"
2192 );
2193 }
2194 }
2195}