use std::collections::VecDeque;
use std::sync::{Arc, Mutex};
use std::time::{SystemTime, UNIX_EPOCH};
use serde::{Deserialize, Serialize};
use thiserror::Error;
use tokio::sync::broadcast;
use super::os_process::PtyForegroundObservation;
pub const DEFAULT_PTY_REPLAY_BYTES: usize = 64 * 1024;
pub const PTY_PROVIDER_PROTOCOL_REVISION: &str = "gate4agent-pty-r1";
const DEFAULT_PTY_REPLAY_EVENTS: usize = 4_096;
const DEFAULT_PTY_SNAPSHOT_SCROLLBACK_ROWS: usize = 1_000;
pub const PTY_TERMINAL_SCROLLBACK_ROWS_MAX: usize = 256;
#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
pub struct PtySize {
pub rows: u16,
pub cols: u16,
}
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq, Serialize, Deserialize)]
#[serde(rename_all = "kebab-case")]
pub enum PtyMouseProtocolEncoding {
#[default]
Default,
Utf8,
Sgr,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
#[serde(rename_all = "kebab-case")]
pub enum PtySignal {
InterruptKey,
EndOfFileKey,
TerminateProcess,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
#[serde(rename_all = "kebab-case")]
pub enum PtySignalOutcome {
ControlWritten,
TerminationRequested,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
#[serde(rename_all = "kebab-case")]
pub enum PtyGapReason {
ReplayEvicted,
SubscriberLagged,
}
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
#[serde(tag = "kind", content = "payload", rename_all = "kebab-case")]
pub enum PtyEvent {
Started,
Output(Vec<u8>),
DataGap {
from_sequence: u64,
to_sequence: u64,
reason: PtyGapReason,
},
Resized(PtySize),
ForegroundProcess(PtyForegroundObservation),
SnapshotAvailable {
snapshot_sequence: u64,
},
ReaderError {
message: String,
},
OperatorActionRequired {
message: String,
},
Exited {
code: i32,
},
}
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
pub struct PtyEventEnvelope {
pub pty_id: String,
pub provider_revision: String,
pub generation: u64,
pub sequence: u64,
pub event: PtyEvent,
}
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
pub struct PtyReplayCursor {
pub provider_revision: String,
pub generation: u64,
pub next_sequence: u64,
}
impl PtyReplayCursor {
pub fn beginning(provider_revision: impl Into<String>, generation: u64) -> Self {
Self {
provider_revision: provider_revision.into(),
generation,
next_sequence: 1,
}
}
}
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
pub struct PtyTerminalSnapshot {
pub pty_id: String,
pub provider_revision: String,
pub generation: u64,
pub sequence: u64,
pub size: PtySize,
pub cursor: (u16, u16),
#[serde(default)]
pub bracketed_paste: bool,
pub contents: String,
pub formatted: Vec<u8>,
#[serde(default)]
pub scrollback_formatted: Vec<Vec<u8>>,
#[serde(default)]
pub alternate_screen: bool,
#[serde(default)]
pub mouse_protocol_enabled: bool,
#[serde(default)]
pub mouse_protocol_encoding: PtyMouseProtocolEncoding,
#[serde(default)]
pub produced_at_unix_ms: u64,
}
pub(crate) fn now_unix_ms() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|duration| duration.as_millis().min(u128::from(u64::MAX)) as u64)
.unwrap_or(0)
}
#[derive(Debug, Error)]
pub enum PtyAttachError {
#[error("requested PTY generation {requested}, but the live generation is {actual}")]
StaleGeneration { requested: u64, actual: u64 },
#[error("requested PTY provider revision '{requested}', but the live revision is '{actual}'")]
ProviderRevisionMismatch { requested: String, actual: String },
#[error("PTY replay cursor sequence starts at 1")]
InvalidCursor,
#[error("PTY event journal mutex poisoned")]
JournalPoisoned,
}
#[derive(Debug, Error)]
pub enum PtyEventRecvError {
#[error("PTY event stream closed")]
Closed,
}
pub struct PtyAttachment {
pub replay: Vec<PtyEventEnvelope>,
pub receiver: PtyEventReceiver,
}
pub struct PtyEventReceiver {
rx: broadcast::Receiver<PtyEventEnvelope>,
pty_id: String,
provider_revision: String,
generation: u64,
expected_sequence: u64,
pending: Option<PtyEventEnvelope>,
}
impl PtyEventReceiver {
pub fn try_recv(&mut self) -> Result<Option<PtyEventEnvelope>, PtyEventRecvError> {
if let Some(event) = self.pending.take() {
self.expected_sequence = event.sequence.saturating_add(1);
return Ok(Some(event));
}
loop {
match self.rx.try_recv() {
Ok(event) if event.generation != self.generation => continue,
Ok(event) if event.sequence < self.expected_sequence => continue,
Ok(event) if event.sequence > self.expected_sequence => {
let gap = PtyEventEnvelope {
pty_id: self.pty_id.clone(),
provider_revision: self.provider_revision.clone(),
generation: self.generation,
sequence: event.sequence - 1,
event: PtyEvent::DataGap {
from_sequence: self.expected_sequence,
to_sequence: event.sequence - 1,
reason: PtyGapReason::SubscriberLagged,
},
};
self.pending = Some(event);
return Ok(Some(gap));
}
Ok(event) => {
self.expected_sequence = event.sequence.saturating_add(1);
return Ok(Some(event));
}
Err(broadcast::error::TryRecvError::Lagged(_)) => continue,
Err(broadcast::error::TryRecvError::Empty) => return Ok(None),
Err(broadcast::error::TryRecvError::Closed) => {
return Err(PtyEventRecvError::Closed);
}
}
}
}
pub async fn recv(&mut self) -> Result<PtyEventEnvelope, PtyEventRecvError> {
if let Some(event) = self.pending.take() {
self.expected_sequence = event.sequence.saturating_add(1);
return Ok(event);
}
loop {
match self.rx.recv().await {
Ok(event) if event.generation != self.generation => continue,
Ok(event) if event.sequence < self.expected_sequence => continue,
Ok(event) if event.sequence > self.expected_sequence => {
let gap = PtyEventEnvelope {
pty_id: self.pty_id.clone(),
provider_revision: self.provider_revision.clone(),
generation: self.generation,
sequence: event.sequence - 1,
event: PtyEvent::DataGap {
from_sequence: self.expected_sequence,
to_sequence: event.sequence - 1,
reason: PtyGapReason::SubscriberLagged,
},
};
self.pending = Some(event);
return Ok(gap);
}
Ok(event) => {
self.expected_sequence = event.sequence.saturating_add(1);
return Ok(event);
}
Err(broadcast::error::RecvError::Lagged(_)) => {
}
Err(broadcast::error::RecvError::Closed) => {
return Err(PtyEventRecvError::Closed);
}
}
}
}
}
pub(crate) struct PtyEventPublisher {
pty_id: String,
provider_revision: String,
generation: u64,
tx: broadcast::Sender<PtyEventEnvelope>,
state: Mutex<PtyEventState>,
}
struct PtyEventState {
sequence: u64,
replay_bytes: usize,
replay_byte_limit: usize,
replay_event_limit: usize,
replay: VecDeque<PtyEventEnvelope>,
terminal: vt100::Parser,
}
impl PtyEventPublisher {
pub(crate) fn new(
pty_id: String,
provider_revision: String,
generation: u64,
channel_capacity: usize,
replay_byte_limit: usize,
rows: u16,
cols: u16,
) -> Arc<Self> {
let (tx, _) = broadcast::channel(channel_capacity);
Arc::new(Self {
pty_id,
provider_revision,
generation,
tx,
state: Mutex::new(PtyEventState {
sequence: 0,
replay_bytes: 0,
replay_byte_limit,
replay_event_limit: DEFAULT_PTY_REPLAY_EVENTS,
replay: VecDeque::new(),
terminal: vt100::Parser::new(rows, cols, DEFAULT_PTY_SNAPSHOT_SCROLLBACK_ROWS),
}),
})
}
pub(crate) fn publish(&self, event: PtyEvent) {
let Ok(mut state) = self.state.lock() else {
return;
};
state.observe_terminal(&event);
state.sequence = state.sequence.saturating_add(1);
let envelope = PtyEventEnvelope {
pty_id: self.pty_id.clone(),
provider_revision: self.provider_revision.clone(),
generation: self.generation,
sequence: state.sequence,
event,
};
state.push(envelope.clone());
let _ = self.tx.send(envelope);
}
pub(crate) fn generation(&self) -> u64 {
self.generation
}
pub(crate) fn provider_revision(&self) -> &str {
&self.provider_revision
}
pub(crate) fn subscribe(&self) -> Result<PtyEventReceiver, PtyAttachError> {
let rx = self.tx.subscribe();
let state = self
.state
.lock()
.map_err(|_| PtyAttachError::JournalPoisoned)?;
Ok(PtyEventReceiver {
rx,
pty_id: self.pty_id.clone(),
provider_revision: self.provider_revision.clone(),
generation: self.generation,
expected_sequence: state.sequence.saturating_add(1),
pending: None,
})
}
pub(crate) fn retained_cursor(&self) -> Result<PtyReplayCursor, PtyAttachError> {
let state = self
.state
.lock()
.map_err(|_| PtyAttachError::JournalPoisoned)?;
Ok(PtyReplayCursor {
provider_revision: self.provider_revision.clone(),
generation: self.generation,
next_sequence: state
.replay
.front()
.map_or_else(|| state.sequence.saturating_add(1), |event| event.sequence),
})
}
pub(crate) fn attach_retained(&self) -> Result<PtyAttachment, PtyAttachError> {
let rx = self.tx.subscribe();
let state = self
.state
.lock()
.map_err(|_| PtyAttachError::JournalPoisoned)?;
let snapshot_sequence = state.sequence;
Ok(PtyAttachment {
replay: state.replay.iter().cloned().collect(),
receiver: PtyEventReceiver {
rx,
pty_id: self.pty_id.clone(),
provider_revision: self.provider_revision.clone(),
generation: self.generation,
expected_sequence: snapshot_sequence.saturating_add(1),
pending: None,
},
})
}
pub(crate) fn attach(&self, cursor: PtyReplayCursor) -> Result<PtyAttachment, PtyAttachError> {
if cursor.next_sequence == 0 {
return Err(PtyAttachError::InvalidCursor);
}
if cursor.provider_revision != self.provider_revision {
return Err(PtyAttachError::ProviderRevisionMismatch {
requested: cursor.provider_revision,
actual: self.provider_revision.clone(),
});
}
if cursor.generation != self.generation {
return Err(PtyAttachError::StaleGeneration {
requested: cursor.generation,
actual: self.generation,
});
}
let rx = self.tx.subscribe();
let state = self
.state
.lock()
.map_err(|_| PtyAttachError::JournalPoisoned)?;
let snapshot_sequence = state.sequence;
let mut replay = Vec::new();
let first_retained = state.replay.front().map(|event| event.sequence);
if cursor.next_sequence <= snapshot_sequence {
match first_retained {
Some(first) if cursor.next_sequence < first => {
replay.push(PtyEventEnvelope {
pty_id: self.pty_id.clone(),
provider_revision: self.provider_revision.clone(),
generation: self.generation,
sequence: first - 1,
event: PtyEvent::DataGap {
from_sequence: cursor.next_sequence,
to_sequence: first - 1,
reason: PtyGapReason::ReplayEvicted,
},
});
}
None => {
replay.push(PtyEventEnvelope {
pty_id: self.pty_id.clone(),
provider_revision: self.provider_revision.clone(),
generation: self.generation,
sequence: snapshot_sequence,
event: PtyEvent::DataGap {
from_sequence: cursor.next_sequence,
to_sequence: snapshot_sequence,
reason: PtyGapReason::ReplayEvicted,
},
});
}
_ => {}
}
replay.extend(
state
.replay
.iter()
.filter(|event| event.sequence >= cursor.next_sequence)
.cloned(),
);
}
Ok(PtyAttachment {
replay,
receiver: PtyEventReceiver {
rx,
pty_id: self.pty_id.clone(),
provider_revision: self.provider_revision.clone(),
generation: self.generation,
expected_sequence: snapshot_sequence.saturating_add(1),
pending: None,
},
})
}
pub(crate) fn terminal_sequence(&self) -> Result<u64, PtyAttachError> {
let state = self
.state
.lock()
.map_err(|_| PtyAttachError::JournalPoisoned)?;
Ok(state.sequence)
}
pub(crate) fn snapshot(&self) -> Result<PtyTerminalSnapshot, PtyAttachError> {
let state = self
.state
.lock()
.map_err(|_| PtyAttachError::JournalPoisoned)?;
let screen = state.terminal.screen();
let scrollback_formatted = recent_scrollback_formatted(screen);
let (rows, cols) = screen.size();
Ok(PtyTerminalSnapshot {
pty_id: self.pty_id.clone(),
provider_revision: self.provider_revision.clone(),
generation: self.generation,
sequence: state.sequence,
size: PtySize { rows, cols },
cursor: screen.cursor_position(),
bracketed_paste: screen.bracketed_paste(),
contents: screen.contents(),
formatted: screen.contents_formatted(),
scrollback_formatted,
alternate_screen: screen.alternate_screen(),
mouse_protocol_enabled: screen.mouse_protocol_mode() != vt100::MouseProtocolMode::None,
mouse_protocol_encoding: match screen.mouse_protocol_encoding() {
vt100::MouseProtocolEncoding::Default => PtyMouseProtocolEncoding::Default,
vt100::MouseProtocolEncoding::Utf8 => PtyMouseProtocolEncoding::Utf8,
vt100::MouseProtocolEncoding::Sgr => PtyMouseProtocolEncoding::Sgr,
},
produced_at_unix_ms: now_unix_ms(),
})
}
}
pub(super) fn recent_scrollback_formatted(screen: &vt100::Screen) -> Vec<Vec<u8>> {
if screen.alternate_screen() {
return Vec::new();
}
let mut view = screen.clone();
view.set_scrollback(usize::MAX);
let available = view.scrollback().min(PTY_TERMINAL_SCROLLBACK_ROWS_MAX);
let columns = view.size().1;
let mut rows = Vec::with_capacity(available);
for offset in (1..=available).rev() {
view.set_scrollback(offset);
if let Some(row) = view.rows_formatted(0, columns).next() {
rows.push(row);
}
}
rows
}
impl PtyEventState {
fn observe_terminal(&mut self, event: &PtyEvent) {
match event {
PtyEvent::Output(data) => self.terminal.process(data),
PtyEvent::Resized(size) => self.terminal.screen_mut().set_size(size.rows, size.cols),
PtyEvent::Started
| PtyEvent::DataGap { .. }
| PtyEvent::ReaderError { .. }
| PtyEvent::OperatorActionRequired { .. }
| PtyEvent::ForegroundProcess(_)
| PtyEvent::SnapshotAvailable { .. }
| PtyEvent::Exited { .. } => {}
}
}
fn push(&mut self, event: PtyEventEnvelope) {
self.replay_bytes = self.replay_bytes.saturating_add(event_size(&event));
self.replay.push_back(event);
while self.replay.len() > self.replay_event_limit
|| self.replay_bytes > self.replay_byte_limit
{
let Some(evicted) = self.replay.pop_front() else {
break;
};
self.replay_bytes = self.replay_bytes.saturating_sub(event_size(&evicted));
}
}
}
fn event_size(event: &PtyEventEnvelope) -> usize {
match &event.event {
PtyEvent::Output(data) => data.len(),
_ => 1,
}
}
#[cfg(test)]
mod tests {
use super::*;
const REVISION: &str = "test-provider:1";
#[test]
fn replay_reports_evicted_sequences() {
let publisher =
PtyEventPublisher::new("pty-test".to_owned(), REVISION.to_owned(), 1, 8, 3, 24, 80);
publisher.publish(PtyEvent::Output(vec![1, 2]));
publisher.publish(PtyEvent::Output(vec![3, 4]));
let attachment = publisher
.attach(PtyReplayCursor::beginning(REVISION, 1))
.expect("attach");
assert!(matches!(
attachment.replay.first().map(|event| &event.event),
Some(PtyEvent::DataGap {
from_sequence: 1,
to_sequence: 1,
reason: PtyGapReason::ReplayEvicted,
})
));
assert_eq!(attachment.replay.last().unwrap().sequence, 2);
}
#[test]
fn retained_cursor_replays_the_available_tail_without_a_synthetic_gap() {
let publisher =
PtyEventPublisher::new("pty-test".to_owned(), REVISION.to_owned(), 1, 8, 3, 24, 80);
publisher.publish(PtyEvent::Output(vec![1, 2]));
publisher.publish(PtyEvent::Output(vec![3, 4]));
let cursor = publisher.retained_cursor().expect("retained cursor");
assert_eq!(cursor.next_sequence, 2);
let attachment = publisher.attach(cursor).expect("tail attach");
assert_eq!(attachment.replay.len(), 1);
assert_eq!(attachment.replay[0].sequence, 2);
assert!(!matches!(attachment.replay[0].event, PtyEvent::DataGap { .. }));
}
#[test]
fn retained_attachment_atomically_replays_available_events_without_a_gap() {
let publisher =
PtyEventPublisher::new("pty-test".to_owned(), REVISION.to_owned(), 1, 8, 3, 24, 80);
publisher.publish(PtyEvent::Output(vec![1, 2]));
publisher.publish(PtyEvent::Output(vec![3, 4]));
let attachment = publisher.attach_retained().expect("retained attach");
assert_eq!(attachment.replay.len(), 1);
assert_eq!(attachment.replay[0].sequence, 2);
assert!(!matches!(attachment.replay[0].event, PtyEvent::DataGap { .. }));
}
#[tokio::test]
async fn receiver_turns_broadcast_lag_into_an_exact_gap() {
let publisher =
PtyEventPublisher::new("pty-test".to_owned(), REVISION.to_owned(), 1, 2, 64, 24, 80);
let mut receiver = publisher.subscribe().expect("subscribe");
for value in 0..4 {
publisher.publish(PtyEvent::Output(vec![value]));
}
let gap = receiver.recv().await.expect("gap");
assert!(matches!(
gap.event,
PtyEvent::DataGap {
from_sequence: 1,
to_sequence: 2,
reason: PtyGapReason::SubscriberLagged,
}
));
assert_eq!(receiver.recv().await.expect("first retained").sequence, 3);
}
#[test]
fn nonblocking_receiver_preserves_exact_gap_and_pending_event() {
let publisher =
PtyEventPublisher::new("pty-test".to_owned(), REVISION.to_owned(), 1, 2, 64, 24, 80);
let mut receiver = publisher.subscribe().expect("subscribe");
for value in 0..4 {
publisher.publish(PtyEvent::Output(vec![value]));
}
let gap = receiver.try_recv().unwrap().unwrap();
assert!(matches!(
gap.event,
PtyEvent::DataGap {
from_sequence: 1,
to_sequence: 2,
reason: PtyGapReason::SubscriberLagged,
}
));
assert_eq!(receiver.try_recv().unwrap().unwrap().sequence, 3);
assert_eq!(receiver.try_recv().unwrap().unwrap().sequence, 4);
assert!(receiver.try_recv().unwrap().is_none());
}
#[test]
fn snapshot_is_pinned_to_the_last_incorporated_sequence() {
let publisher =
PtyEventPublisher::new("pty-test".to_owned(), REVISION.to_owned(), 1, 8, 64, 2, 10);
publisher.publish(PtyEvent::Output(b"\x1b[?2004hhello".to_vec()));
let snapshot = publisher.snapshot().expect("snapshot");
assert_eq!(snapshot.sequence, 1);
assert_eq!(snapshot.contents, "hello");
assert_eq!(snapshot.size, PtySize { rows: 2, cols: 10 });
assert_eq!(snapshot.provider_revision, REVISION);
assert!(snapshot.bracketed_paste);
}
#[test]
fn snapshot_exports_only_the_recent_primary_scrollback_in_chronological_order() {
let publisher =
PtyEventPublisher::new("pty-test".to_owned(), REVISION.to_owned(), 1, 8, 64, 2, 16);
for line in 0..300 {
publisher.publish(PtyEvent::Output(format!("{line:03}\r\n").into_bytes()));
}
let snapshot = publisher.snapshot().expect("snapshot");
assert_eq!(snapshot.scrollback_formatted.len(), PTY_TERMINAL_SCROLLBACK_ROWS_MAX);
let row_text = |formatted: &[u8]| {
let mut parser = vt100::Parser::new(1, 16, 0);
parser.process(formatted);
parser.screen().contents()
};
assert_eq!(row_text(&snapshot.scrollback_formatted[0]), "043");
assert_eq!(row_text(snapshot.scrollback_formatted.last().unwrap()), "298");
}
#[test]
fn snapshot_reports_alternate_screen_and_mouse_protocol_metadata() {
let publisher =
PtyEventPublisher::new("pty-test".to_owned(), REVISION.to_owned(), 1, 8, 64, 2, 16);
publisher.publish(PtyEvent::Output(
b"before\r\nafter\r\n\x1b[?1000h\x1b[?1006h\x1b[?1049halternate".to_vec(),
));
let snapshot = publisher.snapshot().expect("snapshot");
assert!(snapshot.alternate_screen);
assert!(snapshot.mouse_protocol_enabled);
assert_eq!(snapshot.mouse_protocol_encoding, PtyMouseProtocolEncoding::Sgr);
assert!(snapshot.scrollback_formatted.is_empty());
}
#[test]
fn attach_rejects_a_different_provider_revision() {
let publisher =
PtyEventPublisher::new("pty-test".to_owned(), REVISION.to_owned(), 1, 8, 64, 2, 10);
let error = match publisher.attach(PtyReplayCursor::beginning("other-provider:1", 1)) {
Err(error) => error,
Ok(_) => panic!("revision mismatch must reject attach"),
};
assert!(matches!(
error,
PtyAttachError::ProviderRevisionMismatch { .. }
));
}
}