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> {
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,
);
});
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;
};
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 {
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::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;
}
}
}
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
};
}
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();
}
}
}
#[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 HandshakeReader {
release: Receiver<()>,
released: bool,
}
impl TerminalEventReader for HandshakeReader {
fn poll(&mut self, _timeout: Duration) -> io::Result<bool> {
if !self.released {
self.release
.recv()
.map_err(|_| io::Error::other("reader release channel closed"))?;
self.released = true;
}
Err(io::Error::other("reader stopped"))
}
fn read(&mut self) -> io::Result<Event> {
Err(io::Error::other("read should not be called"))
}
}
#[test]
fn spawn_completes_reader_ready_handshake_before_polling() {
let (release_sender, release_receiver) = bounded(1);
let bridge = TerminalInputBridge::spawn_with_reader(HandshakeReader {
release: release_receiver,
released: false,
});
release_sender.send(()).unwrap();
drop(bridge);
}
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 DelayedSingleEventReader {
event: Option<Event>,
}
impl TerminalEventReader for DelayedSingleEventReader {
fn poll(&mut self, _timeout: Duration) -> io::Result<bool> {
if self.event.is_some() {
thread::sleep(Duration::from_millis(20));
Ok(true)
} else {
Ok(false)
}
}
fn read(&mut self) -> io::Result<Event> {
self.event
.take()
.ok_or_else(|| io::Error::other("event already read"))
}
}
struct BurstBoundaryReader {
events: VecDeque<Event>,
collection_started: Sender<()>,
release_collection: Receiver<()>,
collection_released: bool,
}
impl TerminalEventReader for BurstBoundaryReader {
fn poll(&mut self, _timeout: Duration) -> io::Result<bool> {
if self.events.len() == 1 && !self.collection_released {
self.collection_started
.send(())
.map_err(|_| io::Error::other("collection-start signal disconnected"))?;
self.release_collection
.recv()
.map_err(|_| io::Error::other("collection release disconnected"))?;
self.collection_released = true;
}
Ok(!self.events.is_empty())
}
fn read(&mut self) -> io::Result<Event> {
self.events
.pop_front()
.ok_or_else(|| io::Error::other("burst reader has no event"))
}
}
#[test]
fn flush_ack_follows_delayed_terminal_input() {
let bridge = TerminalInputBridge::spawn_with_reader(DelayedSingleEventReader {
event: Some(key_event('x')),
});
let flush_ack = bridge.flush().expect("flush command should be accepted");
flush_ack.recv_timeout(Duration::from_secs(1)).unwrap();
assert_eq!(
bridge.receiver.try_recv().unwrap(),
TerminalInputEvent::Event(key_event('x'))
);
}
#[test]
fn flush_ack_waits_for_complete_logical_burst() {
let (collection_started_sender, collection_started_receiver) = bounded(1);
let (release_sender, release_receiver) = bounded(1);
let bridge = TerminalInputBridge::spawn_with_reader_and_burst_timeout(
BurstBoundaryReader {
events: VecDeque::from([key_event('a'), key_event('b')]),
collection_started: collection_started_sender,
release_collection: release_receiver,
collection_released: false,
},
Some(WINDOWS_KEY_BURST_GAP_TIMEOUT),
);
collection_started_receiver.recv().unwrap();
let flush_ack = match bridge.flush() {
Some(ack) => ack,
None => {
release_sender.send(()).unwrap();
panic!("flush command should be accepted");
}
};
let ack_was_pending = matches!(
flush_ack.try_recv(),
Err(crossbeam_channel::TryRecvError::Empty)
);
release_sender.send(()).unwrap();
assert_eq!(
bridge
.receiver
.recv_timeout(Duration::from_secs(1))
.unwrap(),
TerminalInputEvent::KeyBurst("ab".to_owned())
);
if ack_was_pending {
flush_ack.recv_timeout(Duration::from_secs(1)).unwrap();
}
assert!(
ack_was_pending,
"flush was acknowledged before the complete burst was forwarded"
);
}
#[test]
fn flush_ack_disconnects_when_reader_stops_before_acknowledging() {
let (release_sender, release_receiver) = bounded(1);
let bridge = TerminalInputBridge::spawn_with_reader(HandshakeReader {
release: release_receiver,
released: false,
});
let flush_ack = bridge.flush().expect("flush command should be accepted");
release_sender.send(()).unwrap();
let deadline = Instant::now() + Duration::from_secs(1);
while !bridge.reader_finished() {
assert!(Instant::now() < deadline, "reader did not stop");
thread::yield_now();
}
let result = flush_ack.try_recv();
assert!(
matches!(result, Err(crossbeam_channel::TryRecvError::Disconnected)),
"unexpected flush result: {result:?}"
);
}
#[test]
fn flush_is_nonblocking_across_repeated_invalidations_and_full_command_queue() {
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,
poll_count,
});
let first = bridge
.flush()
.expect("first flush command should be accepted");
drop(first);
let mut saw_full_command_queue = false;
for _ in 0..10_000 {
if bridge.flush().is_none() {
saw_full_command_queue = true;
break;
}
}
assert!(saw_full_command_queue);
let deadline = Instant::now() + Duration::from_secs(1);
let fresh_ack = loop {
if let Some(ack) = bridge.flush() {
break ack;
}
assert!(
Instant::now() < deadline,
"flush retry did not make progress"
);
thread::yield_now();
};
fresh_ack.recv_timeout(Duration::from_secs(1)).unwrap();
}
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'))
}
}
struct ContinuousBurstShutdownReader {
collection_started: Sender<()>,
release_collection: Receiver<()>,
second_key_read: Sender<()>,
release_second_key: Receiver<()>,
extra_poll: Sender<()>,
release_extra_poll: Receiver<()>,
reads: usize,
collection_released: bool,
extra_poll_released: bool,
}
impl TerminalEventReader for ContinuousBurstShutdownReader {
fn poll(&mut self, _timeout: Duration) -> io::Result<bool> {
if self.reads == 1 && !self.collection_released {
self.collection_started
.send(())
.map_err(|_| io::Error::other("collection-start signal disconnected"))?;
self.release_collection
.recv()
.map_err(|_| io::Error::other("collection release disconnected"))?;
self.collection_released = true;
}
if self.reads >= 2 && !self.extra_poll_released {
self.extra_poll
.send(())
.map_err(|_| io::Error::other("extra-poll signal disconnected"))?;
self.release_extra_poll
.recv()
.map_err(|_| io::Error::other("extra-poll release disconnected"))?;
self.extra_poll_released = true;
return Ok(false);
}
Ok(true)
}
fn read(&mut self) -> io::Result<Event> {
self.reads += 1;
match self.reads {
1 => Ok(key_event('a')),
2 => {
self.second_key_read
.send(())
.map_err(|_| io::Error::other("second-key signal disconnected"))?;
self.release_second_key
.recv()
.map_err(|_| io::Error::other("second-key release disconnected"))?;
Ok(key_event('b'))
}
_ => Ok(key_event('c')),
}
}
}
struct ShutdownBurstThreadGuard {
shutdown: Arc<AtomicBool>,
release_collection: Sender<()>,
release_second_key: Sender<()>,
release_extra_poll: Sender<()>,
handle: Option<JoinHandle<()>>,
}
impl Drop for ShutdownBurstThreadGuard {
fn drop(&mut self) {
self.shutdown.store(true, Ordering::SeqCst);
let _ = self.release_collection.try_send(());
let _ = self.release_second_key.try_send(());
let _ = self.release_extra_poll.try_send(());
if let Some(handle) = self.handle.take() {
let _ = handle.join();
}
}
}
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);
let (_command_sender, command_receiver) = bounded(1);
run_input_reader(
ShutdownOnReadReader {
shutdown: Arc::clone(&shutdown),
},
sender,
command_receiver,
shutdown,
None,
);
assert!(receiver.try_recv().is_err());
}
#[test]
fn shutdown_during_continuous_burst_abandons_partial_event() {
let shutdown = Arc::new(AtomicBool::new(false));
let (collection_started_sender, collection_started_receiver) = bounded(1);
let (release_collection_sender, release_collection_receiver) = bounded(1);
let (second_key_read_sender, second_key_read_receiver) = bounded(1);
let (release_second_key_sender, release_second_key_receiver) = bounded(1);
let (extra_poll_sender, extra_poll_receiver) = bounded(1);
let _extra_poll_observer = extra_poll_sender.clone();
let (release_extra_poll_sender, release_extra_poll_receiver) = bounded(1);
let (sender, receiver) = bounded(1);
let (_command_sender, command_receiver) = bounded(1);
let (finished_sender, finished_receiver) = bounded(1);
let thread_shutdown = Arc::clone(&shutdown);
let handle = thread::spawn(move || {
run_input_reader(
ContinuousBurstShutdownReader {
collection_started: collection_started_sender,
release_collection: release_collection_receiver,
second_key_read: second_key_read_sender,
release_second_key: release_second_key_receiver,
extra_poll: extra_poll_sender,
release_extra_poll: release_extra_poll_receiver,
reads: 0,
collection_released: false,
extra_poll_released: false,
},
sender,
command_receiver,
thread_shutdown,
Some(WINDOWS_KEY_BURST_GAP_TIMEOUT),
);
finished_sender.send(()).unwrap();
});
let guard = ShutdownBurstThreadGuard {
shutdown: Arc::clone(&shutdown),
release_collection: release_collection_sender.clone(),
release_second_key: release_second_key_sender.clone(),
release_extra_poll: release_extra_poll_sender.clone(),
handle: Some(handle),
};
collection_started_receiver.recv().unwrap();
release_collection_sender.send(()).unwrap();
second_key_read_receiver.recv().unwrap();
shutdown.store(true, Ordering::SeqCst);
release_second_key_sender.send(()).unwrap();
let mut stopped_at_shutdown_boundary = false;
crossbeam_channel::select! {
recv(finished_receiver) -> result => {
result.unwrap();
stopped_at_shutdown_boundary = true;
}
recv(extra_poll_receiver) -> result => {
result.unwrap();
}
}
if !stopped_at_shutdown_boundary {
release_extra_poll_sender.send(()).unwrap();
finished_receiver.recv().unwrap();
}
drop(guard);
assert!(
stopped_at_shutdown_boundary,
"reader continued accumulating after shutdown"
);
assert!(receiver.try_recv().is_err());
}
fn key_release(event: Event) -> Event {
let Event::Key(mut key) = event else {
panic!("expected a key event");
};
key.kind = KeyEventKind::Release;
Event::Key(key)
}
#[test]
fn windows_press_release_paste_keeps_multiline_text_in_one_burst() {
let observed = Arc::new(Mutex::new(Vec::new()));
let mut steps = VecDeque::new();
for event in [
key_event('a'),
key_event('é'),
enter_event(),
tab_event(),
key_event('b'),
enter_event(),
] {
steps.push_back(Ok(Some(event.clone())));
steps.push_back(Ok(Some(key_release(event))));
}
let bridge = spawn_with_paste_fallback(ScriptedReader::new(steps, observed));
assert_eq!(
finished_events(&bridge),
vec![TerminalInputEvent::KeyBurst("aé\n\tb\n".to_owned())]
);
}
#[test]
fn text_key_releases_do_not_confirm_a_burst_across_idle_gaps() {
let observed = Arc::new(Mutex::new(Vec::new()));
let mut steps = VecDeque::new();
for event in [key_event('a'), enter_event()] {
steps.push_back(Ok(Some(event.clone())));
steps.push_back(Ok(Some(key_release(event))));
steps.push_back(Ok(None));
}
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 shortcuts_still_split_press_release_paste_bursts() {
let observed = Arc::new(Mutex::new(Vec::new()));
let shortcut = Event::Key(KeyEvent::new(KeyCode::Char('c'), KeyModifiers::CONTROL));
let steps = VecDeque::from([
Ok(Some(key_event('a'))),
Ok(Some(key_release(key_event('a')))),
Ok(Some(shortcut.clone())),
Ok(Some(key_release(shortcut.clone()))),
Ok(Some(key_event('b'))),
Ok(Some(key_release(key_event('b')))),
]);
let bridge = spawn_with_paste_fallback(ScriptedReader::new(steps, observed));
assert_eq!(
finished_events(&bridge),
vec![
TerminalInputEvent::Event(key_event('a')),
TerminalInputEvent::Event(shortcut.clone()),
TerminalInputEvent::Event(key_release(shortcut)),
TerminalInputEvent::Event(key_event('b')),
]
);
}
#[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 isolated_tab_remains_a_navigation_event_with_windows_fallback() {
let observed = Arc::new(Mutex::new(Vec::new()));
let steps = VecDeque::from([Ok(Some(tab_event())), Err(io::Error::other("done"))]);
let bridge = spawn_with_paste_fallback(ScriptedReader::new(steps, observed));
assert_eq!(
finished_events(&bridge),
vec![TerminalInputEvent::Event(tab_event())]
);
}
#[test]
fn repeated_tabs_remain_individual_navigation_events_with_windows_fallback() {
let observed = Arc::new(Mutex::new(Vec::new()));
let steps = VecDeque::from([
Ok(Some(tab_event())),
Ok(Some(tab_event())),
Err(io::Error::other("done")),
]);
let bridge = spawn_with_paste_fallback(ScriptedReader::new(steps, observed));
assert_eq!(
finished_events(&bridge),
vec![
TerminalInputEvent::Event(tab_event()),
TerminalInputEvent::Event(tab_event()),
]
);
}
#[test]
fn tab_inside_a_text_burst_remains_literal_text() {
let observed = Arc::new(Mutex::new(Vec::new()));
let steps = VecDeque::from([
Ok(Some(key_event('a'))),
Ok(Some(tab_event())),
Ok(Some(key_event('b'))),
Err(io::Error::other("done")),
]);
let bridge = spawn_with_paste_fallback(ScriptedReader::new(steps, observed));
assert_eq!(
finished_events(&bridge),
vec![TerminalInputEvent::KeyBurst("a\tb".to_owned())]
);
}
#[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(text) = &events[0] else {
panic!("expected key burst");
};
assert_eq!(text, "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(text) if text == "ab"));
assert_eq!(events[1], TerminalInputEvent::Event(resize));
assert!(matches!(&events[2], TerminalInputEvent::KeyBurst(text) if text == "cd"));
}
#[test]
fn large_burst_emits_one_bounded_event_without_loss() {
let observed = Arc::new(Mutex::new(Vec::new()));
let expected = 256 * 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 events = finished_events(&bridge);
assert!(
matches!(&events[..], [TerminalInputEvent::KeyBurst(text)] if text.len() == expected)
);
assert_eq!(
events
.iter()
.filter(|event| matches!(event, TerminalInputEvent::KeyBurst(_)))
.count(),
1
);
}
#[test]
fn exact_normalized_burst_cap_is_accepted() {
let observed = Arc::new(Mutex::new(Vec::new()));
let mut steps = VecDeque::new();
for _ in 0..crate::tui::state::MAX_PROMPT_BYTES {
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));
assert!(matches!(
finished_events(&bridge).as_slice(),
[TerminalInputEvent::KeyBurst(text)] if text.len() == crate::tui::state::MAX_PROMPT_BYTES
));
}
#[test]
fn normalized_burst_over_cap_emits_only_oversized_marker() {
let observed = Arc::new(Mutex::new(Vec::new()));
let mut steps = VecDeque::new();
for _ in 0..=crate::tui::state::MAX_PROMPT_BYTES {
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));
assert_eq!(
finished_events(&bridge),
vec![TerminalInputEvent::KeyBurstTooLarge]
);
}
#[test]
fn oversized_burst_preserves_following_non_text_order() {
let observed = Arc::new(Mutex::new(Vec::new()));
let mut steps = VecDeque::new();
for _ in 0..(crate::tui::state::MAX_PROMPT_BYTES / 2 + 1) {
steps.push_back(Ok(Some(key_event('é'))));
}
steps.push_back(Ok(Some(Event::Resize(100, 30))));
steps.push_back(Ok(Some(key_event('c'))));
steps.push_back(Ok(Some(key_event('d'))));
steps.push_back(Err(io::Error::other("done")));
let bridge = spawn_with_paste_fallback(ScriptedReader::new(steps, observed));
assert!(matches!(
finished_events(&bridge).as_slice(),
[
TerminalInputEvent::KeyBurstTooLarge,
TerminalInputEvent::Event(Event::Resize(100, 30)),
TerminalInputEvent::KeyBurst(text)
] if text == "cd"
));
}
#[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(text)] if text == "ab"));
}
}