use std::collections::{HashMap, VecDeque};
use std::path::Path;
use std::sync::Mutex;
use std::time::{Duration, Instant};
use libtmux::ServerGeneration;
use crate::run_request::{EndpointIdentity, endpoint_identity};
const RECENT_TTL: Duration = Duration::from_secs(10);
const RECENT_PER_PANE: usize = 4;
const MAX_TRACKED_PANES: usize = 256;
#[derive(Clone, Debug, Eq, Hash, PartialEq)]
pub(crate) struct EchoKey {
generation: ServerGeneration,
endpoint: EndpointIdentity,
pane: String,
}
impl EchoKey {
pub(crate) fn new(generation: ServerGeneration, endpoint: &Path, pane: &str) -> Option<Self> {
Some(Self {
generation,
endpoint: endpoint_identity(endpoint).ok()?,
pane: pane.to_owned(),
})
}
}
#[derive(Default)]
struct PaneRecord {
pending: String,
in_flight: u32,
recent: VecDeque<(String, Instant)>,
touched: Option<Instant>,
}
impl PaneRecord {
fn push_recent(&mut self, line: String, now: Instant) {
if let Some((last, at)) = self.recent.back_mut()
&& *last == line
{
*at = now;
return;
}
self.recent.push_back((line, now));
while self.recent.len() > RECENT_PER_PANE {
self.recent.pop_front();
}
}
fn prune(&mut self, now: Instant) {
while self
.recent
.front()
.is_some_and(|(_, at)| now.duration_since(*at) > RECENT_TTL)
{
self.recent.pop_front();
}
}
fn has_pending(&self) -> bool {
self.in_flight > 0 || !self.pending.is_empty()
}
fn is_empty(&self, now: Instant) -> bool {
self.in_flight == 0
&& self.pending.is_empty()
&& self.recent.is_empty()
&& self
.touched
.is_none_or(|touched| now.duration_since(touched) > RECENT_TTL)
}
}
struct DispatchOutcome {
submitted: Option<String>,
}
const SUBMIT_KEYS: [&str; 3] = ["Enter", "C-m", "KPEnter"];
const KILL_LINE_KEYS: [&str; 2] = ["C-u", "C-c"];
const ERASE_KEYS: [&str; 2] = ["BSpace", "C-h"];
fn typed_char(key: &str) -> Option<char> {
if key == "Space" {
return Some(' ');
}
let mut chars = key.chars();
let only = chars.next()?;
if chars.next().is_some() || only.is_control() {
return None;
}
Some(only)
}
fn apply_dispatch(pending: &mut String, text: Option<&str>, keys: &[String]) -> DispatchOutcome {
if let Some(text) = text.filter(|text| !text.is_empty()) {
pending.push_str(text);
}
let mut submitted = None;
for key in keys {
let key = key.as_str();
if SUBMIT_KEYS.contains(&key) {
let line = std::mem::take(pending);
if !line.is_empty() {
submitted = Some(line);
}
} else if KILL_LINE_KEYS.contains(&key) {
pending.clear();
} else if ERASE_KEYS.contains(&key) {
pending.pop();
} else if key == "DC" {
} else if let Some(character) = typed_char(key) {
pending.push(character);
} else {
pending.clear();
}
}
DispatchOutcome { submitted }
}
pub(crate) struct EchoUpdate {
pending: Vec<(EchoKey, String)>,
}
#[derive(Default)]
pub(crate) struct PaneEchoes {
inner: Mutex<HashMap<EchoKey, PaneRecord>>,
}
impl PaneEchoes {
pub(crate) fn new() -> Self {
Self::default()
}
pub(crate) fn apply(
&self,
generation: ServerGeneration,
endpoint: &Path,
panes: &[String],
text: Option<&str>,
keys: &[String],
) -> EchoUpdate {
let Ok(endpoint) = endpoint_identity(endpoint) else {
return EchoUpdate {
pending: Vec::new(),
};
};
let now = Instant::now();
let mut table = self.hold();
let mut pending = Vec::with_capacity(panes.len());
for pane in panes {
let key = EchoKey {
generation,
endpoint,
pane: pane.clone(),
};
let record = table.entry(key.clone()).or_default();
record.in_flight += 1;
record.touched = Some(now);
let mut scratch = record.pending.clone();
let outcome = apply_dispatch(&mut scratch, text, keys);
if let Some(line) = outcome.submitted {
record.push_recent(line, now);
}
pending.push((key, scratch));
}
evict_stale(&mut table, now);
EchoUpdate { pending }
}
pub(crate) fn commit(&self, update: EchoUpdate) {
let now = Instant::now();
let mut table = self.hold();
for (key, computed) in update.pending {
if let Some(record) = table.get_mut(&key) {
record.pending = computed;
record.in_flight = record.in_flight.saturating_sub(1);
record.touched = Some(now);
}
}
evict_stale(&mut table, now);
}
pub(crate) fn abandon(&self, update: EchoUpdate) {
let mut table = self.hold();
for (key, _) in update.pending {
if let Some(record) = table.get_mut(&key) {
record.in_flight = record.in_flight.saturating_sub(1);
}
}
}
pub(crate) fn has_pending(&self, key: &EchoKey) -> bool {
self.hold().get(key).is_some_and(PaneRecord::has_pending)
}
pub(crate) fn snapshot(&self, key: &EchoKey) -> Vec<Vec<u8>> {
let now = Instant::now();
let mut table = self.hold();
let Some(record) = table.get_mut(key) else {
return Vec::new();
};
record.prune(now);
record
.recent
.iter()
.map(|(line, _)| line.clone().into_bytes())
.collect()
}
#[cfg(test)]
pub(crate) fn tracked_panes(&self) -> usize {
self.hold().len()
}
fn hold(&self) -> std::sync::MutexGuard<'_, HashMap<EchoKey, PaneRecord>> {
self.inner
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
}
}
fn evict_stale(table: &mut HashMap<EchoKey, PaneRecord>, now: Instant) {
table.retain(|_, record| {
record.prune(now);
!record.is_empty(now)
});
while table.len() > MAX_TRACKED_PANES {
let Some(stale) = table
.iter()
.min_by_key(|(_, record)| record.touched)
.map(|(key, _)| key.clone())
else {
break;
};
table.remove(&stale);
}
}
const fn is_word_byte(byte: u8) -> bool {
byte.is_ascii_alphanumeric() || byte == b'_' || byte >= 0x80
}
fn without_echo(haystack: &[u8], needle: &[u8]) -> Vec<u8> {
if needle.is_empty() || haystack.len() < needle.len() {
return haystack.to_vec();
}
let mut masked = vec![false; haystack.len()];
let mut from = 0;
while from + needle.len() <= haystack.len() {
let Some(offset) = haystack[from..]
.windows(needle.len())
.position(|window| window == needle)
else {
break;
};
let at = from + offset;
let end = at + needle.len();
let opens = at == 0 || !is_word_byte(haystack[at - 1]);
let closes = end == haystack.len() || !is_word_byte(haystack[end]);
if opens && closes {
for slot in &mut masked[at..end] {
*slot = true;
}
from = end;
} else {
from = at + 1;
}
}
haystack
.iter()
.zip(masked)
.filter_map(|(&byte, hit)| (!hit).then_some(byte))
.collect()
}
pub(crate) fn mask(haystack: &[u8], echoes: &[Vec<u8>]) -> Vec<u8> {
let mut text = haystack.to_vec();
for echo in echoes {
text = without_echo(&text, echo);
}
text
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_submitted_line_is_masked_whole_word_only() {
let haystack = b"echo MARKER\nMARKER\nprefixMARKER\n";
let masked = mask(haystack, &[b"echo MARKER".to_vec()]);
assert_eq!(masked, b"\nMARKER\nprefixMARKER\n");
}
#[test]
fn masking_never_tears_a_longer_word() {
let masked = without_echo(b"prefixMARKERsuffix", b"MARKER");
assert_eq!(masked, b"prefixMARKERsuffix");
}
#[test]
fn repeated_output_still_matches_after_masking_the_echo() {
let haystack = b"sleep 1; echo MARKER\nMARKER\n";
let masked = mask(haystack, &[b"sleep 1; echo MARKER".to_owned().to_vec()]);
assert_eq!(masked, b"\nMARKER\n");
}
#[test]
fn backspaces_reach_an_earlier_calls_pending_text() {
let mut pending = String::new();
apply_dispatch(&mut pending, Some("xMARKER"), &[]);
assert_eq!(pending, "xMARKER");
let backspaces = vec!["BSpace".to_owned(); 7];
apply_dispatch(&mut pending, None, &backspaces);
assert_eq!(pending, "");
let outcome = apply_dispatch(&mut pending, Some("echo MARKER"), &["Enter".to_owned()]);
assert_eq!(outcome.submitted.as_deref(), Some("echo MARKER"));
assert_eq!(pending, "");
}
#[test]
fn an_unmodelable_key_clears_pending_instead_of_keeping_it_stale() {
let mut pending = "MARKER".to_owned();
let outcome = apply_dispatch(&mut pending, None, &["Left".to_owned()]);
assert!(outcome.submitted.is_none());
assert_eq!(
pending, "",
"an unrecognized key stops discounting the line"
);
}
#[test]
fn kill_line_keys_discard_without_recording_an_echo() {
let mut pending = "doomed".to_owned();
let outcome = apply_dispatch(&mut pending, None, &["C-u".to_owned()]);
assert!(outcome.submitted.is_none());
assert_eq!(pending, "");
}
#[tokio::test]
async fn a_killed_panes_record_does_not_outlive_the_table_cap() {
let guard = libtmux::test::TestServer::builder()
.start()
.await
.expect("tmux starts");
let generation = guard
.server()
.generation()
.await
.expect("server generation");
let endpoint = guard.server().socket_path().to_path_buf();
let echoes = PaneEchoes::new();
for index in 0..MAX_TRACKED_PANES + 8 {
let pane = format!("%{index}");
let _ = echoes.apply(
generation,
&endpoint,
std::slice::from_ref(&pane),
Some("x"),
&[],
);
}
assert!(
echoes.tracked_panes() <= MAX_TRACKED_PANES,
"the table stays bounded at {} panes",
echoes.tracked_panes()
);
guard.shutdown().await.expect("tmux fixture shuts down");
}
}