use std::pin::Pin;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::mpsc::Sender;
use std::task::{Context, Poll, Wake, Waker};
use std::thread::{JoinHandle, Thread};
use std::time::Instant;
use crossterm::event::{Event, EventStream};
use futures_core::Stream;
use crate::AppEvent;
use crate::app::terminal_color::{self, ReplyScanner, Scanned};
struct Unpark(Thread);
impl Wake for Unpark {
fn wake(self: Arc<Self>) {
self.0.unpark();
}
fn wake_by_ref(self: &Arc<Self>) {
self.0.unpark();
}
}
pub struct TerminalInput {
stop: Arc<AtomicBool>,
thread: Option<JoinHandle<()>>,
}
impl TerminalInput {
pub fn start(tx: Sender<AppEvent>) -> std::io::Result<Self> {
let stop = Arc::new(AtomicBool::new(false));
let stopping = Arc::clone(&stop);
let thread = std::thread::Builder::new()
.name("datui-input".into())
.spawn(move || read(tx, &stopping))?;
Ok(Self {
stop,
thread: Some(thread),
})
}
pub fn stop(&mut self) {
self.stop.store(true, Ordering::SeqCst);
if let Some(thread) = self.thread.take() {
thread.thread().unpark();
let _ = thread.join();
}
}
}
impl Drop for TerminalInput {
fn drop(&mut self) {
self.stop();
}
}
fn read(tx: Sender<AppEvent>, stop: &AtomicBool) {
let mut stream = EventStream::new();
let waker = Waker::from(Arc::new(Unpark(std::thread::current())));
let mut cx = Context::from_waker(&waker);
let mut replies = ReplyScanner::default();
let mut scanned = Vec::new();
let mut held_since: Option<Instant> = None;
while !stop.load(Ordering::SeqCst) {
match Pin::new(&mut stream).poll_next(&mut cx) {
Poll::Ready(Some(Ok(event))) => {
replies.feed(event, terminal_color::armed(), &mut scanned);
held_since = replies
.holding()
.then(|| held_since.unwrap_or_else(Instant::now));
if pass_on(&tx, &mut scanned).is_err() {
break;
}
}
Poll::Ready(Some(Err(e))) => {
let _ = tx.send(AppEvent::Crash(format!("Cannot read the terminal: {e}")));
break;
}
Poll::Ready(None) => break,
Poll::Pending => match held_since {
None => std::thread::park(),
Some(since) => {
let left = terminal_color::HOLD.saturating_sub(since.elapsed());
if left.is_zero() {
replies.flush(&mut scanned);
held_since = None;
if pass_on(&tx, &mut scanned).is_err() {
break;
}
} else {
std::thread::park_timeout(left);
}
}
},
}
}
drop(stream);
}
fn pass_on(tx: &Sender<AppEvent>, scanned: &mut Vec<Scanned>) -> Result<(), ()> {
for item in scanned.drain(..) {
match item {
Scanned::Event(event) => forward(tx, event)?,
Scanned::Background(mode) => {
terminal_color::disarm();
if let Some(mode) = mode {
tx.send(AppEvent::TerminalBackground(mode))
.map_err(|_| ())?;
}
}
}
}
Ok(())
}
fn forward(tx: &Sender<AppEvent>, event: Event) -> Result<(), ()> {
let event = match event {
Event::Key(key) if !key.is_press() => return Ok(()),
Event::Key(_) | Event::Resize(..) => event,
Event::Mouse(mouse) if crate::app::pointer::wanted(&mouse) => event,
Event::FocusGained => return tx.send(AppEvent::TerminalFocused).map_err(|_| ()),
_ => return Ok(()),
};
tx.send(AppEvent::Terminal(event)).map_err(|_| ())
}
#[cfg(test)]
mod tests;