magi-code 0.96.1

Repository-aware CLI coding agent for terminal work
Documentation
use super::*;
use std::io;

const INPUT_POLL_IDLE_TIMEOUT: Duration = Duration::from_millis(100);
const INPUT_POLL_DRAIN_TIMEOUT: Duration = Duration::ZERO;
const INPUT_SEND_TIMEOUT: Duration = Duration::from_millis(10);
const WINDOWS_KEY_BURST_GAP_TIMEOUT: Duration = Duration::from_millis(10);

fn key_burst_gap_timeout() -> Option<Duration> {
    // ponytail: crossterm 0.29 lacks reliable Windows paste events. Group only Windows key
    // records arriving within 10ms; remove this heuristic after upstream emits Event::Paste.
    cfg!(windows).then_some(WINDOWS_KEY_BURST_GAP_TIMEOUT)
}

#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum TerminalInputEvent {
    Event(crossterm::event::Event),
    KeyBurst(String),
    KeyBurstTooLarge,
}

enum TerminalInputCommand {
    Flush(Sender<()>),
}

pub(crate) struct TerminalInputBridge {
    pub(crate) receiver: Receiver<TerminalInputEvent>,
    commands: Sender<TerminalInputCommand>,
    shutdown: Arc<AtomicBool>,
    handle: Option<JoinHandle<()>>,
}

trait TerminalEventReader: Send + 'static {
    fn poll(&mut self, timeout: Duration) -> io::Result<bool>;
    fn read(&mut self) -> io::Result<crossterm::event::Event>;
}

struct CrosstermEventReader;

impl TerminalEventReader for CrosstermEventReader {
    fn poll(&mut self, timeout: Duration) -> io::Result<bool> {
        crossterm::event::poll(timeout)
    }

    fn read(&mut self) -> io::Result<crossterm::event::Event> {
        crossterm::event::read()
    }
}

impl TerminalInputBridge {
    pub(crate) fn spawn() -> Self {
        Self::spawn_with_reader(CrosstermEventReader)
    }

    fn spawn_with_reader<R>(reader: R) -> Self
    where
        R: TerminalEventReader,
    {
        Self::spawn_with_reader_and_burst_timeout(reader, key_burst_gap_timeout())
    }

    fn spawn_with_reader_and_burst_timeout<R>(
        reader: R,
        burst_gap_timeout: Option<Duration>,
    ) -> Self
    where
        R: TerminalEventReader,
    {
        let (sender, receiver) = bounded::<TerminalInputEvent>(1024);
        let (command_sender, command_receiver) = bounded::<TerminalInputCommand>(1);
        let (ready_sender, ready_receiver) = bounded(1);
        let shutdown = Arc::new(AtomicBool::new(false));
        let thread_shutdown = Arc::clone(&shutdown);
        let handle = thread::spawn(move || {
            let _ = ready_sender.send(());
            run_input_reader(
                reader,
                sender,
                command_receiver,
                thread_shutdown,
                burst_gap_timeout,
            );
        });
        // Do not paint the first interactive frame until the reader thread has
        // taken ownership of terminal input. This closes the startup window in
        // which keys could arrive before the bridge existed.
        let _ = ready_receiver.recv();
        Self {
            receiver,
            commands: command_sender,
            shutdown,
            handle: Some(handle),
        }
    }

    pub(crate) fn flush(&self) -> Option<Receiver<()>> {
        let (sender, receiver) = bounded(1);
        match self.commands.try_send(TerminalInputCommand::Flush(sender)) {
            Ok(()) => Some(receiver),
            Err(_) => None,
        }
    }
    pub(crate) fn reader_finished(&self) -> bool {
        self.handle
            .as_ref()
            .is_some_and(|handle| handle.is_finished())
    }
}

fn is_burst_key(key: crossterm::event::KeyEvent) -> bool {
    crate::tui::input::text_char_for_key_burst(key).is_some()
}

fn is_burst_start_key(key: crossterm::event::KeyEvent) -> bool {
    is_burst_key(key) && key.code != crossterm::event::KeyCode::Tab
}

fn is_text_key_release(mut key: crossterm::event::KeyEvent) -> bool {
    if key.kind != crossterm::event::KeyEventKind::Release {
        return false;
    }
    key.kind = crossterm::event::KeyEventKind::Press;
    is_burst_key(key)
}

fn send_event(
    sender: &crossbeam_channel::Sender<TerminalInputEvent>,
    mut event: TerminalInputEvent,
    shutdown: &AtomicBool,
) -> bool {
    loop {
        if shutdown.load(Ordering::SeqCst) {
            return false;
        }
        match sender.send_timeout(event, INPUT_SEND_TIMEOUT) {
            Ok(()) => return true,
            Err(crossbeam_channel::SendTimeoutError::Timeout(returned)) => event = returned,
            Err(crossbeam_channel::SendTimeoutError::Disconnected(_)) => return false,
        }
    }
}

fn key_run_event(
    first: crossterm::event::KeyEvent,
    confirmed_burst: bool,
    normalizer: crate::tui::state::PromptTextNormalizer,
) -> TerminalInputEvent {
    if !confirmed_burst {
        return TerminalInputEvent::Event(crossterm::event::Event::Key(first));
    }

    match normalizer.finish() {
        Ok(text) => TerminalInputEvent::KeyBurst(text),
        Err(_) => TerminalInputEvent::KeyBurstTooLarge,
    }
}

