use std::sync::mpsc::{Receiver, RecvTimeoutError, Sender};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use ftui::core::geometry::Rect;
use ftui::render::sanitize::sanitize;
use ftui::runtime::subscription::{StopSignal, SubId, Subscription};
use ftui::text::Text;
use ftui::widgets::Widget;
use ftui::widgets::paragraph::Paragraph;
use ftui::widgets::spinner::{DOTS, SpinnerState};
use ftui::widgets::textarea::TextArea;
use ftui::{Cmd, Event, Frame, KeyCode, Model, Modifiers, MouseEventKind};
use crate::ask::{AskAnswer, AskResponse, AskUiRequest, QuestionReply};
use crate::extensions::{ExtensionUiRequest, ExtensionUiResponse};
use crate::interactive::PiMsg;
use crate::interactive::{format_extension_ui_prompt, parse_extension_ui_response};
use crate::keybindings::{AppAction, KeyBinding, KeyBindings};
use std::collections::VecDeque;
#[derive(Debug)]
pub enum PiFtuiMsg {
Term(Event),
Agent(PiMsg),
}
impl From<Event> for PiFtuiMsg {
fn from(event: Event) -> Self {
Self::Term(event)
}
}
const AGENT_EVENTS_SUB_ID: SubId = 0x5049_4147;
pub struct AgentEventSubscription {
rx: Arc<Mutex<Option<Receiver<PiMsg>>>>,
}
impl AgentEventSubscription {
pub fn new(rx: Receiver<PiMsg>) -> Self {
Self::from_shared(Arc::new(Mutex::new(Some(rx))))
}
const fn from_shared(rx: Arc<Mutex<Option<Receiver<PiMsg>>>>) -> Self {
Self { rx }
}
}
const AGENT_EVENT_POLL: Duration = Duration::from_millis(50);
const SPINNER_INTERVAL: Duration = Duration::from_millis(120);
const PICKER_HINT: &str = "↑/↓ j/k navigate · Enter apply · Esc close";
fn drain_agent_events(
rx: &Receiver<PiMsg>,
sender: &Sender<PiFtuiMsg>,
stopped: impl Fn() -> bool,
) {
loop {
if stopped() {
return;
}
match rx.recv_timeout(AGENT_EVENT_POLL) {
Ok(msg) => {
if sender.send(PiFtuiMsg::Agent(msg)).is_err() {
return;
}
}
Err(RecvTimeoutError::Timeout) => {}
Err(RecvTimeoutError::Disconnected) => {
return;
}
}
}
}
impl Subscription<PiFtuiMsg> for AgentEventSubscription {
fn id(&self) -> SubId {
AGENT_EVENTS_SUB_ID
}
fn run(&self, sender: Sender<PiFtuiMsg>, stop: StopSignal) {
let Some(rx) = self.rx.lock().ok().and_then(|mut slot| slot.take()) else {
return;
};
drain_agent_events(&rx, &sender, || stop.is_stopped());
}
}
#[derive(Debug, Clone, Copy)]
pub struct FtuiPalette {
accent: ftui::PackedRgba,
muted: ftui::PackedRgba,
error: ftui::PackedRgba,
warning: ftui::PackedRgba,
}
impl Default for FtuiPalette {
fn default() -> Self {
Self {
accent: ftui::PackedRgba::rgb(97, 175, 239),
muted: ftui::PackedRgba::rgb(130, 137, 151),
error: ftui::PackedRgba::rgb(220, 80, 80),
warning: ftui::PackedRgba::rgb(229, 192, 123),
}
}
}
impl FtuiPalette {
#[must_use]
pub fn from_theme(theme: &crate::theme::Theme) -> Self {
let fallback = Self::default();
let parse = |hex: &str, fallback: ftui::PackedRgba| {
crate::theme::parse_hex_color(hex)
.map_or(fallback, |(r, g, b)| ftui::PackedRgba::rgb(r, g, b))
};
Self {
accent: parse(&theme.colors.accent, fallback.accent),
muted: parse(&theme.colors.muted, fallback.muted),
error: parse(&theme.colors.error, fallback.error),
warning: parse(&theme.colors.warning, fallback.warning),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum EntryRole {
User,
Assistant,
System,
Error,
Ask,
}
impl EntryRole {
const fn prefix(self) -> &'static str {
match self {
Self::User => "› ",
Self::System => "· ",
Self::Error => "✗ ",
Self::Assistant | Self::Ask => "",
}
}
fn style(self, palette: &FtuiPalette) -> ftui::Style {
match self {
Self::User => ftui::Style::new().bold().fg(palette.accent),
Self::Assistant => ftui::Style::new(),
Self::System | Self::Ask => ftui::Style::new().dim().fg(palette.muted),
Self::Error => ftui::Style::new().bold().fg(palette.error),
}
}
}
#[derive(Debug)]
struct TranscriptEntry {
role: EntryRole,
text: String,
}
struct ActiveAsk {
request: AskUiRequest,
question_index: usize,
answers: Vec<AskAnswer>,
}
#[derive(Debug)]
pub struct AskUiReply {
pub request_id: String,
pub response: AskResponse,
}
struct PickerOverlay {
title: String,
items: Vec<String>,
values: Vec<String>,
selected: usize,
kind: PickerKind,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum PickerKind {
Theme,
Model,
Session,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum UiCommand {
Prompt(String),
SetModel { provider: String, model: String },
Bash { command: String, exclude: bool },
ResumeSession { path: String },
Compact,
ExtensionCommand { name: String, args: String },
NewSession,
SessionInfo,
TreeSummary,
SetThinking(Option<crate::model::ThinkingLevel>),
SetName(String),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum AgentUiState {
Ready,
Working,
}
impl AgentUiState {
const fn label(self) -> &'static str {
match self {
Self::Ready => "ready",
Self::Working => "working",
}
}
}
pub struct PiFtuiModel {
state: AgentUiState,
transcript: Vec<TranscriptEntry>,
streaming: String,
current_tool: Option<String>,
todo_summary: Option<String>,
thinking: String,
spinner: SpinnerState,
usage_line: Option<String>,
palette: FtuiPalette,
picker: Option<PickerOverlay>,
available_models: Vec<String>,
pending_quit: bool,
available_sessions: Vec<(String, String)>,
keybindings: KeyBindings,
active_ask: Option<ActiveAsk>,
active_ext: Option<ExtensionUiRequest>,
ext_queue: VecDeque<ExtensionUiRequest>,
ext_reply_tx: Option<Sender<ExtensionUiResponse>>,
ask_reply_tx: Option<Sender<AskUiReply>>,
term: (u16, u16),
scroll_from_tail: usize,
input: TextArea,
submit_tx: Option<Sender<UiCommand>>,
agent_rx: Arc<Mutex<Option<Receiver<PiMsg>>>>,
}
struct Regions {
header: Rect,
body: Rect,
status: Rect,
input: Rect,
footer: Rect,
}
const FIXED_CHROME_ROWS: u16 = 3;
const MAX_INPUT_ROWS: u16 = 5;
fn layout_regions(area: Rect, input_rows: u16) -> Regions {
use ftui::layout::{Constraint, Flex};
let rects = Flex::vertical()
.constraints([
Constraint::Fixed(1), Constraint::Fill, Constraint::Fixed(1), Constraint::Fixed(input_rows), Constraint::Fixed(1), ])
.split(area);
Regions {
header: rects[0],
body: rects[1],
status: rects[2],
input: rects[3],
footer: rects[4],
}
}
impl PiFtuiModel {
pub fn new(agent_rx: Receiver<PiMsg>) -> Self {
Self {
state: AgentUiState::Ready,
transcript: Vec::new(),
streaming: String::new(),
current_tool: None,
todo_summary: None,
thinking: String::new(),
spinner: SpinnerState::default(),
usage_line: None,
palette: FtuiPalette::default(),
picker: None,
available_models: Vec::new(),
pending_quit: false,
available_sessions: Vec::new(),
keybindings: KeyBindings::default(),
active_ask: None,
active_ext: None,
ext_queue: VecDeque::new(),
ext_reply_tx: None,
ask_reply_tx: None,
term: (80, 24),
scroll_from_tail: 0,
input: TextArea::new()
.with_placeholder("Type a message (Enter to send, Alt+Enter for newline)")
.with_focus(true)
.with_soft_wrap(true),
submit_tx: None,
agent_rx: Arc::new(Mutex::new(Some(agent_rx))),
}
}
#[must_use]
pub fn with_submit_channel(mut self, tx: Sender<UiCommand>) -> Self {
self.submit_tx = Some(tx);
self
}
#[must_use]
pub fn with_ask_reply_channel(mut self, tx: Sender<AskUiReply>) -> Self {
self.ask_reply_tx = Some(tx);
self
}
#[must_use]
pub fn with_ext_reply_channel(mut self, tx: Sender<ExtensionUiResponse>) -> Self {
self.ext_reply_tx = Some(tx);
self
}
#[must_use]
pub const fn with_palette(mut self, palette: FtuiPalette) -> Self {
self.palette = palette;
self
}
#[must_use]
pub fn with_available_models(mut self, models: Vec<String>) -> Self {
self.available_models = models;
self
}
#[must_use]
pub fn with_available_sessions(mut self, sessions: Vec<(String, String)>) -> Self {
self.available_sessions = sessions;
self
}
fn input_rows(&self) -> u16 {
let lines = if self.input.is_empty() {
1
} else {
self.input.text().lines().count().max(1)
};
u16::try_from(lines)
.unwrap_or(MAX_INPUT_ROWS)
.min(MAX_INPUT_ROWS)
}
fn body_height(&self) -> usize {
usize::from(
self.term
.1
.saturating_sub(FIXED_CHROME_ROWS + self.input_rows()),
)
.max(1)
}
fn conversation_line_count(&self) -> usize {
let transcript: usize = self
.transcript
.iter()
.map(|e| e.text.lines().count().max(1))
.sum();
let streaming = if self.streaming.is_empty() {
0
} else {
self.streaming.lines().count().max(1)
};
transcript + streaming
}
fn push_entry(&mut self, role: EntryRole, text: String) {
self.transcript.push(TranscriptEntry { role, text });
}
fn max_scroll_from_tail(&self) -> usize {
self.conversation_line_count()
.saturating_sub(self.body_height())
}
fn scroll_up(&mut self, lines: usize) {
self.scroll_from_tail = self
.scroll_from_tail
.saturating_add(lines)
.min(self.max_scroll_from_tail());
}
const fn scroll_down(&mut self, lines: usize) {
self.scroll_from_tail = self.scroll_from_tail.saturating_sub(lines);
}
fn handle_agent(&mut self, msg: PiMsg) -> Cmd<PiFtuiMsg> {
match msg {
PiMsg::AgentStart => {
self.state = AgentUiState::Working;
return Cmd::tick(SPINNER_INTERVAL);
}
PiMsg::TextDelta(delta) => {
self.streaming.push_str(&sanitize(&delta));
}
PiMsg::ThinkingDelta(delta) => {
self.thinking.push_str(&sanitize(&delta));
}
PiMsg::ToolStart { name, .. } => {
self.current_tool = Some(sanitize(&name).into_owned());
}
PiMsg::ToolEnd { name, is_error, .. } => {
let mark = if is_error { "✗" } else { "✓" };
let text = format!("{mark} {}", sanitize(&name));
self.push_entry(EntryRole::System, text);
self.current_tool = None;
}
PiMsg::TodoSummary { summary } => {
self.todo_summary = summary.map(|s| sanitize(&s).into_owned());
}
PiMsg::AgentDone {
usage,
error_message,
..
} => {
if !self.streaming.is_empty() {
let text = std::mem::take(&mut self.streaming);
self.push_entry(EntryRole::Assistant, text);
}
if let Some(err) = error_message {
let text = sanitize(&err).into_owned();
self.push_entry(EntryRole::Error, text);
}
if let Some(usage) = usage {
self.usage_line = Some(format!(
"tokens {}↑ {}↓ · total {}",
usage.input, usage.output, usage.total_tokens
));
}
self.state = AgentUiState::Ready;
self.current_tool = None;
self.thinking.clear();
}
PiMsg::AgentError(err) => {
let text = sanitize(&err).into_owned();
self.push_entry(EntryRole::Error, text);
self.state = AgentUiState::Ready;
self.current_tool = None;
self.thinking.clear();
}
PiMsg::System(text) | PiMsg::SystemNote(text) => {
let text = sanitize(&text).into_owned();
self.push_entry(EntryRole::System, text);
}
PiMsg::ConversationReset {
messages, status, ..
} => {
self.apply_conversation_reset(messages, status);
}
PiMsg::BashResult { display, .. } => {
let text = sanitize(&display).into_owned();
self.push_entry(EntryRole::System, text);
self.current_tool = None;
self.scroll_from_tail = 0;
}
PiMsg::AskUiRequest(request) => {
if request.request.questions.is_empty() {
self.send_ask_reply(request.id, Vec::new(), true);
} else {
self.push_ask_card(&request, 0);
self.active_ask = Some(ActiveAsk {
request,
question_index: 0,
answers: Vec::new(),
});
}
}
PiMsg::ExtensionUiRequest(request) => {
if self.active_ext.is_none() && self.active_ask.is_none() {
self.activate_ext_request(request);
} else {
self.ext_queue.push_back(request);
}
}
PiMsg::UiShutdown => return Cmd::quit(),
_ => {}
}
Cmd::none()
}
fn push_ask_card(&mut self, request: &AskUiRequest, index: usize) {
let total = request.request.questions.len();
let card =
crate::ask::format_question_card(&request.request.questions[index], index, total);
let text = sanitize(card.trim_end()).into_owned();
self.push_entry(EntryRole::Ask, text);
self.scroll_from_tail = 0;
}
fn send_ask_reply(&self, request_id: String, answers: Vec<AskAnswer>, dismissed: bool) {
if let Some(tx) = &self.ask_reply_tx {
let _ = tx.send(AskUiReply {
request_id,
response: AskResponse { answers, dismissed },
});
}
}
fn submit_ask_answer(&mut self) {
let Some(mut ask) = self.active_ask.take() else {
return;
};
let raw = self.input.text();
self.input.set_text("");
let index = ask.question_index;
let question = &ask.request.request.questions[index];
match crate::ask::parse_question_reply(question, &raw) {
Err(err) => {
let text = format!(" ! {}", sanitize(&err));
self.push_entry(EntryRole::Ask, text);
self.scroll_from_tail = 0;
self.active_ask = Some(ask); }
Ok(QuestionReply::Cancel) => {
self.push_entry(EntryRole::Ask, String::from(" (dismissed)"));
self.scroll_from_tail = 0;
self.send_ask_reply(ask.request.id, Vec::new(), true);
self.maybe_activate_queued_ext();
}
Ok(reply) => {
let (selected, other) = match reply {
QuestionReply::Selected(labels) => (labels, None),
QuestionReply::Other(text) => (Vec::new(), Some(text)),
QuestionReply::Cancel => unreachable!("handled above"),
};
let echo = other.as_ref().map_or_else(
|| format!(" → {}", selected.join(", ")),
|text| format!(" → {text}"),
);
let echo = sanitize(&echo).into_owned();
self.push_entry(EntryRole::Ask, echo);
let question_id = question.id.clone().unwrap_or_else(|| index.to_string());
ask.answers.push(AskAnswer {
question_id,
selected,
other,
});
let next = index + 1;
if next < ask.request.request.questions.len() {
self.push_ask_card(&ask.request, next);
ask.question_index = next;
self.active_ask = Some(ask);
} else {
self.scroll_from_tail = 0;
self.send_ask_reply(ask.request.id, ask.answers, false);
self.maybe_activate_queued_ext();
}
}
}
}
fn submit_input(&mut self) {
let text = self.input.text();
let trimmed = text.trim();
if trimmed.is_empty() {
return;
}
let clean = sanitize(trimmed).into_owned();
self.input.set_text("");
self.scroll_from_tail = 0;
self.push_entry(EntryRole::User, clean.clone());
let bang = clean
.strip_prefix("!!")
.map(|rest| (rest.trim(), true))
.or_else(|| clean.strip_prefix('!').map(|rest| (rest.trim(), false)));
if let Some((command, exclude)) = bang {
if command.is_empty() {
self.push_entry(EntryRole::Error, String::from("usage: !<command>"));
} else {
self.send_command(UiCommand::Bash {
command: command.to_string(),
exclude,
});
}
return;
}
if clean.starts_with('/') && self.route_slash_command(&clean) {
return;
}
self.send_command(UiCommand::Prompt(clean));
}
fn route_slash_command(&mut self, clean: &str) -> bool {
if let Some(rest) = clean.strip_prefix("/model") {
self.route_model_command(rest.trim());
return true;
}
self.route_slash_command_tail(clean)
}
fn route_model_command(&mut self, spec: &str) {
{
if spec.is_empty() {
if self.available_models.is_empty() {
self.push_entry(
EntryRole::Error,
String::from("no models available; use /model <provider>/<model>"),
);
} else {
self.picker = Some(PickerOverlay {
title: String::from("Model (Enter to switch, Esc to close)"),
items: self.available_models.clone(),
values: Vec::new(),
selected: 0,
kind: PickerKind::Model,
});
}
} else if let Some((provider, model)) = spec.split_once('/')
&& !provider.is_empty()
&& !model.is_empty()
{
self.push_entry(EntryRole::System, format!("switching model to {spec} ..."));
self.send_command(UiCommand::SetModel {
provider: provider.to_string(),
model: model.to_string(),
});
} else {
self.push_entry(
EntryRole::Error,
String::from("usage: /model <provider>/<model>"),
);
}
}
}
fn route_slash_command_tail(&mut self, clean: &str) -> bool {
if clean == "/exit" || clean == "/quit" {
self.pending_quit = true;
return true;
}
if clean == "/compact" {
self.push_entry(
EntryRole::System,
String::from("compacting conversation ..."),
);
self.send_command(UiCommand::Compact);
return true;
}
if clean == "/theme" {
self.picker = Some(PickerOverlay {
title: String::from("Theme (Enter to apply, Esc to close)"),
items: vec![String::from("dark"), String::from("light")],
values: Vec::new(),
selected: 0,
kind: PickerKind::Theme,
});
return true;
}
if clean == "/resume" {
if self.available_sessions.is_empty() {
self.push_entry(EntryRole::Error, String::from("no saved sessions found"));
} else {
let (items, values) = self
.available_sessions
.iter()
.map(|(label, path)| (label.clone(), path.clone()))
.unzip();
self.picker = Some(PickerOverlay {
title: String::from("Resume session (Enter to load, Esc to close)"),
items,
values,
selected: 0,
kind: PickerKind::Session,
});
}
return true;
}
if clean == "/help" {
self.push_entry(
EntryRole::System,
String::from(
"ftui preview commands: /model [provider/model], /resume, /compact, \
/theme, /new, /clear, /session, /tree, /thinking [level], \
/name <name>, /exit, /help, !<cmd> (runs + sends output to the \
agent), !!<cmd> (display-only)",
),
);
return true;
}
let (cmd_name, cmd_args) = clean.split_once(char::is_whitespace).unwrap_or((clean, ""));
match cmd_name {
"/new" => {
self.send_command(UiCommand::NewSession);
return true;
}
"/clear" | "/cls" => {
self.transcript.clear();
self.streaming.clear();
self.thinking.clear();
self.current_tool = None;
self.scroll_from_tail = 0;
self.push_entry(EntryRole::System, String::from("Conversation cleared"));
return true;
}
"/session" | "/info" => {
self.send_command(UiCommand::SessionInfo);
return true;
}
"/tree" => {
self.send_command(UiCommand::TreeSummary);
return true;
}
"/thinking" | "/think" | "/t" => {
let value = cmd_args.trim();
if value.is_empty() {
self.send_command(UiCommand::SetThinking(None));
return true;
}
match value.parse::<crate::model::ThinkingLevel>() {
Ok(level) => self.send_command(UiCommand::SetThinking(Some(level))),
Err(err) => self.push_entry(EntryRole::Error, err),
}
return true;
}
"/name" => {
let name = cmd_args.trim();
if name.is_empty() {
self.push_entry(EntryRole::Error, String::from("Usage: /name <name>"));
} else {
self.send_command(UiCommand::SetName(name.to_string()));
}
return true;
}
_ => {}
}
if !clean.starts_with("/skill:") {
let body = clean.trim_start_matches('/');
let (name, args) = body.split_once(char::is_whitespace).unwrap_or((body, ""));
if name.is_empty() {
self.push_entry(EntryRole::Error, String::from("Unknown command: /"));
} else {
self.send_command(UiCommand::ExtensionCommand {
name: name.to_string(),
args: args.trim().to_string(),
});
}
return true;
}
false
}
fn handle_picker_key(&mut self, key: &ftui::KeyEvent) {
let Some(picker) = self.picker.as_mut() else {
return;
};
match key.code {
KeyCode::Up | KeyCode::Char('k') => {
picker.selected = picker.selected.saturating_sub(1);
}
KeyCode::Down | KeyCode::Char('j') => {
picker.selected = (picker.selected + 1).min(picker.items.len().saturating_sub(1));
}
KeyCode::Escape => {
self.picker = None;
}
KeyCode::Enter => {
let Some(mut picker) = self.picker.take() else {
return;
};
let choice = if picker.values.is_empty() {
picker.items.swap_remove(picker.selected)
} else {
picker.values.swap_remove(picker.selected)
};
self.apply_picker_choice(picker.kind, &choice);
}
_ => {}
}
}
fn apply_picker_choice(&mut self, kind: PickerKind, choice: &str) {
match kind {
PickerKind::Theme => {
let theme = if choice == "light" {
crate::theme::Theme::light()
} else {
crate::theme::Theme::dark()
};
self.palette = FtuiPalette::from_theme(&theme);
self.push_entry(EntryRole::System, format!("theme set to {choice}"));
self.scroll_from_tail = 0;
}
PickerKind::Model => {
if let Some((provider, model)) = choice.split_once('/') {
self.push_entry(
EntryRole::System,
format!("switching model to {choice} ..."),
);
self.scroll_from_tail = 0;
self.send_command(UiCommand::SetModel {
provider: provider.to_string(),
model: model.to_string(),
});
} else {
self.push_entry(EntryRole::Error, format!("malformed model entry: {choice}"));
}
}
PickerKind::Session => {
self.push_entry(EntryRole::System, String::from("resuming session ..."));
self.scroll_from_tail = 0;
self.send_command(UiCommand::ResumeSession {
path: choice.to_string(),
});
}
}
}
fn send_command(&self, command: UiCommand) {
if let Some(tx) = &self.submit_tx {
let _ = tx.send(command);
}
}
fn handle_term(&mut self, event: &Event) -> Cmd<PiFtuiMsg> {
match event {
Event::Tick => {
if self.state == AgentUiState::Working {
self.spinner.tick();
return Cmd::tick(SPINNER_INTERVAL);
}
return Cmd::none();
}
Event::Key(key) => {
let ctrl_c =
key.code == KeyCode::Char('c') && key.modifiers.contains(Modifiers::CTRL);
if ctrl_c {
return Cmd::quit();
}
if self.picker.is_some() {
self.handle_picker_key(key);
return Cmd::none();
}
let actions = KeyBinding::from_ftui_key(key)
.map(|binding| self.keybindings.matching_actions(&binding))
.unwrap_or_default();
let pick = |wanted: AppAction| actions.contains(&wanted).then_some(wanted);
let action = pick(AppAction::PageUp)
.or_else(|| pick(AppAction::PageDown))
.or_else(|| pick(AppAction::Submit))
.or_else(|| pick(AppAction::NewLine))
.or_else(|| pick(AppAction::Interrupt))
.or_else(|| pick(AppAction::CursorLineEnd))
.or_else(|| {
if self.input.is_empty() {
pick(AppAction::Exit)
} else {
None
}
});
let page = self.body_height().saturating_sub(1).max(1);
match action {
Some(AppAction::PageUp) => return self.consume_scroll(|m| m.scroll_up(page)),
Some(AppAction::PageDown) => {
return self.consume_scroll(|m| m.scroll_down(page));
}
Some(AppAction::Exit) if self.input.is_empty() => return Cmd::quit(),
Some(AppAction::Interrupt) if self.active_ask.is_some() => {
if let Some(ask) = self.active_ask.take() {
self.push_entry(EntryRole::Ask, String::from(" (dismissed)"));
self.scroll_from_tail = 0;
self.send_ask_reply(ask.request.id, Vec::new(), true);
self.maybe_activate_queued_ext();
}
return Cmd::none();
}
Some(AppAction::Interrupt) if self.active_ext.is_some() => {
self.cancel_active_ext();
return Cmd::none();
}
Some(AppAction::Submit) if self.input_active() => {
if self.active_ask.is_some() {
self.submit_ask_answer();
} else if self.active_ext.is_some() {
self.submit_ext_answer();
} else {
self.submit_input();
if self.pending_quit {
return Cmd::quit();
}
}
return Cmd::none();
}
Some(AppAction::NewLine) if self.input_active() => {
self.input.insert_newline();
return Cmd::none();
}
Some(AppAction::CursorLineEnd) if self.input.is_empty() => {
self.scroll_from_tail = 0;
return Cmd::none();
}
_ => {}
}
if self.input_active() {
self.input.handle_event(event);
}
}
Event::Mouse(mouse) => match mouse.kind {
MouseEventKind::ScrollUp => self.scroll_up(3),
MouseEventKind::ScrollDown => self.scroll_down(3),
_ => {}
},
Event::Resize { width, height } => {
self.term = (*width, *height);
self.scroll_from_tail = self.scroll_from_tail.min(self.max_scroll_from_tail());
}
_ => {
if self.input_active() {
self.input.handle_event(event);
}
}
}
Cmd::none()
}
fn input_active(&self) -> bool {
self.state == AgentUiState::Ready || self.active_ask.is_some() || self.active_ext.is_some()
}
fn apply_conversation_reset(
&mut self,
messages: Vec<crate::interactive::ConversationMessage>,
status: Option<String>,
) {
self.transcript.clear();
self.streaming.clear();
for message in messages {
let role = match message.role {
crate::interactive::MessageRole::User => EntryRole::User,
crate::interactive::MessageRole::Assistant => EntryRole::Assistant,
crate::interactive::MessageRole::Tool | crate::interactive::MessageRole::System => {
EntryRole::System
}
};
let text = sanitize(&message.content).into_owned();
self.push_entry(role, text);
}
if let Some(status) = status {
let text = sanitize(&status).into_owned();
self.push_entry(EntryRole::System, text);
}
self.scroll_from_tail = 0;
}
fn activate_ext_request(&mut self, request: ExtensionUiRequest) {
let card = format_extension_ui_prompt(&request);
let text = sanitize(card.trim_end()).into_owned();
self.push_entry(EntryRole::Ask, text);
self.scroll_from_tail = 0;
self.active_ext = Some(request);
}
fn send_ext_reply(&self, response: ExtensionUiResponse) {
if let Some(tx) = &self.ext_reply_tx {
let _ = tx.send(response);
}
}
fn submit_ext_answer(&mut self) {
let Some(request) = self.active_ext.take() else {
return;
};
let raw = self.input.text();
self.input.set_text("");
match parse_extension_ui_response(&request, &raw) {
Err(err) => {
let text = format!(" ! {}", sanitize(&err));
self.push_entry(EntryRole::Ask, text);
self.scroll_from_tail = 0;
self.active_ext = Some(request);
}
Ok(response) => {
let echo = if response.cancelled {
String::from(" (cancelled)")
} else {
format!(" → {}", sanitize(raw.trim()))
};
self.push_entry(EntryRole::Ask, echo);
self.scroll_from_tail = 0;
self.send_ext_reply(response);
if let Some(next) = self.ext_queue.pop_front() {
self.activate_ext_request(next);
}
}
}
}
fn maybe_activate_queued_ext(&mut self) {
if self.active_ask.is_none()
&& self.active_ext.is_none()
&& let Some(next) = self.ext_queue.pop_front()
{
self.activate_ext_request(next);
}
}
fn cancel_active_ext(&mut self) {
if let Some(request) = self.active_ext.take() {
self.push_entry(EntryRole::Ask, String::from(" (cancelled)"));
self.scroll_from_tail = 0;
self.send_ext_reply(ExtensionUiResponse {
id: request.id,
value: None,
cancelled: true,
});
if let Some(next) = self.ext_queue.pop_front() {
self.activate_ext_request(next);
}
}
}
fn consume_scroll(&mut self, scroll: impl FnOnce(&mut Self)) -> Cmd<PiFtuiMsg> {
scroll(self);
Cmd::none()
}
fn conversation_text(&self) -> Text<'static> {
let md = ftui_extras::markdown::MarkdownRenderer::new(
ftui_extras::markdown::MarkdownTheme::default(),
);
let palette = self.palette;
let mut lines: Vec<ftui::text::Line<'static>> =
Vec::with_capacity(self.conversation_line_count());
let mut push_block = |role: EntryRole, content: &str| {
if role == EntryRole::Assistant {
let rendered = md.render(content);
lines.extend(rendered.lines().iter().cloned());
return;
}
let style = role.style(&palette);
let prefix = role.prefix();
let indent = " ".repeat(prefix.chars().count());
for (i, line) in content.lines().enumerate() {
let lead = if i == 0 { prefix } else { indent.as_str() };
let mut rendered = String::with_capacity(lead.len() + line.len());
rendered.push_str(lead);
rendered.push_str(line);
lines.push(ftui::text::Line::styled(rendered, style));
}
if content.is_empty() {
lines.push(ftui::text::Line::styled(prefix.to_string(), style));
}
};
for entry in &self.transcript {
push_block(entry.role, &entry.text);
}
if !self.streaming.is_empty() {
let rendered = md.render_streaming(&self.streaming);
lines.extend(rendered.lines().iter().cloned());
}
Text::from_lines(lines)
}
}
impl Model for PiFtuiModel {
type Message = PiFtuiMsg;
fn update(&mut self, msg: PiFtuiMsg) -> Cmd<PiFtuiMsg> {
match msg {
PiFtuiMsg::Term(event) => self.handle_term(&event),
PiFtuiMsg::Agent(agent) => self.handle_agent(agent),
}
}
fn view(&self, frame: &mut Frame) {
let area = Rect::new(0, 0, frame.width(), frame.height());
let regions = layout_regions(area, self.input_rows());
let header = format!("pi · {}", self.state.label());
let header_style = ftui::Style::new().bold().fg(self.palette.accent);
Paragraph::new(Text::from_lines([ftui::text::Line::styled(
header,
header_style,
)]))
.render(regions.header, frame);
if let Some(picker) = &self.picker {
let mut lines = vec![ftui::text::Line::styled(
picker.title.as_str(),
ftui::Style::new().bold().fg(self.palette.accent),
)];
for (i, item) in picker.items.iter().enumerate() {
let (marker, style) = if i == picker.selected {
("▸ ", ftui::Style::new().bold().fg(self.palette.accent))
} else {
(" ", ftui::Style::new())
};
lines.push(ftui::text::Line::from_spans([
ftui::text::Span::styled(marker, style),
ftui::text::Span::styled(item.as_str(), style),
]));
}
Paragraph::new(Text::from_lines(lines)).render(regions.body, frame);
let footer_style = ftui::Style::new().dim().fg(self.palette.muted);
Paragraph::new(Text::from_lines([ftui::text::Line::styled(
PICKER_HINT,
footer_style,
)]))
.render(regions.footer, frame);
return;
}
let body_text = self.conversation_text();
let total_lines = body_text.lines().len();
let visible = usize::from(regions.body.height).max(1);
let from_tail = self
.scroll_from_tail
.min(total_lines.saturating_sub(visible));
let offset = total_lines.saturating_sub(visible + from_tail);
let offset_u16 = u16::try_from(offset).unwrap_or(u16::MAX);
Paragraph::new(body_text)
.scroll((offset_u16, 0))
.render(regions.body, frame);
let status_line = if self.state == AgentUiState::Working {
let spin = DOTS[self.spinner.current_frame % DOTS.len()];
let activity = self.current_tool.as_ref().map_or_else(
|| {
if self.streaming.is_empty() && !self.thinking.is_empty() {
String::from("thinking ...")
} else {
String::from("responding ...")
}
},
|tool| format!("running {tool} ..."),
);
format!("{spin} {activity}")
} else {
self.todo_summary
.as_ref()
.map_or_else(String::new, |todo| format!("todo {todo}"))
};
if !status_line.is_empty() {
let status_style = if self.state == AgentUiState::Working {
ftui::Style::new().fg(self.palette.warning)
} else {
ftui::Style::new().dim().fg(self.palette.muted)
};
Paragraph::new(Text::from_lines([ftui::text::Line::styled(
status_line,
status_style,
)]))
.render(regions.status, frame);
}
if self.input_active() {
self.input.render(regions.input, frame);
} else {
Paragraph::new(Text::raw("… processing (ctrl+c to quit)")).render(regions.input, frame);
}
let footer = if from_tail > 0 {
format!("[{from_tail} lines up] End to follow")
} else if let Some(usage) = &self.usage_line {
usage.clone()
} else {
String::from("pi — ftui preview")
};
let footer_style = ftui::Style::new().dim().fg(self.palette.muted);
Paragraph::new(Text::from_lines([ftui::text::Line::styled(
footer,
footer_style,
)]))
.render(regions.footer, frame);
}
fn subscriptions(&self) -> Vec<Box<dyn Subscription<PiFtuiMsg>>> {
vec![Box::new(AgentEventSubscription::from_shared(Arc::clone(
&self.agent_rx,
)))]
}
}
pub fn agent_event_to_pi_msgs(event: &crate::agent::AgentEvent) -> Vec<PiMsg> {
use crate::agent::AgentEvent as E;
use crate::model::AssistantMessageEvent as A;
match event {
E::AgentStart { .. } => vec![PiMsg::AgentStart],
E::AgentEnd {
messages, error, ..
} => {
let last_assistant = messages.iter().rev().find_map(|message| match message {
crate::model::Message::Assistant(assistant) => Some(assistant),
_ => None,
});
vec![PiMsg::AgentDone {
usage: last_assistant.map(|a| a.usage.clone()),
stop_reason: last_assistant
.map_or(crate::model::StopReason::Stop, |a| a.stop_reason),
error_message: error.clone(),
}]
}
E::MessageUpdate {
assistant_message_event,
..
} => match assistant_message_event {
A::TextDelta { delta, .. } => vec![PiMsg::TextDelta(delta.clone())],
A::ThinkingDelta { delta, .. } => vec![PiMsg::ThinkingDelta(delta.clone())],
_ => Vec::new(),
},
E::ToolExecutionStart {
tool_call_id,
tool_name,
..
} => vec![PiMsg::ToolStart {
name: tool_name.clone(),
tool_id: tool_call_id.clone(),
}],
E::ToolExecutionEnd {
tool_call_id,
tool_name,
is_error,
..
} => vec![PiMsg::ToolEnd {
name: tool_name.clone(),
tool_id: tool_call_id.clone(),
is_error: *is_error,
}],
E::AutoRetryStart {
attempt,
max_attempts,
error_message,
..
} => vec![PiMsg::SystemNote(format!(
"retry {attempt}/{max_attempts}: {error_message}"
))],
E::AutoCompactionStart { reason } => {
vec![PiMsg::SystemNote(format!("compacting context: {reason}"))]
}
E::AutoCompactionEnd {
aborted,
error_message,
..
} => {
let note = if *aborted {
String::from("compaction aborted")
} else if let Some(err) = error_message {
format!("compaction failed: {err}")
} else {
String::from("compaction complete")
};
vec![PiMsg::SystemNote(note)]
}
E::ExtensionError { event, error, .. } => {
vec![PiMsg::System(format!("extension error ({event}): {error}"))]
}
_ => Vec::new(),
}
}
const SUBMIT_POLL: Duration = Duration::from_millis(50);
const INLINE_MIN_HEIGHT: u16 = 10;
const INLINE_MAX_HEIGHT: u16 = 15;
const EXT_UI_TIMEOUT_MS: u64 = 300_000;
struct FtuiExtensionUiHandler {
agent_tx: Sender<PiMsg>,
pending: Mutex<
std::collections::HashMap<
String,
asupersync::channel::oneshot::Sender<ExtensionUiResponse>,
>,
>,
}
impl FtuiExtensionUiHandler {
fn new(agent_tx: Sender<PiMsg>) -> Self {
Self {
agent_tx,
pending: Mutex::new(std::collections::HashMap::new()),
}
}
fn resolve(&self, response: ExtensionUiResponse) {
let cx = crate::agent_cx::AgentCx::for_current_or_request();
let sender = self
.pending
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.remove(&response.id);
if let Some(sender) = sender {
let _ = sender.send(cx.cx(), response);
}
}
fn drop_pending(&self, id: &str) {
self.pending
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.remove(id);
}
}
#[async_trait::async_trait]
impl crate::sdk::ExtensionUiHandler for FtuiExtensionUiHandler {
async fn request_ui(
&self,
request: ExtensionUiRequest,
) -> crate::error::Result<Option<ExtensionUiResponse>> {
let cx = crate::agent_cx::AgentCx::for_current_or_request();
let id = request.id.clone();
let timeout_ms = request.timeout_ms.unwrap_or(EXT_UI_TIMEOUT_MS);
let (reply_tx, mut reply_rx) = asupersync::channel::oneshot::channel();
self.pending
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.insert(id.clone(), reply_tx);
if self
.agent_tx
.send(PiMsg::ExtensionUiRequest(request))
.is_err()
{
self.drop_pending(&id);
return Ok(None);
}
let waited = asupersync::time::timeout(
asupersync::time::wall_now(),
std::time::Duration::from_millis(timeout_ms),
reply_rx.recv(cx.cx()),
)
.await;
if let Ok(Ok(response)) = waited {
Ok(Some(response))
} else {
self.drop_pending(&id);
Ok(Some(ExtensionUiResponse {
id,
value: None,
cancelled: true,
}))
}
}
}
fn spawn_ext_reply_pump(
handler: Arc<FtuiExtensionUiHandler>,
ext_reply_rx: Receiver<ExtensionUiResponse>,
runtime_handle: &asupersync::runtime::RuntimeHandle,
) {
runtime_handle.spawn(async move {
loop {
match ext_reply_rx.try_recv() {
Ok(response) => handler.resolve(response),
Err(std::sync::mpsc::TryRecvError::Empty) => {
asupersync::time::sleep(asupersync::time::wall_now(), SUBMIT_POLL).await;
}
Err(std::sync::mpsc::TryRecvError::Disconnected) => break,
}
}
});
}
fn install_ask_bridges(
handle: &crate::sdk::AgentSessionHandle,
agent_tx: &Sender<PiMsg>,
ask_reply_rx: Receiver<AskUiReply>,
runtime_handle: &asupersync::runtime::RuntimeHandle,
) -> CurrentAsk {
let current_ask: CurrentAsk = Arc::new(Mutex::new(handle.ask_tool()));
if let Some(ask) = handle.ask_tool() {
install_ask_forwarder(&ask, agent_tx, runtime_handle);
}
spawn_ask_reply_pump(Arc::clone(¤t_ask), ask_reply_rx, runtime_handle);
current_ask
}
type CurrentAsk = Arc<Mutex<Option<crate::ask::AskTool>>>;
fn install_ask_forwarder(
ask: &crate::ask::AskTool,
agent_tx: &Sender<PiMsg>,
runtime_handle: &asupersync::runtime::RuntimeHandle,
) {
let (ask_ui_tx, mut ask_ui_rx) = asupersync::channel::mpsc::channel::<AskUiRequest>(4);
ask.install_channel_ui(ask_ui_tx);
let ask_fwd_tx = agent_tx.clone();
runtime_handle.spawn(async move {
let cx = crate::agent_cx::AgentCx::for_request();
while let Ok(request) = ask_ui_rx.recv(&cx).await {
let _ = ask_fwd_tx.send(PiMsg::AskUiRequest(request));
}
});
}
fn spawn_ask_reply_pump(
current_ask: CurrentAsk,
ask_reply_rx: Receiver<AskUiReply>,
runtime_handle: &asupersync::runtime::RuntimeHandle,
) {
runtime_handle.spawn(async move {
loop {
match ask_reply_rx.try_recv() {
Ok(reply) => {
let guard = current_ask
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if let Some(ask) = guard.as_ref() {
let _ = ask.respond_ui(&reply.request_id, reply.response);
}
}
Err(std::sync::mpsc::TryRecvError::Empty) => {
asupersync::time::sleep(asupersync::time::wall_now(), SUBMIT_POLL).await;
}
Err(std::sync::mpsc::TryRecvError::Disconnected) => break,
}
}
});
}
async fn run_prompt_turn(
handle: &mut crate::sdk::AgentSessionHandle,
prompt: String,
agent_tx: &Sender<PiMsg>,
) {
let tx = agent_tx.clone();
let result = handle
.prompt(prompt, move |event| {
for msg in agent_event_to_pi_msgs(&event) {
let _ = tx.send(msg);
}
})
.await;
if let Err(err) = result {
let _ = agent_tx.send(PiMsg::AgentError(err.to_string()));
}
}
fn resume_template_from(options: &crate::sdk::SessionOptions) -> crate::sdk::SessionOptions {
crate::sdk::SessionOptions {
provider: options.provider.clone(),
model: options.model.clone(),
api_key: options.api_key.clone(),
working_directory: options.working_directory.clone(),
session_dir: options.session_dir.clone(),
extension_paths: options.extension_paths.clone(),
extension_policy: options.extension_policy.clone(),
no_session: false,
..Default::default()
}
}
const EXT_COMMAND_TIMEOUT_MS: u64 = 24 * 60 * 60 * 1000;
async fn run_extension_command(
handle: &crate::sdk::AgentSessionHandle,
cwd: &std::path::Path,
name: &str,
args: &str,
agent_tx: &Sender<PiMsg>,
) {
let manager = handle
.session()
.extensions
.as_ref()
.map(|region| region.manager().clone());
let Some(manager) = manager else {
let _ = agent_tx.send(PiMsg::System(format!(
"Unknown command: /{name} (extensions disabled; try /help)"
)));
return;
};
if !manager.has_command(name) {
let _ = agent_tx.send(PiMsg::System(format!(
"Unknown command: /{name} (try /help)"
)));
return;
}
let Some(runtime) = manager.runtime() else {
let _ = agent_tx.send(PiMsg::System(format!(
"Extension command '/{name}' is not available (runtime not enabled)"
)));
return;
};
let _ = agent_tx.send(PiMsg::ToolStart {
name: format!("/{name}"),
tool_id: String::from("ftui-ext-command"),
});
let ctx_payload = serde_json::json!({
"cwd": cwd.display().to_string(),
"hasUI": true,
});
let result = runtime
.execute_command(
name.to_string(),
args.to_string(),
Arc::new(ctx_payload),
EXT_COMMAND_TIMEOUT_MS,
)
.await;
let msg = match result {
Ok(value) if value.is_null() => PiMsg::SystemNote(format!("/{name} done")),
Ok(value) => PiMsg::SystemNote(format!("/{name} → {value}")),
Err(err) => PiMsg::AgentError(format!("/{name}: {err}")),
};
let _ = agent_tx.send(msg);
let _ = agent_tx.send(PiMsg::ToolEnd {
name: format!("/{name}"),
tool_id: String::from("ftui-ext-command"),
is_error: false,
});
}
async fn run_set_model_command(
handle: &mut crate::sdk::AgentSessionHandle,
provider: &str,
model: &str,
agent_tx: &Sender<PiMsg>,
) {
let msg = match handle.set_model(provider, model).await {
Ok(()) => PiMsg::System(format!("model set to {provider}/{model}")),
Err(err) => PiMsg::AgentError(format!("model switch: {err}")),
};
let _ = agent_tx.send(msg);
}
async fn run_compact_command(
handle: &mut crate::sdk::AgentSessionHandle,
agent_tx: &Sender<PiMsg>,
) {
let tx = agent_tx.clone();
let result = handle
.compact(move |event| {
for msg in agent_event_to_pi_msgs(&event) {
let _ = tx.send(msg);
}
})
.await;
match result {
Ok(()) => {
send_conversation_reset(handle, agent_tx, "conversation compacted").await;
}
Err(err) => {
let _ = agent_tx.send(PiMsg::AgentError(format!("compact: {err}")));
}
}
}
async fn new_session_command(
template: &crate::sdk::SessionOptions,
handle: &crate::sdk::AgentSessionHandle,
current_ask: &CurrentAsk,
ext_handler: &Arc<FtuiExtensionUiHandler>,
agent_tx: &Sender<PiMsg>,
runtime_handle: &asupersync::runtime::RuntimeHandle,
) -> Option<crate::sdk::AgentSessionHandle> {
let (provider, model_id) = handle.model();
let options = crate::sdk::SessionOptions {
provider: Some(provider.clone()),
model: Some(model_id.clone()),
api_key: template.api_key.clone(),
working_directory: template.working_directory.clone(),
session_dir: template.session_dir.clone(),
extension_paths: template.extension_paths.clone(),
extension_policy: template.extension_policy.clone(),
extension_ui_handler: Some(
Arc::clone(ext_handler) as Arc<dyn crate::sdk::ExtensionUiHandler>
),
thinking: Some(crate::model::ThinkingLevel::Off),
no_session: false,
..Default::default()
};
match crate::sdk::create_agent_session(options).await {
Ok(new_handle) => {
*current_ask
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) = new_handle.ask_tool();
if let Some(ask) = new_handle.ask_tool() {
install_ask_forwarder(&ask, agent_tx, runtime_handle);
}
send_conversation_reset(
&new_handle,
agent_tx,
&format!(
"Started new session\nModel set to {provider}/{model_id}\nThinking level: off"
),
)
.await;
Some(new_handle)
}
Err(err) => {
let _ = agent_tx.send(PiMsg::AgentError(format!("new session: {err}")));
None
}
}
}
async fn run_session_info_command(
handle: &crate::sdk::AgentSessionHandle,
agent_tx: &Sender<PiMsg>,
) {
let state = match handle.state().await {
Ok(state) => state,
Err(err) => {
let _ = agent_tx.send(PiMsg::AgentError(format!("session info: {err}")));
return;
}
};
let info = handle
.with_session(|session| {
let file = session.path.as_ref().map_or_else(
|| String::from("(not saved yet)"),
|p| p.display().to_string(),
);
let name = session.get_name().unwrap_or_else(|| String::from("-"));
format!(
"Session info:\n file: {file}\n id: {id}\n name: {name}\n model: {provider}/{model_id}\n thinking: {thinking}\n messageCount: {message_count}",
id = state.session_id.as_deref().unwrap_or("-"),
provider = state.provider,
model_id = state.model_id,
thinking = state
.thinking_level
.as_ref()
.map_or_else(|| String::from("off"), ToString::to_string),
message_count = state.message_count,
)
})
.await;
match info {
Ok(text) => {
let _ = agent_tx.send(PiMsg::System(text));
}
Err(err) => {
let _ = agent_tx.send(PiMsg::AgentError(format!("session info: {err}")));
}
}
}
async fn run_tree_summary_command(
handle: &crate::sdk::AgentSessionHandle,
agent_tx: &Sender<PiMsg>,
) {
let summary = handle.with_session(|session| {
let leaves = session.list_leaves();
let entry_count = session.entries.len();
if leaves.is_empty() {
return format!("Session tree: no branches, {entry_count} entries");
}
let rendered = leaves
.iter()
.map(String::as_str)
.collect::<Vec<_>>()
.join("\n ");
format!(
"Session tree: {} branch(es), {entry_count} entries\nLeaves:\n {rendered}",
leaves.len()
)
});
match summary.await {
Ok(text) => {
let _ = agent_tx.send(PiMsg::System(text));
}
Err(err) => {
let _ = agent_tx.send(PiMsg::AgentError(format!("tree: {err}")));
}
}
}
async fn run_set_thinking_command(
handle: &mut crate::sdk::AgentSessionHandle,
level: Option<crate::model::ThinkingLevel>,
agent_tx: &Sender<PiMsg>,
) {
let msg = match level {
None => match handle.state().await {
Ok(state) => PiMsg::System(format!(
"Thinking level: {}",
state
.thinking_level
.as_ref()
.map_or_else(|| String::from("off"), ToString::to_string)
)),
Err(err) => PiMsg::AgentError(format!("thinking: {err}")),
},
Some(level) => match handle.set_thinking_level(level).await {
Ok(()) => PiMsg::System(format!("Thinking level: {level}")),
Err(err) => PiMsg::AgentError(format!("thinking: {err}")),
},
};
let _ = agent_tx.send(msg);
}
async fn run_set_name_command(
handle: &mut crate::sdk::AgentSessionHandle,
name: &str,
agent_tx: &Sender<PiMsg>,
) {
let msg = match handle.set_session_name(name).await {
Ok(()) => PiMsg::System(format!("Session name: {name}")),
Err(err) => PiMsg::AgentError(format!("name: {err}")),
};
let _ = agent_tx.send(msg);
}
async fn resume_session_command(
path: &str,
template: &crate::sdk::SessionOptions,
current_ask: &CurrentAsk,
ext_handler: &Arc<FtuiExtensionUiHandler>,
agent_tx: &Sender<PiMsg>,
runtime_handle: &asupersync::runtime::RuntimeHandle,
) -> Option<crate::sdk::AgentSessionHandle> {
let options = crate::sdk::SessionOptions {
session_path: Some(std::path::PathBuf::from(path)),
provider: template.provider.clone(),
model: template.model.clone(),
api_key: template.api_key.clone(),
working_directory: template.working_directory.clone(),
session_dir: template.session_dir.clone(),
extension_paths: template.extension_paths.clone(),
extension_policy: template.extension_policy.clone(),
extension_ui_handler: Some(
Arc::clone(ext_handler) as Arc<dyn crate::sdk::ExtensionUiHandler>
),
no_session: false,
..Default::default()
};
match crate::sdk::create_agent_session(options).await {
Ok(handle) => {
*current_ask
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) = handle.ask_tool();
if let Some(ask) = handle.ask_tool() {
install_ask_forwarder(&ask, agent_tx, runtime_handle);
}
send_conversation_reset(&handle, agent_tx, "session resumed").await;
Some(handle)
}
Err(err) => {
let _ = agent_tx.send(PiMsg::AgentError(format!("resume: {err}")));
None
}
}
}
async fn send_conversation_reset(
handle: &crate::sdk::AgentSessionHandle,
agent_tx: &Sender<PiMsg>,
status: &str,
) {
match handle
.with_session(crate::interactive::conversation_from_session)
.await
{
Ok((messages, usage)) => {
let _ = agent_tx.send(PiMsg::ConversationReset {
messages,
usage,
status: Some(status.to_string()),
});
}
Err(err) => {
let _ = agent_tx.send(PiMsg::AgentError(format!("conversation snapshot: {err}")));
}
}
}
async fn run_bash_ui_command(
cwd: &std::path::Path,
command: &str,
exclude: bool,
agent_tx: &Sender<PiMsg>,
) -> Option<String> {
let _ = agent_tx.send(PiMsg::ToolStart {
name: String::from("bash"),
tool_id: String::from("ftui-bash"),
});
let result = crate::tools::run_bash_command(cwd, None, None, command, None, None).await;
let output = match result {
Ok(result) => {
let display = crate::session::bash_execution_to_text(
command,
&result.output,
result.exit_code,
result.cancelled,
result.truncated,
result.full_output_path.as_deref(),
);
let mut shown = display.clone();
if exclude {
shown.push_str("\n\n[Output excluded from model context]");
}
let _ = agent_tx.send(PiMsg::BashResult {
display: shown,
content_for_agent: None,
});
Some(display)
}
Err(err) => {
let _ = agent_tx.send(PiMsg::AgentError(format!("bash: {err}")));
None
}
};
let _ = agent_tx.send(PiMsg::ToolEnd {
name: String::from("bash"),
tool_id: String::from("ftui-bash"),
is_error: false,
});
output
}
pub fn run(
session_options: crate::sdk::SessionOptions,
theme: &crate::theme::Theme,
inline: bool,
available_models: Vec<String>,
available_sessions: Vec<(String, String)>,
) -> std::io::Result<()> {
let (submit_tx, submit_rx) = std::sync::mpsc::channel::<UiCommand>();
let (agent_tx, agent_rx) = std::sync::mpsc::channel::<PiMsg>();
let (ask_reply_tx, ask_reply_rx) = std::sync::mpsc::channel::<AskUiReply>();
let (ext_reply_tx, ext_reply_rx) = std::sync::mpsc::channel::<ExtensionUiResponse>();
let bash_cwd = session_options
.working_directory
.clone()
.or_else(|| std::env::current_dir().ok())
.unwrap_or_else(|| std::path::PathBuf::from("."));
let resume_template = resume_template_from(&session_options);
let driver = std::thread::Builder::new()
.name("pi-ftui-agent-driver".into())
.spawn(move || {
let runtime = match asupersync::runtime::RuntimeBuilder::new().build() {
Ok(runtime) => runtime,
Err(err) => {
let _ = agent_tx.send(PiMsg::AgentError(format!("runtime build: {err}")));
return;
}
};
let runtime_handle = runtime.handle();
runtime.block_on(async move {
let mut session_options = session_options;
let ext_handler = Arc::new(FtuiExtensionUiHandler::new(agent_tx.clone()));
session_options.extension_ui_handler =
Some(Arc::clone(&ext_handler) as Arc<dyn crate::sdk::ExtensionUiHandler>);
spawn_ext_reply_pump(Arc::clone(&ext_handler), ext_reply_rx, &runtime_handle);
let mut handle = match crate::sdk::create_agent_session(session_options).await {
Ok(handle) => handle,
Err(err) => {
let _ = agent_tx.send(PiMsg::AgentError(format!("session: {err}")));
return;
}
};
let current_ask =
install_ask_bridges(&handle, &agent_tx, ask_reply_rx, &runtime_handle);
let _ = agent_tx.send(PiMsg::System(String::from(
"ftui preview stack — experimental (bd-cv653.9.1)",
)));
loop {
match submit_rx.try_recv() {
Ok(UiCommand::Prompt(prompt)) => {
run_prompt_turn(&mut handle, prompt, &agent_tx).await;
}
Ok(UiCommand::SetModel { provider, model }) => {
run_set_model_command(&mut handle, &provider, &model, &agent_tx).await;
}
Ok(UiCommand::Bash { command, exclude }) => {
if let Some(output) =
run_bash_ui_command(&bash_cwd, &command, exclude, &agent_tx).await
&& !exclude
{
run_prompt_turn(&mut handle, output, &agent_tx).await;
}
}
Ok(UiCommand::Compact) => {
run_compact_command(&mut handle, &agent_tx).await;
}
Ok(UiCommand::ExtensionCommand { name, args }) => {
run_extension_command(&handle, &bash_cwd, &name, &args, &agent_tx)
.await;
}
Ok(UiCommand::ResumeSession { path }) => {
if let Some(new_handle) = resume_session_command(
&path,
&resume_template,
¤t_ask,
&ext_handler,
&agent_tx,
&runtime_handle,
)
.await
{
handle = new_handle;
}
}
Ok(UiCommand::NewSession) => {
if let Some(new_handle) = new_session_command(
&resume_template,
&handle,
¤t_ask,
&ext_handler,
&agent_tx,
&runtime_handle,
)
.await
{
handle = new_handle;
}
}
Ok(UiCommand::SessionInfo) => {
run_session_info_command(&handle, &agent_tx).await;
}
Ok(UiCommand::TreeSummary) => {
run_tree_summary_command(&handle, &agent_tx).await;
}
Ok(UiCommand::SetThinking(level)) => {
run_set_thinking_command(&mut handle, level, &agent_tx).await;
}
Ok(UiCommand::SetName(name)) => {
run_set_name_command(&mut handle, &name, &agent_tx).await;
}
Err(std::sync::mpsc::TryRecvError::Empty) => {
asupersync::time::sleep(asupersync::time::wall_now(), SUBMIT_POLL)
.await;
}
Err(std::sync::mpsc::TryRecvError::Disconnected) => break,
}
}
});
})?;
let model = PiFtuiModel::new(agent_rx)
.with_submit_channel(submit_tx)
.with_ask_reply_channel(ask_reply_tx)
.with_palette(FtuiPalette::from_theme(theme))
.with_available_models(available_models)
.with_available_sessions(available_sessions)
.with_ext_reply_channel(ext_reply_tx);
let app = if inline {
ftui::App::inline_auto(model, INLINE_MIN_HEIGHT, INLINE_MAX_HEIGHT)
} else {
ftui::App::fullscreen(model)
};
let log_guard = crate::tui::TuiLogRedirectGuard::begin();
let result = app.with_mouse().run();
drop(log_guard);
let _ = driver.join();
result
}
#[cfg(test)]
mod tests {
use super::*;
use crate::model::StopReason;
use ftui::runtime::simulator::ProgramSimulator;
use ftui::{KeyEvent, KeyEventKind};
use std::sync::mpsc;
fn key(code: KeyCode, modifiers: Modifiers) -> Event {
Event::Key(KeyEvent {
code,
modifiers,
kind: KeyEventKind::Press,
})
}
fn new_model() -> (mpsc::Sender<PiMsg>, PiFtuiModel) {
let (tx, rx) = mpsc::channel();
(tx, PiFtuiModel::new(rx))
}
#[test]
fn streaming_deltas_accumulate_and_flush_on_done() {
let (_tx, model) = new_model();
let mut sim = ProgramSimulator::new(model);
sim.init();
sim.send(PiFtuiMsg::Agent(PiMsg::AgentStart));
assert_eq!(sim.model().state, AgentUiState::Working);
sim.send(PiFtuiMsg::Agent(PiMsg::TextDelta("hello ".into())));
sim.send(PiFtuiMsg::Agent(PiMsg::TextDelta("world".into())));
assert_eq!(sim.model().streaming, "hello world");
sim.send(PiFtuiMsg::Agent(PiMsg::AgentDone {
usage: None,
stop_reason: StopReason::Stop,
error_message: None,
}));
assert_eq!(sim.model().state, AgentUiState::Ready);
let transcript = &sim.model().transcript;
assert_eq!(transcript.len(), 1);
assert_eq!(transcript[0].text, "hello world");
assert_eq!(transcript[0].role, EntryRole::Assistant);
assert!(sim.model().streaming.is_empty());
}
#[test]
fn agent_text_is_sanitized_before_display() {
let (_tx, model) = new_model();
let mut sim = ProgramSimulator::new(model);
sim.init();
sim.send(PiFtuiMsg::Agent(PiMsg::TextDelta(
"safe\x1b]0;pwned\x07 text".into(),
)));
let streamed = sim.model().streaming.clone();
assert!(!streamed.contains('\x1b'), "ESC survived: {streamed:?}");
assert!(!streamed.contains('\x07'), "BEL survived: {streamed:?}");
assert!(streamed.contains("safe"));
assert!(streamed.contains("text"));
}
#[test]
fn ctrl_c_quits() {
let (_tx, model) = new_model();
let mut sim = ProgramSimulator::new(model);
sim.init();
sim.inject_event(key(KeyCode::Char('c'), Modifiers::CTRL));
assert!(!sim.is_running());
}
fn buffer_text(buf: &ftui::Buffer, width: u16, height: u16) -> String {
let mut out = String::new();
for y in 0..height {
for x in 0..width {
let ch = buf
.get(x, y)
.and_then(|cell| cell.content.as_char())
.unwrap_or(' ');
out.push(ch);
}
out.push('\n');
}
out
}
#[test]
fn view_renders_transcript_and_status() {
let (_tx, model) = new_model();
let mut sim = ProgramSimulator::new(model);
sim.init();
sim.send(PiFtuiMsg::Agent(PiMsg::System("session restored".into())));
let rendered = buffer_text(sim.capture_frame(40, 8), 40, 8);
assert!(
rendered.contains("session restored"),
"frame missing transcript line: {rendered:?}"
);
assert!(rendered.contains("pi · ready"), "frame missing header");
assert!(
rendered.contains("Type a message"),
"frame missing input placeholder: {rendered:?}"
);
}
#[test]
fn typing_and_enter_submits_to_channel_and_transcript() {
let (_agent_tx, rx) = mpsc::channel();
let (submit_tx, submit_rx) = mpsc::channel::<UiCommand>();
let model = PiFtuiModel::new(rx).with_submit_channel(submit_tx);
let mut sim = ProgramSimulator::new(model);
sim.init();
for ch in ['h', 'i'] {
sim.inject_event(key(KeyCode::Char(ch), Modifiers::empty()));
}
assert_eq!(sim.model().input.text(), "hi");
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
assert_eq!(
submit_rx.try_recv().expect("submitted"),
UiCommand::Prompt("hi".into())
);
assert!(sim.model().input.is_empty(), "editor not cleared");
let transcript = &sim.model().transcript;
assert_eq!(transcript.len(), 1);
assert_eq!(transcript[0].text, "hi");
assert_eq!(transcript[0].role, EntryRole::User);
}
#[test]
fn alt_enter_inserts_newline_and_grows_input_region() {
let (_tx, model) = new_model();
let mut sim = ProgramSimulator::new(model);
sim.init();
assert_eq!(sim.model().input_rows(), 1);
sim.inject_event(key(KeyCode::Char('a'), Modifiers::empty()));
sim.inject_event(key(KeyCode::Enter, Modifiers::ALT));
sim.inject_event(key(KeyCode::Char('b'), Modifiers::empty()));
assert_eq!(sim.model().input.text(), "a\nb");
assert_eq!(sim.model().input_rows(), 2);
}
#[test]
fn empty_submit_is_a_noop() {
let (_tx, model) = new_model();
let mut sim = ProgramSimulator::new(model);
sim.init();
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
assert!(sim.model().transcript.is_empty());
}
#[test]
fn editor_ignores_keys_while_agent_works() {
let (_tx, model) = new_model();
let mut sim = ProgramSimulator::new(model);
sim.init();
sim.send(PiFtuiMsg::Agent(PiMsg::AgentStart));
sim.inject_event(key(KeyCode::Char('x'), Modifiers::empty()));
assert!(
sim.model().input.is_empty(),
"editor took input while working"
);
let rendered = buffer_text(sim.capture_frame(40, 8), 40, 8);
assert!(rendered.contains("processing"), "missing processing note");
}
#[test]
fn submitted_text_is_sanitized() {
let (_agent_tx, rx) = mpsc::channel();
let (submit_tx, submit_rx) = mpsc::channel::<UiCommand>();
let model = PiFtuiModel::new(rx).with_submit_channel(submit_tx);
let mut sim = ProgramSimulator::new(model);
sim.init();
sim.inject_event(Event::Paste(ftui::PasteEvent::new(
"hello\x1b]0;pwned\x07world",
true,
)));
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
let UiCommand::Prompt(submitted) = submit_rx.try_recv().expect("submitted") else {
panic!("expected a prompt command");
};
assert!(!submitted.contains('\x1b'), "ESC survived: {submitted:?}");
assert!(submitted.contains("hello"));
assert!(submitted.contains("world"));
}
#[test]
fn slash_model_routes_set_model_and_bad_specs_error() {
let (_agent_tx, rx) = mpsc::channel();
let (submit_tx, submit_rx) = mpsc::channel::<UiCommand>();
let model = PiFtuiModel::new(rx).with_submit_channel(submit_tx);
let mut sim = ProgramSimulator::new(model);
sim.init();
type_str(&mut sim, "/model openai/gpt-5");
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
assert_eq!(
submit_rx.try_recv().expect("routed"),
UiCommand::SetModel {
provider: "openai".into(),
model: "gpt-5".into(),
}
);
type_str(&mut sim, "/model nonsense");
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
assert!(submit_rx.try_recv().is_err(), "bad spec reached the driver");
assert!(
sim.model()
.transcript
.iter()
.any(|e| e.role == EntryRole::Error && e.text.contains("usage: /model")),
"usage error missing"
);
}
#[test]
fn non_builtin_slash_commands_route_to_extension_dispatch() {
let (_agent_tx, rx) = mpsc::channel();
let (submit_tx, submit_rx) = mpsc::channel::<UiCommand>();
let model = PiFtuiModel::new(rx).with_submit_channel(submit_tx);
let mut sim = ProgramSimulator::new(model);
sim.init();
type_str(&mut sim, "/deploy --force");
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
assert_eq!(
submit_rx.try_recv().expect("routed"),
UiCommand::ExtensionCommand {
name: "deploy".into(),
args: "--force".into(),
}
);
type_str(&mut sim, "/help");
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
assert!(submit_rx.try_recv().is_err());
assert!(
sim.model()
.transcript
.iter()
.any(|e| e.role == EntryRole::System && e.text.contains("/model")),
"help text missing"
);
}
#[test]
fn tool_status_renders_while_running() {
let (_tx, model) = new_model();
let mut sim = ProgramSimulator::new(model);
sim.init();
sim.send(PiFtuiMsg::Agent(PiMsg::AgentStart));
sim.send(PiFtuiMsg::Agent(PiMsg::ToolStart {
name: "bash".into(),
tool_id: "t1".into(),
}));
let rendered = buffer_text(sim.capture_frame(40, 8), 40, 8);
assert!(
rendered.contains("running bash"),
"missing tool status: {rendered:?}"
);
sim.send(PiFtuiMsg::Agent(PiMsg::ToolEnd {
name: "bash".into(),
tool_id: "t1".into(),
is_error: false,
}));
let rendered = buffer_text(sim.capture_frame(40, 8), 40, 8);
assert!(
!rendered.contains("running bash"),
"tool status not cleared"
);
assert!(
rendered.contains("✓ bash"),
"durable tool trace missing: {rendered:?}"
);
sim.send(PiFtuiMsg::Agent(PiMsg::ToolEnd {
name: "edit".into(),
tool_id: "t2".into(),
is_error: true,
}));
assert!(
sim.model()
.transcript
.iter()
.any(|e| e.text.contains("✗ edit")),
"error trace missing"
);
}
#[test]
fn scroll_pins_view_and_end_resumes_tail_follow() {
let (_tx, model) = new_model();
let mut sim = ProgramSimulator::new(model);
sim.init();
sim.inject_event(Event::Resize {
width: 30,
height: 10,
});
for i in 0..20 {
sim.send(PiFtuiMsg::Agent(PiMsg::System(format!("line-{i}"))));
}
let rendered = buffer_text(sim.capture_frame(30, 10), 30, 10);
assert!(
rendered.contains("line-19"),
"tail not followed: {rendered:?}"
);
assert!(
!rendered.contains("line-0 "),
"oldest line unexpectedly visible"
);
sim.inject_event(key(KeyCode::PageUp, Modifiers::empty()));
let rendered = buffer_text(sim.capture_frame(30, 10), 30, 10);
assert!(!rendered.contains("line-19"), "still at tail after PageUp");
assert!(
rendered.contains("lines up"),
"footer missing scroll indicator"
);
sim.send(PiFtuiMsg::Agent(PiMsg::System("line-20".into())));
let rendered = buffer_text(sim.capture_frame(30, 10), 30, 10);
assert!(
!rendered.contains("line-20"),
"pinned view was yanked to tail"
);
sim.inject_event(key(KeyCode::End, Modifiers::empty()));
let rendered = buffer_text(sim.capture_frame(30, 10), 30, 10);
assert!(
rendered.contains("line-20"),
"End did not resume tail follow"
);
}
#[test]
fn resize_reclamps_scroll_offset() {
let (_tx, model) = new_model();
let mut sim = ProgramSimulator::new(model);
sim.init();
sim.inject_event(Event::Resize {
width: 30,
height: 10,
});
for i in 0..12 {
sim.send(PiFtuiMsg::Agent(PiMsg::System(format!("line-{i}"))));
}
sim.inject_event(key(KeyCode::PageUp, Modifiers::empty()));
assert!(sim.model().scroll_from_tail > 0);
sim.inject_event(Event::Resize {
width: 30,
height: 40,
});
assert_eq!(sim.model().scroll_from_tail, 0);
}
#[test]
fn drain_loop_bridges_agent_channel_until_disconnect() {
let (agent_tx, agent_rx) = mpsc::channel::<PiMsg>();
let (msg_tx, msg_rx) = mpsc::channel::<PiFtuiMsg>();
let handle = std::thread::spawn(move || {
drain_agent_events(&agent_rx, &msg_tx, || false);
});
agent_tx.send(PiMsg::AgentStart).unwrap();
let bridged = msg_rx
.recv_timeout(Duration::from_secs(5))
.expect("bridged message");
assert!(matches!(bridged, PiFtuiMsg::Agent(PiMsg::AgentStart)));
drop(agent_tx);
handle.join().expect("bridge thread exits cleanly");
}
#[test]
fn drain_loop_honors_stop_predicate() {
let (_agent_tx, agent_rx) = mpsc::channel::<PiMsg>();
let (msg_tx, _msg_rx) = mpsc::channel::<PiFtuiMsg>();
drain_agent_events(&agent_rx, &msg_tx, || true);
}
#[test]
fn spinner_ticks_while_working_and_stops_when_idle() {
use ftui::runtime::simulator::CmdRecord;
let (_tx, model) = new_model();
let mut sim = ProgramSimulator::new(model);
sim.init();
sim.send(PiFtuiMsg::Agent(PiMsg::AgentStart));
assert!(
matches!(sim.command_log().last(), Some(CmdRecord::Tick(_))),
"AgentStart did not schedule a tick: {:?}",
sim.command_log().last()
);
let frame_before = sim.model().spinner.current_frame;
sim.inject_event(Event::Tick);
assert_eq!(sim.model().spinner.current_frame, frame_before + 1);
assert!(matches!(sim.command_log().last(), Some(CmdRecord::Tick(_))));
let spin = DOTS[sim.model().spinner.current_frame % DOTS.len()];
let rendered = buffer_text(sim.capture_frame(40, 8), 40, 8);
assert!(
rendered.contains(spin),
"status missing spinner frame {spin:?}: {rendered:?}"
);
sim.send(PiFtuiMsg::Agent(PiMsg::AgentDone {
usage: None,
stop_reason: StopReason::Stop,
error_message: None,
}));
let frame_after_done = sim.model().spinner.current_frame;
sim.inject_event(Event::Tick);
assert_eq!(sim.model().spinner.current_frame, frame_after_done);
assert!(matches!(sim.command_log().last(), Some(CmdRecord::None)));
}
#[test]
fn thinking_status_then_responding_then_usage_footer() {
let (_tx, model) = new_model();
let mut sim = ProgramSimulator::new(model);
sim.init();
sim.send(PiFtuiMsg::Agent(PiMsg::AgentStart));
sim.send(PiFtuiMsg::Agent(PiMsg::ThinkingDelta(
"mull it over".into(),
)));
let rendered = buffer_text(sim.capture_frame(44, 8), 44, 8);
assert!(
rendered.contains("thinking ..."),
"missing thinking: {rendered:?}"
);
sim.send(PiFtuiMsg::Agent(PiMsg::TextDelta("answer".into())));
let rendered = buffer_text(sim.capture_frame(44, 8), 44, 8);
assert!(
rendered.contains("responding ..."),
"missing responding: {rendered:?}"
);
sim.send(PiFtuiMsg::Agent(PiMsg::AgentDone {
usage: Some(crate::model::Usage {
input: 120,
output: 45,
total_tokens: 165,
..Default::default()
}),
stop_reason: StopReason::Stop,
error_message: None,
}));
let rendered = buffer_text(sim.capture_frame(44, 8), 44, 8);
assert!(
rendered.contains("tokens 120↑ 45↓ · total 165"),
"missing usage footer: {rendered:?}"
);
assert!(sim.model().thinking.is_empty(), "thinking not cleared");
}
fn ask_request(id: &str, questions: Vec<crate::ask::AskQuestion>) -> AskUiRequest {
AskUiRequest {
id: id.to_string(),
request: crate::ask::AskRequest { questions },
}
}
fn question(q: &str, options: &[&str], multi: bool) -> crate::ask::AskQuestion {
crate::ask::AskQuestion {
id: None,
question: q.to_string(),
header: None,
options: options
.iter()
.map(|label| crate::ask::AskOption {
label: (*label).to_string(),
description: None,
})
.collect(),
multi,
recommended: None,
}
}
fn type_str(sim: &mut ProgramSimulator<PiFtuiModel>, s: &str) {
for ch in s.chars() {
sim.inject_event(key(KeyCode::Char(ch), Modifiers::empty()));
}
}
#[test]
fn ask_card_collects_answers_across_questions() {
let (agent_tx, agent_rx) = mpsc::channel();
let (reply_tx, reply_rx) = mpsc::channel::<AskUiReply>();
let model = PiFtuiModel::new(agent_rx).with_ask_reply_channel(reply_tx);
drop(agent_tx);
let mut sim = ProgramSimulator::new(model);
sim.init();
sim.send(PiFtuiMsg::Agent(PiMsg::AgentStart));
sim.send(PiFtuiMsg::Agent(PiMsg::AskUiRequest(ask_request(
"ask-1",
vec![
question("Pick a color?", &["red", "blue"], false),
question("Pick tools?", &["hammer", "saw"], true),
],
))));
let rendered = buffer_text(sim.capture_frame(50, 12), 50, 12);
assert!(
rendered.contains("Pick a color?"),
"card not rendered: {rendered:?}"
);
type_str(&mut sim, "2");
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
let rendered = buffer_text(sim.capture_frame(50, 14), 50, 14);
assert!(
rendered.contains("Pick tools?"),
"second card missing: {rendered:?}"
);
type_str(&mut sim, "hammer, saw");
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
let reply = reply_rx.try_recv().expect("ask reply sent");
assert_eq!(reply.request_id, "ask-1");
assert!(!reply.response.dismissed);
assert_eq!(reply.response.answers.len(), 2);
assert_eq!(reply.response.answers[0].selected, vec!["blue".to_string()]);
assert_eq!(
reply.response.answers[1].selected,
vec!["hammer".to_string(), "saw".to_string()]
);
assert!(sim.model().active_ask.is_none(), "ask not cleared");
}
#[test]
fn ask_cancel_dismisses() {
let (agent_tx, agent_rx) = mpsc::channel();
let (reply_tx, reply_rx) = mpsc::channel::<AskUiReply>();
let model = PiFtuiModel::new(agent_rx).with_ask_reply_channel(reply_tx);
drop(agent_tx);
let mut sim = ProgramSimulator::new(model);
sim.init();
sim.send(PiFtuiMsg::Agent(PiMsg::AgentStart));
sim.send(PiFtuiMsg::Agent(PiMsg::AskUiRequest(ask_request(
"ask-2",
vec![question("Sure?", &["yes", "no"], false)],
))));
type_str(&mut sim, "cancel");
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
let reply = reply_rx.try_recv().expect("dismissal sent");
assert!(reply.response.dismissed);
assert!(reply.response.answers.is_empty());
assert!(sim.model().active_ask.is_none());
}
#[test]
fn ask_free_text_becomes_other_answer() {
let (agent_tx, agent_rx) = mpsc::channel();
let (reply_tx, reply_rx) = mpsc::channel::<AskUiReply>();
let model = PiFtuiModel::new(agent_rx).with_ask_reply_channel(reply_tx);
drop(agent_tx);
let mut sim = ProgramSimulator::new(model);
sim.init();
sim.send(PiFtuiMsg::Agent(PiMsg::AgentStart));
sim.send(PiFtuiMsg::Agent(PiMsg::AskUiRequest(ask_request(
"ask-3",
vec![question("Which env?", &["dev", "prod"], false)],
))));
type_str(&mut sim, "staging with canary");
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
let reply = reply_rx.try_recv().expect("reply sent");
assert_eq!(
reply.response.answers[0].other.as_deref(),
Some("staging with canary")
);
assert!(reply.response.answers[0].selected.is_empty());
}
#[test]
fn catalog_routes_shift_enter_newline_and_ctrl_d_exit() {
let (_tx, model) = new_model();
let mut sim = ProgramSimulator::new(model);
sim.init();
sim.inject_event(key(KeyCode::Char('a'), Modifiers::empty()));
sim.inject_event(key(KeyCode::Enter, Modifiers::SHIFT));
sim.inject_event(key(KeyCode::Char('b'), Modifiers::empty()));
assert_eq!(sim.model().input.text(), "a\nb");
sim.inject_event(key(KeyCode::Char('d'), Modifiers::CTRL));
assert!(sim.is_running(), "ctrl+d exited despite editor content");
sim.model_mut().input.set_text("");
sim.inject_event(key(KeyCode::Char('d'), Modifiers::CTRL));
assert!(!sim.is_running(), "ctrl+d on empty editor did not exit");
}
fn ext_request(id: &str, method: &str, payload: serde_json::Value) -> ExtensionUiRequest {
ExtensionUiRequest {
id: id.to_string(),
method: method.to_string(),
payload,
timeout_ms: None,
extension_id: Some(String::from("demo-ext")),
}
}
#[test]
fn extension_confirm_prompt_renders_and_reply_routes() {
let (_agent_tx, rx) = mpsc::channel();
let (ext_tx, ext_rx) = mpsc::channel::<ExtensionUiResponse>();
let model = PiFtuiModel::new(rx).with_ext_reply_channel(ext_tx);
let mut sim = ProgramSimulator::new(model);
sim.init();
sim.send(PiFtuiMsg::Agent(PiMsg::AgentStart));
sim.send(PiFtuiMsg::Agent(PiMsg::ExtensionUiRequest(ext_request(
"ext-1",
"confirm",
serde_json::json!({"title": "Deploy?", "message": "Ship to prod?"}),
))));
let rendered = buffer_text(sim.capture_frame(50, 12), 50, 12);
assert!(rendered.contains("Deploy?"), "prompt missing: {rendered:?}");
assert!(
rendered.contains("demo-ext"),
"provenance missing: {rendered:?}"
);
type_str(&mut sim, "yes");
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
let reply = ext_rx.try_recv().expect("reply routed");
assert_eq!(reply.id, "ext-1");
assert!(!reply.cancelled);
assert_eq!(reply.value, Some(serde_json::Value::Bool(true)));
assert!(sim.model().active_ext.is_none());
}
#[test]
fn extension_prompt_escape_cancels_and_queue_advances() {
let (_agent_tx, rx) = mpsc::channel();
let (ext_tx, ext_rx) = mpsc::channel::<ExtensionUiResponse>();
let model = PiFtuiModel::new(rx).with_ext_reply_channel(ext_tx);
let mut sim = ProgramSimulator::new(model);
sim.init();
sim.send(PiFtuiMsg::Agent(PiMsg::ExtensionUiRequest(ext_request(
"ext-a",
"confirm",
serde_json::json!({"title": "First?"}),
))));
sim.send(PiFtuiMsg::Agent(PiMsg::ExtensionUiRequest(ext_request(
"ext-b",
"confirm",
serde_json::json!({"title": "Second?"}),
))));
assert_eq!(sim.model().ext_queue.len(), 1, "second request not queued");
sim.inject_event(key(KeyCode::Escape, Modifiers::empty()));
let reply = ext_rx.try_recv().expect("cancel routed");
assert_eq!(reply.id, "ext-a");
assert!(reply.cancelled);
assert_eq!(
sim.model().active_ext.as_ref().map(|r| r.id.as_str()),
Some("ext-b")
);
}
#[test]
fn extension_prompt_queues_behind_active_ask() {
let (_agent_tx, rx) = mpsc::channel();
let (ask_tx, _ask_rx) = mpsc::channel::<AskUiReply>();
let (ext_tx, _ext_rx) = mpsc::channel::<ExtensionUiResponse>();
let model = PiFtuiModel::new(rx)
.with_ask_reply_channel(ask_tx)
.with_ext_reply_channel(ext_tx);
let mut sim = ProgramSimulator::new(model);
sim.init();
sim.send(PiFtuiMsg::Agent(PiMsg::AgentStart));
sim.send(PiFtuiMsg::Agent(PiMsg::AskUiRequest(ask_request(
"ask-hold",
vec![question("Pick?", &["a", "b"], false)],
))));
sim.send(PiFtuiMsg::Agent(PiMsg::ExtensionUiRequest(ext_request(
"ext-waiting",
"confirm",
serde_json::json!({"title": "Later?"}),
))));
assert!(sim.model().active_ext.is_none(), "ext jumped the ask");
assert_eq!(sim.model().ext_queue.len(), 1);
type_str(&mut sim, "1");
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
assert_eq!(
sim.model().active_ext.as_ref().map(|r| r.id.as_str()),
Some("ext-waiting")
);
}
#[test]
fn escape_dismisses_active_ask() {
let (agent_tx, agent_rx) = mpsc::channel();
let (reply_tx, reply_rx) = mpsc::channel::<AskUiReply>();
let model = PiFtuiModel::new(agent_rx).with_ask_reply_channel(reply_tx);
drop(agent_tx);
let mut sim = ProgramSimulator::new(model);
sim.init();
sim.send(PiFtuiMsg::Agent(PiMsg::AgentStart));
sim.send(PiFtuiMsg::Agent(PiMsg::AskUiRequest(ask_request(
"ask-esc",
vec![question("Continue?", &["yes", "no"], false)],
))));
sim.inject_event(key(KeyCode::Escape, Modifiers::empty()));
let reply = reply_rx.try_recv().expect("dismissal sent");
assert!(reply.response.dismissed);
assert!(sim.model().active_ask.is_none());
}
#[test]
fn agent_event_translation_covers_lifecycle_stream_and_tools() {
use crate::agent::AgentEvent as E;
use crate::model::{AssistantMessage, AssistantMessageEvent as A, Message, Usage};
use std::sync::Arc;
let msgs = agent_event_to_pi_msgs(&E::AgentStart {
session_id: Arc::from("s1"),
});
assert!(matches!(msgs.as_slice(), [PiMsg::AgentStart]));
let assistant = Arc::new(AssistantMessage {
usage: Usage {
input: 10,
output: 5,
total_tokens: 15,
..Default::default()
},
stop_reason: StopReason::Stop,
..Default::default()
});
let partial = Arc::clone(&assistant);
let msgs = agent_event_to_pi_msgs(&E::MessageUpdate {
message: Message::Assistant(Arc::clone(&assistant)),
assistant_message_event: A::TextDelta {
content_index: 0,
delta: "hi".into(),
partial,
},
});
assert!(matches!(msgs.as_slice(), [PiMsg::TextDelta(d)] if d == "hi"));
let msgs = agent_event_to_pi_msgs(&E::ToolExecutionStart {
tool_call_id: "t1".into(),
tool_name: "bash".into(),
args: serde_json::json!({}),
});
assert!(
matches!(msgs.as_slice(), [PiMsg::ToolStart { name, tool_id }] if name == "bash" && tool_id == "t1")
);
let msgs = agent_event_to_pi_msgs(&E::AgentEnd {
session_id: Arc::from("s1"),
messages: vec![Message::Assistant(assistant)],
error: None,
});
match msgs.as_slice() {
[
PiMsg::AgentDone {
usage: Some(usage),
stop_reason: StopReason::Stop,
error_message: None,
},
] => assert_eq!(usage.total_tokens, 15),
other => panic!("unexpected translation: {other:?}"),
}
}
#[test]
fn assistant_markdown_renders_without_markers() {
let (_tx, model) = new_model();
let mut sim = ProgramSimulator::new(model);
sim.init();
sim.send(PiFtuiMsg::Agent(PiMsg::AgentStart));
sim.send(PiFtuiMsg::Agent(PiMsg::TextDelta(
"# Release Notes\n\nplain body".into(),
)));
sim.send(PiFtuiMsg::Agent(PiMsg::AgentDone {
usage: None,
stop_reason: StopReason::Stop,
error_message: None,
}));
let rendered = buffer_text(sim.capture_frame(50, 10), 50, 10);
assert!(
rendered.contains("Release Notes"),
"heading text missing: {rendered:?}"
);
assert!(
!rendered.contains("# Release Notes"),
"markdown marker leaked into frame: {rendered:?}"
);
assert!(
rendered.contains("plain body"),
"body missing: {rendered:?}"
);
}
#[test]
fn theme_picker_opens_navigates_applies_and_captures_keys() {
let (_tx, model) = new_model();
let mut sim = ProgramSimulator::new(model);
sim.init();
let dark_accent = sim.model().palette.accent;
type_str(&mut sim, "/theme");
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
assert!(sim.model().picker.is_some(), "picker did not open");
let rendered = buffer_text(sim.capture_frame(50, 10), 50, 10);
assert!(
rendered.contains("Theme"),
"picker title missing: {rendered:?}"
);
assert!(rendered.contains("▸ dark"), "selection marker missing");
sim.inject_event(key(KeyCode::Char('j'), Modifiers::empty()));
assert!(sim.model().input.is_empty(), "picker leaked keys to editor");
assert_eq!(sim.model().picker.as_ref().unwrap().selected, 1);
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
assert!(sim.model().picker.is_none(), "picker did not close");
assert_ne!(
sim.model().palette.accent,
dark_accent,
"palette unchanged after applying light theme"
);
assert!(
sim.model()
.transcript
.iter()
.any(|e| e.text.contains("theme set to light")),
"confirmation note missing"
);
}
#[test]
fn bare_model_command_opens_picker_and_selection_routes_set_model() {
let (_agent_tx, rx) = mpsc::channel();
let (submit_tx, submit_rx) = mpsc::channel::<UiCommand>();
let model = PiFtuiModel::new(rx)
.with_submit_channel(submit_tx)
.with_available_models(vec![
String::from("openai/gpt-5"),
String::from("anthropic/claude-opus-5"),
]);
let mut sim = ProgramSimulator::new(model);
sim.init();
type_str(&mut sim, "/model");
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
assert!(sim.model().picker.is_some(), "picker did not open");
let rendered = buffer_text(sim.capture_frame(50, 10), 50, 10);
assert!(
rendered.contains("▸ openai/gpt-5"),
"first entry not selected: {rendered:?}"
);
sim.inject_event(key(KeyCode::Down, Modifiers::empty()));
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
assert_eq!(
submit_rx.try_recv().expect("routed"),
UiCommand::SetModel {
provider: "anthropic".into(),
model: "claude-opus-5".into(),
}
);
assert!(sim.model().picker.is_none());
}
#[test]
fn slash_compact_routes_command() {
let (_agent_tx, rx) = mpsc::channel();
let (submit_tx, submit_rx) = mpsc::channel::<UiCommand>();
let model = PiFtuiModel::new(rx).with_submit_channel(submit_tx);
let mut sim = ProgramSimulator::new(model);
sim.init();
type_str(&mut sim, "/compact");
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
assert_eq!(submit_rx.try_recv().expect("routed"), UiCommand::Compact);
assert!(
sim.model()
.transcript
.iter()
.any(|e| e.text.contains("compacting")),
"compact note missing"
);
}
#[test]
fn slash_exit_quits() {
let (_tx, model) = new_model();
let mut sim = ProgramSimulator::new(model);
sim.init();
type_str(&mut sim, "/exit");
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
assert!(!sim.is_running(), "/exit did not quit");
}
#[test]
fn bare_model_command_errors_without_registry() {
let (_tx, model) = new_model();
let mut sim = ProgramSimulator::new(model);
sim.init();
type_str(&mut sim, "/model");
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
assert!(sim.model().picker.is_none());
assert!(
sim.model()
.transcript
.iter()
.any(|e| e.role == EntryRole::Error && e.text.contains("no models available")),
"empty-registry error missing"
);
}
#[test]
fn resume_picker_shows_labels_and_routes_paths() {
let (_agent_tx, rx) = mpsc::channel();
let (submit_tx, submit_rx) = mpsc::channel::<UiCommand>();
let model = PiFtuiModel::new(rx)
.with_submit_channel(submit_tx)
.with_available_sessions(vec![
(
String::from("fix parser · 12 msgs"),
String::from("/tmp/sessions/a.jsonl"),
),
(
String::from("older run · 3 msgs"),
String::from("/tmp/sessions/b.jsonl"),
),
]);
let mut sim = ProgramSimulator::new(model);
sim.init();
type_str(&mut sim, "/resume");
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
let rendered = buffer_text(sim.capture_frame(50, 10), 50, 10);
assert!(
rendered.contains("▸ fix parser · 12 msgs"),
"labels not shown: {rendered:?}"
);
assert!(
!rendered.contains("/tmp/sessions"),
"paths leaked into display: {rendered:?}"
);
sim.inject_event(key(KeyCode::Char('j'), Modifiers::empty()));
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
assert_eq!(
submit_rx.try_recv().expect("routed"),
UiCommand::ResumeSession {
path: "/tmp/sessions/b.jsonl".into()
}
);
}
#[test]
fn conversation_reset_rebuilds_transcript() {
use crate::interactive::{ConversationMessage, MessageRole};
let (_tx, model) = new_model();
let mut sim = ProgramSimulator::new(model);
sim.init();
sim.send(PiFtuiMsg::Agent(PiMsg::System("old line".into())));
sim.send(PiFtuiMsg::Agent(PiMsg::ConversationReset {
messages: vec![
ConversationMessage {
role: MessageRole::User,
content: "restore me".into(),
thinking: None,
collapsed: false,
},
ConversationMessage {
role: MessageRole::Assistant,
content: "restored reply".into(),
thinking: None,
collapsed: false,
},
],
usage: crate::model::Usage::default(),
status: Some("session resumed".into()),
}));
let transcript = &sim.model().transcript;
assert!(
!transcript.iter().any(|e| e.text.contains("old line")),
"stale transcript survived reset"
);
assert!(
transcript
.iter()
.any(|e| e.role == EntryRole::User && e.text == "restore me")
);
assert!(
transcript
.iter()
.any(|e| e.role == EntryRole::Assistant && e.text == "restored reply")
);
assert!(
transcript
.iter()
.any(|e| e.text.contains("session resumed"))
);
}
#[test]
fn theme_picker_escape_closes_without_change() {
let (_tx, model) = new_model();
let mut sim = ProgramSimulator::new(model);
sim.init();
let accent_before = sim.model().palette.accent;
type_str(&mut sim, "/theme");
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
sim.inject_event(key(KeyCode::Escape, Modifiers::empty()));
assert!(sim.model().picker.is_none());
assert_eq!(sim.model().palette.accent, accent_before);
}
#[test]
fn bang_routes_bash_command_and_result_renders() {
let (_agent_tx, rx) = mpsc::channel();
let (submit_tx, submit_rx) = mpsc::channel::<UiCommand>();
let model = PiFtuiModel::new(rx).with_submit_channel(submit_tx);
let mut sim = ProgramSimulator::new(model);
sim.init();
type_str(&mut sim, "!echo hi");
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
assert_eq!(
submit_rx.try_recv().expect("routed"),
UiCommand::Bash {
command: "echo hi".into(),
exclude: false,
}
);
type_str(&mut sim, "!!ls");
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
assert_eq!(
submit_rx.try_recv().expect("routed"),
UiCommand::Bash {
command: "ls".into(),
exclude: true,
}
);
type_str(&mut sim, "!");
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
assert!(submit_rx.try_recv().is_err());
assert!(
sim.model()
.transcript
.iter()
.any(|e| e.role == EntryRole::Error && e.text.contains("usage: !")),
"bare-bang usage error missing"
);
sim.send(PiFtuiMsg::Agent(PiMsg::BashResult {
display: "$ echo hi\nhi".into(),
content_for_agent: None,
}));
let rendered = buffer_text(sim.capture_frame(40, 10), 40, 10);
assert!(
rendered.contains("echo hi"),
"bash display missing: {rendered:?}"
);
}
#[test]
fn subscription_id_is_stable() {
let (_tx, rx) = mpsc::channel::<PiMsg>();
let sub = AgentEventSubscription::new(rx);
assert_eq!(sub.id(), AGENT_EVENTS_SUB_ID);
}
#[test]
fn session_slash_commands_route_to_driver() {
let (_agent_tx, rx) = mpsc::channel();
let (submit_tx, submit_rx) = mpsc::channel::<UiCommand>();
let model = PiFtuiModel::new(rx).with_submit_channel(submit_tx);
let mut sim = ProgramSimulator::new(model);
sim.init();
type_str(&mut sim, "/new");
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
assert_eq!(submit_rx.try_recv().expect("routed"), UiCommand::NewSession);
type_str(&mut sim, "/session");
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
assert_eq!(
submit_rx.try_recv().expect("routed"),
UiCommand::SessionInfo
);
type_str(&mut sim, "/tree deep --all");
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
assert_eq!(
submit_rx.try_recv().expect("routed"),
UiCommand::TreeSummary
);
type_str(&mut sim, "/thinking medium");
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
assert_eq!(
submit_rx.try_recv().expect("routed"),
UiCommand::SetThinking(Some(crate::model::ThinkingLevel::Medium))
);
type_str(&mut sim, "/t 3");
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
assert_eq!(
submit_rx.try_recv().expect("routed"),
UiCommand::SetThinking(Some(crate::model::ThinkingLevel::High))
);
type_str(&mut sim, "/think");
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
assert_eq!(
submit_rx.try_recv().expect("routed"),
UiCommand::SetThinking(None)
);
type_str(&mut sim, "/thinking bogus");
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
assert!(submit_rx.try_recv().is_err());
assert!(
sim.model()
.transcript
.iter()
.any(|e| e.role == EntryRole::Error && e.text.contains("Invalid thinking level")),
"invalid-level error missing"
);
type_str(&mut sim, "/name");
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
assert!(submit_rx.try_recv().is_err());
type_str(&mut sim, "/name ship-it");
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
assert_eq!(
submit_rx.try_recv().expect("routed"),
UiCommand::SetName(String::from("ship-it"))
);
}
#[test]
fn slash_input_is_gated_while_working() {
let (_agent_tx, rx) = mpsc::channel();
let (submit_tx, submit_rx) = mpsc::channel::<UiCommand>();
let model = PiFtuiModel::new(rx).with_submit_channel(submit_tx);
let mut sim = ProgramSimulator::new(model);
sim.init();
sim.send(PiFtuiMsg::Agent(PiMsg::AgentStart));
type_str(&mut sim, "/new");
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
type_str(&mut sim, "/tree");
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
assert!(submit_rx.try_recv().is_err());
assert!(sim.model().transcript.is_empty());
}
#[test]
fn clear_resets_transcript_locally() {
let (_tx, model) = new_model();
let mut sim = ProgramSimulator::new(model);
sim.init();
sim.send(PiFtuiMsg::Agent(PiMsg::System(String::from(
"earlier note",
))));
type_str(&mut sim, "/cls");
sim.inject_event(key(KeyCode::Enter, Modifiers::empty()));
let transcript = &sim.model().transcript;
assert!(!transcript.iter().any(|e| e.text.contains("earlier note")));
assert!(transcript.iter().any(|e| e.text == "Conversation cleared"));
}
}