use std::time::Duration;
use tokio_util::sync::CancellationToken;
use libtmux::{CaptureOptions, ControlLimits, ControlModeErrorKind, Error, Pane};
use regex::bytes::Regex;
use serde::Serialize;
use crate::echo::{EchoKey, PaneEchoes};
use crate::retained::MAX_BYTES as OUTPUT_LIMIT;
#[cfg(test)]
use crate::retained::{COMPACT_AFTER, RetainedBytes};
use crate::text::TextFilter;
#[cfg(test)]
use crate::text::readable_from;
const MAX_PATTERNS: usize = 32;
const MAX_PATTERN_BYTES: usize = 4_096;
const MAX_TOTAL_PATTERN_BYTES: usize = 16_384;
#[derive(Clone, Copy)]
pub(crate) struct EchoContext<'a> {
pub(crate) echoes: &'a PaneEchoes,
pub(crate) key: Option<&'a EchoKey>,
}
mod run;
#[cfg(test)]
pub(crate) use run::observing_prepared_shutdowns;
#[cfg(test)]
use run::{
FrameError, Scanner, TRAP_DECLARATION_LIMIT, find, frame_path, frame_with_random,
inherited_trap_capture, quote_shell_word, render_payload, stage_frame, staged_line,
};
pub(crate) use run::{
PrepareRunError, RunDispatch, RunProgress, prepare_run, readable, route_is_terminal_safe,
route_path_is_terminal_safe,
};
#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, schemars::JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum RunOutcome {
Completed,
Deadline,
PaneClosed,
Cancelled,
NoShell,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, schemars::JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum WaitOutcome {
PresentAtEntry,
Pending,
Matched,
Stopped,
Deadline,
PaneClosed,
Cancelled,
}
#[derive(Clone, Debug, Serialize, schemars::JsonSchema)]
pub struct RunView {
pub pane: String,
pub outcome: RunOutcome,
pub exit_status: Option<i32>,
pub output: String,
pub bytes: usize,
pub truncated: bool,
}
#[derive(Debug, Serialize, schemars::JsonSchema)]
pub struct WaitView {
pub pane: String,
pub outcome: WaitOutcome,
pub matched_index: Option<usize>,
pub matched_pattern: Option<String>,
pub text: String,
pub bytes: usize,
}
pub(crate) struct Patterns {
compiled: Vec<Regex>,
sources: Vec<String>,
}
impl Patterns {
pub(crate) fn compile(
sources: &[String],
regex: bool,
match_case: bool,
) -> Result<Self, (String, String)> {
if sources.len() > MAX_PATTERNS {
return Err((
"set".to_owned(),
format!("contains more than {MAX_PATTERNS} patterns"),
));
}
let mut compiled = Vec::with_capacity(sources.len());
let mut total_bytes = 0_usize;
for (index, source) in sources.iter().enumerate() {
if source.len() > MAX_PATTERN_BYTES {
return Err((
format!("{}", index + 1),
format!("exceeds {MAX_PATTERN_BYTES} bytes"),
));
}
total_bytes = total_bytes.saturating_add(source.len());
if total_bytes > MAX_TOTAL_PATTERN_BYTES {
return Err((
"set".to_owned(),
format!("exceeds {MAX_TOTAL_PATTERN_BYTES} bytes in total"),
));
}
let body = if regex {
source.clone()
} else {
regex::escape(source)
};
let expression = if match_case {
body
} else {
format!("(?i){body}")
};
match Regex::new(&expression) {
Ok(pattern) => compiled.push(pattern),
Err(error) => return Err((source.clone(), error.to_string())),
}
}
Ok(Self {
compiled,
sources: sources.to_vec(),
})
}
fn is_empty(&self) -> bool {
self.compiled.is_empty()
}
pub(crate) fn first_match(&self, haystack: &[u8]) -> Option<(usize, &str)> {
self.compiled
.iter()
.position(|pattern| pattern.is_match(haystack))
.map(|index| (index, self.sources[index].as_str()))
}
}
pub(crate) async fn wait_for_text(
pane: &Pane,
patterns: &Patterns,
stops: &Patterns,
timeout: Duration,
cancelled: &CancellationToken,
echo: EchoContext<'_>,
) -> Result<WaitView, Error> {
wait_for_text_with_limits(
pane,
patterns,
stops,
timeout,
cancelled,
ControlLimits::default(),
echo,
)
.await
}
pub(crate) async fn wait_for_text_with_limits(
pane: &Pane,
patterns: &Patterns,
stops: &Patterns,
timeout: Duration,
cancelled: &CancellationToken,
limits: ControlLimits,
echo: EchoContext<'_>,
) -> Result<WaitView, Error> {
let output = pane.stream_output_with_limits(limits).await?;
if let Some(view) = read_present_at_entry(pane, patterns, echo).await? {
let _ = output.shutdown().await;
return Ok(view);
}
wait_on_output(pane, output, patterns, stops, timeout, cancelled, echo).await
}
struct Screen {
above: Vec<u8>,
pending: Vec<u8>,
}
impl Screen {
async fn capture(pane: &Pane) -> Option<Self> {
let cursor_row: usize = pane
.format("#{cursor_y}")
.await
.ok()?
.to_string_lossy()
.trim()
.parse()
.ok()?;
let lines = pane.capture_with(CaptureOptions::visible()).await.ok()?;
let pending_row = cursor_row.min(lines.len().saturating_sub(1));
let mut above = Vec::new();
for line in lines.iter().take(pending_row) {
above.extend_from_slice(line.as_bytes());
above.push(b'\n');
}
let mut pending = lines
.get(pending_row)
.map_or_else(Vec::new, |line| line.as_bytes().to_vec());
pending.push(b'\n');
Some(Self { above, pending })
}
fn whole(&self) -> Vec<u8> {
let mut all = self.above.clone();
all.extend_from_slice(&self.pending);
all
}
}
async fn read_present_at_entry(
pane: &Pane,
patterns: &Patterns,
echo: EchoContext<'_>,
) -> Result<Option<WaitView>, Error> {
if patterns.is_empty() {
return Ok(None);
}
let Some(screen) = Screen::capture(pane).await else {
return Ok(None);
};
let recent = echo
.key
.map(|key| echo.echoes.snapshot(key))
.unwrap_or_default();
let masked_above = crate::echo::mask(&screen.above, &recent);
let pending_outcome = if echo.key.is_some_and(|key| echo.echoes.has_pending(key)) {
WaitOutcome::Pending
} else {
WaitOutcome::PresentAtEntry
};
let outcome = patterns
.first_match(&masked_above)
.map(|found| (WaitOutcome::PresentAtEntry, found))
.or_else(|| {
patterns
.first_match(&screen.pending)
.map(|found| (pending_outcome, found))
});
let Some((outcome, (index, source))) = outcome else {
return Ok(None);
};
let whole = screen.whole();
Ok(Some(WaitView {
pane: pane.id().to_string(),
outcome,
matched_index: Some(index),
matched_pattern: Some(source.to_owned()),
text: String::from_utf8_lossy(&whole).into_owned(),
bytes: whole.len(),
}))
}
async fn confirmed_above(
pane: &Pane,
patterns: &Patterns,
echo: EchoContext<'_>,
sticky: &mut Vec<Vec<u8>>,
) -> bool {
if let Some(key) = echo.key {
for line in echo.echoes.snapshot(key) {
if !sticky.contains(&line) {
sticky.push(line);
}
}
}
let Some(screen) = Screen::capture(pane).await else {
return false;
};
let mut haystack = crate::echo::mask(&screen.above, sticky);
if !echo.key.is_some_and(|key| echo.echoes.has_pending(key)) {
haystack.extend_from_slice(&screen.pending);
}
patterns.first_match(&haystack).is_some()
}
async fn wait_on_output(
pane: &Pane,
mut output: libtmux::control::PaneOutput,
patterns: &Patterns,
stops: &Patterns,
timeout: Duration,
cancelled: &CancellationToken,
echo: EchoContext<'_>,
) -> Result<WaitView, Error> {
let mut filter = TextFilter::new();
let mut text: Vec<u8> = Vec::new();
let mut bytes = 0usize;
let mut outcome = WaitOutcome::Deadline;
let mut matched_index = None;
let mut matched_pattern = None;
let mut sticky_echoes: Vec<Vec<u8>> = Vec::new();
let deadline = tokio::time::Instant::now() + timeout;
loop {
let chunk = tokio::select! {
biased;
() = cancelled.cancelled() => {
outcome = WaitOutcome::Cancelled;
break;
}
chunk = tokio::time::timeout_at(deadline, output.next_chunk()) => chunk,
};
match chunk {
Ok(Some(chunk)) => {
bytes = bytes.saturating_add(chunk.len());
filter.push(&chunk, &mut text);
if let Some((index, source)) = stops.first_match(&text) {
outcome = WaitOutcome::Stopped;
matched_index = Some(index);
matched_pattern = Some(source.to_owned());
break;
}
if patterns.is_empty() {
if !text.is_empty() {
outcome = WaitOutcome::Matched;
break;
}
} else if let Some((index, source)) = patterns.first_match(&text) {
let confirmed = confirmed_above(pane, patterns, echo, &mut sticky_echoes).await;
if confirmed {
outcome = WaitOutcome::Matched;
matched_index = Some(index);
matched_pattern = Some(source.to_owned());
break;
}
}
if text.len() > OUTPUT_LIMIT {
let excess = text.len() - OUTPUT_LIMIT;
text.drain(..excess);
}
}
Ok(None) => {
outcome = WaitOutcome::PaneClosed;
break;
}
Err(_) => break,
}
}
let pane_id = output.pane().to_string();
if let Err(error) = output.shutdown().await
&& !matches!(
error,
Error::ControlMode {
kind: ControlModeErrorKind::Closed,
..
}
)
{
return Err(error);
}
let (mut outcome, mut matched_index, mut matched_pattern) =
reconcile_deadline(outcome, matched_index, matched_pattern, patterns, &text);
if matches!(outcome, WaitOutcome::Matched)
&& !confirmed_above(pane, patterns, echo, &mut sticky_echoes).await
{
outcome = WaitOutcome::Deadline;
matched_index = None;
matched_pattern = None;
}
Ok(WaitView {
pane: pane_id,
outcome,
matched_index,
matched_pattern,
text: String::from_utf8_lossy(&text).into_owned(),
bytes,
})
}
fn reconcile_deadline(
outcome: WaitOutcome,
matched_index: Option<usize>,
matched_pattern: Option<String>,
patterns: &Patterns,
text: &[u8],
) -> (WaitOutcome, Option<usize>, Option<String>) {
if !matches!(outcome, WaitOutcome::Deadline) {
return (outcome, matched_index, matched_pattern);
}
match patterns.first_match(text) {
Some((index, source)) => (WaitOutcome::Matched, Some(index), Some(source.to_owned())),
None => (outcome, matched_index, matched_pattern),
}
}
#[cfg(test)]
mod tests;