use crate::{
agent::{mcp, serve},
model::{
media,
session_preferences::{self, Choices},
settings,
},
};
use anyhow::{Result, anyhow};
use cacp::{
AgentConn, Client, Direction, Error, Tap,
client::{HistoryEntry, HistoryFork, HistoryRole},
schema::{
AuthenticateRequest, CancelNotification, ClientCapabilities, ContentBlock, EnvVariable,
FileSystemCapabilities, HttpHeader, ImageContent, InitializeRequest, InitializeResponse,
LoadSessionRequest, McpServer, McpServerHttp, McpServerStdio, NewSessionRequest,
NewSessionResponse, PromptRequest, ReadTextFileRequest, ReadTextFileResponse,
RequestPermissionRequest, RequestPermissionResponse, SessionConfigOptionValue, SessionId,
SessionNotification, SessionUpdate, SetSessionConfigOptionRequest, SetSessionModeRequest,
StopReason, WriteTextFileRequest, WriteTextFileResponse,
},
};
use std::process::Stdio;
use std::{
path::PathBuf,
sync::{Arc, OnceLock},
time::Duration,
};
use tokio::{
io::{AsyncBufReadExt, BufReader},
process::{Child, ChildStderr, Command},
runtime::Runtime,
sync::{Mutex, mpsc, oneshot},
};
const DEBUG: &str = "CYDONIA_DEBUG";
const SERVER: &str = "cydonia";
const SHUTDOWN_GRACE: Duration = Duration::from_millis(250);
pub fn runtime() -> &'static Runtime {
static RUNTIME: OnceLock<Runtime> = OnceLock::new();
RUNTIME.get_or_init(|| Runtime::new().expect("failed to start the tokio runtime"))
}
pub enum Event {
Update(SessionUpdate),
Permission(RequestPermissionRequest, Reply<RequestPermissionResponse>),
Stderr(String),
TurnDone(Result<StopReason, Error>),
Closed,
}
pub type Events = mpsc::UnboundedReceiver<Event>;
pub type Sender = mpsc::UnboundedSender<Event>;
pub fn channel() -> (Sender, Events) {
mpsc::unbounded_channel()
}
pub struct Reply<T>(oneshot::Sender<Result<T, Error>>);
impl<T> Reply<T> {
pub fn send(self, value: T) {
let _ = self.0.send(Ok(value));
}
}
pub struct Session {
conn: Option<AgentConn>,
tx: mpsc::UnboundedSender<Event>,
child: Option<Child>,
pub session_id: SessionId,
pub init: InitializeResponse,
pub response: NewSessionResponse,
pub cwd: PathBuf,
pub loaded: bool,
built_in_mcp: bool,
history_fork: Option<Arc<Mutex<HistoryFork>>>,
}
#[derive(Default)]
pub struct Launch {
pub cwd: PathBuf,
pub previous: Option<String>,
pub history: Option<Vec<HistoryEntry>>,
pub history_pending: bool,
pub choices: Choices,
}
impl Launch {
pub fn new(cwd: PathBuf) -> Self {
Self {
cwd,
..Default::default()
}
}
}
impl Session {
pub async fn spawn(entry: &settings::Agent, launch: Launch, tx: Sender) -> Result<Self> {
let mut command = Command::new(&entry.command);
command.args(&entry.args).envs(&entry.env);
let bypass = loopback_bypass(["NO_PROXY", "no_proxy"].map(|key| {
entry
.env
.get(key)
.cloned()
.or_else(|| std::env::var(key).ok())
}));
command.env("NO_PROXY", &bypass).env("no_proxy", &bypass);
command.stderr(Stdio::piped());
let configured = mcp::servers();
let (conn, mut child) =
cacp::spawn(&mut command, Arc::new(Frontend(tx.clone())), debug_tap())
.map_err(|e| anyhow!("failed to start {}: {}", entry.command, error_text(&e)))?;
if let Some(stderr) = child.stderr.take() {
runtime().spawn(drain(stderr, tx.clone()));
}
Self::open(conn, child, tx, launch, configured).await
}
async fn open(
conn: AgentConn,
child: Child,
tx: mpsc::UnboundedSender<Event>,
launch: Launch,
configured: Vec<mcp::McpServer>,
) -> Result<Self> {
let cwd = launch.cwd.clone();
let init = conn
.initialize(InitializeRequest::new(ClientCapabilities {
fs: FileSystemCapabilities {
read_text_file: true,
write_text_file: true,
meta: None,
},
..Default::default()
}))
.await
.map_err(|e| anyhow!("initialize failed: {}", error_text(&e)))?;
let (mcp_servers, built_in_mcp) = acp_mcp_servers(&configured, &init, &cwd);
let mut loaded = false;
let mut response = None;
if let Some(id) = launch
.previous
.clone()
.filter(|_| init.agent_capabilities.load_session)
{
let load = || LoadSessionRequest {
mcp_servers: mcp_servers.clone(),
..LoadSessionRequest::new(id.clone(), cwd.clone())
};
let result = match conn.load_session(load()).await {
Err(e) if e.is_auth_required() => {
authenticate(&conn, &init).await?;
conn.load_session(load()).await
}
other => other,
};
if let Ok(load_response) = result {
response = Some(NewSessionResponse {
session_id: id.clone().into(),
modes: load_response.modes,
config_options: load_response.config_options,
meta: None,
});
loaded = true;
}
}
let response = match response {
Some(response) => response,
None => {
let new_session = || NewSessionRequest {
mcp_servers: mcp_servers.clone(),
..NewSessionRequest::new(cwd.clone())
};
match conn.new_session(new_session()).await {
Ok(response) => response,
Err(e) if e.is_auth_required() => {
authenticate(&conn, &init).await?;
conn.new_session(new_session()).await.map_err(|e| {
anyhow!(
"session/new failed after authentication: {}",
error_text(&e)
)
})?
}
Err(e) => return Err(anyhow!("session/new failed: {}", error_text(&e))),
}
}
};
let history_fork = if !loaded || launch.history_pending {
if let Some(mut history) = launch.history {
if init.agent_capabilities.prompt_capabilities.image {
history = tokio::task::spawn_blocking(move || {
for entry in &mut history {
if entry.role == HistoryRole::User {
let paths: Vec<_> = entry
.content
.iter()
.flat_map(|block| match block {
ContentBlock::Text(text) => media::attached(&text.text),
_ => Vec::new(),
})
.collect();
entry.content.extend(
paths.iter().filter_map(|path| media::encode(path)).map(
|(data, mime)| {
ContentBlock::Image(ImageContent {
data,
mime_type: mime.to_owned(),
uri: None,
annotations: None,
meta: None,
})
},
),
);
}
}
history
})
.await
.map_err(|e| anyhow!("fork history preparation failed: {e}"))?;
}
Some(Arc::new(Mutex::new(HistoryFork::restore(
conn.clone(),
response.clone(),
history,
))))
} else {
None
}
} else {
None
};
let mut session = Self {
conn: Some(conn),
tx,
child: Some(child),
session_id: response.session_id.clone(),
init,
response,
cwd,
loaded,
built_in_mcp,
history_fork,
};
session.restore_choices(&launch.choices).await?;
Ok(session)
}
async fn restore_choices(&mut self, choices: &Choices) -> Result<()> {
if let Some(mode) = &choices.mode
&& let Some(modes) = &self.response.modes
&& modes
.available_modes
.iter()
.any(|offered| offered.id.to_string() == *mode)
&& modes.current_mode_id.to_string() != *mode
{
self.conn()
.set_session_mode(SetSessionModeRequest {
session_id: self.session_id.clone(),
mode_id: mode.clone().into(),
meta: None,
})
.await
.map_err(|error| {
anyhow!("restoring session mode failed: {}", error_text(&error))
})?;
if let Some(modes) = &mut self.response.modes {
modes.current_mode_id = mode.clone().into();
}
}
let mut keys: Vec<_> = self
.response
.config_options
.as_deref()
.unwrap_or_default()
.iter()
.map(|option| {
(
option.category != Some(cacp::schema::SessionConfigOptionCategory::Model),
option.id.to_string(),
)
})
.collect();
keys.sort();
for (_, id) in keys {
let Some(value) = choices.config.get(&id) else {
continue;
};
let Some(option) = self
.response
.config_options
.as_deref()
.unwrap_or_default()
.iter()
.find(|option| option.id.to_string() == id)
else {
continue;
};
if !session_preferences::supports(option, value)
|| session_preferences::current(option) == *value
{
continue;
}
let response = self
.conn()
.set_session_config_option(SetSessionConfigOptionRequest {
session_id: self.session_id.clone(),
config_id: id.clone().into(),
value: value.clone(),
meta: None,
})
.await
.map_err(|error| {
anyhow!(
"restoring session option {id} failed: {}",
error_text(&error)
)
})?;
self.response.config_options = Some(response.config_options);
}
Ok(())
}
fn conn(&self) -> AgentConn {
self.conn.clone().expect("the session is being dropped")
}
pub fn prompt(&self, content: &str) {
let capabilities = &self.init.agent_capabilities.prompt_capabilities;
let pictures = match capabilities.image {
true => media::attached(content),
false => Vec::new(),
};
let mut blocks = super::context::prompt(
&self.cwd,
self.built_in_mcp,
capabilities.embedded_context,
vec![content.to_owned().into()],
);
let session_id = self.session_id.clone();
let conn = self.conn();
let tx = self.tx.clone();
let history_fork = self.history_fork.clone();
runtime().spawn(async move {
if !pictures.is_empty() {
let images = tokio::task::spawn_blocking(move || {
pictures
.iter()
.filter_map(|path| media::encode(path))
.map(|(data, mime)| {
ContentBlock::Image(ImageContent {
data,
mime_type: mime.to_owned(),
uri: None,
annotations: None,
meta: None,
})
})
.collect::<Vec<_>>()
})
.await
.unwrap_or_default();
blocks.splice(1..1, images);
}
let request = PromptRequest::new(session_id, blocks);
let result = match history_fork {
Some(fork) => fork.lock().await.prompt(request).await,
None => conn.prompt(request).await,
};
let done = result.map(|response| response.stop_reason);
let _ = tx.send(Event::TurnDone(done));
});
}
pub fn cancel(&self) -> Result<(), Error> {
self.conn().cancel(CancelNotification {
session_id: self.session_id.clone(),
meta: None,
})
}
pub fn set_mode(&self, mode_id: &str) {
let request = SetSessionModeRequest {
session_id: self.session_id.clone(),
mode_id: mode_id.into(),
meta: None,
};
let conn = self.conn();
runtime().spawn(async move { conn.set_session_mode(request).await });
}
pub fn set_config_option(&self, config_id: &str, value: SessionConfigOptionValue) {
let request = SetSessionConfigOptionRequest {
session_id: self.session_id.clone(),
config_id: config_id.into(),
value,
meta: None,
};
let conn = self.conn();
runtime().spawn(async move { conn.set_session_config_option(request).await });
}
}
pub fn history(items: &[artifact::session::chat::ChatItem]) -> Vec<HistoryEntry> {
use artifact::session::chat::ChatItem;
items
.iter()
.map(|item| {
let (role, text) = match item {
ChatItem::User(text) => (HistoryRole::User, text.clone()),
ChatItem::Agent(text) => (HistoryRole::Agent, text.clone()),
ChatItem::Thinking { text, .. } => {
(HistoryRole::Agent, format!("Prior reasoning: {text}"))
}
ChatItem::Tool { label, output, .. } => {
(HistoryRole::Tool, format!("{label}\n{output}"))
}
ChatItem::Process { command, output } => {
(HistoryRole::Tool, format!("{command}\n{output}"))
}
ChatItem::Notice { text, .. } => {
(HistoryRole::Tool, format!("Session notice: {text}"))
}
};
HistoryEntry {
role,
content: vec![text.into()],
}
})
.collect()
}
struct Frontend(mpsc::UnboundedSender<Event>);
impl Client for Frontend {
async fn session_update(&self, notification: SessionNotification) {
let _ = self.0.send(Event::Update(notification.update));
}
async fn request_permission(
&self,
request: RequestPermissionRequest,
) -> Result<RequestPermissionResponse, Error> {
let (tx, rx) = oneshot::channel();
self.0
.send(Event::Permission(request, Reply(tx)))
.map_err(|_| Error::internal_error().data("the frontend is gone"))?;
rx.await.unwrap_or_else(|_| Err(Error::method_not_found()))
}
async fn read_text_file(
&self,
request: ReadTextFileRequest,
) -> Result<ReadTextFileResponse, Error> {
read_text_file(&request)
}
async fn write_text_file(
&self,
request: WriteTextFileRequest,
) -> Result<WriteTextFileResponse, Error> {
std::fs::write(&request.path, &request.content)
.map(|()| WriteTextFileResponse::default())
.map_err(|e| io_error(&request.path, &e))
}
}
impl Drop for Session {
fn drop(&mut self) {
let (Some(conn), Some(mut child)) = (self.conn.take(), self.child.take()) else {
return;
};
drop(conn);
runtime().spawn(async move {
let _ = tokio::time::timeout(SHUTDOWN_GRACE, child.wait()).await;
});
}
}
impl Drop for Frontend {
fn drop(&mut self) {
let _ = self.0.send(Event::Closed);
}
}
const _: () = {
const fn assert_send<T: Send>() {}
assert_send::<Session>();
assert_send::<Event>();
};
fn loopback_bypass(existing: [Option<String>; 2]) -> String {
let mut entries = Vec::new();
for value in existing
.iter()
.flatten()
.map(String::as_str)
.chain(["localhost,127.0.0.1,::1"])
{
for entry in value
.split(',')
.map(str::trim)
.filter(|entry| !entry.is_empty())
{
if !entries.contains(&entry) {
entries.push(entry);
}
}
}
entries.join(",")
}
fn acp_mcp_servers(
configured: &[mcp::McpServer],
init: &InitializeResponse,
cwd: &std::path::Path,
) -> (Vec<McpServer>, bool) {
let http = init.agent_capabilities.mcp_capabilities.http;
let ours = http.then(serve::url).flatten().map(|url| {
McpServer::Http(McpServerHttp {
name: SERVER.to_owned(),
url,
headers: vec![{
let (name, value) = serve::project(cwd);
HttpHeader {
name: name.to_owned(),
value,
meta: None,
}
}],
meta: None,
})
});
let available = ours.is_some();
let servers = ours
.into_iter()
.chain(
configured
.iter()
.filter(|server| server.enabled)
.filter_map(|server| match (&server.command, &server.url) {
(Some(command), _) => Some(McpServer::Stdio(McpServerStdio {
name: server.name.clone(),
command: command.into(),
args: server.args.clone(),
env: server
.env
.iter()
.map(|(name, value)| EnvVariable {
name: name.clone(),
value: value.clone(),
meta: None,
})
.collect(),
meta: None,
})),
(None, Some(url)) if http => Some(McpServer::Http(McpServerHttp {
name: server.name.clone(),
url: url.clone(),
headers: Vec::new(),
meta: None,
})),
_ => None,
}),
)
.collect();
(servers, available)
}
async fn authenticate(conn: &AgentConn, init: &InitializeResponse) -> Result<()> {
if init.auth_methods.is_empty() {
return Err(anyhow!(
"authentication required, but the agent advertises no auth methods"
));
}
let mut failures = Vec::new();
for method in &init.auth_methods {
let request = AuthenticateRequest {
method_id: method.id().clone(),
meta: None,
};
match conn.authenticate(request).await {
Ok(_) => return Ok(()),
Err(e) => failures.push(format!("{}: {}", method.name(), error_text(&e))),
}
}
Err(anyhow!("authentication failed — {}", failures.join("; ")))
}
async fn drain(stderr: ChildStderr, tx: Sender) {
let echo = std::env::var_os(DEBUG).is_some();
let mut lines = BufReader::new(stderr).lines();
while let Ok(Some(line)) = lines.next_line().await {
if echo {
eprintln!("{line}");
}
if tx.send(Event::Stderr(line)).is_err() {
return;
}
}
}
pub fn spend_replay(events: &mut Events, tx: &Sender) {
let mut usage = None;
let mut kept = Vec::new();
while let Ok(event) = events.try_recv() {
match event {
Event::Update(SessionUpdate::UsageUpdate(update)) => usage = Some(update),
Event::Update(_) => {}
other => kept.push(other),
}
}
if let Some(update) = usage {
let _ = tx.send(Event::Update(SessionUpdate::UsageUpdate(update)));
}
for event in kept {
let _ = tx.send(event);
}
}
pub fn error_text(e: &Error) -> String {
match &e.data {
Some(data) => {
let detail = data
.as_str()
.map(str::to_owned)
.unwrap_or_else(|| data.to_string());
format!("{} — {detail}", e.message)
}
None => e.message.clone(),
}
}
fn read_text_file(request: &ReadTextFileRequest) -> Result<ReadTextFileResponse, Error> {
let content =
std::fs::read_to_string(&request.path).map_err(|e| io_error(&request.path, &e))?;
let content = match (request.line, request.limit) {
(None, None) => content,
(line, limit) => {
let skip = line.map(|l| l.saturating_sub(1) as usize).unwrap_or(0);
let take = limit.map(|l| l as usize).unwrap_or(usize::MAX);
content
.lines()
.skip(skip)
.take(take)
.collect::<Vec<_>>()
.join("\n")
}
};
Ok(ReadTextFileResponse {
content,
meta: None,
})
}
fn io_error(path: &std::path::Path, e: &std::io::Error) -> Error {
Error::internal_error().data(format!("{}: {e}", path.display()))
}
fn debug_tap() -> Option<Tap> {
let path = std::env::var(DEBUG).ok()?;
Some(Arc::new(move |direction: Direction, line: &str| {
use std::io::Write;
if let Ok(mut f) = std::fs::OpenOptions::new()
.create(true)
.append(true)
.open(&path)
{
let _ = writeln!(f, "{direction:?}: {line}");
}
}))
}
#[cfg(test)]
#[path = "../../tests/unit/acp_proxy.rs"]
mod proxy_tests;