use std::io;
use std::path::{Path, PathBuf};
use std::time::{Duration, Instant};
use chrono::{DateTime, Local, Utc};
use ratatui::Terminal;
use ratatui::backend::CrosstermBackend;
use ratatui::crossterm::SynchronizedUpdate;
use ratatui::crossterm::event::Event as TerminalEvent;
use ratatui::crossterm::event::KeyEventKind;
use ratatui::crossterm::execute;
use ratatui::crossterm::style::Print;
use tokio::io::AsyncWriteExt as _;
use tokio::time::MissedTickBehavior;
use super::TranscriptTone;
use super::TuiState;
use super::clipboard::{
ClipboardPreparation, ClipboardUploads, UploadCandidate, prepare_clipboard,
};
use super::events::{handle_gateway_event, handle_gateway_history};
use super::input::UiAction;
use super::view::render_preview;
use crate::frontend::FrontendExit;
use crate::frontend::bots;
use crate::frontend::catalog::{GatewayAction, UiCatalog};
use crate::frontend::dashboard::render_capability_overlay;
use crate::frontend::extensions;
use crate::frontend::gateway;
use crate::frontend::gateway_actions::{ResponseSeverity, prepare, render_response};
use crate::frontend::setup;
use crate::frontend::terminal::{INPUT_POLL, MAX_INPUT_BATCH, TerminalGuard, poll_event};
use mobius::backend::checkpoint::ExecutionOutcome;
use mobius::protocol::{
EventMsg, FrontendBlockFormat, FrontendEvent, ModelInfo, Op, SessionFileReference, Submission,
};
use mobius::{Error, Result};
use mobius_gateway::client::{GatewayEvents, GatewaySender};
use mobius_gateway::wire::{
BotRecord, ClientMessage, ReadyPayload, ServerMessage, SessionActivityState,
SessionReadyPayload, SessionRecord, WorkspaceFileScope,
};
use uuid::Uuid;
const ELAPSED_INTERVAL: Duration = Duration::from_secs(1);
const CLEAR_SCREEN_AND_SCROLLBACK: &str = "\x1b[r\x1b[0m\x1b[H\x1b[2J\x1b[3J\x1b[H";
const FILE_READ_CHUNK_BYTES: usize = 256 * 1024;
type TuiTerminal = Terminal<CrosstermBackend<io::Stdout>>;
struct ReplayHydration {
pending: bool,
}
#[derive(Debug, PartialEq, Eq)]
struct PendingSessionCreation {
request_id: String,
clear: bool,
}
impl ReplayHydration {
const fn pending() -> Self {
Self { pending: true }
}
const fn allows_draw(&self) -> bool {
!self.pending
}
fn observe(&mut self, message: &ServerMessage, session_id: &str) {
if matches!(
message,
ServerMessage::SessionReplayComplete {
session_id: actual,
..
} if actual == session_id
) {
self.pending = false;
}
}
fn finish(&mut self) {
self.pending = false;
}
}
pub(in crate::frontend) async fn run(
sender: GatewaySender,
mut events: GatewayEvents,
gateway: &mut ReadyPayload,
session: &mut SessionReadyPayload,
mut catalog: UiCatalog,
gateway_endpoint: String,
choose_initial_bot: bool,
) -> Result<(FrontendExit, GatewaySender, GatewayEvents)> {
let mut guard = TerminalGuard::alternate()?;
let mut terminal = TuiTerminal::new(CrosstermBackend::new(io::stdout()))?;
terminal.clear()?;
let session_id = session.session.session_id.clone();
let mut state = TuiState::new(
&catalog,
catalog.workspace().to_path_buf(),
ModelInfo::default(),
String::new(),
String::new(),
);
sync_session(&mut state, session, gateway)?;
if choose_initial_bot {
choose_bot(
gateway,
session,
session.workspace.path.clone(),
false,
&mut state,
);
}
let mut uploads = ClipboardUploads::default();
let mut tick = tokio::time::interval(INPUT_POLL);
tick.set_missed_tick_behavior(MissedTickBehavior::Skip);
let mut elapsed = tokio::time::interval(ELAPSED_INTERVAL);
elapsed.set_missed_tick_behavior(MissedTickBehavior::Skip);
let mut events_open = true;
let mut dirty = true;
let mut replay_hydration = ReplayHydration::pending();
let exit;
let mut clear_on_exit = false;
let mut pending_session_creation = None;
let mut clipboard_preparation = None;
let mut workspace_reference_open = false;
request_workspace_inventory(&sender, &session_id, &mut state).await;
'ui: loop {
draw_if_dirty(
&mut terminal,
&mut state,
&catalog,
&mut dirty,
&replay_hydration,
)?;
tokio::select! {
event = events.next(), if events_open => {
match event {
Ok(Some(frame)) => {
replay_hydration.observe(&frame.message, &session_id);
let message = match frame.message {
ServerMessage::WorkspaceFiles { session_id: actual, files, .. }
if actual == session_id => {
catalog.set_workspace_paths(files.into_iter().map(|file| file.path));
state.reference_cache = None;
dirty = true;
continue 'ui;
}
message => message,
};
if replay_hydration.allows_draw()
&& refresh_workspace_inventory(&message, &session_id, workspace_reference_open)
{
request_workspace_inventory(&sender, &session_id, &mut state).await;
}
if let Some((next_exit, should_clear)) = handle_incoming_message(
message,
&sender,
gateway,
session,
&mut state,
&mut uploads,
&mut pending_session_creation,
)
.await
{
clear_on_exit |= should_clear;
exit = next_exit;
break 'ui;
}
}
Ok(None) => {
disconnect(
&mut state,
&mut uploads,
&mut clipboard_preparation,
&mut replay_hydration,
"gateway disconnected · press q to exit",
);
events_open = false;
}
Err(error) => {
disconnect(
&mut state,
&mut uploads,
&mut clipboard_preparation,
&mut replay_hydration,
error.to_string(),
);
events_open = false;
}
}
dirty = true;
}
_ = tick.tick() => {
if let Some(next_exit) = handle_terminal_input(
&mut terminal,
&sender,
&mut events,
gateway,
session,
&catalog,
&gateway_endpoint,
&session_id,
&mut state,
&mut uploads,
&mut pending_session_creation,
&mut clipboard_preparation,
&mut dirty,
).await? {
exit = next_exit;
break 'ui;
}
let reference_open = !state.reference_menu_dismissed
&& super::references::active_reference_token(&state.input, state.cursor, '@').is_some();
if events_open && reference_open && !workspace_reference_open {
request_workspace_inventory(&sender, &session_id, &mut state).await;
}
workspace_reference_open = reference_open;
}
result = async {
clipboard_preparation
.as_mut()
.expect("clipboard preparation is active")
.await
}, if clipboard_preparation.is_some() => {
clipboard_preparation = None;
handle_clipboard_preparation(
result,
&sender,
&session_id,
&mut state,
&mut uploads,
).await;
dirty = true;
}
_ = elapsed.tick(), if state.active_turn().is_some() => {
dirty = true;
}
}
guard.set_mouse_capture(
state
.preview
.as_ref()
.is_none_or(|preview| matches!(preview.content, super::PreviewContent::Diff(_))),
)?;
}
if let Some(preparation) = clipboard_preparation.take() {
drop(preparation);
}
if let Some(download) = state.file_download.take() {
drop(download.output);
let _ = tokio::fs::remove_file(download.destination).await;
}
drop(terminal);
drop(guard);
if clear_on_exit {
execute!(io::stdout(), Print(CLEAR_SCREEN_AND_SCROLLBACK))?;
}
Ok((exit, sender, events))
}
async fn request_workspace_inventory(
sender: &GatewaySender,
session_id: &str,
state: &mut TuiState,
) {
if let Err(error) = sender
.send(ClientMessage::ListWorkspaceFiles {
request_id: Uuid::new_v4().to_string(),
session_id: session_id.into(),
scope: WorkspaceFileScope::All,
})
.await
{
state.push(error.to_string(), TranscriptTone::Error);
}
}
fn refresh_workspace_inventory(
message: &ServerMessage,
session_id: &str,
reference_open: bool,
) -> bool {
matches!(message, ServerMessage::AgentEvent { session_id: actual, record }
if actual == session_id && (matches!(record.event.msg, EventMsg::TurnComplete(_) | EventMsg::TurnAborted(_))
|| reference_open && matches!(record.event.msg, EventMsg::ToolCallEnd(_))))
}
#[expect(
clippy::too_many_arguments,
reason = "terminal input dispatch keeps the active frontend state explicit"
)]
async fn handle_terminal_input(
terminal: &mut TuiTerminal,
sender: &GatewaySender,
events: &mut GatewayEvents,
gateway: &mut ReadyPayload,
session: &mut SessionReadyPayload,
catalog: &UiCatalog,
gateway_endpoint: &str,
session_id: &str,
state: &mut TuiState,
uploads: &mut ClipboardUploads,
pending_session_creation: &mut Option<PendingSessionCreation>,
clipboard_preparation: &mut Option<ClipboardPreparation>,
dirty: &mut bool,
) -> Result<Option<FrontendExit>> {
for _ in 0..MAX_INPUT_BATCH {
let Some(event) = poll_event()? else { break };
let action = terminal_action(event, state, catalog, dirty);
match action {
UiAction::None => {}
UiAction::PasteClipboard => {
paste_clipboard(gateway, state, uploads, clipboard_preparation);
}
UiAction::Exit => {
interrupt_active_turn(sender, session_id, state).await;
return Ok(Some(FrontendExit::Exit));
}
UiAction::ChooseBot { workspace, clear } => {
choose_bot(gateway, session, workspace, clear, state);
*dirty = true;
}
UiAction::CreateSession {
workspace,
bot_id,
clear,
} => {
*pending_session_creation =
create_session(sender, workspace, bot_id, clear, state).await;
}
UiAction::Submit(op) => send_and_report(sender, session_id, op, state).await,
UiAction::Resume(session_id) => return Ok(Some(FrontendExit::Resume(session_id))),
UiAction::Gateway(action) => {
send_gateway_action(sender, session_id, action, state).await
}
UiAction::Download(file) => {
start_file_download(sender, session_id, file, state).await;
}
UiAction::Events => {
state.open_background_approvals(&gateway.background_approvals, &gateway.bots);
}
UiAction::ReassignBot => {
state.open_reassign_bot_picker(&gateway.bots, &session.session.context.owner_id);
}
UiAction::Branches => state.open_branch_picker(session.git.as_ref()),
UiAction::ConfirmDelete => state.confirm_delete_current_session(),
UiAction::ConfirmDeleteFile(file) => state.confirm_delete_file(file),
UiAction::Queued => state.open_queued_messages(),
UiAction::LoadEarlierHistory => {
request_earlier_history(sender, session_id, state).await;
}
UiAction::GatewaySettings => {
if open_gateway_settings(terminal, gateway_endpoint, state).await {
return Ok(Some(FrontendExit::Reconnect));
}
*dirty = true;
}
UiAction::Extensions => {
open_extensions(terminal, sender, events, gateway, session, state).await;
*dirty = true;
}
UiAction::Bots => {
if let Some(exit) = open_bots(
terminal, sender, events, gateway, session, state, session_id,
)
.await
{
return Ok(Some(exit));
}
*dirty = true;
}
UiAction::Setup { mode, provider } => {
let (result, session_changed) = run_setup(
terminal,
mode,
provider.as_deref(),
sender,
events,
gateway,
session,
)
.await;
if session_changed {
return Ok(Some(FrontendExit::Resume(
session.session.session_id.clone(),
)));
}
if let Err(error) = sync_session_info(state, session, gateway) {
state.push(error.to_string(), TranscriptTone::Error);
}
if let Err(error) = result {
state.push(error.to_string(), TranscriptTone::Error);
}
*dirty = true;
}
}
}
Ok(None)
}
async fn handle_incoming_message(
message: ServerMessage,
sender: &GatewaySender,
gateway: &mut ReadyPayload,
session: &mut SessionReadyPayload,
state: &mut TuiState,
uploads: &mut ClipboardUploads,
pending_session_creation: &mut Option<PendingSessionCreation>,
) -> Option<(FrontendExit, bool)> {
let session_id = session.session.session_id.clone();
if settle_session_deletion(&message, state) {
return Some((FrontendExit::Fresh, true));
}
let session_opened = match &message {
ServerMessage::SessionOpened { request_id, .. } => {
matches_pending_session_creation(request_id, pending_session_creation)
}
_ => false,
};
let should_clear = settle_session_creation(&message, pending_session_creation);
if !handle_file_download_message(&message, sender, state).await
&& !handle_upload_message(&message, sender, gateway, state, uploads, &session_id).await
&& let Some(exit) = handle_server_message(
message,
gateway,
session,
state,
&session_id,
session_opened,
)
{
return Some((exit, should_clear));
}
state
.requested_resume
.take()
.map(|request| (FrontendExit::Resume(request.session_id), true))
}
fn settle_session_deletion(message: &ServerMessage, state: &mut TuiState) -> bool {
let Some(request_id) = state.deletion_request_id.as_deref() else {
return false;
};
match message {
ServerMessage::Accepted { request_id: actual } if actual == request_id => {
state.deletion_request_id = None;
true
}
ServerMessage::Rejected {
request_id: actual, ..
} if actual == request_id => {
state.deletion_request_id = None;
false
}
_ => false,
}
}
fn matches_pending_session_creation(
request_id: &str,
pending: &Option<PendingSessionCreation>,
) -> bool {
pending
.as_ref()
.is_some_and(|request| request.request_id == request_id)
}
fn choose_bot(
gateway: &ReadyPayload,
session: &SessionReadyPayload,
workspace: std::path::PathBuf,
clear: bool,
state: &mut TuiState,
) {
if gateway.bots.is_empty() {
state.push(
"create a Bot before starting another chat",
TranscriptTone::Warning,
);
return;
}
state.open_bot_picker(
&gateway.bots,
workspace,
&session.session.context.owner_id,
clear,
);
}
async fn create_session(
sender: &GatewaySender,
workspace: std::path::PathBuf,
bot_id: String,
clear: bool,
state: &mut TuiState,
) -> Option<PendingSessionCreation> {
let request_id = Uuid::new_v4().to_string();
let result = sender
.send(ClientMessage::CreateSession {
request_id: request_id.clone(),
workspace,
bot_id,
})
.await;
if let Err(error) = result {
state.push(error.to_string(), TranscriptTone::Error);
return None;
}
Some(PendingSessionCreation { request_id, clear })
}
fn settle_session_creation(
message: &ServerMessage,
pending: &mut Option<PendingSessionCreation>,
) -> bool {
let Some(request) = pending.as_ref() else {
return false;
};
let clear = match message {
ServerMessage::SessionOpened { request_id, .. } if request_id == &request.request_id => {
request.clear
}
ServerMessage::Rejected { request_id, .. } if request_id == &request.request_id => false,
_ => return false,
};
*pending = None;
clear
}
fn draw_if_dirty(
terminal: &mut TuiTerminal,
state: &mut TuiState,
catalog: &UiCatalog,
dirty: &mut bool,
replay_hydration: &ReplayHydration,
) -> Result<()> {
if *dirty && replay_hydration.allows_draw() {
draw(terminal, state, catalog)?;
*dirty = false;
}
Ok(())
}
fn draw(terminal: &mut TuiTerminal, state: &mut TuiState, catalog: &UiCatalog) -> Result<()> {
io::stdout().sync_update(|_| -> Result<()> {
terminal.draw(|frame| {
super::view::render(frame, state, catalog);
if state.preview.is_some() {
render_preview(frame, state);
} else if let Some(overlay) = state.capability_overlay.as_mut() {
render_capability_overlay(frame, overlay);
}
})?;
Ok(())
})??;
Ok(())
}
async fn handle_upload_message(
message: &ServerMessage,
sender: &GatewaySender,
gateway: &ReadyPayload,
state: &mut TuiState,
uploads: &mut ClipboardUploads,
session_id: &str,
) -> bool {
let Some(result) = uploads
.handle(message, session_id, &gateway.session_file_limits)
.await
else {
return false;
};
match result {
Ok(advance) => {
if let Some(attachment) = advance.attachment {
state.attachments.push(attachment);
}
if let Some(message) = advance.message
&& let Err(error) = sender.send(message).await
{
uploads.abort();
state.push(error.to_string(), TranscriptTone::Error);
}
}
Err(error) => state.push(error, TranscriptTone::Error),
}
state.upload_in_progress = uploads.is_active();
true
}
fn handle_server_message(
message: ServerMessage,
gateway: &mut ReadyPayload,
session: &mut SessionReadyPayload,
state: &mut TuiState,
session_id: &str,
session_opened: bool,
) -> Option<FrontendExit> {
match message {
ServerMessage::AgentEvent {
session_id: actual,
mut record,
} if actual == session_id => {
enrich_resume_picker(
&mut record.event.msg,
&gateway.sessions,
&gateway.bots,
&session.workspace.id,
);
let live = record.sequence > session.latest_sequence;
handle_gateway_event(state, record, live);
}
ServerMessage::SessionHistory {
request_id,
session_id: actual,
mut records,
next_before_sequence,
} if actual == session_id
&& state.history_request_id.as_deref() == Some(request_id.as_str()) =>
{
state.history_request_id = None;
state.next_before_sequence = next_before_sequence;
session.next_before_sequence = next_before_sequence;
for record in &mut records {
enrich_resume_picker(
&mut record.event.msg,
&gateway.sessions,
&gateway.bots,
&session.workspace.id,
);
}
handle_gateway_history(state, records);
}
ServerMessage::SessionOpened { payload, .. } if session_opened => {
*session = payload;
return Some(FrontendExit::Reload);
}
ServerMessage::Ready { payload } => {
*gateway = payload;
if let Err(error) = sync_session_info(state, session, gateway) {
state.push(error.to_string(), TranscriptTone::Error);
}
}
ServerMessage::Sessions { sessions, .. } => gateway.sessions = sessions,
ServerMessage::BackgroundApprovals { approvals } => {
gateway.background_approvals = approvals;
state.background_approval_count = gateway.background_approvals.len();
}
ServerMessage::Bots { bots, .. } => {
gateway.bots = bots;
if let Err(error) = sync_session_info(state, session, gateway) {
state.push(error.to_string(), TranscriptTone::Error);
}
}
ServerMessage::SessionChanged { payload } if payload.session.session_id == session_id => {
if payload.workspace.id == session.workspace.id
&& payload.contributions == session.contributions
{
if let Err(error) = refresh_session(state, session, payload, gateway) {
state.push(error.to_string(), TranscriptTone::Error);
}
} else {
*session = payload;
return Some(FrontendExit::Resume(session_id.into()));
}
}
ServerMessage::SessionFiles {
session_id: actual,
files,
..
} if actual == session_id => state.open_session_files(files),
ServerMessage::GitDiff {
session_id: actual,
scope,
diff,
..
} if actual == session_id => {
let title = format!("{scope:?} diff").to_ascii_lowercase();
state.open_diff_preview(title, diff);
}
ServerMessage::Rejected {
request_id,
message,
..
} if state.history_request_id.as_deref() == Some(request_id.as_str()) => {
state.history_request_id = None;
state.push(message, TranscriptTone::Error);
}
message => {
if let Some(response) = render_response(&message, &gateway.provider_instances) {
let tone = match response.severity {
ResponseSeverity::Neutral => TranscriptTone::Neutral,
ResponseSeverity::Error | ResponseSeverity::Fatal => TranscriptTone::Error,
};
if response.severity == ResponseSeverity::Fatal {
state.disconnected = true;
}
state.push(response.text, tone);
}
}
}
None
}
fn disconnect(
state: &mut TuiState,
uploads: &mut ClipboardUploads,
clipboard_preparation: &mut Option<ClipboardPreparation>,
replay_hydration: &mut ReplayHydration,
message: impl AsRef<str>,
) {
if let Some(preparation) = clipboard_preparation.take() {
drop(preparation);
}
uploads.abort();
state.upload_in_progress = false;
replay_hydration.finish();
state.disconnected = true;
for turn_id in state.active_turns.clone() {
state.finish_turn(&turn_id);
}
state.push(message, TranscriptTone::Error);
}
fn terminal_action(
event: TerminalEvent,
state: &mut TuiState,
catalog: &UiCatalog,
dirty: &mut bool,
) -> UiAction {
match event {
TerminalEvent::Key(key) => {
*dirty |= matches!(key.kind, KeyEventKind::Press | KeyEventKind::Repeat);
state.handle_key(key, catalog)
}
TerminalEvent::Paste(text) => {
if state.capability_overlay.is_some() {
*dirty |= state.insert_capability_overlay_paste(&text);
} else if state.preview.is_none() && state.picker.is_none() {
let before = (state.input.len(), state.input_limit_reached);
state.insert_paste(&text);
*dirty |= before != (state.input.len(), state.input_limit_reached);
}
UiAction::None
}
TerminalEvent::Resize(_, _) => {
*dirty = true;
UiAction::None
}
TerminalEvent::Mouse(mouse) => {
*dirty |= state.handle_mouse(mouse);
UiAction::None
}
TerminalEvent::FocusGained | TerminalEvent::FocusLost => UiAction::None,
}
}
fn paste_clipboard(
gateway: &ReadyPayload,
state: &mut TuiState,
uploads: &mut ClipboardUploads,
clipboard_preparation: &mut Option<ClipboardPreparation>,
) {
if state.composer_target_turn().is_some() {
state.push(
"files can be pasted when the agent is idle",
TranscriptTone::Warning,
);
return;
}
if state.disconnected {
state.push("gateway is disconnected", TranscriptTone::Error);
return;
}
if state.upload_in_progress || uploads.is_active() || clipboard_preparation.is_some() {
state.push(
"an attachment upload is already in progress",
TranscriptTone::Warning,
);
return;
}
let existing = state.attachments.clone();
let limits = gateway.session_file_limits;
match prepare_clipboard(existing, limits) {
Ok(preparation) => {
*clipboard_preparation = Some(preparation);
state.upload_in_progress = true;
}
Err(error) => state.push(error, TranscriptTone::Error),
}
}
async fn handle_clipboard_preparation(
result: std::result::Result<
std::result::Result<Vec<UploadCandidate>, String>,
tokio::sync::oneshot::error::RecvError,
>,
sender: &GatewaySender,
session_id: &str,
state: &mut TuiState,
uploads: &mut ClipboardUploads,
) {
let candidates = match result {
Ok(Ok(candidates)) => candidates,
Ok(Err(error)) => {
state.upload_in_progress = false;
state.push(error, TranscriptTone::Error);
return;
}
Err(error) => {
state.upload_in_progress = false;
state.push(
format!("clipboard preparation worker stopped: {error}"),
TranscriptTone::Error,
);
return;
}
};
let message = match uploads.start(candidates, session_id) {
Ok(message) => message,
Err(error) => {
state.upload_in_progress = false;
state.push(error, TranscriptTone::Error);
return;
}
};
if let Err(error) = sender.send(message).await {
uploads.abort();
state.upload_in_progress = false;
state.push(error.to_string(), TranscriptTone::Error);
}
}
async fn interrupt_active_turn(sender: &GatewaySender, session_id: &str, state: &TuiState) {
if let Some(turn_id) = state.active_turn().map(str::to_owned) {
let _ = send_op(sender, session_id, Op::Interrupt { turn_id }).await;
}
}
async fn send_and_report(sender: &GatewaySender, session_id: &str, op: Op, state: &mut TuiState) {
let requests_preview = matches!(op, Op::CapabilityCommand { .. });
if requests_preview {
state.preview_request_id = None;
}
match send_op(sender, session_id, op).await {
Ok(id) if requests_preview => state.preview_request_id = Some(id),
Ok(_) => {}
Err(error) => state.push(error.to_string(), TranscriptTone::Error),
}
}
async fn send_gateway_action(
sender: &GatewaySender,
session_id: &str,
action: GatewayAction,
state: &mut TuiState,
) {
match action {
GatewayAction::DeleteCurrent => {
let request_id = Uuid::new_v4().to_string();
let result = sender
.send(ClientMessage::DeleteSessions {
request_id: request_id.clone(),
session_ids: vec![session_id.into()],
})
.await;
match result {
Ok(()) => state.deletion_request_id = Some(request_id),
Err(error) => state.push(error.to_string(), TranscriptTone::Error),
}
}
action => {
if let Err(error) = sender.send(*prepare(action, session_id)).await {
state.push(error.to_string(), TranscriptTone::Error);
}
}
}
}
async fn request_earlier_history(sender: &GatewaySender, session_id: &str, state: &mut TuiState) {
let Some(before_sequence) = state.next_before_sequence else {
return;
};
let request_id = Uuid::new_v4().to_string();
match sender
.send(ClientMessage::GetSessionHistory {
request_id: request_id.clone(),
session_id: session_id.into(),
before_sequence: Some(before_sequence),
})
.await
{
Ok(()) => state.history_request_id = Some(request_id),
Err(error) => state.push(error.to_string(), TranscriptTone::Error),
}
}
async fn start_file_download(
sender: &GatewaySender,
session_id: &str,
file: SessionFileReference,
state: &mut TuiState,
) {
if state.file_download.is_some() {
state.push(
"wait for the current file download to finish",
TranscriptTone::Warning,
);
return;
}
let directory = match std::env::current_dir() {
Ok(directory) => directory,
Err(error) => {
state.push(error.to_string(), TranscriptTone::Error);
return;
}
};
let (destination, output) = match create_download_file(&directory, &file.name).await {
Ok(download) => download,
Err(error) => {
state.push(error.to_string(), TranscriptTone::Error);
return;
}
};
let request_id = Uuid::new_v4().to_string();
let request = ClientMessage::ReadSessionFile {
request_id: request_id.clone(),
session_id: session_id.into(),
file_id: file.id.clone(),
offset: 0,
max_bytes: FILE_READ_CHUNK_BYTES,
};
if let Err(error) = sender.send(request).await {
drop(output);
let _ = tokio::fs::remove_file(destination).await;
state.push(error.to_string(), TranscriptTone::Error);
return;
}
state.file_download = Some(super::PendingFileDownload {
request_id,
session_id: session_id.into(),
file,
destination,
output,
offset: 0,
preview: Vec::new(),
preview_truncated: false,
});
}
async fn create_download_file(
directory: &Path,
requested_name: &str,
) -> io::Result<(PathBuf, tokio::fs::File)> {
let requested = Path::new(requested_name)
.file_name()
.filter(|name| !name.is_empty())
.unwrap_or_else(|| std::ffi::OsStr::new("download"));
let path = Path::new(requested);
let stem = path.file_stem().unwrap_or(requested).to_string_lossy();
let extension = path.extension().map(|value| value.to_string_lossy());
for suffix in 0_u32..=u32::MAX {
let name = if suffix == 0 {
requested.to_os_string()
} else if let Some(extension) = &extension {
format!("{stem}-{suffix}.{extension}").into()
} else {
format!("{stem}-{suffix}").into()
};
let destination = directory.join(name);
match tokio::fs::OpenOptions::new()
.create_new(true)
.write(true)
.open(&destination)
.await
{
Ok(file) => return Ok((destination, file)),
Err(error) if error.kind() == io::ErrorKind::AlreadyExists => {}
Err(error) => return Err(error),
}
}
Err(io::Error::new(
io::ErrorKind::AlreadyExists,
"no available download filename",
))
}
async fn handle_file_download_message(
message: &ServerMessage,
sender: &GatewaySender,
state: &mut TuiState,
) -> bool {
let Some(download) = state.file_download.as_ref() else {
return false;
};
let matches_request = match message {
ServerMessage::SessionFileChunk {
request_id,
session_id,
file_id,
..
} => {
request_id == &download.request_id
&& session_id == &download.session_id
&& file_id == &download.file.id
}
ServerMessage::Rejected { request_id, .. } => request_id == &download.request_id,
_ => false,
};
if !matches_request {
return false;
}
match message {
ServerMessage::Rejected { message, .. } => {
fail_file_download(state, message).await;
}
ServerMessage::SessionFileChunk {
offset,
data,
next_offset,
..
} => {
if let Err(error) = write_file_chunk(state, *offset, data, *next_offset).await {
fail_file_download(state, &error.to_string()).await;
return true;
}
if let Some(offset) = next_offset {
let request_id = Uuid::new_v4().to_string();
let Some(download) = state.file_download.as_mut() else {
return true;
};
download.request_id.clone_from(&request_id);
let request = ClientMessage::ReadSessionFile {
request_id,
session_id: download.session_id.clone(),
file_id: download.file.id.clone(),
offset: *offset,
max_bytes: FILE_READ_CHUNK_BYTES,
};
if let Err(error) = sender.send(request).await {
fail_file_download(state, &error.to_string()).await;
}
} else {
finish_file_download(state).await;
}
}
_ => unreachable!("matched file response"),
}
true
}
async fn write_file_chunk(
state: &mut TuiState,
offset: u64,
data: &[u8],
next_offset: Option<u64>,
) -> io::Result<()> {
let download = state.file_download.as_mut().ok_or_else(|| {
io::Error::new(
io::ErrorKind::NotFound,
"session file download is not active",
)
})?;
if offset != download.offset {
return Err(io::Error::new(
io::ErrorKind::InvalidData,
"session file chunk offset did not match the requested offset",
));
}
let end = offset
.checked_add(u64::try_from(data.len()).map_err(io::Error::other)?)
.ok_or_else(|| io::Error::new(io::ErrorKind::InvalidData, "file offset overflow"))?;
if data.is_empty() && next_offset.is_some() {
return Err(io::Error::new(
io::ErrorKind::InvalidData,
"session file chunk did not advance the download",
));
}
if end > download.file.size || next_offset.is_some_and(|next| next != end) {
return Err(io::Error::new(
io::ErrorKind::InvalidData,
"session file chunk exceeded the advertised file size",
));
}
if next_offset.is_none() && end != download.file.size {
return Err(io::Error::new(
io::ErrorKind::UnexpectedEof,
"session file download ended before the advertised file size",
));
}
download.output.write_all(data).await?;
let available = super::MAX_ENTRY_BYTES.saturating_sub(download.preview.len());
download
.preview
.extend_from_slice(&data[..data.len().min(available)]);
download.preview_truncated |= data.len() > available;
download.offset = end;
Ok(())
}
async fn finish_file_download(state: &mut TuiState) {
let Some(mut download) = state.file_download.take() else {
return;
};
if let Err(error) = download.output.flush().await {
drop(download.output);
let _ = tokio::fs::remove_file(&download.destination).await;
state.push(error.to_string(), TranscriptTone::Error);
return;
}
drop(download.output);
state.push(
format!("saved {}", download.destination.display()),
TranscriptTone::Success,
);
if let Ok(mut text) = String::from_utf8(download.preview) {
if download.preview_truncated {
text.push_str("\n\n[preview truncated]");
}
state.open_text_preview(
download.file.name,
text,
FrontendBlockFormat::PlainText,
TranscriptTone::Neutral,
);
}
}
async fn fail_file_download(state: &mut TuiState, message: &str) {
if let Some(download) = state.file_download.take() {
drop(download.output);
let _ = tokio::fs::remove_file(download.destination).await;
}
state.push(message, TranscriptTone::Error);
}
async fn open_gateway_settings(
terminal: &mut TuiTerminal,
gateway_endpoint: &str,
state: &mut TuiState,
) -> bool {
match gateway::run(terminal, gateway_endpoint).await {
Ok(reconnect) => reconnect,
Err(error) => {
state.push(error.to_string(), TranscriptTone::Error);
false
}
}
}
async fn open_extensions(
terminal: &mut TuiTerminal,
sender: &GatewaySender,
events: &mut GatewayEvents,
gateway: &mut ReadyPayload,
session: &SessionReadyPayload,
state: &mut TuiState,
) {
let result = extensions::run(terminal, sender, events, gateway).await;
if let Err(error) = sync_session_info(state, session, gateway) {
state.push(error.to_string(), TranscriptTone::Error);
}
if let Err(error) = result {
state.push(error.to_string(), TranscriptTone::Error);
}
}
async fn open_bots(
terminal: &mut TuiTerminal,
sender: &GatewaySender,
events: &mut GatewayEvents,
gateway: &mut ReadyPayload,
session: &SessionReadyPayload,
state: &mut TuiState,
session_id: &str,
) -> Option<FrontendExit> {
let bot_id = session.session.context.owner_id.clone();
let result = bots::run(
terminal,
sender,
events,
gateway,
Some(&bot_id),
Some(&bot_id),
)
.await;
if !gateway
.sessions
.iter()
.any(|candidate| candidate.session_id == session_id)
{
return Some(FrontendExit::Fresh);
}
if let Err(error) = sync_session_info(state, session, gateway) {
state.push(error.to_string(), TranscriptTone::Error);
}
match result {
Ok(Some(session_id)) => return Some(FrontendExit::Resume(session_id)),
Ok(None) => {}
Err(error) => state.push(error.to_string(), TranscriptTone::Error),
}
None
}
async fn run_setup(
terminal: &mut TuiTerminal,
mode: setup::SetupMode,
provider: Option<&str>,
sender: &GatewaySender,
events: &mut GatewayEvents,
gateway: &mut ReadyPayload,
session: &mut SessionReadyPayload,
) -> (Result<()>, bool) {
let workspace = session.workspace.id.clone();
let selected = session.session.session_id.clone();
let contributions = session.contributions.clone();
let result = setup::run(terminal, mode, provider, sender, events, gateway, session).await;
let changed = session.workspace.id != workspace
|| session.session.session_id != selected
|| session.contributions != contributions;
(result, changed)
}
fn refresh_session(
state: &mut TuiState,
session: &mut SessionReadyPayload,
payload: SessionReadyPayload,
gateway: &ReadyPayload,
) -> Result<()> {
sync_session(state, &payload, gateway)?;
*session = payload;
Ok(())
}
fn sync_session(
state: &mut TuiState,
session: &SessionReadyPayload,
gateway: &ReadyPayload,
) -> Result<()> {
sync_session_info(state, session, gateway)?;
state.active_turns.clone_from(&session.active_turn_ids);
state.approvals.clear();
state.restore_draft();
state.turn_started_at = (!state.active_turns.is_empty()).then(Instant::now);
for request in &session.pending_approvals {
state.begin_approval(request.clone());
}
Ok(())
}
fn sync_session_info(
state: &mut TuiState,
session: &SessionReadyPayload,
gateway: &ReadyPayload,
) -> Result<()> {
let bot = session_bot(gateway, session)?;
state.model.model = super::terminal_text(&session.session.model.model);
state.model.reasoning_effort = session
.session
.model
.reasoning_effort
.as_deref()
.map(super::terminal_text);
state.model_route.clone_from(&session.session.model.route);
state.agent_summary = agent_summary(gateway, session, bot);
state.active_message_delivery = Some(session.active_message_delivery);
state.context_limit = session.context_limit_tokens;
state.usage.apply_context_limit(state.context_limit);
state.background_approval_count = gateway.background_approvals.len();
state.next_before_sequence = session.next_before_sequence;
Ok(())
}
fn enrich_resume_picker(
event: &mut EventMsg,
sessions: &[SessionRecord],
bots: &[BotRecord],
current_workspace_id: &str,
) {
let EventMsg::Frontend(FrontendEvent::Picker { options, .. }) = event else {
return;
};
for option in options {
let Op::ResumeSession { session_id } = &option.op else {
continue;
};
let Some(session) = sessions
.iter()
.find(|session| session.session_id == *session_id)
else {
continue;
};
if let Some(title) = &session.title {
option.label.clone_from(title);
}
let mut details = vec![session_status(session).into()];
details.push(
if session.session_context.workspace_id.as_deref() == Some(current_workspace_id) {
"this workspace"
} else {
session
.session_context
.workspace_label
.as_deref()
.unwrap_or("other workspace")
}
.into(),
);
if let Some(origin) = &session.session_context.origin_label {
details.push(origin.clone());
}
if let Some(handle) = bot_handle(&session.session_context.owner_id, bots) {
details.push(format!("@{handle}"));
}
details.push(format!("started {}", human_time(session.created_at)));
option.description = details.join(" · ");
}
}
fn session_status(session: &SessionRecord) -> &'static str {
match session.activity.state {
SessionActivityState::Running => "running",
SessionActivityState::AwaitingApproval => "awaiting approval",
SessionActivityState::Idle => match session.activity.last_outcome {
Some(ExecutionOutcome::Completed) => "done",
Some(ExecutionOutcome::Aborted) => "aborted",
Some(ExecutionOutcome::Failed) => "failed",
None if session.execution_stats.run_count == 0 => "new",
None => "idle",
},
}
}
fn bot_handle<'a>(bot_id: &str, bots: &'a [BotRecord]) -> Option<&'a str> {
bots.iter()
.find(|bot| bot.id == bot_id)
.map(|bot| bot.handle.as_str())
}
fn session_bot<'a>(
gateway: &'a ReadyPayload,
session: &SessionReadyPayload,
) -> Result<&'a BotRecord> {
gateway
.bots
.iter()
.find(|bot| bot.id == session.session.context.owner_id)
.ok_or_else(|| {
Error::Config(format!(
"session {} references unknown Bot {}",
session.session.session_id, session.session.context.owner_id
))
})
}
fn human_time(timestamp_ms: i64) -> String {
DateTime::<Utc>::from_timestamp_millis(timestamp_ms).map_or_else(
|| "unknown time".into(),
|time| {
time.with_timezone(&Local)
.format("%Y-%m-%d %H:%M %Z")
.to_string()
},
)
}
fn agent_summary(gateway: &ReadyPayload, session: &SessionReadyPayload, bot: &BotRecord) -> String {
let bot_label = format!("@{}", bot.handle);
let providers = gateway
.provider_instances
.iter()
.filter(|entry| entry.configured)
.map(|entry| entry.label.as_str())
.collect::<Vec<_>>()
.join(", ");
let middleware = gateway
.middleware_features
.iter()
.filter(|feature| feature.required || bot.config.config.middleware.enabled(&feature.id))
.map(|feature| feature.label.as_str())
.collect::<Vec<_>>()
.join(", ");
let counts = session
.contributions
.iter()
.filter_map(|contribution| {
contribution.count.map(|count| {
let label = gateway
.middleware_features
.iter()
.find(|feature| feature.id == contribution.capability)
.map_or(contribution.capability.as_str(), |feature| {
feature.label.as_str()
});
format!("{label}: {count}")
})
})
.collect::<Vec<_>>()
.join(" · ");
let reasoning = session
.session
.model
.reasoning_effort
.as_deref()
.unwrap_or("default");
format!(
"MÖBIUS v{} · {bot_label}\nmodel: {} · {reasoning}\nproviders: {}\nmiddleware: {}\n{}tools: {}\nworkspace: {}",
env!("CARGO_PKG_VERSION"),
super::terminal_text(&session.session.model.model),
if providers.is_empty() {
"none"
} else {
&providers
},
if middleware.is_empty() {
"none"
} else {
&middleware
},
if counts.is_empty() {
String::new()
} else {
format!("{counts} · ")
},
session.tool_count,
super::terminal_text(&session.workspace.path.display().to_string()),
)
}
async fn send_op(
sender: &mobius_gateway::client::GatewaySender,
session_id: &str,
op: Op,
) -> Result<String> {
let id = Uuid::new_v4().to_string();
sender
.send(ClientMessage::Submit {
session_id: session_id.into(),
submission: Submission { id: id.clone(), op },
})
.await
.map(|()| id)
.map_err(|error| Error::Stopped(error.to_string()))
}
#[cfg(test)]
mod tests {
use super::*;
use mobius::protocol::{
ActiveMessageDelivery, Event, EventMsg, FrontendPickerOption, ModelChangedEvent,
SessionConfiguredEvent, SessionContext, SessionFileLimits, SessionFileOrigin,
SessionFileRecord,
};
use mobius_gateway::wire::{
BackgroundApproval, GitDiffScope, ReadyPayload, RecordedEvent, RoutineInteractionPolicy,
RunStats, SessionActivity, SessionReadyPayload, VersionedAgentConfig, WorkspaceInfo,
};
fn replay_event(sequence: u64) -> ServerMessage {
ServerMessage::AgentEvent {
session_id: "session-a".into(),
record: RecordedEvent {
sequence,
recorded_at_ms: 0,
event: Event {
submission_id: None,
msg: EventMsg::ContextCompacted,
},
stream_metrics: Vec::new(),
blocks: Vec::new(),
preview: None,
},
}
}
#[test]
fn workspace_inventory_refreshes_current_session_changes() {
use mobius::protocol::{ToolCallEndEvent, TurnAbortedEvent, TurnCompleteEvent};
let mut message = replay_event(1);
assert!(!refresh_workspace_inventory(&message, "session-a", true));
for event in [
EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: "turn-1".into(),
}),
EventMsg::TurnAborted(TurnAbortedEvent {
turn_id: "turn-1".into(),
reason: "cancelled".into(),
}),
EventMsg::ToolCallEnd(ToolCallEndEvent {
turn_id: "turn-1".into(),
call_id: "call-1".into(),
name: "write".into(),
output: String::new().into(),
is_error: false,
}),
] {
let ServerMessage::AgentEvent { record, .. } = &mut message else {
unreachable!()
};
let terminal = !matches!(event, EventMsg::ToolCallEnd(_));
record.event.msg = event;
assert!(refresh_workspace_inventory(&message, "session-a", true));
assert_eq!(
refresh_workspace_inventory(&message, "session-a", false),
terminal
);
assert!(!refresh_workspace_inventory(&message, "session-b", true));
}
}
#[test]
fn replay_hydration_blocks_draw_until_current_session_completes() {
let mut hydration = ReplayHydration::pending();
for sequence in 1..=3 {
hydration.observe(&replay_event(sequence), "session-a");
assert!(!hydration.allows_draw());
}
hydration.observe(
&ServerMessage::SessionReplayComplete {
request_id: "other-request".into(),
session_id: "session-b".into(),
},
"session-a",
);
assert!(!hydration.allows_draw());
hydration.observe(
&ServerMessage::SessionReplayComplete {
request_id: "open-request".into(),
session_id: "session-a".into(),
},
"session-a",
);
assert!(hydration.allows_draw());
hydration.observe(&replay_event(4), "session-a");
assert!(hydration.allows_draw());
}
#[tokio::test]
async fn downloads_do_not_overwrite_and_reject_invalid_chunks() {
let directory = tempfile::tempdir().expect("download directory");
let (first_path, first) = create_download_file(directory.path(), "report.txt")
.await
.expect("first download");
drop(first);
let (second_path, second) = create_download_file(directory.path(), "report.txt")
.await
.expect("second download");
assert_ne!(first_path, second_path);
let mut state = TuiState {
file_download: Some(super::super::PendingFileDownload {
request_id: "request".into(),
session_id: "session".into(),
file: SessionFileReference {
id: "file".into(),
name: "report.txt".into(),
size: 2,
media_type: "text/plain".into(),
},
destination: second_path.clone(),
output: second,
offset: 0,
preview: Vec::new(),
preview_truncated: false,
}),
..TuiState::default()
};
assert!(
write_file_chunk(&mut state, 0, b"too long", None)
.await
.is_err()
);
assert!(write_file_chunk(&mut state, 0, b"", Some(0)).await.is_err());
write_file_chunk(&mut state, 0, b"ok", None)
.await
.expect("valid chunk");
let download = state.file_download.take().expect("download state");
drop(download.output);
assert_eq!(
tokio::fs::read(second_path).await.expect("saved file"),
b"ok"
);
}
#[test]
fn rejected_session_creation_does_not_commit_clear() {
let mut pending = Some(PendingSessionCreation {
request_id: "create-1".into(),
clear: true,
});
let rejection = ServerMessage::Rejected {
request_id: "create-1".into(),
code: "invalid_bot".into(),
message: "Bot unavailable".into(),
fatal: false,
};
assert!(!settle_session_creation(&rejection, &mut pending));
assert_eq!(pending, None);
}
#[test]
fn stale_session_opened_does_not_reload_or_settle_new_creation() {
let mut pending = Some(PendingSessionCreation {
request_id: "create-new".into(),
clear: true,
});
let stale = ServerMessage::SessionOpened {
request_id: "create-old".into(),
payload: session_payload("session-old-response"),
};
let session_opened = match &stale {
ServerMessage::SessionOpened { request_id, .. } => {
matches_pending_session_creation(request_id, &pending)
}
_ => false,
};
let mut session = session_payload("session-current");
let mut gateway = ready_payload();
let catalog = UiCatalog::build(&[], std::path::Path::new("/tmp")).expect("catalog");
let mut state = TuiState::new(
&catalog,
"/tmp".into(),
ModelInfo::default(),
String::new(),
String::new(),
);
assert!(!session_opened);
assert!(!settle_session_creation(&stale, &mut pending));
assert!(
handle_server_message(
stale,
&mut gateway,
&mut session,
&mut state,
"session-current",
session_opened,
)
.is_none()
);
assert_eq!(session.session.session_id, "session-current");
assert_eq!(
pending.as_ref().map(|request| request.request_id.as_str()),
Some("create-new")
);
}
#[test]
fn session_snapshot_restores_current_work_and_approval() {
use mobius::protocol::{ExecApprovalRequestEvent, TurnCompleteEvent};
let catalog = UiCatalog::build(&[], std::path::Path::new("/tmp")).unwrap();
let mut state = TuiState::new(
&catalog,
"/tmp".into(),
ModelInfo::default(),
String::new(),
String::new(),
);
let mut session = session_payload("session-a");
session.next_before_sequence = Some(7);
session.active_turn_ids = vec!["turn-a".into()];
session.pending_approvals = vec![ExecApprovalRequestEvent {
id: "approval-a".into(),
turn_id: "turn-a".into(),
calls: Vec::new(),
reason: "approve work".into(),
}];
let mut gateway = ready_payload();
gateway.background_approvals.push(BackgroundApproval {
session_id: "background".into(),
bot_id: "bot-a".into(),
turn_id: "turn-b".into(),
request_id: "approval-b".into(),
});
sync_session(&mut state, &session, &gateway).unwrap();
assert!(state.active_turn().is_some());
assert_eq!(state.next_before_sequence, Some(7));
assert_eq!(state.background_approval_count, 1);
assert_eq!(
state.approval().map(|request| request.id.as_str()),
Some("approval-a")
);
state.handle_agent_event(
EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: "turn-a".into(),
}),
Vec::new(),
);
assert!(state.active_turn().is_none());
assert!(state.approval().is_none());
}
#[test]
fn history_and_phase_responses_update_the_current_chat() {
let mut state = TuiState {
history_request_id: Some("history".into()),
..TuiState::default()
};
let mut session = session_payload("session-a");
let mut gateway = ready_payload();
let ServerMessage::AgentEvent { record, .. } = replay_event(1) else {
unreachable!()
};
handle_server_message(
ServerMessage::SessionHistory {
request_id: "history".into(),
session_id: "session-a".into(),
records: vec![record],
next_before_sequence: Some(4),
},
&mut gateway,
&mut session,
&mut state,
"session-a",
false,
);
handle_server_message(
ServerMessage::BackgroundApprovals {
approvals: vec![BackgroundApproval {
session_id: "hidden".into(),
bot_id: "bot-a".into(),
turn_id: "turn".into(),
request_id: "approval".into(),
}],
},
&mut gateway,
&mut session,
&mut state,
"session-a",
false,
);
handle_server_message(
ServerMessage::GitDiff {
request_id: "diff".into(),
session_id: "session-a".into(),
scope: GitDiffScope::Staged,
diff: "+added".into(),
},
&mut gateway,
&mut session,
&mut state,
"session-a",
false,
);
assert_eq!(state.history_request_id, None);
assert_eq!(state.next_before_sequence, Some(4));
assert_eq!(session.next_before_sequence, Some(4));
assert_eq!(state.background_approval_count, 1);
assert!(matches!(
state.preview.as_ref().map(|preview| &preview.content),
Some(super::super::PreviewContent::Diff(_))
));
handle_server_message(
ServerMessage::SessionFiles {
request_id: "files".into(),
session_id: "session-a".into(),
files: vec![SessionFileRecord {
origin: SessionFileOrigin::Agent,
file: SessionFileReference {
id: "file".into(),
name: "notes.txt".into(),
size: 5,
media_type: "text/plain".into(),
},
}],
},
&mut gateway,
&mut session,
&mut state,
"session-a",
false,
);
assert_eq!(
state
.picker
.as_ref()
.map(super::super::PickerState::match_count),
Some(1)
);
}
#[test]
fn replay_preserves_current_activity_and_live_completion_clears_only_its_turn() {
use mobius::protocol::{ExecApprovalRequestEvent, TurnCompleteEvent};
let mut state = TuiState::default();
let mut gateway = ready_payload();
let mut session = session_payload("session-a");
session.latest_sequence = 2;
session.active_turn_ids = vec!["current-turn".into()];
session.pending_approvals = vec![ExecApprovalRequestEvent {
id: "current-approval".into(),
turn_id: "current-turn".into(),
calls: Vec::new(),
reason: "Approve current work".into(),
}];
sync_session(&mut state, &session, &gateway).unwrap();
let mut message = replay_event(1);
let ServerMessage::AgentEvent { record, .. } = &mut message else {
unreachable!()
};
record.event.msg = EventMsg::ExecApprovalRequest(ExecApprovalRequestEvent {
id: "old-approval".into(),
turn_id: "old-turn".into(),
calls: Vec::new(),
reason: "Already decided".into(),
});
handle_server_message(
message,
&mut gateway,
&mut session,
&mut state,
"session-a",
false,
);
let mut completion = replay_event(2);
let ServerMessage::AgentEvent { record, .. } = &mut completion else {
unreachable!()
};
record.event.msg = EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: "current-turn".into(),
});
handle_server_message(
completion.clone(),
&mut gateway,
&mut session,
&mut state,
"session-a",
false,
);
assert_eq!(state.active_turn(), Some("current-turn"));
assert_eq!(state.approvals.len(), 1);
assert_eq!(state.approval().unwrap().id, "current-approval");
handle_server_message(
ServerMessage::Ready {
payload: ready_payload(),
},
&mut gateway,
&mut session,
&mut state,
"session-a",
false,
);
assert_eq!(state.approval().unwrap().id, "current-approval");
let ServerMessage::AgentEvent { record, .. } = &mut completion else {
unreachable!()
};
record.sequence = 3;
handle_server_message(
completion,
&mut gateway,
&mut session,
&mut state,
"session-a",
false,
);
assert!(state.active_turn().is_none());
assert!(state.approval().is_none());
handle_server_message(
ServerMessage::Ready {
payload: ready_payload(),
},
&mut gateway,
&mut session,
&mut state,
"session-a",
false,
);
assert!(state.active_turn().is_none());
assert!(state.approval().is_none());
}
fn session_payload(session_id: &str) -> SessionReadyPayload {
SessionReadyPayload {
active_turn_ids: Vec::new(),
pending_approvals: Vec::new(),
latest_sequence: 0,
next_before_sequence: None,
attached_folders: Vec::new(),
workspace: WorkspaceInfo {
id: "workspace".into(),
path: "/tmp".into(),
},
git: None,
session: SessionConfiguredEvent {
session_id: session_id.into(),
context: SessionContext {
owner_id: "bot-a".into(),
..SessionContext::default()
},
model: ModelChangedEvent {
route: String::new(),
model: String::new(),
reasoning_effort: None,
model_context_window: None,
},
},
contributions: Vec::new(),
widgets: Vec::new(),
tool_count: 0,
compaction_count: 0,
context_limit_tokens: None,
active_message_delivery: ActiveMessageDelivery::Steer,
run_stats: RunStats::default(),
}
}
fn ready_payload() -> ReadyPayload {
ReadyPayload {
gateway_version: env!("CARGO_PKG_VERSION").into(),
machine_name: String::new(),
bots: vec![BotRecord {
id: "bot-a".into(),
handle: "ada".into(),
name: "Ada".into(),
description: String::new(),
tint: Default::default(),
config: VersionedAgentConfig {
revision: 1,
config: Default::default(),
},
accepts_file_attachments: false,
routine_interaction_policy: RoutineInteractionPolicy::Unattended,
}],
sessions: Vec::new(),
background_approvals: Vec::new(),
providers: Vec::new(),
provider_instances: Vec::new(),
bot_defaults: None,
models: Vec::new(),
model_providers: Default::default(),
middleware_features: Vec::new(),
extensions: Vec::new(),
contributions: Vec::new(),
max_active_sessions: 0,
session_file_limits: SessionFileLimits {
max_attachment_references: 1,
max_file_bytes: 1,
max_session_files: 1,
max_session_bytes: 1,
max_upload_chunk_bytes: 1,
},
}
}
#[test]
fn resume_picker_uses_live_session_and_bot_metadata() {
let mut event = EventMsg::Frontend(FrontendEvent::Picker {
title: "Resume chat".into(),
options: vec![FrontendPickerOption {
label: "old label".into(),
description: "created at Unix time 1700000000000".into(),
detail: String::new(),
symbol: None,
shows_detail: false,
op: Op::ResumeSession {
session_id: "session-a".into(),
},
}],
});
let sessions = [SessionRecord {
session_id: "session-a".into(),
session_context: SessionContext {
owner_id: "bot-a".into(),
workspace_id: Some("workspace-a".into()),
workspace_label: Some("Project A".into()),
..SessionContext::default()
},
parent_session_id: None,
parent_sequence: None,
sequence: 1,
first_user_message: Some("First message".into()),
execution_stats: Default::default(),
title: Some("Named chat".into()),
pinned: false,
activity: SessionActivity {
state: SessionActivityState::Running,
..SessionActivity::default()
},
created_at: 1_700_000_000_000,
updated_at: 1_700_000_000_000,
}];
let bots = [BotRecord {
id: "bot-a".into(),
handle: "curie".into(),
name: "Curie".into(),
description: "Own research.".into(),
tint: Default::default(),
config: VersionedAgentConfig {
revision: 1,
config: Default::default(),
},
accepts_file_attachments: false,
routine_interaction_policy: RoutineInteractionPolicy::Unattended,
}];
enrich_resume_picker(&mut event, &sessions, &bots, "workspace-a");
let EventMsg::Frontend(FrontendEvent::Picker { options, .. }) = event else {
panic!("resume picker");
};
assert_eq!(options[0].label, "Named chat");
assert!(
options[0]
.description
.starts_with("running · this workspace")
);
assert!(options[0].description.contains("@curie"));
assert!(options[0].description.contains("started "));
assert!(!options[0].description.contains("Unix"));
}
}