datui_lib/app/
terminal_input.rs1use std::pin::Pin;
16use std::sync::Arc;
17use std::sync::atomic::{AtomicBool, Ordering};
18use std::sync::mpsc::Sender;
19use std::task::{Context, Poll, Wake, Waker};
20use std::thread::{JoinHandle, Thread};
21use std::time::Instant;
22
23use crossterm::event::{Event, EventStream};
24use futures_core::Stream;
25
26use crate::AppEvent;
27use crate::app::terminal_color::{self, ReplyScanner, Scanned};
28
29struct Unpark(Thread);
31
32impl Wake for Unpark {
33 fn wake(self: Arc<Self>) {
34 self.0.unpark();
35 }
36
37 fn wake_by_ref(self: &Arc<Self>) {
38 self.0.unpark();
39 }
40}
41
42pub struct TerminalInput {
44 stop: Arc<AtomicBool>,
45 thread: Option<JoinHandle<()>>,
46}
47
48impl TerminalInput {
49 pub fn start(tx: Sender<AppEvent>) -> std::io::Result<Self> {
53 let stop = Arc::new(AtomicBool::new(false));
54 let stopping = Arc::clone(&stop);
55 let thread = std::thread::Builder::new()
56 .name("datui-input".into())
57 .spawn(move || read(tx, &stopping))?;
58 Ok(Self {
59 stop,
60 thread: Some(thread),
61 })
62 }
63
64 pub fn stop(&mut self) {
67 self.stop.store(true, Ordering::SeqCst);
68 if let Some(thread) = self.thread.take() {
69 thread.thread().unpark();
70 let _ = thread.join();
71 }
72 }
73}
74
75impl Drop for TerminalInput {
76 fn drop(&mut self) {
77 self.stop();
78 }
79}
80
81fn read(tx: Sender<AppEvent>, stop: &AtomicBool) {
82 let mut stream = EventStream::new();
83 let waker = Waker::from(Arc::new(Unpark(std::thread::current())));
84 let mut cx = Context::from_waker(&waker);
85 let mut replies = ReplyScanner::default();
87 let mut scanned = Vec::new();
88 let mut held_since: Option<Instant> = None;
89 while !stop.load(Ordering::SeqCst) {
90 match Pin::new(&mut stream).poll_next(&mut cx) {
91 Poll::Ready(Some(Ok(event))) => {
92 replies.feed(event, terminal_color::armed(), &mut scanned);
93 held_since = replies
94 .holding()
95 .then(|| held_since.unwrap_or_else(Instant::now));
96 if pass_on(&tx, &mut scanned).is_err() {
97 break;
98 }
99 }
100 Poll::Ready(Some(Err(e))) => {
101 let _ = tx.send(AppEvent::Crash(format!("Cannot read the terminal: {e}")));
102 break;
103 }
104 Poll::Ready(None) => break,
105 Poll::Pending => match held_since {
107 None => std::thread::park(),
108 Some(since) => {
109 let left = terminal_color::HOLD.saturating_sub(since.elapsed());
110 if left.is_zero() {
111 replies.flush(&mut scanned);
113 held_since = None;
114 if pass_on(&tx, &mut scanned).is_err() {
115 break;
116 }
117 } else {
118 std::thread::park_timeout(left);
119 }
120 }
121 },
122 }
123 }
124 drop(stream);
127}
128
129fn pass_on(tx: &Sender<AppEvent>, scanned: &mut Vec<Scanned>) -> Result<(), ()> {
132 for item in scanned.drain(..) {
133 match item {
134 Scanned::Event(event) => forward(tx, event)?,
135 Scanned::Background(mode) => {
136 terminal_color::disarm();
137 if let Some(mode) = mode {
138 tx.send(AppEvent::TerminalBackground(mode))
139 .map_err(|_| ())?;
140 }
141 }
142 }
143 }
144 Ok(())
145}
146
147fn forward(tx: &Sender<AppEvent>, event: Event) -> Result<(), ()> {
151 let event = match event {
152 Event::Key(key) if !key.is_press() => return Ok(()),
153 Event::Key(_) | Event::Resize(..) => event,
154 Event::Mouse(mouse) if crate::app::pointer::wanted(&mouse) => event,
155 Event::FocusGained => return tx.send(AppEvent::TerminalFocused).map_err(|_| ()),
157 _ => return Ok(()),
158 };
159 tx.send(AppEvent::Terminal(event)).map_err(|_| ())
160}
161
162#[cfg(test)]
163mod tests;