mod drain_processing;
mod frame;
mod key_bursts;
use drain_processing::DrainLoopBudgetState;
use drain_processing::PostDrainContext;
use drain_processing::PostDrainFlow;
use drain_processing::PostDrainProcessing;
use frame::draw_frame;
use frame::tui_perf_enabled;
use key_bursts::classify_prompt_key_burst_event;
mod subagent_prompt;
use super::*;
#[derive(Debug)]
struct PendingThemeSave {
request_id: u64,
theme_id: String,
revision: u64,
}
#[derive(Debug)]
struct ThemePickerRuntime {
opening_appearance: crate::appearance::RuntimeAppearance,
opening_revision: u64,
catalog: Option<crate::appearance::ThemeCatalog>,
preview_appearance: crate::appearance::RuntimeAppearance,
preview_revision: u64,
}
#[derive(Debug)]
struct ThemeWorker {
request_id: u64,
kind: ThemeWorkerKind,
handle: JoinHandle<ThemeWorkerResult>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum ThemeWorkerKind {
Catalog,
Save,
}
#[derive(Debug)]
enum ThemeWorkerResult {
Catalog(Result<crate::appearance::ThemeCatalog, String>),
CatalogDelivered,
Saved {
theme_id: String,
revision: u64,
result: Result<(), String>,
},
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct PendingModelSelection {
request_id: u64,
session_id: Option<String>,
generation: u64,
provider: String,
model: String,
}
#[derive(Debug)]
struct ModelSelectionWorkerOutcome {
result: Option<Result<TuiModelSelectionResult, String>>,
}
#[derive(Debug)]
struct ModelSelectionWorker {
request_id: u64,
cancel: Arc<AtomicBool>,
handle: JoinHandle<ModelSelectionWorkerOutcome>,
reconciled: bool,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum RewindRequestKind {
Changes,
Plan,
Execute,
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct PendingRewind {
request_id: u64,
session_id: String,
generation: u64,
kind: RewindRequestKind,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum RewindWorkerDelivery {
Delivered,
Canceled,
Failed,
}
#[derive(Debug)]
struct RewindWorker {
request_id: u64,
kind: RewindRequestKind,
cancel: Arc<AtomicBool>,
handle: JoinHandle<RewindWorkerOutcome>,
}
#[derive(Debug)]
struct RewindWorkerOutcome {
result: Option<Result<TuiRewindWorkerResult, String>>,
delivery: RewindWorkerDelivery,
}
#[derive(Debug)]
struct ExportWorker {
request_id: u64,
cancel: Arc<AtomicBool>,
handle: JoinHandle<ExportWorkerOutcome>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum ExportWorkerDelivery {
Delivered,
Canceled,
Failed,
}
#[derive(Debug)]
struct ExportWorkerOutcome {
result: Option<Result<TuiExportWorkerResult, String>>,
delivery: ExportWorkerDelivery,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct StartupInputBoundary {
remaining: usize,
}
impl StartupInputBoundary {
fn capture(pending_input_events: usize, queued_input_events: usize) -> Self {
Self {
remaining: pending_input_events.saturating_add(queued_input_events),
}
}
fn is_complete(self) -> bool {
self.remaining == 0
}
fn consume(&mut self) {
self.remaining = self.remaining.saturating_sub(1);
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct PendingSessionPreview {
request_id: u64,
session_id: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct PendingSessionSwitch {
request_id: u64,
session_id: String,
generation: u64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct PendingExport {
request_id: u64,
session_id: String,
generation: u64,
}
enum PendingModelCatalogConsumer {
ModelPicker,
}
impl MissionControlApp {
fn poll_summarizer(&mut self, ui_state: &mut state::MissionControlState) -> bool {
self.summarizer.refresh_config(&self.config, &self.settings);
self.summarizer
.set_session(self.state.current_session.clone());
let mut changed = false;
if ui_state.summary.session_id.as_deref() != self.state.active_session_id() {
ui_state.summary = Default::default();
ui_state.summary.session_id = self.state.active_session_id().map(str::to_string);
changed = true;
}
if let Some(enabled) = ui_state.summary.take_requested_enabled() {
if self.state.current_session.is_some() {
self.summarizer.set_enabled(enabled);
} else {
ui_state.summary.enabled = false;
ui_state.summary.error =
Some("Start a session before using summary controls.".into());
}
changed = true;
}
if let Some(snapshot) = self.summarizer.poll()
&& Some(snapshot.session_id.as_str()) == self.state.active_session_id()
{
ui_state.summary.apply(snapshot);
changed = true;
}
changed
}
}
pub(crate) struct MissionControlApp {
side: side::SideRuntime,
release_notes: crate::tui::release_notes::WelcomeReleaseNotes,
release_notes_eligible: bool,
update_restart_launch: bool,
update_checks: crate::updates::UpdateChecks,
pub(crate) update_restart: Option<crate::updates::RestartContext>,
pub(crate) config: EffectiveConfig,
default_config: EffectiveConfig,
pub(crate) settings: crate::config::Settings,
summarizer: crate::summarizer::SummarizerRuntime,
pub(crate) appearance: crate::appearance::RuntimeAppearance,
pub(crate) theme_cli_override: bool,
pub(crate) instructions: Vec<InstructionFile>,
pub(crate) discovered_skills: SkillDiscovery,
pub(crate) subagent_profile_discovery: crate::subagents::profiles::SubagentProfileDiscovery,
pub(crate) skills: SkillDiscovery,
pub(crate) commands: CommandRegistry,
pub(crate) state: ShellState,
pub(crate) events: Sender<TuiEvent>,
pub(crate) active_run: bool,
pub(crate) last_steering_pending_count: usize,
pub(in crate::tui) worker: Option<WorkerState>,
startup_cancel: Arc<AtomicBool>,
startup_critical_event_sent: Arc<AtomicBool>,
pub(crate) model_catalog_loading: bool,
pub(crate) next_model_catalog_request_id: u64,
pub(crate) next_usage_request_id: u64,
codex_quota: codex_quota::CodexQuotaRefresh,
claude_quota: claude_quota::ClaudeQuotaRefresh,
model_catalog_cache: Option<crate::model_catalog::CatalogForUi>,
pending_model_catalog_consumer: Option<PendingModelCatalogConsumer>,
next_model_selection_request_id: u64,
pending_model_selection: Option<PendingModelSelection>,
model_selection_workers: Vec<ModelSelectionWorker>,
pub(crate) steering: crate::agent::steering::AgentSteering,
pub(crate) disabled_tools: std::sync::Arc<std::sync::Mutex<std::collections::HashSet<String>>>,
pub(crate) disabled_subagent_profiles:
std::sync::Arc<std::sync::Mutex<std::collections::HashSet<String>>>,
pub(crate) mcp: Option<std::sync::Arc<std::sync::Mutex<crate::mcp::manager::McpManager>>>,
pub(crate) startup_readiness: StartupReadiness,
pub(crate) startup_failure: Option<String>,
pub(crate) startup_warning: Option<String>,
startup_critical_worker: Option<JoinHandle<()>>,
startup_decorative_worker: Option<JoinHandle<()>>,
pub(crate) next_startup_request_id: u64,
pending_startup_prompt: Option<PendingStartupPrompt>,
pending_initial_prompt: Option<String>,
initial_prompt_rejected: bool,
initial_prompt_flush_ack: Option<Receiver<()>>,
startup_input_boundary: Option<StartupInputBoundary>,
deferred_startup_critical_loaded: Vec<(u64, Result<TuiStartupCritical, String>)>,
session_generation: u64,
session_preferences: Option<crate::sessions::preferences::SessionPreferences>,
preference_save_worker: Option<completion_worker::CompletionWorker<()>>,
next_session_preview_request_id: u64,
active_session_preview: Option<PendingSessionPreview>,
next_session_switch_request_id: u64,
queued_session_preview: Option<String>,
pending_session_switch: Option<PendingSessionSwitch>,
fast_mode_persistence: fast_mode_persistence::FastModePersistenceRuntime,
settings_persistence: settings_persistence::SettingsPersistenceRuntime,
panel_layout_persistence: panel_layout_persistence::PanelLayoutPersistenceRuntime,
session_maintenance: session_maintenance::SessionMaintenanceRuntime,
system_prompt: system_prompt::SystemPromptRuntime,
theme_catalog: Option<crate::appearance::ThemeCatalog>,
theme_catalog_request_id: u64,
theme_picker_runtime: Option<ThemePickerRuntime>,
theme_workers: Vec<ThemeWorker>,
auth_worker: Option<auth::AuthWorker>,
pending_theme_save: Option<PendingThemeSave>,
next_theme_request_id: u64,
theme_revision: u64,
startup_prompt_launch_in_progress: bool,
next_rewind_request_id: u64,
pending_rewind: Option<PendingRewind>,
rewind_workers: Vec<RewindWorker>,
next_export_request_id: u64,
pending_export: Option<PendingExport>,
export_workers: Vec<ExportWorker>,
next_footer_git_branch_request_id: u64,
pending_footer_git_branch_request_id: Option<u64>,
footer_git_branch_worker: Option<completion_worker::CompletionWorker<TuiEvent>>,
pending_critical_events: VecDeque<TuiEvent>,
#[cfg(any(test, debug_assertions))]
tui_perf_enabled: bool,
}
mod actions;
mod auth;
mod catalog;
mod claude_quota;
mod codex_quota;
mod compaction;
mod completion_worker;
mod fast_mode_persistence;
mod panel_layout_persistence;
mod session_maintenance;
mod settings_persistence;
mod system_prompt;
use fast_mode_persistence::FastModePersistenceRuntime;
mod display_reducer;
mod export;
mod modals;
mod preferences;
mod reducer;
mod rewind;
mod run_loop;
mod side;
mod startup;
mod submit;
mod theme;
mod updates;
mod workers;
use startup::PendingStartupPrompt;
impl MissionControlApp {
fn apply_ui_error(ui_state: &mut state::MissionControlState, error: impl std::fmt::Display) {
let message = error.to_string();
let _ = apply_worker_final_event(ui_state, WorkerFinalEvent::Error(message));
}
pub(crate) fn new(config: TuiSessionConfig, events: Sender<TuiEvent>) -> Self {
let config = config;
let appearance = config.appearance;
let theme_cli_override = config.theme_cli_override;
let (initial_prompt, initial_prompt_rejected) = match config.initial_prompt.as_deref() {
Some(prompt) => match state::normalize_prompt_text_bounded(prompt) {
Ok(prompt) => (Some(prompt), false),
Err(_) => (None, true),
},
None => (None, false),
};
let herdr_reporter = config.herdr_reporter;
let side = side::SideRuntime::new(
config.settings.side.clone(),
config.active_session.is_some(),
);
let model = config
.config
.model
.clone()
.unwrap_or_else(|| crate::providers::DEFAULT_CODEX_MODEL.to_string());
let state = ShellState::new(
config.manager,
config.active_session,
config.cwd,
model,
config.config.auth_state(),
)
.with_config(config.config.clone())
.with_herdr_reporter(herdr_reporter);
Self {
release_notes: Default::default(),
release_notes_eligible: crate::tui::release_notes::eligible_for_launch(
config.resumed_session,
config.update_restart,
config.initial_prompt.as_deref(),
),
update_restart_launch: config.update_restart && config.initial_prompt.is_none(),
update_checks: crate::updates::UpdateChecks::default(),
update_restart: None,
summarizer: crate::summarizer::SummarizerRuntime::new(
config.config.clone(),
config.settings.clone(),
),
default_config: config.config.clone(),
config: config.config,
settings: config.settings,
appearance,
theme_cli_override,
discovered_skills: config.discovered_skills,
instructions: config.instructions,
subagent_profile_discovery:
crate::subagents::profiles::SubagentProfileDiscovery::default(),
skills: config.skills,
commands: config.commands,
state,
events,
side,
active_run: false,
last_steering_pending_count: 0,
worker: None,
model_catalog_loading: false,
next_model_catalog_request_id: 0,
next_usage_request_id: 0,
codex_quota: Default::default(),
claude_quota: Default::default(),
model_catalog_cache: None,
pending_model_catalog_consumer: None,
next_model_selection_request_id: 0,
pending_model_selection: None,
model_selection_workers: Vec::new(),
steering: crate::agent::steering::AgentSteering::new(),
disabled_tools: std::sync::Arc::new(std::sync::Mutex::new(
std::collections::HashSet::new(),
)),
disabled_subagent_profiles: std::sync::Arc::new(std::sync::Mutex::new(
std::collections::HashSet::new(),
)),
mcp: None,
startup_readiness: StartupReadiness::Loading,
startup_failure: None,
startup_warning: None,
startup_cancel: Arc::new(AtomicBool::new(false)),
startup_critical_event_sent: Arc::new(AtomicBool::new(false)),
startup_critical_worker: None,
startup_decorative_worker: None,
next_startup_request_id: 0,
pending_startup_prompt: None,
pending_initial_prompt: initial_prompt,
initial_prompt_rejected,
initial_prompt_flush_ack: None,
deferred_startup_critical_loaded: Vec::new(),
startup_input_boundary: None,
session_generation: 0,
session_preferences: None,
preference_save_worker: None,
next_session_preview_request_id: 0,
active_session_preview: None,
queued_session_preview: None,
next_rewind_request_id: 0,
pending_rewind: None,
rewind_workers: Vec::new(),
next_session_switch_request_id: 0,
pending_session_switch: None,
fast_mode_persistence: FastModePersistenceRuntime::new(),
settings_persistence: settings_persistence::SettingsPersistenceRuntime::default(),
panel_layout_persistence:
panel_layout_persistence::PanelLayoutPersistenceRuntime::default(),
session_maintenance: session_maintenance::SessionMaintenanceRuntime::default(),
system_prompt: system_prompt::SystemPromptRuntime::default(),
theme_catalog: None,
theme_catalog_request_id: 0,
theme_picker_runtime: None,
theme_workers: Vec::new(),
auth_worker: None,
pending_theme_save: None,
next_theme_request_id: 0,
theme_revision: 0,
startup_prompt_launch_in_progress: false,
next_export_request_id: 0,
pending_export: None,
export_workers: Vec::new(),
next_footer_git_branch_request_id: 0,
pending_footer_git_branch_request_id: None,
footer_git_branch_worker: None,
pending_critical_events: VecDeque::new(),
#[cfg(any(test, debug_assertions))]
tui_perf_enabled: crate::tui::perf::TuiPerfCounters::from_env().enabled(),
}
}
fn request_footer_git_branch_refresh(&mut self) {
if self.pending_footer_git_branch_request_id.is_some()
|| self.footer_git_branch_worker.is_some()
{
return;
}
self.next_footer_git_branch_request_id =
self.next_footer_git_branch_request_id.saturating_add(1);
let request_id = self.next_footer_git_branch_request_id;
self.pending_footer_git_branch_request_id = Some(request_id);
let cwd = self.state.cwd.clone();
match completion_worker::CompletionWorker::spawn(
"magi-footer-branch",
self.events.clone(),
move || {
Ok(TuiEvent::FooterGitBranchLoaded {
request_id,
branch: crate::tui::sessions::commands::current_git_branch(&cwd),
})
},
) {
Ok(worker) => self.footer_git_branch_worker = Some(worker),
Err(_) => self.pending_footer_git_branch_request_id = None,
}
}
fn handle_footer_git_branch_drain(
&mut self,
ui_state: &mut state::MissionControlState,
drain_result: &DrainResult,
) -> bool {
let mut changed = false;
let mut owned = DrainResult::default();
if let Some(worker) = self.footer_git_branch_worker.as_mut() {
if let Some(result) = worker.take_result() {
match result {
Ok(event) => apply_control_event_to_state(ui_state, event, &mut owned),
Err(_) => self.pending_footer_git_branch_request_id = None,
}
}
if worker.ready_to_reap() {
let _ = self.footer_git_branch_worker.take().unwrap().join();
}
}
for (request_id, branch) in owned
.footer_git_branch_loaded
.iter()
.chain(&drain_result.footer_git_branch_loaded)
{
if self.pending_footer_git_branch_request_id != Some(*request_id) {
continue;
}
self.pending_footer_git_branch_request_id = None;
if ui_state.footer_git_branch != *branch {
ui_state.footer_git_branch = branch.clone();
changed = true;
}
}
changed
}
fn handle_background_session_drain(
&mut self,
ui_state: &mut state::MissionControlState,
drain_result: &DrainResult,
receiver: &Receiver<TuiEvent>,
) -> bool {
let mut changed = self.handle_session_preview_drain(ui_state, drain_result);
let (switch_changed, selected_title) =
self.handle_session_switch_drain_with_title(ui_state, drain_result, receiver);
changed |= switch_changed;
changed |= self.handle_session_title_drain_with_skip(
ui_state,
drain_result,
selected_title.as_ref(),
);
changed |= self.handle_footer_git_branch_drain(ui_state, drain_result);
changed
}
fn handle_session_switch_drain_with_title(
&mut self,
ui_state: &mut state::MissionControlState,
drain_result: &DrainResult,
receiver: &Receiver<TuiEvent>,
) -> (bool, Option<(String, Option<String>)>) {
let mut changed = false;
let mut selected_title = None;
for (request_id, session_id, result) in &drain_result.session_switch_loaded {
let Some(pending) = self.pending_session_switch.clone() else {
continue;
};
if pending.request_id != *request_id || pending.session_id != *session_id {
continue;
}
self.pending_session_switch = None;
if self.pending_export.is_some() {
ui_state.status =
"session export is in progress; session switch canceled".to_string();
changed = true;
continue;
}
match result {
Ok(loaded) => {
if self.active_run
|| self.session_generation != pending.generation
|| self.state.active_session_id() == Some(session_id.as_str())
{
ui_state.status =
"session switch result ignored; state changed".to_string();
changed = true;
continue;
}
preserve_critical_tui_events(receiver, &mut self.pending_critical_events);
self.state.current_session = Some(loaded.session.clone());
self.state.clear_fast_observations();
self.session_generation = self.session_generation.saturating_add(1);
ui_state.reset_for_new_session();
apply_session_hydration_snapshot(ui_state, loaded.snapshot.clone());
self.restore_session_preferences(ui_state, loaded.preferences.clone());
for diagnostic in &loaded.diagnostics {
ui_state.apply_output_event(&OutputEvent::Diagnostic {
level: "warning".to_string(),
message: diagnostic.clone(),
});
}
apply_footer_context(
ui_state,
self.state.current_session.as_ref(),
loaded.latest_title.as_deref(),
&self.state.cwd,
);
ui_state.close_session_picker();
ui_state.status = loaded.status.clone();
if let Some(reporter) = self.state.herdr_reporter.as_ref() {
reporter.report_agent_session_with_source(
loaded.session.id(),
crate::herdr::HerdrSessionStartSource::Select,
);
reporter.report_title(loaded.latest_title.as_deref());
}
selected_title =
Some((loaded.session.id().to_string(), loaded.latest_title.clone()));
changed = true;
}
Err(error) => {
ui_state.status = error.clone();
changed = true;
}
}
}
(changed, selected_title)
}
fn handle_session_title_drain_with_skip(
&mut self,
_ui_state: &mut state::MissionControlState,
drain_result: &DrainResult,
selected_title: Option<&(String, Option<String>)>,
) -> bool {
let Some(reporter) = self.state.herdr_reporter.as_ref() else {
return false;
};
let active_session_id = self.state.active_session_id();
for (session_id, title) in &drain_result.session_title_updated {
if selected_title.is_some_and(|(selected_id, selected)| {
selected_id == session_id && selected.as_deref() == Some(title.as_str())
}) {
continue;
}
if active_session_id == Some(session_id.as_str()) {
reporter.report_title(Some(title));
}
}
false
}
fn input_reader_disconnected_error(&self) -> anyhow::Error {
if self.pending_startup_prompt.is_some() {
anyhow::anyhow!(
"terminal input reader disconnected while waiting for startup prompt input fence"
)
} else {
anyhow::anyhow!("terminal input reader disconnected unexpectedly")
}
}
fn consume_startup_input_boundary_events(&mut self, count: usize) {
if let Some(boundary) = self.startup_input_boundary.as_mut() {
boundary.remaining = boundary.remaining.saturating_sub(count);
}
}
fn consume_startup_input_boundary_event(&mut self) {
if let Some(boundary) = self.startup_input_boundary.as_mut() {
boundary.consume();
}
}
pub(crate) fn sync_steering_feedback(
&mut self,
ui_state: &mut state::MissionControlState,
) -> bool {
let pending = self.steering.pending_count();
let mut changed = ui_state.set_pending_steering_count(pending);
if pending == 0 && self.last_steering_pending_count > 0 {
ui_state.show_steering_injected_feedback(Instant::now());
changed = true;
}
self.last_steering_pending_count = pending;
changed
}
fn current_fast_request_state(
&self,
ui_state: &state::MissionControlState,
) -> crate::fast::FastRequestState {
crate::fast::resolve_fast_capability(
&self.settings,
&ui_state.provider,
&ui_state.model,
&self.config.paths,
self.config.custom_providers.get(&ui_state.provider),
crate::fast::FastWorkload::Primary,
)
}
pub(crate) fn fast_mode_status(&self, ui_state: &state::MissionControlState) -> String {
let capability = self.current_fast_request_state(ui_state);
let display_state = ui_state.fast_display_state(self.settings.fast.enabled, &capability);
let observation = ui_state
.selected_fast_observation()
.map(state::StoredFastObservation::to_fast_observation);
match display_state {
state::FastDisplayState::Off => crate::commands::runtime::fast_mode_status(
false,
&ui_state.provider,
&ui_state.model,
false,
),
state::FastDisplayState::Unavailable => crate::commands::runtime::fast_mode_status(
true,
&ui_state.provider,
&ui_state.model,
false,
),
state::FastDisplayState::WillRequest => crate::commands::runtime::fast_mode_status(
true,
&ui_state.provider,
&ui_state.model,
true,
),
state::FastDisplayState::Confirmed
| state::FastDisplayState::Different(_)
| state::FastDisplayState::Unconfirmed => crate::fast::fast_status_message(
self.settings.fast.enabled,
&ui_state.provider,
&ui_state.model,
&capability,
observation.as_ref(),
),
}
}
pub(crate) fn refresh_fast_mode_state(&self, ui_state: &mut state::MissionControlState) {
let capability = self.current_fast_request_state(ui_state);
ui_state.refresh_fast_mode_state(self.settings.fast.enabled, &capability);
}
}
#[cfg(test)]
pub(super) mod profiling;