use std::iter;
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex, OnceLock};
use std::time::Duration;
use agent_client_protocol::schema::ProtocolVersion;
use agent_client_protocol::schema::v1::{
AgentCapabilities, AuthenticateRequest, CancelNotification, ClientCapabilities, ContentBlock,
FileSystemCapabilities, InitializeRequest, LoadSessionRequest, LoadSessionResponse,
NewSessionRequest, PermissionOptionKind, PromptRequest, ReadTextFileRequest,
ReadTextFileResponse, RequestPermissionOutcome, RequestPermissionRequest,
RequestPermissionResponse, ResumeSessionRequest, ResumeSessionResponse,
SelectedPermissionOutcome, SessionConfigId, SessionConfigKind, SessionConfigOption,
SessionConfigSelectOptions, SessionConfigValueId, SessionId, SessionNotification,
SetSessionConfigOptionRequest, StopReason, TextContent, WriteTextFileRequest,
WriteTextFileResponse,
};
use agent_client_protocol::{
AcpAgent, Agent as AcpRole, Client, ConnectionTo, Error as AcpError, LineDirection,
on_receive_notification, on_receive_request,
};
use super::events::SubagentEvent;
#[derive(Debug, Clone, Copy)]
pub struct Agent {
pub name: &'static str,
pub command: &'static str,
pub args: &'static [&'static str],
pub requires: &'static [&'static str],
pub install_hint: &'static str,
pub model: &'static str,
}
impl Agent {
#[must_use]
pub fn invocation(&self) -> String {
let mut parts: Vec<&str> = Vec::with_capacity(self.args.len() + 1);
parts.push(self.command);
parts.extend_from_slice(self.args);
parts.join(" ")
}
}
pub const ACP_AGENTS: &[Agent] = &[
Agent {
name: "gemini",
command: "gemini",
args: &["--experimental-acp"],
requires: &["gemini"],
install_hint: "Install: npm install -g @google/gemini-cli. Then: gemini auth login. More: https://github.com/google-gemini/gemini-cli",
model: "",
},
Agent {
name: "goose",
command: "goose",
args: &["acp"],
requires: &["goose"],
install_hint: "Install: curl -fsSL https://github.com/block/goose/releases/download/stable/download_cli.sh | bash. More: https://block.github.io/goose",
model: "",
},
Agent {
name: "opencode",
command: "opencode",
args: &["acp"],
requires: &["opencode"],
install_hint: "Install: curl -fsSL https://opencode.ai/install | bash (or: npm install -g opencode-ai). More: https://opencode.ai/docs",
model: "",
},
Agent {
name: "kilo",
command: "kilo",
args: &["acp"],
requires: &["kilo"],
install_hint: "Install: curl -fsSL https://kilo.ai/install.sh | sh. More: https://kilo.ai/docs/code-with-ai/platforms/cli",
model: "kilo/nvidia/nemotron-3-ultra-550b-a55b:free",
},
Agent {
name: "cline",
command: "cline",
args: &["--acp"],
requires: &["cline"],
install_hint: "Install: npm install -g cline. Then: cline auth (or sign in from the first session). More: https://docs.cline.bot/usage/acp",
model: "",
},
Agent {
name: "devin",
command: "devin",
args: &["acp"],
requires: &["devin"],
install_hint: "Install: curl -fsSL https://cli.devin.ai/install.sh | bash. Then: devin auth login (or set WINDSURF_API_KEY). More: https://docs.devin.ai/cli",
model: "",
},
Agent {
name: "claude",
command: "npx",
args: &["-y", "@agentclientprotocol/claude-agent-acp@latest"],
requires: &["claude", "npx"],
install_hint: "Install: npm install -g @anthropic-ai/claude-code (the ACP adapter runs via npx). More: https://code.claude.com/docs",
model: "",
},
Agent {
name: "codex",
command: "npx",
args: &["-y", "@agentclientprotocol/codex-acp@latest"],
requires: &["codex", "npx"],
install_hint: "Install: npm install -g @openai/codex (the ACP adapter runs via npx). More: https://developers.openai.com/codex",
model: "",
},
];
pub(crate) fn agent_launcher(entry: &Agent) -> Result<AcpAgent, String> {
let args: Vec<String> = iter::once(entry.command)
.chain(entry.args.iter().copied())
.map(String::from)
.collect();
let args = windows_script_launcher(&args).unwrap_or(args);
AcpAgent::from_args(args)
.map_err(|e| format!("invalid ACP launch command for '{}': {e}", entry.name))
}
pub(crate) fn acp_error(message: impl Into<String>) -> AcpError {
let mut error = AcpError::internal_error();
error.message = message.into();
error
}
fn binary_in_path(binary: &str) -> bool {
std::env::var_os("PATH").is_some_and(|path| {
std::env::split_paths(&path).any(|dir| {
dir.join(binary).is_file()
|| (cfg!(windows)
&& ["exe", "cmd", "bat"]
.into_iter()
.any(|ext| dir.join(format!("{binary}.{ext}")).is_file()))
})
})
}
#[cfg(windows)]
fn exe_in_path(binary: &str) -> bool {
std::env::var_os("PATH").is_some_and(|path| {
std::env::split_paths(&path).any(|dir| dir.join(format!("{binary}.exe")).is_file())
})
}
#[cfg(windows)]
fn script_shim_in_path(binary: &str) -> bool {
std::env::var_os("PATH").is_some_and(|path| {
std::env::split_paths(&path).any(|dir| {
["cmd", "bat"]
.into_iter()
.any(|ext| dir.join(format!("{binary}.{ext}")).is_file())
})
})
}
#[cfg(windows)]
fn windows_script_launcher(args: &[String]) -> Option<Vec<String>> {
const CMD_METACHARACTERS: [char; 7] = ['&', '|', '^', '%', '<', '>', '"'];
if args
.iter()
.any(|arg| arg.chars().any(|c| CMD_METACHARACTERS.contains(&c)))
{
return None;
}
let command = args.first()?;
if exe_in_path(command) || !script_shim_in_path(command) {
return None;
}
Some(vec![
"cmd".to_string(),
"/d".to_string(),
"/s".to_string(),
"/c".to_string(),
args.join(" "),
])
}
#[cfg(not(windows))]
fn windows_script_launcher(_args: &[String]) -> Option<Vec<String>> {
None
}
#[must_use]
pub fn install_hint(agent: &str) -> &'static str {
ACP_AGENTS
.iter()
.find(|entry| entry.name == agent)
.map_or("", |entry| entry.install_hint)
}
pub fn detect_installed() -> &'static Vec<&'static str> {
static INSTALLED: OnceLock<Vec<&'static str>> = OnceLock::new();
INSTALLED.get_or_init(|| {
ACP_AGENTS
.iter()
.filter(|agent| agent.requires.iter().all(|bin| binary_in_path(bin)))
.map(|agent| agent.name)
.collect()
})
}
#[cfg(any(not(unix), test))]
pub(crate) fn ensure_path_within(root: &Path, requested: &Path) -> Result<PathBuf, AcpError> {
let canonical_root = root.canonicalize().map_err(|e| {
acp_error(format!(
"session workspace '{}' is not accessible: {e}",
root.display()
))
})?;
let inside = requested
.strip_prefix(root)
.ok()
.or_else(|| requested.strip_prefix(&canonical_root).ok());
let Some(relative) = inside else {
return Err(outside_workspace_error(requested));
};
if relative
.components()
.any(|c| c == std::path::Component::ParentDir)
{
return Err(outside_workspace_error(requested));
}
let mut resolved = canonical_root.clone();
for component in relative.components() {
let candidate = resolved.join(component);
match candidate.symlink_metadata() {
Ok(meta) if meta.file_type().is_symlink() => {
let target = candidate
.canonicalize()
.map_err(|_| outside_workspace_error(requested))?;
if !target.starts_with(&canonical_root) {
return Err(outside_workspace_error(requested));
}
resolved = target;
}
Ok(_) => resolved = candidate,
Err(_) => resolved = candidate,
}
}
if !resolved.starts_with(&canonical_root) {
return Err(outside_workspace_error(requested));
}
Ok(resolved)
}
pub(crate) fn outside_workspace_error(requested: &Path) -> AcpError {
acp_error(format!(
"path '{}' resolves outside the session workspace; refusing to serve it",
requested.display()
))
}
#[must_use]
pub(crate) fn slice_lines(content: &str, line: Option<u32>, limit: Option<u32>) -> String {
let start = line.map_or(0, |line| line.saturating_sub(1) as usize);
match limit {
None if start == 0 => content.to_string(),
limit => content
.lines()
.skip(start)
.take(limit.map_or(usize::MAX, |l| l as usize))
.collect::<Vec<_>>()
.join("\n"),
}
}
#[cfg(unix)]
fn serve_read(
root: &super::sandbox::PinnedRoot,
requested: &Path,
line: Option<u32>,
limit: Option<u32>,
) -> Result<String, AcpError> {
super::sandbox::read(root, requested, line, limit)
}
#[cfg(not(unix))]
fn serve_read(
root: &Path,
requested: &Path,
line: Option<u32>,
limit: Option<u32>,
) -> Result<String, AcpError> {
let path = ensure_path_within(root, requested)?;
let content = std::fs::read_to_string(&path)
.map_err(|e| acp_error(format!("failed to read {}: {e}", requested.display())))?;
Ok(slice_lines(&content, line, limit))
}
#[cfg(unix)]
fn serve_write(
root: &super::sandbox::PinnedRoot,
requested: &Path,
content: &str,
) -> Result<(), AcpError> {
super::sandbox::write(root, requested, content)
}
#[cfg(not(unix))]
fn serve_write(root: &Path, requested: &Path, content: &str) -> Result<(), AcpError> {
let path = ensure_path_within(root, requested)?;
if let Some(parent) = path.parent() {
let _ = std::fs::create_dir_all(parent);
}
std::fs::write(&path, content)
.map_err(|e| acp_error(format!("failed to write {}: {e}", requested.display())))
}
#[must_use]
pub(crate) fn stop_reason_str(reason: StopReason) -> String {
match reason {
StopReason::EndTurn => "end_turn".to_string(),
StopReason::MaxTokens => "max_tokens".to_string(),
StopReason::MaxTurnRequests => "max_turn_requests".to_string(),
StopReason::Refusal => "refusal".to_string(),
StopReason::Cancelled => "cancelled".to_string(),
other => format!("{other:?}"),
}
}
pub fn validate_agent(agent: &str) -> Result<(), String> {
if ACP_AGENTS.iter().any(|entry| entry.name == agent) {
Ok(())
} else {
let supported: Vec<&str> = ACP_AGENTS.iter().map(|a| a.name).collect();
Err(format!(
"Unsupported agent '{agent}'. Supported agents (ACP): {}. Use bash_run for shell commands.",
supported.join(", "),
))
}
}
#[allow(clippy::unwrap_used)]
pub async fn call(
agent: &str,
input: &str,
resume: Option<String>,
cwd: PathBuf,
stop_signal: Arc<AtomicBool>,
chunk_tx: tokio::sync::mpsc::UnboundedSender<SubagentEvent>,
) -> Result<(String, String, Option<String>), String> {
validate_agent(agent)?;
let entry = ACP_AGENTS
.iter()
.find(|a| a.name == agent)
.ok_or_else(|| format!("unknown agent '{agent}'"))?;
let launcher = agent_launcher(entry)?
.with_debug(|line, direction| {
if matches!(direction, LineDirection::Stderr) {
log::debug!("subagent ACP stderr: {line}");
}
});
let accumulated: Arc<Mutex<super::closure::TurnClosure>> =
Arc::new(Mutex::new(super::closure::TurnClosure::default()));
let input = input.to_string();
let turn = {
let accumulated = accumulated.clone();
tokio::task::spawn_blocking(move || {
tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.map_err(|e| format!("failed to build ACP runtime: {e}"))?
.block_on(run_session(
launcher,
input,
resume,
cwd,
entry.model,
accumulated,
stop_signal,
chunk_tx,
))
})
.await
.map_err(|e| format!("sub-agent task failed: {e}"))?
};
let final_output = accumulated.lock().unwrap().output();
match turn {
Ok((stop_reason, session_id)) => Ok((final_output, stop_reason, Some(session_id))),
Err(session_error) => {
if final_output.is_empty() {
Err(session_error)
} else {
log::warn!("sub-agent '{agent}' ACP turn failed: {session_error}");
Ok((final_output, "error".to_string(), None))
}
}
}
}
fn session_resume_support(capabilities: &AgentCapabilities) -> SessionResumeSupport {
if capabilities.session_capabilities.resume.is_some() {
SessionResumeSupport::Resume
} else if capabilities.load_session {
SessionResumeSupport::Load
} else {
SessionResumeSupport::None
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum SessionResumeSupport {
Resume,
Load,
None,
}
async fn open_or_resume_session(
connection: &ConnectionTo<AcpRole>,
resume: Option<String>,
cwd: &Path,
support: SessionResumeSupport,
) -> Result<(SessionId, Option<Vec<SessionConfigOption>>), AcpError> {
match (resume, support) {
(Some(id), SessionResumeSupport::Resume) => {
let response = connection
.send_request(ResumeSessionRequest::new(
SessionId::new(id.as_str()),
cwd.to_path_buf(),
))
.block_task()
.await;
match response.map(|r: ResumeSessionResponse| (r.config_options,)) {
Ok((config_options,)) => {
log::debug!("sub-agent ACP session resumed: {id}");
return Ok((SessionId::new(id), config_options));
}
Err(e) => {
log::warn!(
"sub-agent ACP resume of session '{id}' failed ({e}); \
falling back to session/new"
);
}
}
}
(Some(id), SessionResumeSupport::Load) => {
let response = connection
.send_request(LoadSessionRequest::new(
SessionId::new(id.as_str()),
cwd.to_path_buf(),
))
.block_task()
.await;
match response.map(|r: LoadSessionResponse| (r.config_options,)) {
Ok((config_options,)) => {
log::debug!("sub-agent ACP session loaded: {id}");
return Ok((SessionId::new(id), config_options));
}
Err(e) => {
log::warn!(
"sub-agent ACP load of session '{id}' failed ({e}); \
falling back to session/new"
);
}
}
}
(Some(_), SessionResumeSupport::None) => {
log::debug!(
"sub-agent harness advertises no session resume support; \
starting a fresh session"
);
}
(None, _) => {}
}
let session = connection
.send_request(NewSessionRequest::new(cwd.to_path_buf()))
.block_task()
.await
.map_err(|e| acp_error(format!("ACP session/new failed: {e}")))?;
Ok((session.session_id, session.config_options))
}
async fn select_session_model(
connection: &ConnectionTo<AcpRole>,
session_id: &SessionId,
config_options: Option<&[SessionConfigOption]>,
preferred_model: &str,
) -> Result<(), String> {
if preferred_model.is_empty() {
return Ok(());
}
let Some(option) = config_options
.unwrap_or_default()
.iter()
.find(|option| option.id.0.as_ref() == "model")
else {
return Ok(());
};
let SessionConfigKind::Select(select) = &option.kind else {
return Ok(());
};
if select.current_value.0.as_ref() == preferred_model {
return Ok(());
}
let offered: Vec<&str> = match &select.options {
SessionConfigSelectOptions::Ungrouped(options) => {
options.iter().map(|o| o.value.0.as_ref()).collect()
}
SessionConfigSelectOptions::Grouped(groups) => groups
.iter()
.flat_map(|g| g.options.iter().map(|o| o.value.0.as_ref()))
.collect(),
_ => Vec::new(),
};
if !offered.contains(&preferred_model) {
return Err(format!(
"ACP session model '{preferred_model}' is not offered by this harness; \
update the agent's `model` entry in the registry",
));
}
connection
.send_request(SetSessionConfigOptionRequest::new(
session_id.clone(),
SessionConfigId::new("model"),
SessionConfigValueId::new(preferred_model),
))
.block_task()
.await
.map_err(|e| format!("ACP set model failed: {e}"))?;
Ok(())
}
#[allow(clippy::unwrap_used)]
#[allow(clippy::too_many_arguments)]
pub(crate) async fn run_session<T>(
transport: T,
input: String,
resume: Option<String>,
cwd: PathBuf,
preferred_model: &'static str,
accumulated: Arc<Mutex<super::closure::TurnClosure>>,
stop_signal: Arc<AtomicBool>,
chunk_tx: tokio::sync::mpsc::UnboundedSender<SubagentEvent>,
) -> Result<(String, String), String>
where
T: agent_client_protocol::ConnectTo<agent_client_protocol::Client> + 'static,
{
#[cfg(unix)]
let read_root = super::sandbox::PinnedRoot::acquire(&cwd).map_err(|e| e.to_string())?;
#[cfg(unix)]
let write_root = super::sandbox::PinnedRoot::acquire(&cwd).map_err(|e| e.to_string())?;
#[cfg(not(unix))]
let read_root = cwd.clone();
#[cfg(not(unix))]
let write_root = cwd.clone();
Client
.builder()
.name("cosh")
.on_receive_notification(
async move |notification: SessionNotification, _cx| {
match SubagentEvent::from_session_update(notification.update) {
Some(event) => {
accumulated.lock().unwrap().observe(&event);
let _ = chunk_tx.send(event);
}
None => {
log::debug!("sub-agent session update without a display mapping");
}
}
Ok(())
},
on_receive_notification!(),
)
.on_receive_request(
async move |request: RequestPermissionRequest, responder, _connection| {
let selected = request
.options
.iter()
.find(|option| option.kind == PermissionOptionKind::AllowAlways)
.or_else(|| {
request
.options
.iter()
.find(|option| option.kind == PermissionOptionKind::AllowOnce)
})
.or_else(|| request.options.first());
match selected {
Some(option) => {
log::debug!(
"sub-agent permission auto-approved: option '{}' (kind {:?}) for tool call {}",
option.option_id.0,
option.kind,
request.tool_call.tool_call_id.0
);
responder.respond(RequestPermissionResponse::new(
RequestPermissionOutcome::Selected(SelectedPermissionOutcome::new(
option.option_id.clone(),
)),
))
}
None => responder.respond(RequestPermissionResponse::new(
RequestPermissionOutcome::Cancelled,
)),
}
},
on_receive_request!(),
)
.on_receive_request(
async move |request: ReadTextFileRequest, responder, _connection| {
let content = serve_read(
&read_root,
&request.path,
request.line,
request.limit,
)?;
responder.respond(ReadTextFileResponse::new(content))
},
on_receive_request!(),
)
.on_receive_request(
async move |request: WriteTextFileRequest, responder, _connection| {
serve_write(&write_root, &request.path, &request.content)?;
responder.respond(WriteTextFileResponse::new())
},
on_receive_request!(),
)
.connect_with(transport, async move |connection: ConnectionTo<AcpRole>| {
let capabilities = ClientCapabilities::new().fs(FileSystemCapabilities::new()
.read_text_file(true)
.write_text_file(true));
let init = connection
.send_request(
InitializeRequest::new(ProtocolVersion::V1).client_capabilities(capabilities),
)
.block_task()
.await
.map_err(|e| acp_error(format!("ACP initialize failed: {e}")))?;
let mut auth_error = None;
for method in &init.auth_methods {
match connection
.send_request(AuthenticateRequest::new(method.id().clone()))
.block_task()
.await
{
Ok(_) => {
auth_error = None;
break;
}
Err(e) => {
log::debug!(
"sub-agent ACP authenticate with method '{}' failed: {e}",
method.id().0
);
auth_error = Some(format!("ACP authenticate failed: {e}"));
}
}
}
if let Some(e) = auth_error {
return Err(acp_error(e));
}
let (session_id, config_options) =
open_or_resume_session(
&connection,
resume,
&cwd,
session_resume_support(&init.agent_capabilities),
)
.await?;
select_session_model(
&connection,
&session_id,
config_options.as_deref(),
preferred_model,
)
.await
.map_err(acp_error)?;
let prompt_request = connection
.send_request(PromptRequest::new(
session_id.clone(),
vec![ContentBlock::Text(TextContent::new(input))],
))
.block_task();
tokio::pin!(prompt_request);
let stop_wait = async {
while !stop_signal.load(Ordering::Relaxed) {
tokio::time::sleep(Duration::from_millis(50)).await;
}
};
let prompt = tokio::select! {
prompt = &mut prompt_request => prompt,
() = stop_wait => {
log::debug!("sub-agent ACP turn cancelled by stop signal");
if let Err(e) =
connection.send_notification(CancelNotification::new(session_id.clone()))
{
log::warn!("sub-agent ACP cancel notification failed: {e}");
}
(&mut prompt_request).await
}
};
let prompt = prompt
.map_err(|e| acp_error(format!("ACP session/prompt failed: {e}")))?;
Ok((
stop_reason_str(prompt.stop_reason),
session_id.0.to_string(),
))
})
.await
.map_err(|e| format!("ACP connection failed: {e}"))
}