fn run_input_reader<R>(
    mut reader: R,
    sender: crossbeam_channel::Sender<TerminalInputEvent>,
    commands: Receiver<TerminalInputCommand>,
    shutdown: Arc<AtomicBool>,
    burst_gap_timeout: Option<Duration>,
) where
    R: TerminalEventReader,
{
    let mut poll_timeout = INPUT_POLL_IDLE_TIMEOUT;
    let mut flush_ack = None;
    let mut pending = None;
    'input: while !shutdown.load(Ordering::SeqCst) {
        if flush_ack.is_none()
            && let Ok(TerminalInputCommand::Flush(ack)) = commands.try_recv()
        {
            flush_ack = Some(ack);
            poll_timeout = INPUT_POLL_DRAIN_TIMEOUT;
        }
        let event = match pending.take() {
            Some(event) => event,
            None => match reader.poll(poll_timeout) {
                Ok(true) => match reader.read() {
                    Ok(event) => event,
                    Err(_) => break,
                },
                Ok(false) => {
                    if let Some(ack) = flush_ack.take() {
                        let _ = ack.send(());
                    }
                    poll_timeout = INPUT_POLL_IDLE_TIMEOUT;
                    continue;
                }
                Err(_) => break,
            },
        };

        if shutdown.load(Ordering::SeqCst) {
            break;
        }
        let crossterm::event::Event::Key(first) = event else {
            if !send_event(&sender, TerminalInputEvent::Event(event), &shutdown) {
                break;
            }
            poll_timeout = INPUT_POLL_DRAIN_TIMEOUT;
            continue;
        };
        // A leading Tab is navigation even when printable text follows immediately; the next
        // key may start its own burst, but the navigation contract wins when Tab is first.
        if burst_gap_timeout.is_none() || !is_burst_start_key(first) {
            if !send_event(
                &sender,
                TerminalInputEvent::Event(crossterm::event::Event::Key(first)),
                &shutdown,
            ) {
                break;
            }
            poll_timeout = INPUT_POLL_DRAIN_TIMEOUT;
            continue;
        }

        let mut normalizer = crate::tui::state::PromptTextNormalizer::new();
        normalizer.push_char(
            crate::tui::input::text_char_for_key_burst(first)
                .expect("is_burst_start_key checked the first key"),
        );
        let mut confirmed_burst = false;
        let mut queue_drained = false;
        let mut reader_failed = false;

        loop {
            // A logical burst is collected as one atomic event. Check shutdown
            // between every reader operation so a continuous stream cannot
            // keep Drop waiting for the burst to end.
            if shutdown.load(Ordering::SeqCst) {
                break 'input;
            }
            match reader.poll(burst_gap_timeout.unwrap_or(INPUT_POLL_DRAIN_TIMEOUT)) {
                Ok(true) => {
                    if shutdown.load(Ordering::SeqCst) {
                        break 'input;
                    }
                    let next = match reader.read() {
                        Ok(event) => event,
                        Err(_) => {
                            reader_failed = true;
                            break;
                        }
                    };
                    match next {
                        crossterm::event::Event::Key(key) if is_burst_key(key) => {
                            confirmed_burst = true;
                            normalizer.push_char(
                                crate::tui::input::text_char_for_key_burst(key)
                                    .expect("is_burst_key checked the next key"),
                            );
                        }
                        // Crossterm 0.29 forwards Windows key-up records. They must neither
                        // split pasted text nor count as another character confirming a burst.
                        // Keep shortcut/modifier releases on the normal event path.
                        crossterm::event::Event::Key(key) if is_text_key_release(key) => {}
                        other => {
                            pending = Some(other);
                            break;
                        }
                    }
                }
                Ok(false) => {
                    queue_drained = true;
                    break;
                }
                Err(_) => {
                    reader_failed = true;
                    break;
                }
            }
        }

        // Do not turn the accumulated prefix into an event after shutdown.
        // The reader-failure path below intentionally still forwards a
        // complete burst when shutdown was not requested.
        if shutdown.load(Ordering::SeqCst) {
            break 'input;
        }

        if !send_event(
            &sender,
            key_run_event(first, confirmed_burst, normalizer),
            &shutdown,
        ) {
            break 'input;
        }
        if reader_failed {
            break 'input;
        }
        poll_timeout = if queue_drained {
            INPUT_POLL_IDLE_TIMEOUT
        } else {
            INPUT_POLL_DRAIN_TIMEOUT
        };
    }
    // A command may have been queued while the reader was blocked in poll/read.
    // Drop those acknowledgements explicitly; keeping the bridge sender alive
    // must not leave a startup fence waiting forever after reader failure.
    while let Ok(command) = commands.try_recv() {
        drop(command);
    }
}

impl Drop for TerminalInputBridge {
    fn drop(&mut self) {
        self.shutdown.store(true, Ordering::SeqCst);
        if let Some(handle) = self.handle.take() {
            let _ = handle.join();
        }
    }
}