use crate::term::theme_notify::{OscColorKind, parse_color_scheme_report, parse_osc_color_reply};
use crate::theme::Appearance;
use crate::ui::key::Key;
#[derive(Clone, PartialEq, Debug)]
pub enum TapEvent {
Key(Key),
ThemeNotification(Appearance),
OscColor(OscColorKind, xterm_color::Color),
}
const MAX_ESCAPE_LEN: usize = 128;
pub const ESC_HOLD: std::time::Duration = std::time::Duration::from_millis(50);
pub struct TapScanner {
buf: Vec<u8>,
silent: std::time::Duration,
}
impl TapScanner {
pub fn new() -> TapScanner {
TapScanner {
buf: Vec::new(),
silent: std::time::Duration::ZERO,
}
}
pub fn feed(&mut self, chunk: &[u8]) -> Vec<TapEvent> {
self.silent = std::time::Duration::ZERO;
let mut events = Vec::new();
for &byte in chunk {
if self.buf.is_empty() {
if byte == 0x1b {
self.buf.push(byte);
} else if let Some(key) = decode_key(byte) {
events.push(TapEvent::Key(key));
}
continue;
}
self.buf.push(byte);
if self.buf.len() == 2 {
if !matches!(self.buf[1], b'[' | b']' | b'O') {
self.buf.clear();
if let Some(key) = decode_key(byte) {
events.push(TapEvent::Key(key));
}
}
continue;
}
if self.buf.len() > MAX_ESCAPE_LEN {
self.buf.clear();
continue;
}
if self.buf[1] == b'[' {
if (0x40..=0x7e).contains(&byte) {
if let Some(appearance) = parse_color_scheme_report(&self.buf) {
events.push(TapEvent::ThemeNotification(appearance));
} else if let Some(key) = decode_csi(&self.buf) {
events.push(TapEvent::Key(key));
}
self.buf.clear();
}
} else if self.buf[1] == b'O' {
if let Some(key) = decode_ss3(byte) {
events.push(TapEvent::Key(key));
}
self.buf.clear();
} else {
debug_assert_eq!(self.buf[1], b']');
if self.buf.ends_with(b"\x07") || self.buf.ends_with(b"\x1b\\") {
if let Some((kind, color)) = parse_osc_color_reply(&self.buf) {
events.push(TapEvent::OscColor(kind, color));
}
self.buf.clear();
}
}
}
events
}
pub fn idle(&mut self, silence: std::time::Duration) -> Vec<TapEvent> {
if self.buf != [0x1b] {
return Vec::new();
}
self.silent += silence;
if self.silent < ESC_HOLD {
return Vec::new();
}
self.buf.clear();
self.silent = std::time::Duration::ZERO;
vec![TapEvent::Key(Key::Esc)]
}
}
impl Default for TapScanner {
fn default() -> Self {
TapScanner::new()
}
}
pub fn decode_key(byte: u8) -> Option<Key> {
match byte {
0x03 => Some(Key::CtrlC),
b'\r' | b'\n' => Some(Key::Enter),
0x20..=0x7e => Some(Key::Char(byte as char)),
_ => None,
}
}
pub fn decode_csi(seq: &[u8]) -> Option<Key> {
match seq {
b"\x1b[A" => Some(Key::Up),
b"\x1b[B" => Some(Key::Down),
b"\x1b[C" => Some(Key::Right),
b"\x1b[D" => Some(Key::Left),
b"\x1b[H" | b"\x1b[1~" | b"\x1b[7~" => Some(Key::Home),
b"\x1b[F" | b"\x1b[4~" | b"\x1b[8~" => Some(Key::End),
b"\x1b[5~" => Some(Key::PageUp),
b"\x1b[6~" => Some(Key::PageDown),
_ => None,
}
}
pub fn decode_ss3(final_byte: u8) -> Option<Key> {
match final_byte {
b'A' => Some(Key::Up),
b'B' => Some(Key::Down),
b'C' => Some(Key::Right),
b'D' => Some(Key::Left),
b'H' => Some(Key::Home),
b'F' => Some(Key::End),
_ => None,
}
}
#[cfg(unix)]
const READ_SLICE: std::time::Duration = std::time::Duration::from_millis(50);
#[cfg(unix)]
const PARK_ACK_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(1);
#[cfg(unix)]
#[derive(Default)]
struct TapControl {
pause: std::sync::atomic::AtomicBool,
parked: std::sync::atomic::AtomicBool,
shutdown: std::sync::atomic::AtomicBool,
}
#[cfg(unix)]
#[derive(Clone, PartialEq, Debug)]
pub enum TapChunk {
Tty(Vec<u8>),
Trigger,
}
#[cfg(unix)]
pub struct TtyTap {
rx: std::sync::mpsc::Receiver<TapChunk>,
tx: std::sync::mpsc::Sender<TapChunk>,
control: std::sync::Arc<TapControl>,
reader: Option<std::thread::JoinHandle<()>>,
}
#[cfg(unix)]
impl TtyTap {
pub fn spawn() -> std::io::Result<TtyTap> {
let tty = std::fs::File::open("/dev/tty")?;
let (tx, rx) = std::sync::mpsc::channel();
let control = std::sync::Arc::new(TapControl::default());
let reader_control = std::sync::Arc::clone(&control);
let reader_tx = tx.clone();
let reader = std::thread::Builder::new()
.name("rat-tty-tap".to_string())
.spawn(move || read_loop(&tty, &reader_tx, &reader_control))?;
Ok(TtyTap {
rx,
tx,
control,
reader: Some(reader),
})
}
pub fn sender(&self) -> std::sync::mpsc::Sender<TapChunk> {
self.tx.clone()
}
pub fn recv_timeout(&self, timeout: std::time::Duration) -> Option<TapChunk> {
use std::sync::mpsc::RecvTimeoutError;
match self.rx.recv_timeout(timeout) {
Ok(chunk) => Some(chunk),
Err(RecvTimeoutError::Timeout) => None,
Err(RecvTimeoutError::Disconnected) => {
std::thread::sleep(timeout);
None
}
}
}
pub fn pause(&self) -> bool {
use std::sync::atomic::Ordering;
self.control.pause.store(true, Ordering::SeqCst);
let deadline = std::time::Instant::now() + PARK_ACK_TIMEOUT;
loop {
if self.control.parked.load(Ordering::SeqCst) {
return true;
}
if self
.reader
.as_ref()
.is_none_or(|reader| reader.is_finished())
{
return true;
}
if std::time::Instant::now() >= deadline {
return false;
}
std::thread::sleep(std::time::Duration::from_millis(2));
}
}
pub fn resume(&self) {
use std::sync::atomic::Ordering;
self.control.parked.store(false, Ordering::SeqCst);
self.control.pause.store(false, Ordering::SeqCst);
}
}
#[cfg(unix)]
impl Drop for TtyTap {
fn drop(&mut self) {
use std::sync::atomic::Ordering;
self.control.shutdown.store(true, Ordering::SeqCst);
self.control.pause.store(false, Ordering::SeqCst);
if let Some(reader) = self.reader.take() {
let _ = reader.join();
}
}
}
#[cfg(unix)]
fn read_loop(tty: &std::fs::File, tx: &std::sync::mpsc::Sender<TapChunk>, control: &TapControl) {
use std::os::unix::io::AsRawFd;
use std::sync::atomic::Ordering;
let fd = tty.as_raw_fd();
let mut buf = [0u8; 256];
loop {
if control.shutdown.load(Ordering::SeqCst) {
return;
}
if control.pause.load(Ordering::SeqCst) {
control.parked.store(true, Ordering::SeqCst);
std::thread::sleep(std::time::Duration::from_millis(2));
continue;
}
let mut read_set: libc::fd_set = unsafe { std::mem::zeroed() };
unsafe {
libc::FD_ZERO(&mut read_set);
libc::FD_SET(fd, &mut read_set);
}
let mut timeout = libc::timeval {
tv_sec: 0,
tv_usec: READ_SLICE.subsec_micros() as libc::suseconds_t,
};
let ready = unsafe {
libc::select(
fd + 1,
&mut read_set,
std::ptr::null_mut(),
std::ptr::null_mut(),
&mut timeout,
)
};
if ready < 0 {
let err = std::io::Error::last_os_error();
if err.kind() == std::io::ErrorKind::Interrupted {
continue;
}
return;
}
if ready == 0 {
continue;
}
if control.pause.load(Ordering::SeqCst) {
continue;
}
let read = unsafe { libc::read(fd, buf.as_mut_ptr().cast::<libc::c_void>(), buf.len()) };
if read <= 0 {
return; }
if tx
.send(TapChunk::Tty(buf[..read as usize].to_vec()))
.is_err()
{
return; }
}
}
#[cfg(test)]
mod tests {
use super::*;
#[cfg(unix)]
#[test]
fn a_posted_trigger_wakes_the_receiver_early() {
let Ok(tap) = TtyTap::spawn() else { return };
let sender = tap.sender();
std::thread::spawn(move || {
let _ = sender.send(TapChunk::Trigger);
});
let start = std::time::Instant::now();
let got = tap.recv_timeout(std::time::Duration::from_secs(5));
assert_eq!(got, Some(TapChunk::Trigger));
assert!(
start.elapsed() < std::time::Duration::from_secs(4),
"the trigger did not wake the receiver early"
);
}
#[test]
fn a_split_report_reassembles_across_feeds() {
let mut scanner = TapScanner::new();
assert_eq!(scanner.feed(b"\x1b[?997"), vec![]);
assert_eq!(
scanner.feed(b";2n"),
vec![TapEvent::ThemeNotification(Appearance::Light)]
);
}
#[test]
fn a_report_sandwiched_between_keys_yields_all_three_in_order() {
let mut scanner = TapScanner::new();
let events = scanner.feed(b"a\x1b[?997;2nb");
assert_eq!(
events,
vec![
TapEvent::Key(Key::Char('a')),
TapEvent::ThemeNotification(Appearance::Light),
TapEvent::Key(Key::Char('b')),
]
);
}
#[test]
fn an_unrecognized_private_csi_is_dropped_without_wedging() {
let mut scanner = TapScanner::new();
assert_eq!(scanner.feed(b"\x1b[?123;4x"), vec![]);
assert_eq!(scanner.feed(b"z"), vec![TapEvent::Key(Key::Char('z'))]);
}
#[test]
fn an_unfinished_sequence_past_the_cap_is_discarded_wholesale() {
let mut scanner = TapScanner::new();
let mut long_run = b"\x1b[".to_vec();
long_run.resize(long_run.len() + 200, 0u8);
assert_eq!(scanner.feed(&long_run), vec![]);
assert_eq!(scanner.feed(b"z"), vec![TapEvent::Key(Key::Char('z'))]);
}
#[test]
fn arrow_keys_decode_through_the_scanner() {
let mut scanner = TapScanner::new();
assert_eq!(scanner.feed(b"\x1b[A"), vec![TapEvent::Key(Key::Up)]);
assert_eq!(scanner.feed(b"\x1b[B"), vec![TapEvent::Key(Key::Down)]);
assert_eq!(scanner.feed(b"\x1b[C"), vec![TapEvent::Key(Key::Right)]);
assert_eq!(scanner.feed(b"\x1b[D"), vec![TapEvent::Key(Key::Left)]);
}
#[test]
fn page_and_home_end_sequences_decode() {
let mut scanner = TapScanner::new();
assert_eq!(scanner.feed(b"\x1b[5~"), vec![TapEvent::Key(Key::PageUp)]);
assert_eq!(scanner.feed(b"\x1b[6~"), vec![TapEvent::Key(Key::PageDown)]);
assert_eq!(scanner.feed(b"\x1b[H"), vec![TapEvent::Key(Key::Home)]);
assert_eq!(scanner.feed(b"\x1b[F"), vec![TapEvent::Key(Key::End)]);
assert_eq!(scanner.feed(b"\x1b[1~"), vec![TapEvent::Key(Key::Home)]);
assert_eq!(scanner.feed(b"\x1b[4~"), vec![TapEvent::Key(Key::End)]);
assert_eq!(scanner.feed(b"\x1b[7~"), vec![TapEvent::Key(Key::Home)]);
assert_eq!(scanner.feed(b"\x1b[8~"), vec![TapEvent::Key(Key::End)]);
}
#[test]
fn a_theme_report_is_still_a_report_not_a_key() {
let mut scanner = TapScanner::new();
assert_eq!(
scanner.feed(b"\x1b[?997;2n"),
vec![TapEvent::ThemeNotification(Appearance::Light)]
);
}
#[test]
fn a_split_arrow_reassembles_across_feeds() {
let mut scanner = TapScanner::new();
assert_eq!(scanner.feed(b"\x1b["), vec![]);
assert_eq!(scanner.feed(b"B"), vec![TapEvent::Key(Key::Down)]);
}
#[test]
fn an_application_cursor_arrow_decodes() {
let mut scanner = TapScanner::new();
assert_eq!(scanner.feed(b"\x1bOA"), vec![TapEvent::Key(Key::Up)]);
assert_eq!(scanner.feed(b"\x1bOS"), vec![]);
}
#[test]
fn an_osc_color_reply_reassembles_across_feeds() {
let mut scanner = TapScanner::new();
assert_eq!(scanner.feed(b"\x1b]11;rgb:1e1e/1e1e/"), vec![]);
assert_eq!(
scanner.feed(b"2e2e\x07"),
vec![TapEvent::OscColor(
OscColorKind::Background,
xterm_color::Color::rgb(0x1e1e, 0x1e1e, 0x2e2e)
)]
);
}
#[test]
fn a_lone_escape_with_no_recognized_introducer_does_not_eat_the_next_byte() {
let mut scanner = TapScanner::new();
assert_eq!(scanner.feed(b"\x1b"), vec![]);
assert_eq!(scanner.feed(b"q"), vec![TapEvent::Key(Key::Char('q'))]);
}
#[test]
fn a_lone_escape_resolves_after_the_hold() {
let mut scanner = TapScanner::new();
assert_eq!(scanner.feed(b"\x1b"), vec![]);
assert_eq!(scanner.idle(std::time::Duration::from_millis(20)), vec![]);
assert_eq!(
scanner.idle(std::time::Duration::from_millis(40)),
vec![TapEvent::Key(Key::Esc)]
);
assert_eq!(scanner.idle(std::time::Duration::from_millis(50)), vec![]);
}
#[test]
fn bytes_cancel_a_pending_escape() {
let mut scanner = TapScanner::new();
assert_eq!(scanner.feed(b"\x1b"), vec![]);
assert_eq!(scanner.idle(std::time::Duration::from_millis(30)), vec![]);
assert_eq!(scanner.feed(b"["), vec![]);
assert_eq!(scanner.idle(std::time::Duration::from_millis(50)), vec![]);
assert_eq!(scanner.feed(b"A"), vec![TapEvent::Key(Key::Up)]);
}
#[test]
fn an_escape_followed_by_a_plain_byte_keeps_todays_behavior() {
let mut scanner = TapScanner::new();
assert_eq!(scanner.feed(b"\x1bq"), vec![TapEvent::Key(Key::Char('q'))]);
}
#[test]
fn idle_never_flushes_a_partial_sequence() {
let mut scanner = TapScanner::new();
assert_eq!(scanner.feed(b"\x1b[?997"), vec![]);
assert_eq!(scanner.idle(std::time::Duration::from_millis(200)), vec![]);
assert_eq!(
scanner.feed(b";2n"),
vec![TapEvent::ThemeNotification(Appearance::Light)]
);
}
#[test]
fn decode_key_maps_the_five_recognized_bytes() {
assert_eq!(decode_key(0x03), Some(Key::CtrlC));
assert_eq!(decode_key(b'\r'), Some(Key::Enter));
assert_eq!(decode_key(b'\n'), Some(Key::Enter));
assert_eq!(decode_key(b'q'), Some(Key::Char('q')));
assert_eq!(decode_key(b'v'), Some(Key::Char('v')));
}
#[test]
fn decode_key_has_no_verdict_for_escape_or_delete() {
assert_eq!(decode_key(0x1b), None);
assert_eq!(decode_key(0x7f), None);
}
}
#[cfg(unix)]
#[derive(Debug)]
pub struct TriggerReader {
fired: std::sync::Arc<std::sync::atomic::AtomicBool>,
ended: std::sync::Arc<std::sync::atomic::AtomicBool>,
shutdown: std::sync::Arc<std::sync::atomic::AtomicBool>,
reader: Option<std::thread::JoinHandle<()>>,
}
#[cfg(unix)]
impl TriggerReader {
pub fn open(
spec: &crate::core::trigger::TriggerSpec,
wake: Option<std::sync::mpsc::Sender<TapChunk>>,
) -> anyhow::Result<TriggerReader> {
use std::os::unix::fs::{FileTypeExt, OpenOptionsExt};
use std::os::unix::io::AsRawFd;
use anyhow::{anyhow, bail};
use crate::core::trigger::TriggerSpec;
let (fd, keep_alive) = match spec {
TriggerSpec::Fifo(path) => {
let read_end = std::fs::OpenOptions::new()
.read(true)
.custom_flags(libc::O_NONBLOCK)
.open(path)
.map_err(|err| match err.kind() {
std::io::ErrorKind::NotFound => anyhow!(
"trigger fifo {} does not exist; create it with: mkfifo {}",
path.display(),
path.display()
),
_ => anyhow!("opening trigger fifo {}: {err}", path.display()),
})?;
if !read_end.metadata()?.file_type().is_fifo() {
bail!(
"fifo:{} is not a named pipe; use file:{} for plain paths",
path.display(),
path.display()
);
}
let write_end = std::fs::OpenOptions::new()
.write(true)
.custom_flags(libc::O_NONBLOCK)
.open(path)?;
(read_end.as_raw_fd(), vec![read_end, write_end])
}
TriggerSpec::Fd(fd) => {
let mut stat: libc::stat = unsafe { std::mem::zeroed() };
if unsafe { libc::fstat(*fd, &mut stat) } != 0 {
bail!("fd:{fd} is not an open descriptor");
}
if stat.st_mode & libc::S_IFMT == libc::S_IFREG {
bail!(
"fd:{fd} is a regular file, which select(2) always reports \
ready; use file:PATH to watch a file"
);
}
(*fd, Vec::new())
}
TriggerSpec::File(_) => bail!("file: triggers are polled, not read"),
};
let fired = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
let ended = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
let shutdown = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
let thread_fired = std::sync::Arc::clone(&fired);
let thread_ended = std::sync::Arc::clone(&ended);
let thread_shutdown = std::sync::Arc::clone(&shutdown);
let reader = std::thread::Builder::new()
.name("rat-trigger".to_string())
.spawn(move || {
let _keep_alive = keep_alive;
trigger_read_loop(fd, &thread_fired, &thread_ended, &thread_shutdown, wake);
})?;
Ok(TriggerReader {
fired,
ended,
shutdown,
reader: Some(reader),
})
}
pub fn fired(&self) -> &std::sync::atomic::AtomicBool {
&self.fired
}
pub fn ended(&self) -> &std::sync::atomic::AtomicBool {
&self.ended
}
}
#[cfg(unix)]
impl Drop for TriggerReader {
fn drop(&mut self) {
use std::sync::atomic::Ordering;
self.shutdown.store(true, Ordering::SeqCst);
if let Some(reader) = self.reader.take() {
let _ = reader.join();
}
}
}
#[cfg(unix)]
fn trigger_read_loop(
fd: i32,
fired: &std::sync::atomic::AtomicBool,
ended: &std::sync::atomic::AtomicBool,
shutdown: &std::sync::atomic::AtomicBool,
wake: Option<std::sync::mpsc::Sender<TapChunk>>,
) {
use std::sync::atomic::Ordering;
let mut buf = [0u8; 256];
loop {
if shutdown.load(Ordering::SeqCst) {
return;
}
let mut read_set: libc::fd_set = unsafe { std::mem::zeroed() };
unsafe {
libc::FD_ZERO(&mut read_set);
libc::FD_SET(fd, &mut read_set);
}
let mut timeout = libc::timeval {
tv_sec: 0,
tv_usec: READ_SLICE.subsec_micros() as libc::suseconds_t,
};
let ready = unsafe {
libc::select(
fd + 1,
&mut read_set,
std::ptr::null_mut(),
std::ptr::null_mut(),
&mut timeout,
)
};
if ready < 0 {
let err = std::io::Error::last_os_error();
if err.kind() == std::io::ErrorKind::Interrupted {
continue;
}
ended.store(true, Ordering::SeqCst);
return;
}
if ready == 0 {
continue;
}
let read = unsafe { libc::read(fd, buf.as_mut_ptr().cast::<libc::c_void>(), buf.len()) };
if read == 0 {
ended.store(true, Ordering::SeqCst);
return; }
if read < 0 {
let err = std::io::Error::last_os_error();
if matches!(
err.kind(),
std::io::ErrorKind::Interrupted | std::io::ErrorKind::WouldBlock
) {
continue;
}
ended.store(true, Ordering::SeqCst);
return;
}
if !fired.swap(true, Ordering::SeqCst)
&& let Some(wake) = wake.as_ref()
{
let _ = wake.send(TapChunk::Trigger);
}
}
}
#[cfg(all(test, unix))]
mod trigger_reader_tests {
use std::io::Write;
use std::os::unix::io::AsRawFd;
use std::sync::atomic::Ordering;
use std::time::Duration;
use super::*;
use crate::core::trigger::TriggerSpec;
fn mkfifo(path: &std::path::Path) {
let cpath = std::ffi::CString::new(path.as_os_str().as_encoded_bytes().to_vec()).unwrap();
assert_eq!(unsafe { libc::mkfifo(cpath.as_ptr(), 0o600) }, 0, "mkfifo");
}
fn wait_until(mut cond: impl FnMut() -> bool) -> bool {
let deadline = std::time::Instant::now() + Duration::from_secs(3);
while std::time::Instant::now() < deadline {
if cond() {
return true;
}
std::thread::sleep(Duration::from_millis(10));
}
false
}
#[test]
fn a_fifo_write_raises_the_fired_flag() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("t.fifo");
mkfifo(&path);
let reader = TriggerReader::open(&TriggerSpec::Fifo(path.clone()), None).unwrap();
let mut writer = std::fs::OpenOptions::new().write(true).open(&path).unwrap();
writer.write_all(b"x").unwrap();
assert!(
wait_until(|| reader.fired().swap(false, Ordering::SeqCst)),
"the write never raised the flag"
);
}
#[test]
fn a_writerless_fifo_does_not_spin_or_end() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("t.fifo");
mkfifo(&path);
let reader = TriggerReader::open(&TriggerSpec::Fifo(path), None).unwrap();
std::thread::sleep(Duration::from_millis(200));
assert!(!reader.ended().load(Ordering::SeqCst));
assert!(!reader.fired().load(Ordering::SeqCst));
}
#[test]
fn a_regular_file_fd_is_rejected_at_open_with_the_teaching_error() {
let dir = tempfile::tempdir().unwrap();
let f = dir.path().join("reg");
std::fs::write(&f, b"x").unwrap();
let file = std::fs::File::open(&f).unwrap();
let err = TriggerReader::open(&TriggerSpec::Fd(file.as_raw_fd()), None)
.unwrap_err()
.to_string();
assert!(err.contains("file:"), "{err}"); }
#[test]
fn fd_eof_sets_ended_and_the_thread_exits() {
let mut fds = [0i32; 2];
assert_eq!(unsafe { libc::pipe(fds.as_mut_ptr()) }, 0, "pipe");
let (r, w) = (fds[0], fds[1]);
let reader = TriggerReader::open(&TriggerSpec::Fd(r), None).unwrap();
unsafe { libc::close(w) };
assert!(
wait_until(|| reader.ended().load(Ordering::SeqCst)),
"EOF never set ended"
);
}
}