use std::io::{self, Stdout};
use std::path::PathBuf;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};
use std::thread::{self, JoinHandle};
use std::time::{Duration, Instant};
use crossbeam_channel::{bounded, Sender};
use crossterm::event::{self as ctevent, Event as CtEvent, KeyEventKind};
use crossterm::terminal::{
disable_raw_mode, enable_raw_mode, EnterAlternateScreen, LeaveAlternateScreen,
};
use crossterm::ExecutableCommand;
#[cfg(unix)]
use nix::sys::signal::{raise, Signal};
use ratatui::backend::CrosstermBackend;
use ratatui::Terminal;
use crate::rhei_tui::dashboard::{GateTransitionSink, InterveneSink, PlanLoader};
use crate::rhei_tui::event::{EventSink, MessageLevel, RunEvent};
mod derive;
mod input;
mod render;
mod state;
mod text;
mod theme;
mod views;
use input::{handle_key_event, InputAction};
use state::UiState;
const CHANNEL_CAPACITY: usize = 1024;
const JOURNAL_BUFFER: usize = 400;
const SLOT_TRAFFIC_BUFFER: usize = 50;
pub type StopRequested = Arc<dyn Fn() -> bool + Send + Sync>;
pub struct TuiContext {
pub workspace: PathBuf,
pub plan_loader: Option<PlanLoader>,
pub intervene: Option<Arc<dyn InterveneSink>>,
pub gate: Option<Arc<dyn GateTransitionSink>>,
pub stop_requested: StopRequested,
pub attached: bool,
}
impl TuiContext {
pub fn driving(
workspace: PathBuf,
plan_loader: Option<PlanLoader>,
intervene: Option<Arc<dyn InterveneSink>>,
gate: Option<Arc<dyn GateTransitionSink>>,
stop_requested: StopRequested,
) -> Self {
Self { workspace, plan_loader, intervene, gate, stop_requested, attached: false }
}
}
pub struct TuiSink {
tx: Sender<Msg>,
join: Mutex<Option<JoinHandle<()>>>,
screen_restored: Arc<AtomicBool>,
}
fn message_goes_to_stderr(screen_restored: bool, event: &RunEvent) -> bool {
screen_restored
&& matches!(
event,
RunEvent::Message { level: MessageLevel::Warn | MessageLevel::Error, .. }
)
}
enum Msg {
Event(Box<RunEvent>),
Shutdown,
}
impl TuiSink {
pub fn start(parallel: u16, total_tasks: usize, context: TuiContext) -> io::Result<Self> {
enable_raw_mode()?;
let mut stdout = io::stdout();
stdout.execute(EnterAlternateScreen)?;
let screen_restored = Arc::new(AtomicBool::new(false));
let prev_hook = std::panic::take_hook();
let panic_restored = Arc::clone(&screen_restored);
std::panic::set_hook(Box::new(move |info| {
let _ = disable_raw_mode();
let _ = io::stdout().execute(LeaveAlternateScreen);
panic_restored.store(true, Ordering::SeqCst);
prev_hook(info);
}));
let backend = CrosstermBackend::new(stdout);
let terminal = Terminal::new(backend)?;
let (tx, rx) = bounded::<Msg>(CHANNEL_CAPACITY);
let state = UiState::with_context(
context.workspace,
parallel.max(1),
total_tasks,
context.plan_loader,
context.intervene,
context.gate,
context.attached,
);
let loop_restored = Arc::clone(&screen_restored);
let stop_requested = context.stop_requested;
let handle = thread::spawn(move || {
render_loop(terminal, rx, state, &loop_restored, stop_requested.as_ref())
});
Ok(Self { tx, join: Mutex::new(Some(handle)), screen_restored })
}
pub fn screen_restored(&self) -> bool {
self.screen_restored.load(Ordering::SeqCst)
}
pub fn finish(&self) {
let _ = self.tx.send(Msg::Shutdown);
let mut guard = match self.join.lock() {
Ok(g) => g,
Err(p) => p.into_inner(),
};
if let Some(handle) = guard.take() {
let _ = handle.join();
}
}
}
impl Drop for TuiSink {
fn drop(&mut self) {
self.finish();
}
}
impl EventSink for TuiSink {
fn emit(&self, event: RunEvent) {
if message_goes_to_stderr(self.screen_restored.load(Ordering::SeqCst), &event) {
if let RunEvent::Message { text, .. } = event {
eprintln!("{text}");
}
return;
}
if matches!(event, RunEvent::AgentOutput { .. }) {
let _ = self.tx.try_send(Msg::Event(Box::new(event)));
} else {
let _ = self.tx.send(Msg::Event(Box::new(event)));
}
}
}
fn render_loop(
mut terminal: Terminal<CrosstermBackend<Stdout>>,
rx: crossbeam_channel::Receiver<Msg>,
mut state: UiState,
screen_restored: &AtomicBool,
stop_requested: &(dyn Fn() -> bool + Send + Sync),
) {
let tick = Duration::from_millis(250);
let mut last_draw = Instant::now().checked_sub(tick).unwrap_or_else(Instant::now);
loop {
let deadline = Instant::now() + tick;
loop {
let remaining = deadline.saturating_duration_since(Instant::now());
match rx.recv_timeout(remaining) {
Ok(Msg::Event(event)) => state.apply(&event),
Ok(Msg::Shutdown) | Err(crossbeam_channel::RecvTimeoutError::Disconnected) => {
state.finished = true;
if !stop_requested() {
stay_until_quit(&mut terminal, &mut state, screen_restored, stop_requested);
}
break_out(terminal, screen_restored);
return;
}
Err(crossbeam_channel::RecvTimeoutError::Timeout) => break,
}
}
if drain_input(&mut terminal, &mut state, screen_restored) {
return;
}
if last_draw.elapsed() >= tick {
state.refresh_plan();
state.tick_spinner();
draw(&mut terminal, &state);
last_draw = Instant::now();
}
}
}
fn leave_finished_screen(stop_requested: bool, poll: &io::Result<bool>) -> bool {
stop_requested || poll.is_err()
}
fn stay_until_quit(
terminal: &mut Terminal<CrosstermBackend<Stdout>>,
state: &mut UiState,
screen_restored: &AtomicBool,
stop_requested: &(dyn Fn() -> bool + Send + Sync),
) {
let tick = Duration::from_millis(250);
state.refresh_plan();
draw(terminal, state);
loop {
let poll = ctevent::poll(tick);
if leave_finished_screen(stop_requested(), &poll) {
break_out_ref(terminal, screen_restored);
return;
}
if poll.unwrap_or(false) {
match ctevent::read() {
Err(_) => {
break_out_ref(terminal, screen_restored);
return;
}
Ok(CtEvent::Key(key)) if key.kind != KeyEventKind::Release => {
match handle_key_event(state, key.code, key.modifiers) {
InputAction::Quit => return,
InputAction::ForwardSigint => {
break_out_ref(terminal, screen_restored);
let _ = forward_sigint_to_self();
return;
}
InputAction::Continue => {}
}
}
Ok(_) => {}
}
}
state.refresh_plan();
state.tick_spinner();
draw(terminal, state);
}
}
fn drain_input(
terminal: &mut Terminal<CrosstermBackend<Stdout>>,
state: &mut UiState,
screen_restored: &AtomicBool,
) -> bool {
loop {
match ctevent::poll(Duration::from_millis(0)) {
Ok(true) => {}
Ok(false) => return false,
Err(_) => {
break_out_ref(terminal, screen_restored);
return true;
}
}
match ctevent::read() {
Err(_) => {
break_out_ref(terminal, screen_restored);
return true;
}
Ok(CtEvent::Key(key)) if key.kind != KeyEventKind::Release => {
match handle_key_event(state, key.code, key.modifiers) {
InputAction::ForwardSigint => {
draw(terminal, state);
break_out_ref(terminal, screen_restored);
let _ = forward_sigint_to_self();
return true;
}
InputAction::Quit => {
break_out_ref(terminal, screen_restored);
return true;
}
InputAction::Continue => {}
}
}
Ok(CtEvent::Resize(_, _)) => draw(terminal, state),
Ok(_) => {}
}
}
}
fn draw(terminal: &mut Terminal<CrosstermBackend<Stdout>>, state: &UiState) {
let _ = terminal.draw(|f| render::draw(f, state));
}
fn break_out(mut terminal: Terminal<CrosstermBackend<Stdout>>, screen_restored: &AtomicBool) {
break_out_ref(&mut terminal, screen_restored);
}
fn break_out_ref(terminal: &mut Terminal<CrosstermBackend<Stdout>>, screen_restored: &AtomicBool) {
let _ = terminal.show_cursor();
let _ = disable_raw_mode();
let _ = io::stdout().execute(LeaveAlternateScreen);
screen_restored.store(true, Ordering::SeqCst);
}
#[cfg(unix)]
fn forward_sigint_to_self() -> nix::Result<()> {
raise(Signal::SIGINT)
}
#[cfg(not(unix))]
fn forward_sigint_to_self() -> io::Result<()> {
Err(io::Error::new(io::ErrorKind::Unsupported, "SIGINT forwarding is Unix-only"))
}
#[cfg(test)]
mod tests;