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(),
);
if let Some(capabilities) = tools.web_capabilities()
&& let Some(text) = web_capabilities_notice(&config, capabilities)
{
state
.ui
.pending_msgs
.push_back(Msg::TransientStatus { text });
}
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
}
fn web_capabilities_notice(
config: &Config,
capabilities: &crate::providers::tool::web::WebCapabilities,
) -> Option<String> {
use crate::providers::tool::web::Egress;
if config.safety.network == crate::app::NetworkPolicy::Deny {
return Some(format!(
"Web egress disabled by safety.network = \"deny\" (selected fetch backend: {}; selected search backend: {}).",
capabilities.fetch.backend, capabilities.search.backend
));
}
let all = [
("fetch", &capabilities.fetch),
("search", &capabilities.search),
];
let degraded = all
.into_iter()
.filter(|(_, status)| !status.available)
.collect::<Vec<_>>();
let leaves_machine = all
.iter()
.any(|(_, status)| status.egress == Egress::OffMachine);
if degraded.is_empty() && !leaves_machine {
return None;
}
let headline = |name: &str, status: &crate::providers::tool::web::WebCapabilityStatus| {
if status.available {
format!(
"{name}: {} (available; {})",
status.backend, status.trust_destination
)
} else {
format!("{name}: {} (unavailable)", status.backend)
}
};
let mut notice = format!(
"Web capabilities - {}; {}.",
headline("fetch", &capabilities.fetch),
headline("search", &capabilities.search)
);
for (name, status) in degraded {
let reason = status
.reason
.as_deref()
.map(crate::utils::redact_secrets)
.unwrap_or_else(|| "backend initialization failed".to_string());
let reason = reason.split_whitespace().collect::<Vec<_>>().join(" ");
let reason = crate::utils::truncate_middle_bytes(&reason, 240)
.split_whitespace()
.collect::<Vec<_>>()
.join(" ");
notice.push_str(&format!("\n- {name}: {reason}"));
}
Some(notice)
}
#[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(_))));
}
fn capabilities(
fetch: crate::providers::tool::web::WebCapabilityStatus,
search: crate::providers::tool::web::WebCapabilityStatus,
) -> crate::providers::tool::web::WebCapabilities {
crate::providers::tool::web::WebCapabilities::from_statuses_for_test(fetch, search)
}
fn available(
backend: &'static str,
trust_destination: &'static str,
egress: crate::providers::tool::web::Egress,
) -> crate::providers::tool::web::WebCapabilityStatus {
crate::providers::tool::web::WebCapabilityStatus {
available: true,
backend,
trust_destination,
egress,
reason: None,
}
}
fn unavailable(
backend: &'static str,
trust_destination: &'static str,
egress: crate::providers::tool::web::Egress,
reason: &str,
) -> crate::providers::tool::web::WebCapabilityStatus {
crate::providers::tool::web::WebCapabilityStatus {
available: false,
backend,
trust_destination,
egress,
reason: Some(reason.to_string()),
}
}
fn local_fetch() -> crate::providers::tool::web::WebCapabilityStatus {
available(
"native",
"direct from this machine",
crate::providers::tool::web::Egress::OnMachine,
)
}
fn local_search() -> crate::providers::tool::web::WebCapabilityStatus {
available(
"managed_searxng",
"local managed process",
crate::providers::tool::web::Egress::OnMachine,
)
}
#[test]
fn web_capability_notice_stays_silent_when_everything_resolved_and_local() {
let config = Config::default();
let capabilities = capabilities(local_fetch(), local_search());
assert_eq!(web_capabilities_notice(&config, &capabilities), None);
}
#[test]
fn web_capability_notice_discloses_working_cloud_egress() {
let config = Config::default();
let capabilities = capabilities(
local_fetch(),
available(
"ollama_cloud",
"Ollama Cloud",
crate::providers::tool::web::Egress::OffMachine,
),
);
let notice =
web_capabilities_notice(&config, &capabilities).expect("cloud egress must disclose");
assert!(
notice.contains("search: ollama_cloud (available; Ollama Cloud)"),
"{notice}"
);
assert!(!notice.contains('\n'), "{notice}");
}
#[test]
fn web_capability_notice_discloses_configured_searxng_endpoint() {
let config = Config::default();
let capabilities = capabilities(
local_fetch(),
available(
"searxng",
"configured SearXNG instance",
crate::providers::tool::web::Egress::OffMachine,
),
);
let notice =
web_capabilities_notice(&config, &capabilities).expect("configured endpoint discloses");
assert!(notice.contains("configured SearXNG instance"), "{notice}");
}
#[test]
fn web_capability_notice_gives_every_degraded_capability_its_own_line() {
let config = Config::default();
let capabilities = capabilities(
unavailable(
"native",
"direct from this machine",
crate::providers::tool::web::Egress::OnMachine,
"TLS backend failed to initialize",
),
unavailable(
"managed_searxng",
"local managed process",
crate::providers::tool::web::Egress::OnMachine,
"no sovereign SearXNG bundle is available for this platform",
),
);
let notice = web_capabilities_notice(&config, &capabilities).expect("both degraded");
let lines = notice.lines().collect::<Vec<_>>();
assert_eq!(lines.len(), 3, "{notice}");
assert!(lines[1].starts_with("- fetch: TLS backend"), "{notice}");
assert!(lines[2].starts_with("- search: no sovereign"), "{notice}");
}
#[test]
fn web_capability_notice_discloses_shared_backend_and_trust_routing() {
let config = Config::default();
let capabilities = capabilities(
local_fetch(),
unavailable(
"managed_searxng",
"local managed process",
crate::providers::tool::web::Egress::OnMachine,
"no sovereign SearXNG bundle is available for this platform",
),
);
let notice = web_capabilities_notice(&config, &capabilities).expect("degraded search");
assert!(notice.contains("fetch: native (available"), "{notice}");
assert!(notice.contains("direct from this machine"), "{notice}");
assert!(
notice.contains("search: managed_searxng (unavailable)"),
"{notice}"
);
}
#[test]
fn web_capability_notice_moves_remediation_off_the_headline() {
let config = Config::default();
let capabilities = capabilities(
local_fetch(),
unavailable(
"managed_searxng",
"local managed process",
crate::providers::tool::web::Egress::OnMachine,
"no sovereign SearXNG bundle is available for this platform (windows/x86_64).\n Configure `[web] search_backend = \"ollama\"`.",
),
);
let notice = web_capabilities_notice(&config, &capabilities).expect("degraded search");
let (headline, detail) = notice.split_once('\n').expect("detail line");
assert!(!headline.contains("local managed process"), "{headline}");
assert!(!headline.contains("SearXNG bundle"), "{headline}");
assert_eq!(
detail,
"- search: no sovereign SearXNG bundle is available for this platform (windows/x86_64). Configure `[web] search_backend = \"ollama\"`."
);
}
#[test]
fn web_capability_notice_honors_global_network_denial() {
let mut config = Config::default();
config.safety.network = crate::app::NetworkPolicy::Deny;
let capabilities = capabilities(local_fetch(), local_search());
let notice = web_capabilities_notice(&config, &capabilities).expect("denial always shows");
assert!(notice.contains("Web egress disabled"), "{notice}");
assert!(notice.contains("fetch backend: native"), "{notice}");
assert!(
notice.contains("search backend: managed_searxng"),
"{notice}"
);
}
}