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 KEY_BURST_CHUNK_SIZE: usize = 256;
const WINDOWS_KEY_BURST_GAP_TIMEOUT: Duration = Duration::from_millis(10);
fn key_burst_gap_timeout() -> Option<Duration> {
cfg!(windows).then_some(WINDOWS_KEY_BURST_GAP_TIMEOUT)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum TerminalInputEvent {
Event(crossterm::event::Event),
KeyBurst(Vec<crossterm::event::KeyEvent>),
}
pub(crate) struct TerminalInputBridge {
pub(crate) receiver: Receiver<TerminalInputEvent>,
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 shutdown = Arc::new(AtomicBool::new(false));
let thread_shutdown = Arc::clone(&shutdown);
let handle = thread::spawn(move || {
run_input_reader(reader, sender, thread_shutdown, burst_gap_timeout);
});
Self {
receiver,
shutdown,
handle: Some(handle),
}
}
}
fn is_burst_key(key: crossterm::event::KeyEvent) -> bool {
crate::tui::input::text_char_for_key_burst(key).is_some()
}
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(
keys: Vec<crossterm::event::KeyEvent>,
confirmed_burst: bool,
) -> TerminalInputEvent {
if confirmed_burst {
TerminalInputEvent::KeyBurst(keys)
} else {
TerminalInputEvent::Event(crossterm::event::Event::Key(keys[0]))
}
}
fn run_input_reader<R>(
mut reader: R,
sender: crossbeam_channel::Sender<TerminalInputEvent>,
shutdown: Arc<AtomicBool>,
burst_gap_timeout: Option<Duration>,
) where
R: TerminalEventReader,
{
let mut poll_timeout = INPUT_POLL_IDLE_TIMEOUT;
let mut pending = None;
'input: while !shutdown.load(Ordering::SeqCst) {
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) => {
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;
};
if burst_gap_timeout.is_none() || !is_burst_key(first) {
if !send_event(
&sender,
TerminalInputEvent::Event(crossterm::event::Event::Key(first)),
&shutdown,
) {
break;
}
poll_timeout = INPUT_POLL_DRAIN_TIMEOUT;
continue;
}
let mut keys = vec![first];
let mut confirmed_burst = false;
let mut queue_drained = false;
loop {
while keys.len() < KEY_BURST_CHUNK_SIZE {
match reader.poll(burst_gap_timeout.unwrap_or(INPUT_POLL_DRAIN_TIMEOUT)) {
Ok(true) => {
let next = match reader.read() {
Ok(event) => event,
Err(_) => {
if !send_event(
&sender,
key_run_event(keys, confirmed_burst),
&shutdown,
) {
break 'input;
}
break 'input;
}
};
match next {
crossterm::event::Event::Key(key) if is_burst_key(key) => {
confirmed_burst = true;
keys.push(key);
}
other => {
pending = Some(other);
break;
}
}
}
Ok(false) => {
queue_drained = true;
break;
}
Err(_) => {
if !send_event(&sender, key_run_event(keys, confirmed_burst), &shutdown) {
break 'input;
}
break 'input;
}
}
}
let chunk_full = keys.len() == KEY_BURST_CHUNK_SIZE;
if !send_event(&sender, key_run_event(keys, confirmed_burst), &shutdown) {
break 'input;
}
if pending.is_some() || queue_drained || !chunk_full {
break;
}
match reader.poll(burst_gap_timeout.unwrap_or(INPUT_POLL_DRAIN_TIMEOUT)) {
Ok(true) => match reader.read() {
Ok(crossterm::event::Event::Key(key)) if is_burst_key(key) => {
keys = vec![key];
confirmed_burst = true;
}
Ok(other) => {
pending = Some(other);
break;
}
Err(_) => break 'input,
},
Ok(false) => {
queue_drained = true;
break;
}
Err(_) => break 'input,
}
}
poll_timeout = if queue_drained {
INPUT_POLL_IDLE_TIMEOUT
} else {
INPUT_POLL_DRAIN_TIMEOUT
};
}
}
impl Drop for TerminalInputBridge {
fn drop(&mut self) {
self.shutdown.store(true, Ordering::SeqCst);
if let Some(handle) = self.handle.take() {
let _ = handle.join();
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crossterm::event::{Event, KeyCode, KeyEvent, KeyEventKind, KeyModifiers};
use std::sync::Mutex;
fn key_event(ch: char) -> Event {
Event::Key(KeyEvent::new_with_kind(
KeyCode::Char(ch),
KeyModifiers::NONE,
KeyEventKind::Press,
))
}
fn enter_event() -> Event {
Event::Key(KeyEvent::new(KeyCode::Enter, KeyModifiers::NONE))
}
fn tab_event() -> Event {
Event::Key(KeyEvent::new(KeyCode::Tab, KeyModifiers::NONE))
}
struct ScriptedReader {
steps: VecDeque<io::Result<Option<Event>>>,
observed_timeouts: Arc<Mutex<Vec<Duration>>>,
}
impl ScriptedReader {
fn new(
steps: impl Into<VecDeque<io::Result<Option<Event>>>>,
observed_timeouts: Arc<Mutex<Vec<Duration>>>,
) -> Self {
Self {
steps: steps.into(),
observed_timeouts,
}
}
}
impl TerminalEventReader for ScriptedReader {
fn poll(&mut self, timeout: Duration) -> io::Result<bool> {
self.observed_timeouts.lock().unwrap().push(timeout);
match self.steps.front() {
Some(Ok(Some(_))) => Ok(true),
Some(Ok(None)) => {
self.steps.pop_front();
Ok(false)
}
Some(Err(_)) | None => Err(io::Error::other("reader stopped")),
}
}
fn read(&mut self) -> io::Result<Event> {
match self.steps.pop_front() {
Some(Ok(Some(event))) => Ok(event),
Some(Ok(None)) => Err(io::Error::other("no event ready")),
Some(Err(error)) => Err(error),
None => Err(io::Error::other("reader stopped")),
}
}
}
struct TimedPollReader {
first_poll_started: Arc<AtomicBool>,
poll_count: Arc<AtomicU64>,
}
impl TerminalEventReader for TimedPollReader {
fn poll(&mut self, _timeout: Duration) -> io::Result<bool> {
self.first_poll_started.store(true, Ordering::SeqCst);
self.poll_count.fetch_add(1, Ordering::SeqCst);
thread::sleep(Duration::from_millis(10));
Ok(false)
}
fn read(&mut self) -> io::Result<Event> {
Err(io::Error::other("read should not be called"))
}
}
struct ShutdownOnReadReader {
shutdown: Arc<AtomicBool>,
}
impl TerminalEventReader for ShutdownOnReadReader {
fn poll(&mut self, _timeout: Duration) -> io::Result<bool> {
Ok(true)
}
fn read(&mut self) -> io::Result<Event> {
self.shutdown.store(true, Ordering::SeqCst);
Ok(key_event('q'))
}
}
fn finished_events(bridge: &TerminalInputBridge) -> Vec<TerminalInputEvent> {
while !bridge.handle.as_ref().unwrap().is_finished() {
thread::yield_now();
}
bridge.receiver.try_iter().collect()
}
fn spawn_with_paste_fallback(reader: ScriptedReader) -> TerminalInputBridge {
TerminalInputBridge::spawn_with_reader_and_burst_timeout(
reader,
Some(WINDOWS_KEY_BURST_GAP_TIMEOUT),
)
}
#[test]
fn drop_sets_shutdown_and_joins_reader_thread() {
let first_poll_started = Arc::new(AtomicBool::new(false));
let poll_count = Arc::new(AtomicU64::new(0));
let bridge = TerminalInputBridge::spawn_with_reader(TimedPollReader {
first_poll_started: Arc::clone(&first_poll_started),
poll_count: Arc::clone(&poll_count),
});
while !first_poll_started.load(Ordering::SeqCst) {
thread::sleep(Duration::from_millis(1));
}
let before_drop = poll_count.load(Ordering::SeqCst);
let started = Instant::now();
drop(bridge);
assert!(started.elapsed() < Duration::from_secs(1));
assert!(poll_count.load(Ordering::SeqCst) <= before_drop + 1);
}
#[test]
fn shutdown_after_read_drops_event_before_forwarding() {
let shutdown = Arc::new(AtomicBool::new(false));
let (sender, receiver) = bounded(1);
run_input_reader(
ShutdownOnReadReader {
shutdown: Arc::clone(&shutdown),
},
sender,
shutdown,
None,
);
assert!(receiver.try_recv().is_err());
}
#[test]
fn single_enter_remains_a_submit_key_event() {
let observed = Arc::new(Mutex::new(Vec::new()));
let steps = VecDeque::from([
Ok(Some(enter_event())),
Ok(None),
Err(io::Error::other("done")),
]);
let bridge = spawn_with_paste_fallback(ScriptedReader::new(steps, observed));
assert_eq!(
finished_events(&bridge),
vec![TerminalInputEvent::Event(enter_event())]
);
}
#[test]
fn disabled_fallback_preserves_queued_character_and_enter_events() {
let observed = Arc::new(Mutex::new(Vec::new()));
let steps = VecDeque::from([
Ok(Some(key_event('a'))),
Ok(Some(enter_event())),
Err(io::Error::other("done")),
]);
let bridge = TerminalInputBridge::spawn_with_reader_and_burst_timeout(
ScriptedReader::new(steps, observed),
None,
);
assert_eq!(
finished_events(&bridge),
vec![
TerminalInputEvent::Event(key_event('a')),
TerminalInputEvent::Event(enter_event()),
]
);
}
#[test]
fn idle_gap_preserves_character_then_enter_as_normal_input() {
let observed = Arc::new(Mutex::new(Vec::new()));
let steps = VecDeque::from([
Ok(Some(key_event('a'))),
Ok(None),
Ok(Some(enter_event())),
Ok(None),
Err(io::Error::other("done")),
]);
let bridge = spawn_with_paste_fallback(ScriptedReader::new(steps, observed));
assert_eq!(
finished_events(&bridge),
vec![
TerminalInputEvent::Event(key_event('a')),
TerminalInputEvent::Event(enter_event()),
]
);
}
#[test]
fn text_keys_batch_in_exact_order_with_unicode_tab_and_newline() {
let observed = Arc::new(Mutex::new(Vec::new()));
let mut steps = VecDeque::new();
for event in [
key_event('a'),
key_event('é'),
tab_event(),
enter_event(),
key_event('🙂'),
] {
steps.push_back(Ok(Some(event)));
}
steps.push_back(Err(io::Error::other("done")));
let bridge = spawn_with_paste_fallback(ScriptedReader::new(steps, observed));
let events = finished_events(&bridge);
let TerminalInputEvent::KeyBurst(keys) = &events[0] else {
panic!("expected key burst");
};
assert_eq!(
crate::tui::input::text_for_key_burst(keys).as_deref(),
Some("aé\t\n🙂")
);
assert_eq!(events.len(), 1);
}
#[test]
fn non_text_event_splits_key_bursts_without_reordering() {
let observed = Arc::new(Mutex::new(Vec::new()));
let resize = Event::Resize(100, 30);
let steps = VecDeque::from([
Ok(Some(key_event('a'))),
Ok(Some(key_event('b'))),
Ok(Some(resize.clone())),
Ok(Some(key_event('c'))),
Ok(Some(key_event('d'))),
Err(io::Error::other("done")),
]);
let bridge = spawn_with_paste_fallback(ScriptedReader::new(steps, observed));
let events = finished_events(&bridge);
assert!(
matches!(&events[0], TerminalInputEvent::KeyBurst(keys) if crate::tui::input::text_for_key_burst(keys).as_deref() == Some("ab"))
);
assert_eq!(events[1], TerminalInputEvent::Event(resize));
assert!(
matches!(&events[2], TerminalInputEvent::KeyBurst(keys) if crate::tui::input::text_for_key_burst(keys).as_deref() == Some("cd"))
);
}
#[test]
fn large_burst_chunks_without_loss() {
let observed = Arc::new(Mutex::new(Vec::new()));
let expected = KEY_BURST_CHUNK_SIZE * 4 + 7;
let mut steps = VecDeque::new();
for _ in 0..expected {
steps.push_back(Ok(Some(key_event('x'))));
}
steps.push_back(Err(io::Error::other("done")));
let bridge = spawn_with_paste_fallback(ScriptedReader::new(steps, observed));
let text: String = finished_events(&bridge)
.iter()
.map(|event| match event {
TerminalInputEvent::KeyBurst(keys) => {
crate::tui::input::text_for_key_burst(keys).unwrap()
}
TerminalInputEvent::Event(_) => String::new(),
})
.collect();
assert_eq!(text.len(), expected);
assert!(text.bytes().all(|byte| byte == b'x'));
}
#[test]
fn backpressure_preserves_more_events_than_channel_capacity() {
let observed = Arc::new(Mutex::new(Vec::new()));
let expected = 1030;
let mut steps = VecDeque::new();
for index in 0..expected {
steps.push_back(Ok(Some(Event::Resize(index as u16, 1))));
}
steps.push_back(Err(io::Error::other("done")));
let bridge = TerminalInputBridge::spawn_with_reader(ScriptedReader::new(steps, observed));
let mut received = Vec::new();
while !bridge.handle.as_ref().unwrap().is_finished() || !bridge.receiver.is_empty() {
if let Ok(event) = bridge.receiver.recv_timeout(Duration::from_millis(10)) {
received.push(event);
}
}
assert_eq!(received.len(), expected);
for (index, event) in received.into_iter().enumerate() {
assert_eq!(
event,
TerminalInputEvent::Event(Event::Resize(index as u16, 1))
);
}
}
#[test]
fn adaptive_poll_uses_gap_timeout_then_idle_timeout_to_drain_burst() {
let observed_timeouts = Arc::new(Mutex::new(Vec::new()));
let steps = VecDeque::from([
Ok(None),
Ok(Some(key_event('a'))),
Ok(Some(key_event('b'))),
Ok(None),
Err(io::Error::other("done")),
]);
let bridge =
spawn_with_paste_fallback(ScriptedReader::new(steps, Arc::clone(&observed_timeouts)));
let events = finished_events(&bridge);
assert_eq!(
*observed_timeouts.lock().unwrap(),
vec![
INPUT_POLL_IDLE_TIMEOUT,
INPUT_POLL_IDLE_TIMEOUT,
WINDOWS_KEY_BURST_GAP_TIMEOUT,
WINDOWS_KEY_BURST_GAP_TIMEOUT,
INPUT_POLL_IDLE_TIMEOUT,
]
);
assert!(
matches!(&events[..], [TerminalInputEvent::KeyBurst(keys)] if crate::tui::input::text_for_key_burst(keys).as_deref() == Some("ab"))
);
}
}