use super::worker::{Response, spawn_worker};
use super::*;
use crate::cli::Args;
use ratatui::crossterm::cursor::Show;
use ratatui::crossterm::event::{
self, DisableBracketedPaste, DisableMouseCapture, EnableBracketedPaste, Event,
};
use ratatui::crossterm::execute;
use std::sync::mpsc::{Receiver, TryRecvError, channel};
const QUOTA_INTERVAL_SECS: u64 = 300;
const QUOTA_TICK: Duration = Duration::from_secs(10);
const FULL_WALK_INTERVAL: Duration = Duration::from_secs(60);
const PENDING_WALK_INTERVAL: Duration = Duration::from_secs(3);
pub fn run(args: &Args, hosted: Option<crate::shim::Hosted>) -> anyhow::Result<i32> {
theme::init_from_env();
let (req_tx, req_rx) = channel::<Request>();
let (res_tx, res_rx) = channel::<Response>();
let worker = spawn_worker(args.plan, req_rx, res_tx.clone());
{
let tx = res_tx.clone();
std::thread::spawn(move || {
crate::pricing::refresh_pricing_blocking();
let _ = tx.send(Response::PricingReady);
});
}
spawn_quota_poller(res_tx.clone());
let hosts = crate::fleet::Host::collect(&args.hosts);
for host in &hosts {
spawn_host_poller(host.clone(), res_tx.clone());
}
std::thread::spawn(move || {
if let Some(version) = crate::update::available_update() {
let _ = res_tx.send(Response::UpdateAvailable(version));
}
});
let mut app = App::new(args.plan, req_tx.clone());
app.refresh_secs = args.delay;
if hosts.is_empty() {
app.hidden_columns.push(ColumnId::Host);
}
if crate::config::OTHER_HOMES.is_empty() {
app.hidden_columns.push(ColumnId::User);
}
if crate::config::profile_count() <= 1 {
app.hidden_columns.push(ColumnId::Profile);
}
if !app.hook_pids.is_empty() {
let _ = req_tx.send(Request::HookClaims(app.hook_pids.clone()));
}
let _ = req_tx.send(Request::Refresh);
app.restore_running_tabs();
let mut hosted = hosted;
if let Some(hosted) = hosted.as_ref() {
app.hosted = Some((hosted.pid, hosted.label.clone()));
app.attach_hosted();
}
let mut terminal = ratatui::init();
if theme::variant() == theme::Variant::Light {
let _ = execute!(
std::io::stdout(),
ratatui::crossterm::style::SetColors(ratatui::crossterm::style::Colors::new(
ratatui::crossterm::style::Color::Black,
ratatui::crossterm::style::Color::White,
)),
);
}
let previous_hook = std::panic::take_hook();
std::panic::set_hook(Box::new(move |info| {
restore_terminal();
previous_hook(info);
}));
let _ = execute!(std::io::stdout(), crossterm::style::Print(MOUSE_ON));
let _ = execute!(std::io::stdout(), EnableBracketedPaste);
if crossterm::terminal::supports_keyboard_enhancement().unwrap_or(false) {
let _ = execute!(
std::io::stdout(),
event::PushKeyboardEnhancementFlags(
event::KeyboardEnhancementFlags::DISAMBIGUATE_ESCAPE_CODES
)
);
}
let watch = crate::watch::Watch::start();
app.listener = crate::hook::Listener::start();
for fixed in crate::hook::repair(app.hook_project().as_deref()) {
app.set_status(&fixed);
}
if app
.hook_status()
.entries
.iter()
.any(|s| s.health.is_problem())
{
app.set_status("Agent hooks need attention — press h");
}
let result = event_loop(
&mut app,
&mut terminal,
&res_rx,
&req_tx,
watch.as_ref(),
hosted.as_mut(),
);
let had_tabs = !app.open_rmux().is_empty();
app.tabs.clear();
drop(hosted);
let left_running = match had_tabs {
true => crate::rmux::sessions(),
false => Vec::new(),
};
restore_terminal();
if !left_running.is_empty() {
println!(
"{} agent{} still running in rmux; `cctop` then `t` to get back to {}.",
left_running.len(),
if left_running.len() == 1 { "" } else { "s" },
if left_running.len() == 1 {
"it"
} else {
"them"
},
);
}
let _ = req_tx.send(Request::Shutdown);
let _ = worker.join();
app.save_prefs();
result
}
const MOUSE_ON: &str = "\x1b[?1000h\x1b[?1002h\x1b[?1006h";
fn restore_terminal() {
let _ = execute!(std::io::stdout(), event::PopKeyboardEnhancementFlags);
let _ = execute!(std::io::stdout(), DisableBracketedPaste);
let _ = execute!(std::io::stdout(), DisableMouseCapture);
let _ = execute!(std::io::stdout(), Show);
let _ = execute!(std::io::stdout(), ratatui::crossterm::style::ResetColor);
ratatui::restore();
let _ = execute!(std::io::stdout(), Show);
}
fn spawn_host_poller(host: crate::fleet::Host, tx: Sender<Response>) {
std::thread::spawn(move || {
loop {
let snapshot = host.poll();
if tx
.send(Response::Remote {
host: host.target.clone(),
snapshot,
})
.is_err()
{
return;
}
std::thread::sleep(crate::fleet::POLL);
}
});
}
fn record_burn(log: &mut crate::burn::Log, quota: &Quota) -> bool {
let at = crate::util::now_ms() / 1000;
let mut stored = false;
for (provider, profiles) in [("claude", "a.claude), ("codex", "a.codex)] {
for profile in profiles {
let crate::quota::ProviderStatus::Ok(q) = &profile.status else {
continue;
};
for window in &q.windows {
stored |= log.record(
crate::burn::key(provider, &profile.profile, window.label),
crate::burn::Sample {
at,
pct: window.pct,
resets_at: window.resets_at,
plan: q.plan.clone(),
},
);
}
}
}
stored
}
fn spawn_quota_poller(tx: Sender<Response>) {
std::thread::spawn(move || {
let mut quota = Quota::default();
let (mut claude_due, mut codex_due) = (Instant::now(), Instant::now());
let mut burn = crate::burn::Log::load();
loop {
let now = Instant::now();
let mut changed = false;
if now >= claude_due {
quota.claude = crate::config::accounts_for(Provider::Claude)
.iter()
.map(|profile| crate::quota::ProfileQuota {
profile: profile.name.clone(),
status: crate::quota::fetch_claude(profile),
source: profile.source,
})
.collect();
let delay = quota
.claude
.iter()
.map(|q| q.status.retry_delay_secs(QUOTA_INTERVAL_SECS))
.max()
.unwrap_or(QUOTA_INTERVAL_SECS);
claude_due = now + Duration::from_secs(delay);
changed = true;
}
if now >= codex_due {
quota.codex = crate::config::accounts_for(Provider::Codex)
.iter()
.map(|profile| crate::quota::ProfileQuota {
profile: profile.name.clone(),
status: crate::quota::fetch_codex(profile),
source: profile.source,
})
.collect();
let delay = quota
.codex
.iter()
.map(|q| q.status.retry_delay_secs(QUOTA_INTERVAL_SECS))
.max()
.unwrap_or(QUOTA_INTERVAL_SECS);
codex_due = now + Duration::from_secs(delay);
changed = true;
}
if changed {
quota.fetched = true;
if record_burn(&mut burn, "a) {
burn.save();
}
if tx.send(Response::Quota(Box::new(quota.clone()))).is_err() {
break;
}
}
std::thread::sleep(QUOTA_TICK);
}
});
}
fn event_loop(
app: &mut App,
terminal: &mut ratatui::DefaultTerminal,
res_rx: &Receiver<Response>,
req_tx: &Sender<Request>,
watch: Option<&crate::watch::Watch>,
mut hosted: Option<&mut crate::shim::Hosted>,
) -> anyhow::Result<i32> {
let mut last_refresh = Instant::now();
let mut last_full_walk = Instant::now();
let mut layout = render::Layout::default();
let mut refresh_in_flight = true;
let mut last_blink = true;
loop {
let mut annotated_rows_changed = false;
let mut rows_changed = false;
loop {
match res_rx.try_recv() {
Ok(Response::Discovered(sessions)) => {
app.sessions = sessions;
app.loaded = true;
app.stats = crate::loader::compute_stats(&app.sessions);
app.refilter();
app.merge_remotes();
rows_changed = true;
}
Ok(Response::Annotated(session)) => {
let found = app.sessions.iter_mut().find(|s| {
s.remote.is_none()
&& s.provider == session.provider
&& s.session_id == session.session_id
});
if let Some(existing) = found {
*existing = *session;
} else {
app.sessions.push(*session);
}
annotated_rows_changed = true;
}
Ok(Response::Sessions(payload)) => {
let (sessions, stats) = *payload;
app.sessions = sessions;
app.loaded = true;
app.stats = stats;
app.merge_remotes();
app.push_history();
app.refilter();
refresh_in_flight = false;
annotated_rows_changed = false;
rows_changed = true;
}
Ok(Response::LiveRows(payload)) => {
let (rows, stats) = *payload;
for row in rows {
let found = app.sessions.iter_mut().find(|s| {
s.remote.is_none()
&& s.provider == row.provider
&& s.session_id == row.session_id
});
match found {
Some(existing) => *existing = row,
None => app.sessions.push(row),
}
}
app.adopt_stats(stats);
app.loaded = true;
app.push_history();
app.refilter();
refresh_in_flight = false;
rows_changed = true;
}
Ok(Response::Data(key, data)) => {
if key == app.panel_key {
app.panel_data = Some(*data);
app.needs_redraw = true;
}
}
Ok(Response::Quota(q)) => {
app.quota = *q;
app.burn = crate::burn::Log::load();
app.needs_redraw = true;
}
Ok(Response::UpdateAvailable(version)) => {
app.update_available = Some(version);
app.needs_redraw = true;
}
Ok(Response::PricingReady) => {
let _ = req_tx.send(Request::Refresh);
refresh_in_flight = true;
}
Ok(Response::Terminated {
session_key,
result,
}) => match result {
Ok(()) => {
app.set_status("Termination signal sent");
let _ = req_tx.send(Request::Refresh);
refresh_in_flight = true;
}
Err(error) => {
app.set_status(format!("Could not stop {session_key}: {error}"));
}
},
Ok(Response::Deleted {
session_key,
result,
}) => {
app.deleting.remove(&session_key);
match result {
Ok(()) => {
app.sessions.retain(|session| session.key() != session_key);
app.marked.remove(&session_key);
app.stats = crate::loader::compute_stats(&app.sessions);
app.refilter();
app.set_status("Deleted session");
}
Err(error) => {
app.set_status(format!("Could not delete {session_key}: {error}"))
}
}
}
Ok(Response::KeysSent { result }) => match result {
Ok(()) => app.set_status("Sent to the session's terminal"),
Err(error) => app.set_status(error),
},
Ok(Response::Remote { host, snapshot }) => {
match snapshot {
crate::fleet::Snapshot::Rows(rows) => {
app.remote_errors.remove(&host);
app.remotes.insert(host, rows);
}
crate::fleet::Snapshot::Failed(why) => {
app.remote_errors.insert(host, why);
}
}
app.merge_remotes();
rows_changed = true;
}
Ok(Response::Scanned { query, hits }) => app.scanned(query, hits),
Ok(Response::Insight(text)) => {
app.insight = Some(text);
app.insight_scroll = 0;
}
Err(TryRecvError::Empty) | Err(TryRecvError::Disconnected) => break,
}
}
rows_changed |= app.promote_matured_prompts();
if annotated_rows_changed || rows_changed {
app.apply_finished_agents();
app.apply_reports();
app.collisions = crate::collide::apply(&mut app.sessions);
}
if annotated_rows_changed {
app.stats = crate::loader::compute_stats(&app.sessions);
app.refilter();
}
if rows_changed {
app.check_bells();
}
if rows_changed || annotated_rows_changed {
app.feed_serving();
}
app.sync_panel_data();
app.tick_scan();
let mut drawn = false;
for tab in &mut app.tabs {
drawn |= tab.pump();
}
if let Some(note) = app.focused_pane().and_then(tabs::Pane::answer_bell) {
app.set_status(note);
}
let closed = app.tabs.iter_mut().fold(false, |any, tab| tab.reap() | any);
if closed {
app.drop_empty_tabs();
}
app.poll_rmux_install();
app.sync_shared_tabs();
if drawn || closed {
app.needs_redraw = true;
}
if let Some(events) = app.listener.as_ref().map(crate::hook::Listener::drain) {
let (changed, lifecycle) = app.apply_hooks(events);
app.needs_redraw |= changed;
if lifecycle && !refresh_in_flight {
let _ = req_tx.send(Request::Refresh);
refresh_in_flight = true;
last_refresh = Instant::now();
}
}
app.tick_handoff();
app.needs_redraw |= app.tick_share();
let phase = app.blink_on();
if phase != last_blink && app.any_attention() {
app.needs_redraw = true;
}
last_blink = phase;
if let Some((_, at)) = &app.status
&& at.elapsed() > Duration::from_secs(3)
{
app.status = None;
app.needs_redraw = true;
}
if app.needs_redraw {
terminal.draw(|frame| layout = render::draw(frame, app))?;
app.needs_redraw = false;
}
let refresh_every = Duration::from_secs_f64(app.refresh_secs);
let idle_wait = match app.tab {
0 => Duration::from_millis(200),
_ => Duration::from_millis(16),
};
let idle_wait = match app.share_opening.is_some() {
true => idle_wait.min(Duration::from_millis(100)),
false => idle_wait,
};
let wait = refresh_every
.checked_sub(last_refresh.elapsed())
.unwrap_or(Duration::ZERO)
.min(idle_wait);
if event::poll(wait)? {
match event::read()? {
Event::Key(key) => app.on_key(key),
Event::Paste(text) => app.on_paste(&text),
Event::Mouse(m) => app.on_mouse(m, &layout),
Event::Resize(_, _) => app.needs_redraw = true,
_ => {}
}
}
if let Some(hosted) = hosted.as_mut()
&& let Some(code) = hosted.finished()
{
return Ok(code);
}
if app.should_quit {
break;
}
let refresh_every = Duration::from_secs_f64(app.refresh_secs);
if last_refresh.elapsed() >= refresh_every && !refresh_in_flight {
last_refresh = Instant::now();
refresh_in_flight = true;
let watched_change = watch.is_some_and(crate::watch::Watch::took_structural_change);
let awaiting = !watched_change
&& last_full_walk.elapsed() >= PENDING_WALK_INTERVAL
&& watch.is_some_and(|w| {
w.awaiting_discovery(|path| {
app.sessions
.iter()
.any(|s| s.data_file.as_deref() == Some(path))
})
});
let full_due =
watched_change || awaiting || last_full_walk.elapsed() >= FULL_WALK_INTERVAL;
if full_due {
last_full_walk = Instant::now();
}
let _ = req_tx.send(if full_due {
Request::Refresh
} else {
Request::RefreshLive
});
}
}
Ok(0)
}