Skip to main content

kimun_notes/app/
events.rs

1//! The **Input source** seam (CONTEXT.md § App shell), merged with the app
2//! channel into the one `next()` the **App loop** awaits.
3//!
4//! Two kinds of event reach the loop. Terminal-originated ones — key, mouse,
5//! paste, resize — come from the input source: crossterm's `EventStream` in
6//! the app, any `Stream<Item = AppEvent>` in tests. App-originated ones —
7//! autosave done, indexing done, a server answer — come from spawned tasks
8//! through the app channel. The loop needs to tell them apart: app messages
9//! are peeked without blocking and coalesced into one frame; a real input
10//! event always forces a fresh await and its own draw. That is why the seam
11//! is the input *stream* and not one merged stream handed to the loop.
12//!
13//! The input source ends only when the terminal is gone (or a script ran
14//! out), and the loop treats its end as quit.
15
16use std::io;
17use std::pin::Pin;
18
19use crossterm::event::{Event as CrosstermEvent, EventStream, KeyEventKind};
20use futures::{Stream, StreamExt};
21use tokio::sync::mpsc;
22use tokio::sync::mpsc::error::TryRecvError;
23
24use crate::components::events::{AppEvent, AppTx, InputEvent};
25
26/// Terminal-originated events, already decoded. Boxed rather than generic so
27/// the loop, the app and `main` carry no type parameter for it (the same line
28/// ADR-0009 draws for the editor backend).
29pub type InputSource = Pin<Box<dyn Stream<Item = AppEvent> + Send>>;
30
31/// Owns the app channel and the input source. Exposes a single `next()`
32/// await point for the loop.
33pub struct EventHandler {
34    tx: AppTx,
35    rx: mpsc::UnboundedReceiver<AppEvent>,
36    input: InputSource,
37}
38
39impl Default for EventHandler {
40    fn default() -> Self {
41        Self::new()
42    }
43}
44
45impl EventHandler {
46    /// The app's handler: crossterm reads the terminal.
47    pub fn new() -> Self {
48        Self::from_input(crossterm_input())
49    }
50
51    /// A handler over any input source — a `futures::stream::iter` of scripted
52    /// events in tests, a replayed recording, anything that yields `AppEvent`.
53    /// The stream is fused here, so an adapter need not be.
54    pub fn from_input(input: impl Stream<Item = AppEvent> + Send + 'static) -> Self {
55        let (tx, rx) = mpsc::unbounded_channel();
56        Self {
57            tx,
58            rx,
59            input: Box::pin(input.fuse()),
60        }
61    }
62
63    /// Returns a cloned sender. Pass this to screens and components as `&AppTx`.
64    pub fn app_sender(&self) -> AppTx {
65        self.tx.clone()
66    }
67
68    /// Non-blocking peek of the app channel only. Input is never polled here:
69    /// the loop coalesces queued app messages between blocking awaits, and a
70    /// real input event must always get its own `next()` and its own draw.
71    ///
72    /// `Disconnected` is structurally unreachable: `self.tx` is owned by this
73    /// handler and live across the `&mut self` borrow, so at least one sender
74    /// always exists while `try_next` runs.
75    pub fn try_next(&mut self) -> Option<AppEvent> {
76        match self.rx.try_recv() {
77            Ok(msg) => Some(msg),
78            Err(TryRecvError::Empty) => None,
79            Err(TryRecvError::Disconnected) => {
80                unreachable!(
81                    "EventHandler::tx is owned by this struct and the `&mut self` borrow \
82                     guarantees it outlives this call; channel cannot be Disconnected here"
83                )
84            }
85        }
86    }
87
88    /// Wait for the next event. App messages first (`biased`), then input.
89    /// An exhausted input source yields `Quit`, every time it is asked.
90    pub async fn next(&mut self) -> AppEvent {
91        tokio::select! {
92            biased;
93            Some(msg) = self.rx.recv() => msg,
94            event = self.input.next() => event.unwrap_or(AppEvent::Quit),
95        }
96    }
97}
98
99/// The crossterm adapter: the terminal's event stream, decoded.
100fn crossterm_input() -> impl Stream<Item = AppEvent> + Send {
101    EventStream::new().filter_map(|event| {
102        tracing::debug!("RAW EVENT: {:?}", event);
103        futures::future::ready(decode(event))
104    })
105}
106
107/// What one crossterm event means to the loop. Key releases are dropped (the
108/// kitty protocol reports them; the app acts on presses), a resize is a
109/// redraw, focus and unknown events are nothing, and a read error is logged
110/// and skipped rather than ending the source.
111pub(crate) fn decode(event: io::Result<CrosstermEvent>) -> Option<AppEvent> {
112    match event {
113        Ok(CrosstermEvent::Key(key)) if key.kind != KeyEventKind::Release => {
114            Some(AppEvent::Input(InputEvent::Key(key)))
115        }
116        Ok(CrosstermEvent::Mouse(mouse)) => Some(AppEvent::Input(InputEvent::Mouse(mouse))),
117        Ok(CrosstermEvent::Paste(text)) => Some(AppEvent::Input(InputEvent::Paste(text))),
118        Ok(CrosstermEvent::Resize(_, _)) => Some(AppEvent::Redraw),
119        Ok(_) => None,
120        Err(e) => {
121            tracing::warn!("terminal input error: {e}");
122            None
123        }
124    }
125}
126
127#[cfg(test)]
128mod tests {
129    use super::*;
130    use futures::stream;
131    use ratatui::crossterm::event::{
132        Event as CrosstermEvent, KeyCode, KeyEvent, KeyEventKind, KeyEventState, KeyModifiers,
133    };
134
135    fn key(code: KeyCode, kind: KeyEventKind) -> CrosstermEvent {
136        CrosstermEvent::Key(KeyEvent {
137            code,
138            modifiers: KeyModifiers::NONE,
139            kind,
140            state: KeyEventState::NONE,
141        })
142    }
143
144    #[test]
145    fn decode_keeps_presses_and_drops_releases() {
146        assert!(matches!(
147            decode(Ok(key(KeyCode::Char('a'), KeyEventKind::Press))),
148            Some(AppEvent::Input(InputEvent::Key(k))) if k.code == KeyCode::Char('a')
149        ));
150        assert!(decode(Ok(key(KeyCode::Char('a'), KeyEventKind::Release))).is_none());
151    }
152
153    #[test]
154    fn decode_turns_a_resize_into_a_redraw_and_skips_errors() {
155        assert!(matches!(
156            decode(Ok(CrosstermEvent::Resize(80, 24))),
157            Some(AppEvent::Redraw)
158        ));
159        assert!(decode(Err(std::io::Error::other("hangup"))).is_none());
160    }
161
162    #[tokio::test]
163    async fn a_scripted_input_source_is_delivered_in_order() {
164        let mut events = EventHandler::from_input(stream::iter([
165            AppEvent::Redraw,
166            AppEvent::Input(InputEvent::Paste("p".into())),
167        ]));
168        assert!(matches!(events.next().await, AppEvent::Redraw));
169        assert!(matches!(
170            events.next().await,
171            AppEvent::Input(InputEvent::Paste(s)) if s == "p"
172        ));
173    }
174
175    #[tokio::test]
176    async fn app_messages_are_drained_before_input() {
177        let mut events = EventHandler::from_input(stream::iter([AppEvent::Redraw]));
178        events.app_sender().send(AppEvent::Quit).unwrap();
179        // Biased: the queued app message wins even though input is ready.
180        assert!(matches!(events.next().await, AppEvent::Quit));
181        assert!(matches!(events.next().await, AppEvent::Redraw));
182    }
183
184    #[tokio::test]
185    async fn an_exhausted_input_source_yields_quit() {
186        let mut events = EventHandler::from_input(stream::empty());
187        assert!(matches!(events.next().await, AppEvent::Quit));
188        // …and keeps yielding it: the loop may ask once more on its way out.
189        assert!(matches!(events.next().await, AppEvent::Quit));
190    }
191
192    #[test]
193    fn try_next_only_peeks_the_app_channel() {
194        let mut events = EventHandler::from_input(stream::iter([AppEvent::Redraw]));
195        assert!(events.try_next().is_none());
196        events.app_sender().send(AppEvent::Redraw).unwrap();
197        assert!(matches!(events.try_next(), Some(AppEvent::Redraw)));
198    }
199}