use std::collections::VecDeque;
use std::path::PathBuf;
use anyhow::Result;
use crossterm::event::EventStream;
use futures::{FutureExt, StreamExt};
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(
config: Config,
cwd: PathBuf,
model_id: String,
mut opts: InteractiveOptions,
) -> Result<()> {
let startup_now = chrono::Local::now();
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);
}
if state.session.conversation.git_branch.is_none() {
state.session.conversation.git_branch = crate::session::detect_git_branch(&cwd);
}
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();
runner.spawn_config_watcher(cwd.clone(), config.memory.clone());
let mut terminal = Some(TerminalGuard::setup()?);
let mut rstate = RenderCache::new();
let mut events = 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) {
runner.dispatch(cmd);
}
#[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(());
loop {
if let Err(err) = terminal
.as_mut()
.expect("terminal guard is alive while the render loop runs")
.inner_mut()
.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.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,
}),
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.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 _ = r.record_msg(now, &msg);
}
state.now = now;
let (new_state, cmds) = update(state, msg);
state = new_state;
for cmd in cmds {
runner.dispatch(cmd);
}
if state.should_exit {
break;
}
}
if let Some(r) = recorder.as_mut() {
let _ = r.record_trailer(chrono::Local::now(), &state.session);
}
drop(events);
if let Some(mut terminal) = terminal.take() {
terminal.restore_now();
}
runner.shutdown().await;
exit_result
}
fn bootstrap_cmds(config: &Config) -> Vec<Cmd> {
let mut cmds = Vec::new();
if !config.mcp_servers.is_empty() {
cmds.push(Cmd::InitMcpServers(config.mcp_servers.clone()));
}
cmds
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn bootstrap_is_empty_without_mcp_servers() {
let cmds = bootstrap_cmds(&Config::default());
assert!(cmds.is_empty());
}
#[test]
fn bootstrap_skips_mcp_init_when_no_servers_configured() {
let cmds = bootstrap_cmds(&Config::default());
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(),
},
);
let cmds = bootstrap_cmds(&cfg);
assert!(cmds.iter().any(|c| matches!(c, Cmd::InitMcpServers(_))));
}
}