use std::collections::VecDeque;
use std::path::PathBuf;
use anyhow::Result;
use crossterm::event::EventStream;
use futures::{FutureExt, StreamExt};
use ratatui::layout::Rect;
use tokio::time::{Duration, interval};
use crate::app::Config;
use crate::app::event_source::coalesce_key_burst;
use crate::app::lifecycle::RuntimeLifecycle;
use crate::app::recorder::{RECORDING_FORMAT_VERSION, Recorder, SessionHeader};
use crate::app::terminal::TerminalGuard;
use crate::domain::{Cmd, Msg, RuntimeSignal, State, update};
use crate::effect::EffectRunner;
use crate::providers::ToolRegistry;
use crate::render::{RenderCache, render};
use crate::session::ConversationHistory;
#[derive(Default)]
pub struct InteractiveOptions {
pub recorder: Option<Recorder>,
pub seed_conversation: Option<ConversationHistory>,
}
pub async fn run_interactive_with(
mut config: Config,
cwd: PathBuf,
model_id: String,
mut opts: InteractiveOptions,
) -> Result<()> {
let startup_now = chrono::Local::now();
let plugin_assets = crate::app::plugin_assets::load();
let plugin_warnings = crate::app::plugin_assets::apply(&mut config, &plugin_assets);
let mut state = State::new(config.clone(), cwd.clone(), model_id.clone(), startup_now);
let seed = opts.seed_conversation.take();
if let Some(r) = opts.recorder.as_mut() {
r.record_header(&SessionHeader {
format: RECORDING_FORMAT_VERSION,
ts: startup_now,
model_id: model_id.clone(),
cwd: cwd.clone(),
config: config.clone(),
seed_conversation: seed.clone(),
})?;
}
if let Some(history) = seed {
state.seed_conversation(history);
}
crate::app::stamp_session_provenance(&mut state, &cwd);
state.ui.no_color = std::env::var_os("NO_COLOR").is_some_and(|v| !v.is_empty());
state.skills = crate::app::skills::load(&cwd);
state.plugin_commands = plugin_assets.commands;
for warning in plugin_warnings {
state
.ui
.pending_msgs
.push_back(Msg::TransientStatus { text: warning });
}
let providers = std::sync::Arc::new(crate::providers::ProviderFactory::new(config.clone()));
let tools = ToolRegistry::build(
&config,
crate::providers::TuiMode::Interactive,
providers.clone(),
);
let (runner, mut msg_rx) = EffectRunner::pair_from(cwd.clone(), providers, tools);
let mut runner = runner
.with_interactive_approvals()
.with_interactive_questions();
runner.spawn_config_watcher(cwd.clone(), config.memory.clone());
let mut terminal = Some(TerminalGuard::setup()?);
let mut rstate = RenderCache::new();
let mut events = Some(EventStream::new());
let mut lifecycle = RuntimeLifecycle::new();
let mut tick = interval(Duration::from_millis(16));
let mut recorder = opts.recorder;
for cmd in bootstrap_cmds(&config, &state.session.conversation.id) {
runner.dispatch(cmd);
}
if !state.session.conversation.tasks.tasks.is_empty() {
runner.dispatch(crate::domain::Cmd::SyncTaskStore(
state.session.conversation.tasks.clone(),
));
}
#[allow(clippy::large_enum_variant)]
enum Sel {
Msg(Option<Msg>),
Term(Option<Result<crossterm::event::Event, std::io::Error>>),
}
let mut pending_msgs: VecDeque<Msg> = VecDeque::new();
let mut exit_result: Result<()> = Ok(());
let mut seen_redraw_seq = state.ui.full_redraw_seq;
loop {
{
let term = terminal
.as_mut()
.expect("terminal guard is alive while the render loop runs")
.inner_mut();
if state.ui.full_redraw_seq != seen_redraw_seq {
seen_redraw_seq = state.ui.full_redraw_seq;
let repaint = term
.size()
.and_then(|size| term.resize(Rect::new(0, 0, size.width, size.height)));
if let Err(err) = repaint {
exit_result = Err(err.into());
break;
}
}
if let Err(err) = term.draw(|f| render(&state, &mut rstate, f)) {
exit_result = Err(err.into());
break;
}
}
let msg = if let Some(queued) = pending_msgs.pop_front() {
Some(queued)
} else {
let selected = tokio::select! {
m = msg_rx.recv() => Sel::Msg(m),
e = events.as_mut().expect("event stream is alive while the loop runs").next() => Sel::Term(e),
s = lifecycle.next_msg() => Sel::Msg(s),
_ = tick.tick() => Sel::Msg(Some(Msg::Tick)),
};
match selected {
Sel::Msg(m) => m,
Sel::Term(Some(Ok(evt))) => {
if let crossterm::event::Event::Mouse(m) = &evt {
use crossterm::event::{KeyModifiers, MouseButton, MouseEventKind as MEK};
let ctrl = m.modifiers.contains(KeyModifiers::CONTROL);
match m.kind {
MEK::Down(MouseButton::Left) if ctrl => rstate
.chat
.find_image_at_screen_pos(m.row)
.map(|target| Msg::OpenImageAt {
message_index: target.message_index,
image_index: target.image_index,
image_number: target.image_number,
}),
MEK::Down(MouseButton::Left) => {
rstate.chat.begin_selection(m.row, m.column);
None
},
MEK::Drag(MouseButton::Left) => {
rstate.chat.update_selection(m.row, m.column);
None
},
MEK::Up(MouseButton::Left) => {
None
},
MEK::ScrollUp => Some(Msg::MouseScroll {
delta: crate::constants::UI_MOUSE_SCROLL_LINES as i16,
}),
MEK::ScrollDown => Some(Msg::MouseScroll {
delta: -(crate::constants::UI_MOUSE_SCROLL_LINES as i16),
}),
_ => None,
}
} else {
if let crossterm::event::Event::Key(k) = &evt
&& k.kind == crossterm::event::KeyEventKind::Press
&& k.modifiers
.contains(crossterm::event::KeyModifiers::CONTROL)
&& k.modifiers.contains(crossterm::event::KeyModifiers::SHIFT)
&& matches!(k.code, crossterm::event::KeyCode::Char(c) if c.eq_ignore_ascii_case(&'c'))
{
rstate
.chat
.selected_text()
.filter(|t| !t.is_empty())
.map(Msg::CopySelection)
} else {
let (primary, trailing) = coalesce_key_burst(evt, || {
events
.as_mut()
.expect("event stream is alive while the loop runs")
.next()
.now_or_never()
.flatten()
.and_then(|r| r.ok())
});
for queued in trailing {
pending_msgs.push_back(queued);
}
primary
}
}
},
Sel::Term(Some(Err(error))) => {
tracing::warn!(error = %error, "terminal event stream failed");
None
},
Sel::Term(None) => Some(Msg::RuntimeSignal(RuntimeSignal::Hangup)),
}
};
let Some(msg) = msg else { continue };
let now = chrono::Local::now();
if let Some(r) = recorder.as_mut()
&& let Err(err) = r.record_msg(now, &msg)
{
tracing::warn!(error = %err, "recorder: failed to record message; --replay may be non-deterministic");
}
state.now = now;
let (new_state, cmds) = update(state, msg);
state = new_state;
let mut compose_draft: Option<String> = None;
for cmd in cmds {
if let Cmd::ComposeInEditor { text } = cmd {
compose_draft = Some(text);
} else {
runner.dispatch(cmd);
}
}
if let Some(draft) = compose_draft {
match crate::app::editor::compose_in_editor(&mut terminal, &mut events, draft).await {
Ok(msg) => pending_msgs.push_back(msg),
Err(err) => {
exit_result = Err(err);
break;
},
}
}
if state.should_exit {
break;
}
}
if let Some(r) = recorder.as_mut()
&& let Err(err) = r.record_trailer(chrono::Local::now(), &state.session)
{
tracing::warn!(error = %err, "recorder: failed to write replay trailer");
}
drop(events);
if let Some(mut terminal) = terminal.take() {
terminal.restore_now();
}
runner.shutdown().await;
exit_result
}
fn bootstrap_cmds(config: &Config, session_id: &str) -> Vec<Cmd> {
let mut cmds = Vec::new();
if !config.mcp_servers.is_empty() {
cmds.push(Cmd::InitMcpServers(config.mcp_servers.clone()));
}
cmds.push(Cmd::EnsureScratchpad {
session_id: session_id.to_string(),
});
cmds
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn bootstrap_always_ensures_the_session_scratchpad() {
let cmds = bootstrap_cmds(&Config::default(), "sess-1");
assert_eq!(cmds.len(), 1);
assert!(
cmds.iter().any(
|c| matches!(c, Cmd::EnsureScratchpad { session_id } if session_id == "sess-1")
)
);
}
#[test]
fn bootstrap_skips_mcp_init_when_no_servers_configured() {
let cmds = bootstrap_cmds(&Config::default(), "sess-1");
assert!(!cmds.iter().any(|c| matches!(c, Cmd::InitMcpServers(_))));
}
#[test]
fn bootstrap_includes_mcp_init_when_servers_configured() {
let mut cfg = Config::default();
cfg.mcp_servers.insert(
"example".to_string(),
crate::app::McpServerConfig {
command: "echo".to_string(),
args: vec![],
env: std::collections::HashMap::new(),
..Default::default()
},
);
let cmds = bootstrap_cmds(&cfg, "sess-1");
assert!(cmds.iter().any(|c| matches!(c, Cmd::InitMcpServers(_))));
}
}