use std::sync::{
LazyLock,
atomic::{AtomicBool, Ordering},
};
use parking_lot::Mutex;
use crate::{
input::{InputEvent, Key, Mouse, MouseButton, MouseReport},
pump::{DebugOp, DebugQuery, TerminalEvent},
};
pub const DEBUG_ENV: &str = "OMP_TUI_DEBUG";
static SCREEN: Mutex<Option<ScreenSnapshot>> = Mutex::new(None);
static RESPONSES: Mutex<Option<flume::Sender<(u64, serde_json::Value)>>> = Mutex::new(None);
pub fn enabled() -> bool {
static ENABLED: LazyLock<bool> =
LazyLock::new(|| std::env::var_os(DEBUG_ENV).is_some_and(|value| !value.is_empty()));
*ENABLED
}
pub fn publishing() -> bool {
enabled()
}
pub fn respond_debug_query(id: u64, response: serde_json::Value) {
let sender = RESPONSES.lock().clone();
if let Some(sender) = sender {
let _ = sender.send((id, response));
}
}
pub fn publish_screen(snapshot: ScreenSnapshot) {
*SCREEN.lock() = Some(snapshot);
}
pub fn screen_snapshot() -> Option<ScreenSnapshot> {
SCREEN.lock().clone()
}
#[derive(Clone)]
pub struct ScreenSnapshot {
pub lines: Vec<String>,
pub cursor: Option<(u16, u16)>,
pub window_top: u16,
pub cols: u16,
pub rows: u16,
pub doc_height: u16,
pub overlay: bool,
}
fn direct_response(request: DebugRequest) -> Result<serde_json::Value, DebugOp> {
use serde_json::json;
Ok(match request {
DebugRequest::Info => return Err(DebugOp::Info),
DebugRequest::Text => return Err(DebugOp::Text),
DebugRequest::Frame => return Err(DebugOp::Frame),
DebugRequest::Tree => return Err(DebugOp::Tree),
DebugRequest::Values => return Err(DebugOp::Values),
DebugRequest::Resize => return Err(DebugOp::Resize),
DebugRequest::Quit => return Err(DebugOp::Quit),
DebugRequest::Inject(events) => {
let injected = events.len();
if events
.into_iter()
.all(|event| crate::pump::send_event(TerminalEvent::Input(event)))
{
json!({ "ok": true, "injected": injected })
} else {
json!({ "ok": false, "error": "no live terminal to inject into" })
}
},
DebugRequest::Events(events) => {
let injected = events.len();
if events.into_iter().all(crate::pump::send_event) {
json!({ "ok": true, "injected": injected })
} else {
json!({ "ok": false, "error": "no live terminal to inject into" })
}
},
DebugRequest::Bytes(bytes) => {
let fed = bytes.len();
if crate::pump::inject_bytes(bytes) {
json!({ "ok": true, "fed": fed })
} else {
json!({ "ok": false, "error": "no live terminal to inject into" })
}
},
})
}
pub fn terminal_response(op: DebugOp) -> Option<serde_json::Value> {
use serde_json::json;
let snapshot = |build: fn(ScreenSnapshot) -> serde_json::Value| {
screen_snapshot()
.map_or_else(|| json!({ "ok": false, "error": "no frame painted yet" }), build)
};
match op {
DebugOp::Info => Some(snapshot(|snapshot| {
json!({
"ok": true,
"cols": snapshot.cols,
"rows": snapshot.rows,
"height": snapshot.doc_height,
"window_top": snapshot.window_top,
"alt_screen": crate::terminal::alt_screen_active(),
"overlay": snapshot.overlay,
})
})),
DebugOp::Text => Some(snapshot(|snapshot| {
json!({
"ok": true,
"lines": snapshot.lines,
"cursor": snapshot.cursor.map(|(row, col)| vec![row, col]),
"window_top": snapshot.window_top,
"alt_screen": crate::terminal::alt_screen_active(),
})
})),
DebugOp::Resize => {
crate::terminal::simulate_resize_signal();
Some(json!({ "ok": true, "signalled": true }))
},
DebugOp::Quit => Some(json!({ "ok": true, "injected": "C-c" })),
DebugOp::Frame | DebugOp::Tree | DebugOp::Values => None,
}
}
pub enum DebugRequest {
Info,
Text,
Frame,
Tree,
Values,
Inject(Vec<InputEvent>),
Events(Vec<TerminalEvent>),
Bytes(Vec<u8>),
Resize,
Quit,
}
pub fn parse_request(line: &[u8]) -> Result<DebugRequest, String> {
let value: serde_json::Value =
serde_json::from_slice(line).map_err(|error| format!("malformed request: {error}"))?;
let op = value
.get("op")
.and_then(serde_json::Value::as_str)
.ok_or_else(|| "missing \"op\"".to_owned())?;
match op {
"info" => Ok(DebugRequest::Info),
"text" => Ok(DebugRequest::Text),
"frame" => Ok(DebugRequest::Frame),
"tree" => Ok(DebugRequest::Tree),
"values" => Ok(DebugRequest::Values),
"keys" => {
let spec = value
.get("keys")
.and_then(serde_json::Value::as_str)
.ok_or_else(|| "keys op needs a \"keys\" string".to_owned())?;
Ok(DebugRequest::Inject(parse_keys(spec)?.into_iter().map(InputEvent::Key).collect()))
},
"bytes" => {
let data = value
.get("data")
.and_then(serde_json::Value::as_str)
.ok_or_else(|| "bytes op needs a \"data\" string".to_owned())?;
Ok(DebugRequest::Bytes(data.as_bytes().to_vec()))
},
"paste" => {
let text = value
.get("text")
.and_then(serde_json::Value::as_str)
.ok_or_else(|| "paste op needs a \"text\" string".to_owned())?;
Ok(DebugRequest::Inject(vec![InputEvent::Paste(text.into())]))
},
"mouse" => Ok(DebugRequest::Inject(vec![InputEvent::Mouse(parse_mouse(&value)?)])),
"event" | "events" => {
let payload = value
.get("event")
.or_else(|| value.get("events"))
.ok_or_else(|| "event op needs an \"event\" (or \"events\") field".to_owned())?;
let events = if payload.is_array() {
serde_json::from_value::<Vec<TerminalEvent>>(payload.clone())
} else {
serde_json::from_value::<TerminalEvent>(payload.clone()).map(|event| vec![event])
}
.map_err(|error| format!("malformed terminal event: {error}"))?;
Ok(DebugRequest::Events(events))
},
"resize" => Ok(DebugRequest::Resize),
"quit" => Ok(DebugRequest::Quit),
other => Err(format!("unknown op {other:?}")),
}
}
pub fn parse_keys(spec: &str) -> Result<Vec<Key>, String> {
let mut keys = Vec::new();
let mut rest = spec.trim_start();
while !rest.is_empty() {
if let Some(quote) = rest.chars().next().filter(|ch| matches!(ch, '\'' | '"')) {
let body = &rest[quote.len_utf8()..];
let end = body
.find(quote)
.ok_or_else(|| format!("unterminated quote in key spec: {rest:?}"))?;
keys.extend(body[..end].chars().map(literal_key));
rest = body[end + quote.len_utf8()..].trim_start();
continue;
}
let token = rest
.split_whitespace()
.next()
.expect("non-empty trimmed spec");
keys.push(parse_token(token)?);
rest = rest[token.len()..].trim_start();
}
Ok(keys)
}
const fn literal_key(ch: char) -> Key {
match ch {
' ' => Key::Space,
_ => Key::Char(ch),
}
}
fn parse_token(token: &str) -> Result<Key, String> {
if token.chars().count() > 1 {
let lower = token.to_ascii_lowercase();
if let Some(ch) = strip_chord(&lower, &["c-m-", "m-c-", "ctrl-alt-"]) {
return Ok(Key::CtrlAlt(ch));
}
if let Some(ch) = strip_chord(&lower, &["c-", "ctrl-", "ctrl+"]) {
return Ok(Key::Ctrl(ch));
}
if let Some(ch) = strip_chord(&lower, &["m-", "a-", "alt-", "alt+"]) {
return Ok(Key::Alt(ch));
}
}
let named = match token.to_ascii_lowercase().as_str() {
"up" => Key::Up,
"down" => Key::Down,
"left" => Key::Left,
"right" => Key::Right,
"tab" => Key::Tab,
"backtab" | "shift-tab" => Key::BackTab,
"enter" | "return" | "cr" => Key::Enter,
"space" => Key::Space,
"esc" | "escape" => Key::Esc,
"backspace" | "bs" => Key::Backspace,
"delete" | "del" => Key::Delete,
"insert" => Key::Insert,
"home" => Key::Home,
"end" => Key::End,
"pgup" | "pageup" => Key::PageUp,
"pgdn" | "pagedown" => Key::PageDown,
"shift-enter" => Key::ShiftEnter,
"word-left" => Key::WordLeft,
"word-right" => Key::WordRight,
"word-delete" => Key::WordDelete,
other => {
if let Some(number) = other.strip_prefix('f')
&& let Ok(number) = number.parse::<u8>()
&& (1..=12).contains(&number)
{
return Ok(Key::Function(number));
}
let mut chars = token.chars();
return match (chars.next(), chars.next()) {
(Some(ch), None) => Ok(literal_key(ch)),
_ => Err(format!("unknown key token {token:?}")),
};
},
};
Ok(named)
}
fn strip_chord(token: &str, prefixes: &[&str]) -> Option<char> {
prefixes.iter().find_map(|prefix| {
let rest = token.strip_prefix(prefix)?;
let mut chars = rest.chars();
match (chars.next(), chars.next()) {
(Some(ch), None) => Some(ch),
_ => None,
}
})
}
fn parse_mouse(value: &serde_json::Value) -> Result<MouseReport, String> {
let coordinate = |name: &str| -> Result<u16, String> {
value
.get(name)
.and_then(serde_json::Value::as_u64)
.and_then(|number| u16::try_from(number).ok())
.ok_or_else(|| format!("mouse op needs a numeric \"{name}\""))
};
let col = coordinate("x")?;
let row = coordinate("y")?;
let action = value
.get("action")
.and_then(serde_json::Value::as_str)
.unwrap_or("click");
let (kind, button, pressed) = match action {
"click" | "press" => (Mouse::Click, MouseButton::Left, true),
"right-click" => (Mouse::RightClick, MouseButton::Right, true),
"middle-click" => (Mouse::MiddleClick, MouseButton::Middle, true),
"move" => (Mouse::Move, MouseButton::None, false),
"drag" => (Mouse::Drag, MouseButton::Left, true),
"release" => (Mouse::Release, MouseButton::Left, false),
"wheel-up" => (Mouse::WheelUp, MouseButton::WheelUp, true),
"wheel-down" => (Mouse::WheelDown, MouseButton::WheelDown, true),
"wheel-left" => (Mouse::WheelLeft, MouseButton::WheelLeft, true),
"wheel-right" => (Mouse::WheelRight, MouseButton::WheelRight, true),
other => return Err(format!("unknown mouse action {other:?}")),
};
Ok(MouseReport { kind, col, row, button, mods: Default::default(), pressed })
}
#[cfg(unix)]
static SERVER_STARTED: AtomicBool = AtomicBool::new(false);
#[cfg(unix)]
pub fn ensure_server() -> std::io::Result<()> {
if !enabled() || SERVER_STARTED.swap(true, Ordering::AcqRel) {
return Ok(());
}
server::spawn_thread()
}
#[cfg(not(unix))]
pub(crate) fn ensure_server() -> std::io::Result<()> {
Ok(())
}
#[cfg(unix)]
mod server {
use std::{io, path::PathBuf, task::Poll, time::Duration};
use tokio::net::{UnixListener, UnixStream};
use super::{
DEBUG_ENV, DebugQuery, DebugRequest, RESPONSES, TerminalEvent, direct_response, parse_request,
};
const QUERY_TIMEOUT: Duration = Duration::from_secs(2);
pub(super) fn spawn_thread() -> io::Result<()> {
let path = PathBuf::from(
std::env::var_os(DEBUG_ENV).expect("enabled() checked the variable before spawning"),
);
match std::fs::remove_file(&path) {
Ok(()) => {},
Err(error) if error.kind() == io::ErrorKind::NotFound => {},
Err(error) => return Err(error),
}
let listener = std::os::unix::net::UnixListener::bind(&path)?;
listener.set_nonblocking(true)?;
let (responses_tx, responses_rx) = flume::unbounded();
*RESPONSES.lock() = Some(responses_tx);
std::thread::Builder::new()
.name("omp-tui-debug".into())
.spawn(move || {
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_io()
.enable_time()
.build()
.expect("debug server runtime builds");
runtime.block_on(async move {
let Ok(listener) = UnixListener::from_std(listener) else {
return;
};
let mut server = DebugServer::new(listener);
serve_loop(&mut server, responses_rx).await;
});
})?;
Ok(())
}
struct PendingQuery {
id: u64,
client: u64,
expires: tokio::time::Instant,
}
async fn serve_loop(
server: &mut DebugServer,
responses: flume::Receiver<(u64, serde_json::Value)>,
) {
let mut pending: Vec<PendingQuery> = Vec::new();
let mut next_id = 1_u64;
loop {
let expiry = pending.iter().map(|query| query.expires).min();
tokio::select! {
received = server.recv() => {
let (client, request) = received;
let request = match request {
Err(error) => {
server
.respond(client, &serde_json::json!({ "ok": false, "error": error }));
continue;
},
Ok(request) => request,
};
match direct_response(request) {
Ok(response) => server.respond(client, &response),
Err(op) => {
let id = next_id;
next_id += 1;
if crate::pump::send_event(TerminalEvent::Debug(DebugQuery { id, op })) {
pending.push(PendingQuery {
id,
client,
expires: tokio::time::Instant::now() + QUERY_TIMEOUT,
});
} else {
server.respond(client, &serde_json::json!({
"ok": false,
"error": "no live terminal to query",
}));
}
},
}
},
response = responses.recv_async() => {
let Ok((id, response)) = response else {
return;
};
if let Some(index) = pending.iter().position(|query| query.id == id) {
let query = pending.swap_remove(index);
server.respond(query.client, &response);
}
},
() = expire(expiry) => {
let now = tokio::time::Instant::now();
let mut index = 0;
while index < pending.len() {
if pending[index].expires <= now {
let query = pending.swap_remove(index);
server.respond(query.client, &serde_json::json!({
"ok": false,
"error": "no retained host answered; use `text` (omp_tui::App answers frame/tree/values)",
}));
} else {
index += 1;
}
}
},
}
}
}
async fn expire(at: Option<tokio::time::Instant>) {
match at {
Some(at) => tokio::time::sleep_until(at).await,
None => std::future::pending().await,
}
}
struct DebugServer {
listener: UnixListener,
conns: Vec<Conn>,
next_conn: u64,
}
struct Conn {
id: u64,
stream: UnixStream,
buf: Vec<u8>,
out: Vec<u8>,
dead: bool,
}
impl Conn {
fn take_line(&mut self) -> Option<Vec<u8>> {
let end = self.buf.iter().position(|byte| *byte == b'\n')?;
let line = self.buf[..end].to_vec();
self.buf.drain(..=end);
Some(line)
}
fn fill(&mut self) {
let mut bytes = [0_u8; 4096];
loop {
match self.stream.try_read(&mut bytes) {
Ok(0) => {
self.dead = true;
return;
},
Ok(read) => self.buf.extend_from_slice(&bytes[..read]),
Err(error) if error.kind() == io::ErrorKind::WouldBlock => return,
Err(_) => {
self.dead = true;
return;
},
}
}
}
fn flush(&mut self) {
while !self.out.is_empty() {
match self.stream.try_write(&self.out) {
Ok(0) => {
self.dead = true;
return;
},
Ok(written) => {
self.out.drain(..written);
},
Err(error) if error.kind() == io::ErrorKind::WouldBlock => return,
Err(_) => {
self.dead = true;
return;
},
}
}
}
}
impl DebugServer {
const fn new(listener: UnixListener) -> Self {
Self { listener, conns: Vec::new(), next_conn: 1 }
}
async fn recv(&mut self) -> (u64, Result<DebugRequest, String>) {
loop {
self.conns.retain(|conn| !conn.dead);
for conn in &mut self.conns {
if let Some(line) = conn.take_line() {
return (conn.id, parse_request(&line));
}
}
let Self { listener, conns, next_conn } = self;
tokio::select! {
accepted = listener.accept() => {
if let Ok((stream, _)) = accepted {
let id = *next_conn;
*next_conn += 1;
conns.push(Conn {
id,
stream,
buf: Vec::new(),
out: Vec::new(),
dead: false,
});
}
},
index = ready(conns) => {
conns[index].fill();
conns[index].flush();
},
}
}
}
fn respond(&mut self, client: u64, response: &serde_json::Value) {
let Some(conn) = self.conns.iter_mut().find(|conn| conn.id == client) else {
return;
};
serde_json::to_writer(&mut conn.out, response).expect("JSON responses serialize");
conn.out.push(b'\n');
conn.flush();
}
}
async fn ready(conns: &[Conn]) -> usize {
std::future::poll_fn(|cx| {
for (index, conn) in conns.iter().enumerate() {
if conn.stream.poll_read_ready(cx).is_ready()
|| (!conn.out.is_empty() && conn.stream.poll_write_ready(cx).is_ready())
{
return Poll::Ready(index);
}
}
Poll::Pending
})
.await
}
#[cfg(test)]
mod tests {
use std::time::Duration;
use tokio::io::{AsyncBufReadExt as _, AsyncWriteExt as _, BufReader};
use super::{DebugServer, serve_loop};
#[tokio::test]
async fn pending_query_replies_follow_stable_client_ids() {
let ingress = crate::pump::publish_ingress_for_test();
let (responses_tx, responses_rx) = flume::unbounded();
let path =
std::env::temp_dir().join(format!("omp-tui-debug-idtest-{}.sock", std::process::id()));
let _ = std::fs::remove_file(&path);
let listener = std::os::unix::net::UnixListener::bind(&path).expect("test socket binds");
listener
.set_nonblocking(true)
.expect("nonblocking listener");
let listener = tokio::net::UnixListener::from_std(listener).expect("listener registers");
let mut server = DebugServer::new(listener);
let serve = tokio::spawn(async move { serve_loop(&mut server, responses_rx).await });
let mut first = tokio::net::UnixStream::connect(&path)
.await
.expect("first client connects");
first
.write_all(b"{\"op\":\"tree\"}\n")
.await
.expect("query sends");
tokio::time::timeout(Duration::from_secs(1), ingress.recv_async())
.await
.expect("query reaches the ingress")
.expect("ingress lives");
drop(first);
let second = tokio::net::UnixStream::connect(&path)
.await
.expect("second client connects");
let mut second = BufReader::new(second);
second
.get_mut()
.write_all(b"{\"op\":\"keys\",\"keys\":\"x\"}\n")
.await
.expect("injection sends");
let mut line = String::new();
tokio::time::timeout(Duration::from_secs(1), second.read_line(&mut line))
.await
.expect("injection is acknowledged")
.expect("ack line reads");
assert!(line.contains("\"injected\""), "unexpected ack: {line:?}");
responses_tx
.send((1, serde_json::json!({ "leak": "wrong-client", "ok": true })))
.expect("reply channel lives");
line.clear();
let stray =
tokio::time::timeout(Duration::from_millis(200), second.read_line(&mut line)).await;
assert!(
stray.is_err() || line.is_empty(),
"reply for a disconnected client reached the survivor: {line:?}"
);
serve.abort();
let _ = std::fs::remove_file(&path);
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn key_spec_tokens_chords_and_literals() {
let keys = parse_keys("tab C-c M-y 'hi there' x pgdn f5").expect("valid spec");
assert_eq!(keys, vec![
Key::Tab,
Key::Ctrl('c'),
Key::Alt('y'),
Key::Char('h'),
Key::Char('i'),
Key::Space,
Key::Char('t'),
Key::Char('h'),
Key::Char('e'),
Key::Char('r'),
Key::Char('e'),
Key::Char('x'),
Key::PageDown,
Key::Function(5),
]);
}
#[test]
fn key_spec_rejects_unknown_and_unterminated() {
assert!(parse_keys("bogus-token").is_err());
assert!(parse_keys("'open").is_err());
}
#[test]
fn request_lines_parse_by_op() {
assert!(matches!(parse_request(br#"{"op":"text"}"#), Ok(DebugRequest::Text)));
assert!(matches!(
parse_request(br#"{"op":"keys","keys":"enter"}"#),
Ok(DebugRequest::Inject(events)) if events == vec![InputEvent::Key(Key::Enter)]
));
assert!(parse_request(br#"{"op":"warp"}"#).is_err());
let mouse = parse_request(br#"{"op":"mouse","x":3,"y":7,"action":"wheel-down"}"#);
assert!(matches!(
mouse,
Ok(DebugRequest::Inject(events))
if matches!(&events[..], [InputEvent::Mouse(report)]
if report.kind == Mouse::WheelDown && report.col == 3 && report.row == 7)
));
}
}