use crate::{
agent::acp::{self, Event, Launch, Reply, Session},
model::{
notify,
session_preferences::{self, Choices},
settings,
workspace::Workspace,
},
view::component::transcript,
};
use anyhow::anyhow;
use artifact::{
project::{Project as _, fs},
session::{
chat::{ChatItem, PlanStatus, ToolStatus},
record::{ForkOrigin, 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::{BTreeMap, VecDeque},
path::PathBuf,
time::{Duration, Instant, 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>,
history_unloaded: bool,
pub fork: Option<ForkOrigin>,
pub draft: String,
pub(crate) draft_save: Option<Task<()>>,
pub sent_at: BTreeMap<usize, u64>,
pub plan: Vec<(String, PlanStatus)>,
pub permission: Option<PermissionPrompt>,
pub commands: Vec<Command>,
pub modes: Option<SessionModeState>,
pub config: Vec<SessionConfigOption>,
preferences: Choices,
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>,
written: bool,
pub closed: bool,
pub streaming: bool,
pub(crate) last_activity: Instant,
pub(crate) turn_started: Option<Instant>,
pub(crate) tool_started: BTreeMap<String, Instant>,
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 preferences = session_preferences::load(&cwd, &entry, None);
let record = fs::Project::new(&cwd).create_session();
let pump = pump(
id,
&entry,
Launch {
choices: preferences.clone(),
record: record.clone(),
..Launch::new(cwd.clone())
},
cx,
);
Self {
id,
entry,
cwd,
connection: Connection::Connecting,
items: Vec::new(),
history_unloaded: false,
fork: None,
draft: String::new(),
draft_save: None,
sent_at: BTreeMap::new(),
plan: Vec::new(),
permission: None,
commands: Vec::new(),
modes: None,
config: Vec::new(),
preferences,
usage: None,
flight: None,
title: String::new(),
name: None,
updated: SystemTime::now(),
agent_session: None,
record,
written: false,
number: None,
closed: false,
streaming: false,
last_activity: Instant::now(),
turn_started: None,
tool_started: BTreeMap::new(),
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();
let preferences = session_preferences::load(&cwd, &entry, Some(&record.id));
Self {
id,
entry,
cwd,
connection: Connection::Idle,
items: record.items,
history_unloaded: false,
fork: record.fork,
draft: record.draft,
draft_save: None,
sent_at: record.sent_at,
plan: Vec::new(),
permission: None,
commands: Vec::new(),
modes: None,
config: Vec::new(),
preferences,
usage: None,
flight: None,
title: record.title,
name: record.name,
updated,
agent_session: record.session,
record: Some(record.id),
written: true,
number: record.number,
closed: record.closed,
streaming: false,
last_activity: Instant::now(),
turn_started: None,
tool_started: BTreeMap::new(),
queue: VecDeque::new(),
transcript: transcript::State::default(),
_pump: Task::ready(()),
}
}
pub fn load_history(&mut self) -> bool {
if !self.history_unloaded {
return true;
}
let Some(record) = self
.record
.as_deref()
.and_then(|id| fs::Project::new(&self.cwd).session(id))
else {
return false;
};
self.items = record.items;
self.sent_at = record.sent_at;
self.fork = record.fork;
self.draft = record.draft;
self.history_unloaded = false;
true
}
pub fn unload_history(&mut self) {
if !self.closed || self.history_unloaded || self.draft_save.is_some() {
return;
}
let Some(record) = self
.record
.as_deref()
.and_then(|id| fs::Project::new(&self.cwd).session(id))
else {
return;
};
if serde_json::to_value(&record).ok() != serde_json::to_value(self.to_record()).ok() {
return;
}
self.items = Vec::new();
self.sent_at = BTreeMap::new();
self.fork = None;
self.draft = String::new();
self.transcript = transcript::State::default();
self.plan = Vec::new();
self.commands = Vec::new();
self.config = Vec::new();
self.modes = None;
self.permission = None;
self.usage = None;
self.tool_started = BTreeMap::new();
self.history_unloaded = true;
}
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(),
fork: self.fork.clone(),
draft: self.draft.clone(),
sent_at: self.sent_at.clone(),
}
}
pub fn fork_at(&self, id: u64, before: usize) -> Option<Self> {
let source = if self.history_unloaded {
fs::Project::new(&self.cwd).session(self.record.as_deref()?)?
} else {
self.to_record()
};
let record = source.fork_at(before)?;
let mut chat = Self::restore(id, self.cwd.clone(), self.entry.clone(), record);
chat.record = None;
chat.written = false;
chat.preferences = self.preferences.clone();
Some(chat)
}
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 filed(&self) -> Option<&str> {
self.record.as_deref().filter(|_| self.written)
}
pub fn unsaid(&self) -> bool {
!self.history_unloaded && nothing_said(&self.items)
}
pub(crate) fn retain_panel(&mut self) -> Option<String> {
if !self.load_history() {
return None;
}
self.mint_record()?;
fs::Project::new(&self.cwd).save_session(&self.to_record());
self.written = true;
self.save_preferences();
self.record.clone()
}
pub fn flush(&mut self) {
if !self.load_history() {
return;
}
if self.unsaid() && self.fork.is_none() {
return;
}
let store = fs::Project::new(&self.cwd);
self.mint_record();
if self.record.is_some() {
store.save_session(&self.to_record());
self.written = true;
self.save_preferences();
}
}
pub fn resume(&mut self, cx: &mut Context<Workspace>) {
if self.closed || !self.load_history() {
return;
}
if self.record.is_none() {
self.record = fs::Project::new(&self.cwd).create_session();
}
self._pump = pump(
self.id,
&self.entry,
Launch {
previous: self.agent_session.clone(),
history: self.fork.as_ref().map(|fork| {
let end = if fork.pending {
fork.before.min(self.items.len())
} else {
self.items.len()
};
acp::history(&self.items[..end])
}),
history_pending: self.fork.as_ref().is_some_and(|fork| fork.pending),
choices: self.preferences.clone(),
record: self.record.clone(),
..Launch::new(self.cwd.clone())
},
cx,
);
self.connection = Connection::Connecting;
self.last_activity = Instant::now();
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.closed {
return;
}
self.updated = SystemTime::now();
if self.streaming || !self.live() {
self.queue.push_back(content);
self.flush();
} else {
self.prompt(content);
}
}
pub fn drain(&mut self) {
if self.closed || 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.last_activity = Instant::now();
self.turn_started = Some(self.last_activity);
self.tool_started.clear();
self.flight = Some(Flight {
at: SystemTime::now(),
used: self.usage.map_or(0, |usage| usage.used),
});
self.sent_at.insert(
self.items.len(),
SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_secs(),
);
self.items.push(ChatItem::User(content));
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;
};
if !self.modes.as_ref().is_some_and(|modes| {
modes
.available_modes
.iter()
.any(|mode| &*mode.id == mode_id)
}) {
return;
}
session.set_mode(mode_id);
if let Some(modes) = &mut self.modes {
modes.current_mode_id = mode_id.into();
}
self.preferences.mode = Some(mode_id.into());
if let Err(error) = session_preferences::remember_mode(&self.cwd, &self.entry, mode_id) {
self.notice(true, &format!("Could not save session default: {error}"));
}
self.save_preferences();
}
pub fn set_config(&mut self, config_id: &str, value: SessionConfigOptionValue) {
let Connection::Live(session) = &self.connection else {
return;
};
if !self
.config
.iter()
.any(|option| &*option.id == config_id && session_preferences::supports(option, &value))
{
return;
}
session.set_config_option(config_id, value.clone());
self.preferences
.config
.insert(config_id.into(), value.clone());
let saved =
session_preferences::remember_config(&self.cwd, &self.entry, 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;
}
_ => {}
}
if let Err(error) = saved {
self.notice(true, &format!("Could not save session default: {error}"));
}
self.save_preferences();
}
fn save_preferences(&mut self) {
if let Some(record) = &self.record
&& let Err(error) = session_preferences::remember_session(
&self.cwd,
&self.entry,
record,
&self.preferences,
)
{
self.notice(true, &format!("Could not save session choices: {error}"));
}
}
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) {
if self.closed {
return;
}
self.last_activity = Instant::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) => {
if result.is_ok()
&& let Some(fork) = &mut self.fork
{
fork.pending = false;
}
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.tool_started
.entry(call.tool_call_id.to_string())
.or_insert_with(Instant::now);
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;
}
self.preferences = Choices::capture(self.modes.as_ref(), &self.config);
self.save_preferences();
}
SessionUpdate::ConfigOptionUpdate(update) => {
self.config = update.config_options;
self.preferences = Choices::capture(self.modes.as_ref(), &self.config);
self.save_preferences();
}
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 use artifact::session::chat::nothing_said;
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, launch: Launch, cx: &mut Context<Workspace>) -> Task<()> {
let entry = entry.clone();
let (tx, mut events) = acp::channel();
let echo = tx.clone();
let mut conn = PendingConnection(
acp::runtime().spawn(async move { Session::spawn(&entry, launch, tx).await }),
);
cx.spawn(async move |this, cx| {
let opened = (&mut conn.0)
.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| {
if chat.closed {
return;
}
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.preferences = Choices::capture(chat.modes.as_ref(), &chat.config);
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| {
let was_streaming = workspace.session(id).is_some_and(|chat| chat.streaming);
workspace.with_session(id, cx, |chat| {
for event in batch {
chat.apply(event);
}
});
let streaming = workspace.session(id).is_some_and(|chat| chat.streaming);
if was_streaming && !streaming {
let notify = workspace.settings.notify_turns;
if let Some(chat) = workspace.session(id) {
notify::turn_finished(chat, notify, cx);
}
}
streaming
});
match streaming {
Ok(true) => cx.background_executor().timer(STREAM_FRAME).await,
Ok(false) => {}
Err(_) => return,
}
}
})
}
struct PendingConnection(tokio::task::JoinHandle<anyhow::Result<Session>>);
impl Drop for PendingConnection {
fn drop(&mut self) {
self.0.abort();
}
}
#[cfg(test)]
#[path = "../../tests/unit/session.rs"]
mod tests;