use crate::{
agent::acp::{self, Event, Launch, Reply, Session},
model::{settings, workspace::Workspace},
view::component::transcript,
};
use anyhow::anyhow;
use artifact::{
project::{Project as _, fs},
session::{
chat::{ChatItem, PlanStatus, ToolStatus},
record::Record,
},
};
use bezel::gpui::{Context, Task};
use cacp::schema::{
ContentBlock, MaybeUndefined, PermissionOptionKind, PlanEntryStatus, RequestPermissionRequest,
RequestPermissionResponse, SessionConfigKind, SessionConfigOption, SessionConfigOptionValue,
SessionModeState, SessionUpdate, StopReason, ToolCallContent, ToolCallStatus,
};
use std::{
collections::VecDeque,
path::PathBuf,
time::{Duration, SystemTime, UNIX_EPOCH},
};
const STREAM_FRAME: Duration = Duration::from_millis(120);
const STDERR_KEEP: usize = 16 * 1024;
#[derive(Clone, Copy)]
pub struct Usage {
pub used: u64,
pub size: u64,
}
impl Usage {
pub fn fraction(&self) -> Option<f32> {
(self.size > 0).then(|| (self.used as f32 / self.size as f32).clamp(0., 1.))
}
}
#[derive(Clone, Copy)]
struct Flight {
at: SystemTime,
used: u64,
}
pub struct Choice {
pub id: String,
pub name: String,
pub kind: PermissionOptionKind,
}
pub enum Connection {
Idle,
Connecting,
Live(Box<Session>),
Lost,
}
#[derive(Clone, PartialEq, Eq)]
pub struct Command {
pub name: String,
pub description: String,
}
pub struct PermissionPrompt {
pub title: String,
pub options: Vec<Choice>,
pub always: bool,
reply: Reply<RequestPermissionResponse>,
}
pub struct ChatSession {
pub number: Option<u64>,
pub id: u64,
pub entry: settings::Agent,
pub cwd: PathBuf,
pub connection: Connection,
pub items: Vec<ChatItem>,
pub plan: Vec<(String, PlanStatus)>,
pub permission: Option<PermissionPrompt>,
pub commands: Vec<Command>,
pub modes: Option<SessionModeState>,
pub config: Vec<SessionConfigOption>,
pub usage: Option<Usage>,
flight: Option<Flight>,
pub title: String,
pub name: Option<String>,
pub updated: SystemTime,
pub agent_session: Option<String>,
pub record: Option<String>,
pub closed: bool,
pub streaming: bool,
pub queue: VecDeque<String>,
pub transcript: transcript::State,
_pump: Task<()>,
}
impl ChatSession {
pub fn connect(
id: u64,
entry: settings::Agent,
cwd: PathBuf,
seed: Option<String>,
cx: &mut Context<Workspace>,
) -> Self {
let pump = pump(id, &entry, cwd.clone(), None, cx);
Self {
id,
entry,
cwd,
connection: Connection::Connecting,
items: Vec::new(),
plan: Vec::new(),
permission: None,
commands: Vec::new(),
modes: None,
config: Vec::new(),
usage: None,
flight: None,
title: String::new(),
name: None,
updated: SystemTime::now(),
agent_session: None,
record: None,
number: None,
closed: false,
streaming: false,
queue: seed.into_iter().collect(),
transcript: transcript::State::default(),
_pump: pump,
}
}
pub fn restore(id: u64, cwd: PathBuf, entry: settings::Agent, record: Record) -> Self {
let updated = record.at();
Self {
id,
entry,
cwd,
connection: Connection::Idle,
items: record.items,
plan: Vec::new(),
permission: None,
commands: Vec::new(),
modes: None,
config: Vec::new(),
usage: None,
flight: None,
title: record.title,
name: record.name,
updated,
agent_session: record.session,
record: Some(record.id),
number: record.number,
closed: record.closed,
streaming: false,
queue: VecDeque::new(),
transcript: transcript::State::default(),
_pump: Task::ready(()),
}
}
pub fn touched(&self) -> u128 {
self.updated
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_millis()
}
fn to_record(&self) -> Record {
Record {
number: self.number,
id: self.record.clone().unwrap_or_default(),
agent: self.entry.name.clone(),
agent_id: self.entry.id.clone(),
session: self.agent_session.clone(),
title: self.title.clone(),
name: self.name.clone(),
updated: self
.updated
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_secs(),
closed: self.closed,
items: self.items.clone(),
}
}
pub fn mint_record(&mut self) -> Option<&str> {
if self.record.is_none() {
self.record = fs::Project::new(&self.cwd).create_session();
}
if self.number.is_none() {
self.number = self
.record
.as_deref()
.and_then(|id| artifact::entry::number(&self.cwd, "session", id).ok());
}
self.record.as_deref()
}
pub fn unsaid(&self) -> bool {
nothing_said(&self.items)
}
pub fn flush(&mut self) {
if self.unsaid() {
return;
}
let store = fs::Project::new(&self.cwd);
self.mint_record();
if self.record.is_some() {
store.save_session(&self.to_record());
}
}
pub fn resume(&mut self, cx: &mut Context<Workspace>) {
self._pump = pump(
self.id,
&self.entry,
self.cwd.clone(),
self.agent_session.clone(),
cx,
);
self.connection = Connection::Connecting;
self.closed = false;
}
pub fn close(&mut self) {
self.cancel();
self.connection = Connection::Idle;
self._pump = Task::ready(());
self.streaming = false;
self.closed = true;
self.queue.clear();
self.flush();
}
pub fn elapsed(&self) -> Option<Duration> {
self.flight?.at.elapsed().ok()
}
pub fn spent(&self) -> Option<u64> {
Some(self.usage?.used.saturating_sub(self.flight?.used))
}
pub fn live(&self) -> bool {
matches!(self.connection, Connection::Live(_))
}
pub fn idle(&self) -> bool {
matches!(self.connection, Connection::Idle | Connection::Lost)
}
pub fn resumable(&self) -> bool {
!self.entry.command.is_empty()
}
pub fn label(&self) -> String {
match (&self.name, self.title.is_empty()) {
(Some(name), _) => name.clone(),
(None, false) => self.title.clone(),
(None, true) => self.entry.name.clone(),
}
}
pub fn send(&mut self, content: String) {
if self.streaming || !self.live() {
self.queue.push_back(content);
} else {
self.prompt(content);
}
}
pub fn drain(&mut self) {
if self.streaming {
return;
}
if let Some(next) = self.queue.pop_front() {
self.prompt(next);
}
}
fn prompt(&mut self, content: String) {
let Connection::Live(session) = &self.connection else {
self.queue.push_front(content);
return;
};
session.prompt(&content);
self.flight = Some(Flight {
at: SystemTime::now(),
used: self.usage.map_or(0, |usage| usage.used),
});
self.items.push(ChatItem::User(content));
self.updated = SystemTime::now();
self.streaming = true;
self.flush();
}
pub fn cancel(&mut self) {
if let Some(prompt) = self.permission.take() {
prompt.reply.send(RequestPermissionResponse::cancelled());
}
if let Connection::Live(session) = &self.connection
&& let Err(e) = session.cancel()
{
self.notice(true, &format!("cancel failed: {}", acp::error_text(&e)));
}
}
pub fn set_mode(&mut self, mode_id: &str) {
let Connection::Live(session) = &self.connection else {
return;
};
session.set_mode(mode_id);
if let Some(modes) = &mut self.modes {
modes.current_mode_id = mode_id.into();
}
}
pub fn set_config(&mut self, config_id: &str, value: SessionConfigOptionValue) {
let Connection::Live(session) = &self.connection else {
return;
};
session.set_config_option(config_id, value.clone());
let found = self
.config
.iter_mut()
.find(|option| &*option.id == config_id);
let Some(option) = found else {
return;
};
match (&mut option.kind, value) {
(SessionConfigKind::Select(select), SessionConfigOptionValue::ValueId { value }) => {
select.current_value = value;
}
(SessionConfigKind::Boolean(flag), SessionConfigOptionValue::Boolean { value }) => {
flag.current_value = value;
}
_ => {}
}
}
pub fn respond_permission(&mut self, option_id: String) {
if let Some(prompt) = self.permission.take() {
prompt
.reply
.send(RequestPermissionResponse::selected(option_id));
}
}
pub fn toggle_permission_always(&mut self) {
if let Some(prompt) = &mut self.permission {
prompt.always = !prompt.always;
}
}
pub fn unqueue(&mut self, ix: usize) {
self.queue.remove(ix);
}
fn apply(&mut self, event: Event) {
self.updated = SystemTime::now();
match event {
Event::Update(update) => self.apply_update(update),
Event::Permission(request, reply) => self.open_permission(request, reply),
Event::Stderr(line) => self.stderr(line),
Event::TurnDone(result) => {
self.finish_thinking();
self.streaming = false;
match result {
Ok(StopReason::EndTurn) => {}
Ok(StopReason::Cancelled) => {
self.fail_running_tools();
self.notice(false, "cancelled");
}
Ok(StopReason::Refusal) => self.notice(false, "the agent refused to continue"),
Ok(StopReason::MaxTokens) => self.notice(false, "stopped: max tokens"),
Ok(StopReason::MaxTurnRequests) => {
self.notice(false, "stopped: max turn requests")
}
Ok(other) => self.notice(false, &format!("stopped: {other:?}")),
Err(e) => {
self.fail_running_tools();
self.notice(true, &format!("turn failed: {}", acp::error_text(&e)));
}
}
self.flush();
self.drain();
}
Event::Closed => {
self.connection = Connection::Lost;
self.streaming = false;
self.fail_running_tools();
self.notice(true, "agent connection lost");
self.flush();
}
}
}
fn apply_update(&mut self, update: SessionUpdate) {
match update {
SessionUpdate::AgentMessageChunk(chunk) => {
self.finish_thinking();
let text = content_text(&chunk.content);
if let Some(ChatItem::Agent(body)) = self.items.last_mut() {
body.push_str(&text);
} else {
self.items.push(ChatItem::Agent(text));
}
}
SessionUpdate::AgentThoughtChunk(chunk) => {
let text = content_text(&chunk.content);
if let Some(ChatItem::Thinking {
text: body,
done: false,
}) = self.items.last_mut()
{
body.push_str(&text);
} else {
self.items.push(ChatItem::Thinking { text, done: false });
}
}
SessionUpdate::ToolCall(call) => {
self.finish_thinking();
self.items.push(ChatItem::Tool {
id: call.tool_call_id.to_string(),
kind: call.kind,
label: call.title,
status: tool_status(call.status),
output: tool_content_text(&call.content),
});
}
SessionUpdate::ToolCallUpdate(update) => {
let id = update.tool_call_id.to_string();
let Some(ix) = self.items.iter().rposition(
|item| matches!(item, ChatItem::Tool { id: tool, .. } if *tool == id),
) else {
return;
};
let ChatItem::Tool {
kind,
label,
status,
output,
..
} = &mut self.items[ix]
else {
unreachable!("rposition matched a Tool item");
};
if let Some(title) = update.fields.title {
*label = title;
}
if let Some(new_kind) = update.fields.kind {
*kind = new_kind;
}
if let Some(content) = update.fields.content {
let text = tool_content_text(&content);
if !text.is_empty() {
if !output.is_empty() {
output.push('\n');
}
output.push_str(&text);
}
}
if let Some(new_status) = update.fields.status {
*status = tool_status(new_status);
}
}
SessionUpdate::SessionInfoUpdate(info) => {
let title = match info.title {
MaybeUndefined::Value(title) => title,
MaybeUndefined::Null => String::new(),
MaybeUndefined::Undefined => return,
};
if title != self.title {
self.title = title;
self.flush();
}
}
SessionUpdate::AvailableCommandsUpdate(cmds) => {
self.commands = cmds
.available_commands
.into_iter()
.map(|c| Command {
name: c.name,
description: c.description,
})
.collect();
}
SessionUpdate::Plan(plan) => {
self.plan = plan
.entries
.into_iter()
.map(|entry| {
let status = match entry.status {
PlanEntryStatus::Completed => PlanStatus::Done,
PlanEntryStatus::InProgress => PlanStatus::Active,
_ => PlanStatus::Pending,
};
(entry.content, status)
})
.collect();
}
SessionUpdate::CurrentModeUpdate(update) => {
if let Some(modes) = &mut self.modes {
modes.current_mode_id = update.current_mode_id;
}
}
SessionUpdate::ConfigOptionUpdate(update) => self.config = update.config_options,
SessionUpdate::UsageUpdate(update) => {
self.usage = Some(Usage {
used: update.used,
size: update.size,
});
}
SessionUpdate::UserMessageChunk(_) => {}
_ => {}
}
}
fn open_permission(
&mut self,
request: RequestPermissionRequest,
reply: Reply<RequestPermissionResponse>,
) {
let options: Vec<Choice> = request
.options
.into_iter()
.map(|opt| Choice {
id: opt.option_id.to_string(),
name: opt.name,
kind: opt.kind,
})
.collect();
if options.is_empty() {
reply.send(RequestPermissionResponse::cancelled());
return;
}
if let Some(previous) = self.permission.take() {
previous.reply.send(RequestPermissionResponse::cancelled());
}
let title = request
.tool_call
.fields
.title
.clone()
.unwrap_or_else(|| "Permission required".to_owned());
self.permission = Some(PermissionPrompt {
title,
options,
always: false,
reply,
});
}
pub(crate) fn stderr(&mut self, line: String) {
if let Some(ChatItem::Process { output, .. }) = self.items.last_mut() {
output.push('\n');
output.push_str(&line);
trim_front(output, STDERR_KEEP);
return;
}
self.items.push(ChatItem::Process {
command: command_line(&self.entry),
output: line,
});
}
pub(crate) fn notice(&mut self, failed: bool, text: &str) {
self.items.push(ChatItem::Notice {
text: text.to_owned(),
failed,
});
}
fn finish_thinking(&mut self) {
if let Some(ChatItem::Thinking { done, .. }) = self.items.last_mut() {
*done = true;
}
}
fn fail_running_tools(&mut self) {
for item in &mut self.items {
if let ChatItem::Tool { status, .. } = item
&& *status == ToolStatus::Running
{
*status = ToolStatus::Failure;
}
}
}
}
pub fn nothing_said(items: &[ChatItem]) -> bool {
items
.iter()
.all(|item| matches!(item, ChatItem::Process { .. }))
}
pub fn command_line(entry: &settings::Agent) -> String {
std::iter::once(entry.command.as_str())
.chain(entry.args.iter().map(String::as_str))
.collect::<Vec<_>>()
.join(" ")
}
pub fn trim_front(text: &mut String, keep: usize) {
if text.len() <= keep {
return;
}
let mut from = text.len() - keep;
while !text.is_char_boundary(from) {
from += 1;
}
let cut = match text[from..].find('\n') {
Some(at) => from + at + 1,
None => from,
};
text.drain(..cut);
}
fn tool_status(status: ToolCallStatus) -> ToolStatus {
match status {
ToolCallStatus::Completed => ToolStatus::Success,
ToolCallStatus::Failed => ToolStatus::Failure,
_ => ToolStatus::Running,
}
}
fn tool_content_text(content: &[ToolCallContent]) -> String {
content
.iter()
.map(|c| match c {
ToolCallContent::Content { content } => content_text(content),
ToolCallContent::Diff(diff) => format!("edited {}", diff.path.display()),
ToolCallContent::Terminal { .. } => "[terminal]".to_owned(),
_ => "[content]".to_owned(),
})
.collect::<Vec<_>>()
.join("\n")
}
fn content_text(block: &ContentBlock) -> String {
match block {
ContentBlock::Text(text) => text.text.clone(),
_ => "[non-text content]".to_owned(),
}
}
fn pump(
id: u64,
entry: &settings::Agent,
cwd: PathBuf,
previous: Option<String>,
cx: &mut Context<Workspace>,
) -> Task<()> {
let entry = entry.clone();
let (tx, mut events) = acp::channel();
let echo = tx.clone();
let conn = acp::runtime().spawn(async move {
Session::spawn(
&entry,
Launch {
previous,
..Launch::new(cwd)
},
tx,
)
.await
});
cx.spawn(async move |this, cx| {
let opened = conn
.await
.unwrap_or_else(|e| Err(anyhow!("the connection task panicked: {e}")));
let session = match opened {
Ok(session) => session,
Err(e) => {
cx.background_executor().timer(STREAM_FRAME).await;
let mut said = Vec::new();
while let Ok(event) = events.try_recv() {
if let Event::Stderr(line) = event {
said.push(line);
}
}
let _ = this.update(cx, |workspace, cx| {
workspace.with_session(id, cx, |chat| {
chat.connection = Connection::Lost;
for line in said {
chat.stderr(line);
}
chat.notice(true, &format!("connection failed: {e:#}"));
});
});
return;
}
};
if session.loaded {
acp::spend_replay(&mut events, &echo);
}
drop(echo);
if this
.update(cx, |workspace, cx| {
workspace.with_session(id, cx, |chat| {
chat.agent_session = Some(session.session_id.to_string());
chat.modes = session.response.modes.clone();
chat.config = session.response.config_options.clone().unwrap_or_default();
chat.connection = Connection::Live(Box::new(session));
});
workspace.session_connected(id, cx);
})
.is_err()
{
return;
}
while let Some(event) = events.recv().await {
let mut batch = vec![event];
while let Ok(event) = events.try_recv() {
batch.push(event);
}
let streaming = this.update(cx, |workspace, cx| {
workspace.with_session(id, cx, |chat| {
for event in batch {
chat.apply(event);
}
});
workspace.session(id).is_some_and(|chat| chat.streaming)
});
match streaming {
Ok(true) => cx.background_executor().timer(STREAM_FRAME).await,
Ok(false) => {}
Err(_) => return,
}
}
})
}