use std::collections::{HashMap, HashSet};
use std::fs::{self, File, OpenOptions};
use std::io::{BufRead, BufReader, Write};
use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex};
use base64::{engine::general_purpose::URL_SAFE_NO_PAD as BASE64URL, Engine};
use serde::{Deserialize, Serialize};
use crate::terminal_cursor::CursorModel;
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct CommandEntry {
pub seq: u64,
pub command: String,
pub output: String,
#[serde(rename = "exitCode")]
pub exit_code: Option<i32>,
#[serde(rename = "startedAt")]
pub started_at: i64,
#[serde(rename = "endedAt")]
pub ended_at: i64,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct RawEntry {
pub seq: u64,
pub raw: String,
pub ts: i64,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(untagged)]
pub enum HistoryEntry {
Command(CommandEntry),
Raw(RawEntry),
}
#[derive(Debug, Clone, PartialEq)]
pub enum PendingEntry {
Command {
command: String,
output: String,
exit_code: Option<i32>,
started_at: i64,
ended_at: i64,
},
Raw {
raw: String,
ts: i64,
},
}
impl PendingEntry {
fn into_entry(self, seq: u64) -> HistoryEntry {
match self {
PendingEntry::Command {
command,
output,
exit_code,
started_at,
ended_at,
} => HistoryEntry::Command(CommandEntry {
seq,
command,
output,
exit_code,
started_at,
ended_at,
}),
PendingEntry::Raw { raw, ts } => HistoryEntry::Raw(RawEntry { seq, raw, ts }),
}
}
}
pub fn now_ms() -> i64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as i64)
.unwrap_or(0)
}
const RAW_FLUSH_THRESHOLD: usize = 4096;
const MAX_COMMAND_OUTPUT_BYTES: usize = 256 * 1024;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Marker {
A,
B,
C,
D(Option<i32>),
}
struct OpenCommand {
command: String,
output: Vec<u8>,
started_at: i64,
}
impl OpenCommand {
fn finish(self, exit_code: Option<i32>, ended_at: i64) -> PendingEntry {
PendingEntry::Command {
command: self.command,
output: String::from_utf8_lossy(&self.output).into_owned(),
exit_code,
started_at: self.started_at,
ended_at,
}
}
}
pub struct Segmenter {
scan_buf: Vec<u8>,
pending: Vec<u8>,
open: Option<OpenCommand>,
cursor: CursorModel,
seen_marker: bool,
}
impl Default for Segmenter {
fn default() -> Self {
Self::new()
}
}
impl Segmenter {
pub fn new() -> Self {
Self {
scan_buf: Vec::new(),
pending: Vec::new(),
open: None,
cursor: CursorModel::new(),
seen_marker: false,
}
}
pub fn feed(&mut self, chunk: &[u8], now_ms: i64) -> Vec<PendingEntry> {
self.scan_buf.extend_from_slice(chunk);
let mut events = Vec::new();
loop {
match self.scan_buf.iter().position(|&b| b == 0x1b) {
None => {
let n = self.scan_buf.len();
if n > 0 {
self.consume_plain(n, now_ms, &mut events);
}
break;
}
Some(0) => match try_match_marker(&self.scan_buf) {
MarkerMatch::Complete { markers, len } => {
self.scan_buf.drain(0..len);
for m in markers {
self.handle_marker(m, now_ms, &mut events);
}
}
MarkerMatch::Incomplete => break,
MarkerMatch::NotAMarker => {
self.consume_plain(1, now_ms, &mut events);
}
},
Some(pos) => self.consume_plain(pos, now_ms, &mut events),
}
}
events
}
pub fn flush(&mut self, now_ms: i64) -> Option<PendingEntry> {
if let Some(open) = self.open.take() {
return Some(open.finish(None, now_ms));
}
if self.seen_marker {
self.pending.clear();
return None;
}
if !self.pending.is_empty() {
let raw = String::from_utf8_lossy(&self.pending).into_owned();
self.pending.clear();
return Some(PendingEntry::Raw { raw, ts: now_ms });
}
None
}
fn consume_plain(&mut self, n: usize, now_ms: i64, events: &mut Vec<PendingEntry>) {
let bytes: Vec<u8> = self.scan_buf.drain(0..n).collect();
self.cursor.feed(&bytes);
if let Some(open) = &mut self.open {
if open.output.len() < MAX_COMMAND_OUTPUT_BYTES {
let room = MAX_COMMAND_OUTPUT_BYTES - open.output.len();
let take = room.min(bytes.len());
open.output.extend_from_slice(&bytes[..take]);
}
return;
}
self.pending.extend_from_slice(&bytes);
while self.pending.len() >= RAW_FLUSH_THRESHOLD {
let chunk: Vec<u8> = self.pending.drain(0..RAW_FLUSH_THRESHOLD).collect();
events.push(PendingEntry::Raw {
raw: String::from_utf8_lossy(&chunk).into_owned(),
ts: now_ms,
});
}
}
fn handle_marker(&mut self, marker: Marker, now_ms: i64, events: &mut Vec<PendingEntry>) {
self.seen_marker = true;
match marker {
Marker::C => {
if let Some(open) = self.open.take() {
events.push(open.finish(None, now_ms));
}
let command = self.cursor.take_command_line();
self.pending.clear();
self.open = Some(OpenCommand {
command,
output: Vec::new(),
started_at: now_ms,
});
}
Marker::D(code) => {
if let Some(open) = self.open.take() {
events.push(open.finish(code, now_ms));
}
self.pending.clear();
self.cursor.mark_output_boundary();
}
Marker::A | Marker::B => {
}
}
}
}
enum MarkerMatch {
Complete { markers: Vec<Marker>, len: usize },
Incomplete,
NotAMarker,
}
fn try_match_marker(buf: &[u8]) -> MarkerMatch {
debug_assert_eq!(buf.first(), Some(&0x1b));
if buf.len() < 2 {
return MarkerMatch::Incomplete;
}
match buf[1] {
b']' => try_match_bare(buf),
b'P' => try_match_wrapped(buf),
_ => MarkerMatch::NotAMarker,
}
}
fn build_marker(kind: u8, code: Option<i32>) -> Option<Marker> {
match kind {
b'A' => Some(Marker::A),
b'B' => Some(Marker::B),
b'C' => Some(Marker::C),
b'D' => Some(Marker::D(code)),
_ => None,
}
}
fn try_match_bare(buf: &[u8]) -> MarkerMatch {
const PREFIX: &[u8] = b"\x1b]133;";
if buf.len() < PREFIX.len() {
return if PREFIX.starts_with(buf) {
MarkerMatch::Incomplete
} else {
MarkerMatch::NotAMarker
};
}
if &buf[..PREFIX.len()] != PREFIX {
return MarkerMatch::NotAMarker;
}
let mut i = PREFIX.len();
if i >= buf.len() {
return MarkerMatch::Incomplete;
}
let kind = buf[i];
i += 1;
let mut code: Option<i32> = None;
if i < buf.len() && buf[i] == b';' {
let start = i + 1;
let mut j = start;
while j < buf.len() && buf[j].is_ascii_digit() {
j += 1;
}
if j == buf.len() {
return MarkerMatch::Incomplete;
}
if j > start {
code = std::str::from_utf8(&buf[start..j])
.ok()
.and_then(|s| s.parse().ok());
}
i = j;
}
if i >= buf.len() {
return MarkerMatch::Incomplete;
}
let (terminator_len, ok) = match buf[i] {
0x07 => (1, true),
0x1b => {
if i + 1 >= buf.len() {
return MarkerMatch::Incomplete;
}
(2, buf[i + 1] == b'\\')
}
_ => (0, false),
};
if !ok {
return MarkerMatch::NotAMarker;
}
match build_marker(kind, code) {
Some(marker) => MarkerMatch::Complete {
markers: vec![marker],
len: i + terminator_len,
},
None => MarkerMatch::NotAMarker,
}
}
fn try_match_wrapped(buf: &[u8]) -> MarkerMatch {
const PREFIX: &[u8] = b"\x1bPtmux;";
if buf.len() < PREFIX.len() {
return if PREFIX.starts_with(buf) {
MarkerMatch::Incomplete
} else {
MarkerMatch::NotAMarker
};
}
if &buf[..PREFIX.len()] != PREFIX {
return MarkerMatch::NotAMarker;
}
let mut payload = Vec::new();
let mut i = PREFIX.len();
loop {
if i >= buf.len() {
return MarkerMatch::Incomplete;
}
if buf[i] == 0x1b {
if i + 1 >= buf.len() {
return MarkerMatch::Incomplete;
}
match buf[i + 1] {
0x1b => {
payload.push(0x1b);
i += 2;
}
b'\\' => {
i += 2;
break;
}
_ => return MarkerMatch::NotAMarker,
}
} else {
payload.push(buf[i]);
i += 1;
}
}
let markers = parse_all_bare_markers(&payload);
if markers.is_empty() {
return MarkerMatch::NotAMarker;
}
MarkerMatch::Complete { markers, len: i }
}
fn parse_all_bare_markers(payload: &[u8]) -> Vec<Marker> {
let mut markers = Vec::new();
let mut pos = 0;
while pos < payload.len() {
if payload[pos] == 0x1b {
match try_match_bare(&payload[pos..]) {
MarkerMatch::Complete {
markers: mut m,
len,
} => {
markers.append(&mut m);
pos += len;
continue;
}
_ => {
pos += 1;
continue;
}
}
}
pos += 1;
}
markers
}
pub const MAX_ENTRIES_PER_SESSION: u64 = 2000;
const TRIM_MARGIN: u64 = 100;
pub const DEFAULT_LIMIT: usize = 50;
pub const MAX_LIMIT: usize = 500;
struct SessionSlot {
last_seq: u64,
entry_count: u64,
}
pub struct SessionHistoryStore {
root: PathBuf,
slots: Mutex<HashMap<String, Arc<Mutex<SessionSlot>>>>,
active_feeders: Mutex<HashSet<String>>,
}
pub struct FeederGuard {
store: Arc<SessionHistoryStore>,
session: String,
}
impl Drop for FeederGuard {
fn drop(&mut self) {
if let Ok(mut active) = self.store.active_feeders.lock() {
active.remove(&self.session);
}
}
}
impl SessionHistoryStore {
pub fn new(data_dir: &Path) -> Self {
Self {
root: data_dir.join("history"),
slots: Mutex::new(HashMap::new()),
active_feeders: Mutex::new(HashSet::new()),
}
}
pub fn try_acquire_feeder(self: &Arc<Self>, session: &str) -> Option<FeederGuard> {
let mut active = self.active_feeders.lock().ok()?;
if active.contains(session) {
return None;
}
active.insert(session.to_string());
Some(FeederGuard {
store: self.clone(),
session: session.to_string(),
})
}
fn file_path(&self, session: &str) -> PathBuf {
self.root.join(format!("{session}.jsonl"))
}
fn slot(&self, session: &str) -> anyhow::Result<Arc<Mutex<SessionSlot>>> {
let mut slots = self
.slots
.lock()
.map_err(|_| anyhow::anyhow!("session history slots lock poisoned"))?;
if let Some(existing) = slots.get(session) {
return Ok(existing.clone());
}
let hydrated = Self::hydrate(&self.file_path(session));
let arc = Arc::new(Mutex::new(hydrated));
slots.insert(session.to_string(), arc.clone());
Ok(arc)
}
fn hydrate(path: &Path) -> SessionSlot {
let mut last_seq = 0u64;
let mut entry_count = 0u64;
if let Ok(f) = File::open(path) {
for line in BufReader::new(f).lines().map_while(Result::ok) {
if line.trim().is_empty() {
continue;
}
entry_count += 1;
if let Ok(v) = serde_json::from_str::<serde_json::Value>(&line) {
if let Some(seq) = v.get("seq").and_then(|s| s.as_u64()) {
last_seq = last_seq.max(seq);
}
}
}
}
SessionSlot {
last_seq,
entry_count,
}
}
pub fn append(&self, session: &str, entry: PendingEntry) -> anyhow::Result<HistoryEntry> {
let slot = self.slot(session)?;
let mut s = slot
.lock()
.map_err(|_| anyhow::anyhow!("session history slot lock poisoned"))?;
s.last_seq += 1;
let full = entry.into_entry(s.last_seq);
fs::create_dir_all(&self.root)?;
let line = serde_json::to_string(&full)?;
let mut f = OpenOptions::new()
.create(true)
.append(true)
.open(self.file_path(session))?;
writeln!(f, "{line}")?;
s.entry_count += 1;
if s.entry_count > MAX_ENTRIES_PER_SESSION + TRIM_MARGIN {
self.trim(session, &mut s)?;
}
Ok(full)
}
fn trim(&self, session: &str, s: &mut SessionSlot) -> anyhow::Result<()> {
let path = self.file_path(session);
let lines: Vec<String> = BufReader::new(File::open(&path)?)
.lines()
.map_while(Result::ok)
.filter(|l| !l.trim().is_empty())
.collect();
let keep_from = lines.len().saturating_sub(MAX_ENTRIES_PER_SESSION as usize);
let kept = &lines[keep_from..];
let tmp_path = path.with_extension("jsonl.tmp");
{
let mut tmp = File::create(&tmp_path)?;
for l in kept {
writeln!(tmp, "{l}")?;
}
}
fs::rename(&tmp_path, &path)?;
s.entry_count = kept.len() as u64;
Ok(())
}
pub fn read_page(
&self,
session: &str,
cursor: Option<u64>,
limit: usize,
) -> anyhow::Result<(Vec<serde_json::Value>, u64)> {
let path = self.file_path(session);
let floor = cursor.unwrap_or(0);
let mut entries = Vec::new();
let mut last_seq = floor;
if let Ok(f) = File::open(&path) {
for line in BufReader::new(f).lines().map_while(Result::ok) {
if line.trim().is_empty() {
continue;
}
let Ok(v) = serde_json::from_str::<serde_json::Value>(&line) else {
continue;
};
let seq = v.get("seq").and_then(|s| s.as_u64()).unwrap_or(0);
if seq <= floor {
continue;
}
last_seq = seq;
entries.push(v);
if entries.len() >= limit {
break;
}
}
}
Ok((entries, last_seq))
}
}
pub fn encode_cursor(seq: u64) -> String {
BASE64URL.encode(format!("v1:{seq}"))
}
pub fn decode_cursor(cursor: &str) -> Option<u64> {
let bytes = BASE64URL.decode(cursor).ok()?;
let text = String::from_utf8(bytes).ok()?;
text.strip_prefix("v1:")?.parse::<u64>().ok()
}
#[cfg(test)]
mod tests {
use super::*;
fn feed_str(seg: &mut Segmenter, s: &str, t: i64) -> Vec<PendingEntry> {
seg.feed(s.as_bytes(), t)
}
const REAL_ECHO_HELLO: &str =
"echo hello\r\n\x1b]133;C\x07hello\r\n\x1b[?2004h\x1b]133;D;0\x07\x1b]133;A\x07mvhenten@sandbox:~$ ";
#[test]
fn real_command_cycle_produces_one_command_entry() {
let mut seg = Segmenter::new();
let events = feed_str(&mut seg, REAL_ECHO_HELLO, 1000);
assert_eq!(events.len(), 1);
match &events[0] {
PendingEntry::Command {
command,
output,
exit_code,
started_at,
ended_at,
} => {
assert_eq!(command, "echo hello");
assert_eq!(output, "hello\r\n\x1b[?2004h");
assert_eq!(*exit_code, Some(0));
assert_eq!(*started_at, 1000);
assert_eq!(*ended_at, 1000);
}
other => panic!("expected a command entry, got {other:?}"),
}
}
#[test]
fn command_text_keeps_the_prompt_prefix_like_the_client_reader_does() {
let mut seg = Segmenter::new();
let events = feed_str(
&mut seg,
"mvhenten@host:~$ ls -la\r\n\x1b]133;C\x07f\r\n\x1b]133;D;0\x07\x1b]133;A\x07",
1,
);
let PendingEntry::Command { command, .. } = &events[0] else {
panic!("expected command entry");
};
assert_eq!(command, "mvhenten@host:~$ ls -la");
}
#[test]
fn command_survives_a_leading_unmarked_repaint_burst() {
let repaint_burst = "\x1b[?1049h\x1b[?1h\x1b=\x1b[H\x1b[2J\x1b[?25h\r\n\
Welcome to Ubuntu — motd banner line one\r\n\
Last login: Tue Jul 21 07:00:00 2026\r\n";
let mut seg = Segmenter::new();
let events = feed_str(
&mut seg,
&format!(
"{repaint_burst}mvhenten@sandbox:~$ echo hello\r\n\x1b]133;C\x07hello\r\n\x1b]133;D;0\x07\x1b]133;A\x07"
),
1,
);
assert_eq!(events.len(), 1);
let PendingEntry::Command {
command,
output,
exit_code,
..
} = &events[0]
else {
panic!("expected a command entry, got {events:?}");
};
assert_eq!(command, "mvhenten@sandbox:~$ echo hello");
assert_eq!(output, "hello\r\n");
assert_eq!(*exit_code, Some(0));
}
#[test]
fn command_survives_a_status_bar_redraw_between_echo_and_c() {
let mut seg = Segmenter::new();
let events = feed_str(
&mut seg,
"mvhenten@sandbox:~$ echo hi\r\n\
\x1b[24;1H\x1b[K[ 0:bash* ]\x1b[2;1H\
\x1b]133;C\x07hi\r\n\x1b]133;D;0\x07\x1b]133;A\x07",
1,
);
assert_eq!(events.len(), 1);
let PendingEntry::Command {
command,
output,
exit_code,
..
} = &events[0]
else {
panic!("expected a command entry, got {events:?}");
};
assert_eq!(command, "mvhenten@sandbox:~$ echo hi");
assert_eq!(output, "hi\r\n");
assert_eq!(*exit_code, Some(0));
}
#[test]
fn command_never_carries_a_mode_toggle_escape_riding_the_same_line() {
let mut seg = Segmenter::new();
let events = feed_str(
&mut seg,
"mvhenten@sandbox:~$ echo hi\x1b[?2004l\r\n\x1b]133;C\x07hi\r\n\x1b]133;D;0\x07\x1b]133;A\x07",
1,
);
let PendingEntry::Command { command, .. } = &events[0] else {
panic!("expected command entry");
};
assert_eq!(command, "mvhenten@sandbox:~$ echo hi");
}
#[test]
fn two_commands_including_a_failing_one() {
let mut seg = Segmenter::new();
let mut events = feed_str(&mut seg, REAL_ECHO_HELLO, 1);
events.extend(feed_str(
&mut seg,
"false\r\n\x1b]133;C\x07\x1b[?2004h\x1b]133;D;1\x07\x1b]133;A\x07mvhenten@sandbox:~$ ",
2,
));
assert_eq!(events.len(), 2);
let PendingEntry::Command { exit_code, .. } = &events[1] else {
panic!("expected command entry");
};
assert_eq!(*exit_code, Some(1));
}
#[test]
fn marker_split_across_two_feed_calls_still_resolves() {
let mut seg = Segmenter::new();
let whole = REAL_ECHO_HELLO;
let split = whole.len() / 2;
let mut events = feed_str(&mut seg, &whole[..split], 5);
assert!(
events.is_empty(),
"no entry should complete before the D marker arrives, got {events:?}"
);
events.extend(feed_str(&mut seg, &whole[split..], 6));
assert_eq!(events.len(), 1);
assert!(matches!(events[0], PendingEntry::Command { .. }));
}
#[test]
fn marker_split_byte_by_byte_still_resolves() {
let mut seg = Segmenter::new();
let mut events = Vec::new();
for b in REAL_ECHO_HELLO.as_bytes() {
events.extend(seg.feed(&[*b], 9));
}
assert_eq!(events.len(), 1);
let PendingEntry::Command {
command,
output,
exit_code,
..
} = &events[0]
else {
panic!("expected command entry");
};
assert_eq!(command, "echo hello");
assert_eq!(output, "hello\r\n\x1b[?2004h");
assert_eq!(*exit_code, Some(0));
}
#[test]
fn tmux_dcs_passthrough_wrapped_markers_are_unwrapped() {
let mut seg = Segmenter::new();
feed_str(
&mut seg,
"prompt$ true\r\n\x1bPtmux;\x1b\x1b]133;C\x07\x1b\\",
1,
);
let events = feed_str(
&mut seg,
"ok\r\n\x1bPtmux;\x1b\x1b]133;D;0\x07\x1b\x1b]133;A\x07\x1b\\",
2,
);
assert_eq!(events.len(), 1);
let PendingEntry::Command {
command,
output,
exit_code,
..
} = &events[0]
else {
panic!("expected command entry, got {events:?}");
};
assert_eq!(command, "prompt$ true");
assert_eq!(output, "ok\r\n");
assert_eq!(*exit_code, Some(0));
}
#[test]
fn wrapped_marker_split_across_feed_calls() {
let mut seg = Segmenter::new();
let envelope = "x\r\n\x1bPtmux;\x1b\x1b]133;C\x07\x1b\\out\x1bPtmux;\x1b\x1b]133;D;7\x07\x1b\x1b]133;A\x07\x1b\\";
let split = envelope.len() / 2;
let mut events = feed_str(&mut seg, &envelope[..split], 1);
events.extend(feed_str(&mut seg, &envelope[split..], 2));
assert_eq!(events.len(), 1);
let PendingEntry::Command { exit_code, .. } = &events[0] else {
panic!("expected command entry");
};
assert_eq!(*exit_code, Some(7));
}
#[test]
fn uninstrumented_session_falls_back_to_raw_chunks_on_flush() {
let mut seg = Segmenter::new();
let events = feed_str(&mut seg, "plain output, no OSC 133 at all\n", 1);
assert!(events.is_empty(), "must not flush eagerly mid-session");
let flushed = seg.flush(2);
match flushed {
Some(PendingEntry::Raw { raw, ts }) => {
assert_eq!(raw, "plain output, no OSC 133 at all\n");
assert_eq!(ts, 2);
}
other => panic!("expected a raw entry, got {other:?}"),
}
}
#[test]
fn large_uninstrumented_burst_flushes_in_bounded_chunks() {
let mut seg = Segmenter::new();
let burst = "x".repeat(RAW_FLUSH_THRESHOLD * 2 + 10);
let events = seg.feed(burst.as_bytes(), 1);
assert_eq!(events.len(), 2, "two full RAW_FLUSH_THRESHOLD chunks");
for e in &events {
let PendingEntry::Raw { raw, .. } = e else {
panic!("expected raw entry");
};
assert_eq!(raw.len(), RAW_FLUSH_THRESHOLD);
}
let flushed = seg.flush(2);
let Some(PendingEntry::Raw { raw, .. }) = flushed else {
panic!("expected leftover raw entry");
};
assert_eq!(raw.len(), 10);
}
#[test]
fn open_command_at_disconnect_flushes_best_effort_with_no_exit_code() {
let mut seg = Segmenter::new();
feed_str(&mut seg, "cmd\r\n\x1b]133;C\x07partial out", 1);
let flushed = seg.flush(5);
match flushed {
Some(PendingEntry::Command {
command,
output,
exit_code,
ended_at,
..
}) => {
assert_eq!(command, "cmd");
assert_eq!(output, "partial out");
assert_eq!(exit_code, None);
assert_eq!(ended_at, 5);
}
other => panic!("expected a command entry, got {other:?}"),
}
}
#[test]
fn instrumented_session_drops_the_trailing_prompt_remainder_on_flush() {
let mut seg = Segmenter::new();
feed_str(&mut seg, "cmd\r\n\x1b]133;C\x07out", 1);
let events = feed_str(&mut seg, "\x1b]133;D;0\x07\x1b]133;A\x07", 2);
assert_eq!(events.len(), 1);
feed_str(&mut seg, "mvhenten@sandbox:~$ ", 3);
assert_eq!(seg.flush(4), None);
}
#[test]
fn command_output_is_capped_not_unbounded() {
let mut seg = Segmenter::new();
feed_str(&mut seg, "cmd\r\n\x1b]133;C\x07", 1);
let huge = "y".repeat(MAX_COMMAND_OUTPUT_BYTES + 1000);
seg.feed(huge.as_bytes(), 1);
let events = seg.feed(b"\x1b]133;D;0\x07\x1b]133;A\x07", 2);
let PendingEntry::Command { output, .. } = &events[0] else {
panic!("expected command entry");
};
assert_eq!(output.len(), MAX_COMMAND_OUTPUT_BYTES);
}
#[test]
fn d_without_open_command_is_a_defensive_no_op() {
let mut seg = Segmenter::new();
let events = feed_str(&mut seg, "\x1b]133;D;0\x07\x1b]133;A\x07", 1);
assert!(events.is_empty());
}
#[test]
fn cursor_roundtrips_and_is_not_a_bare_integer_string() {
let c = encode_cursor(42);
assert_ne!(c, "42");
assert_eq!(decode_cursor(&c), Some(42));
assert_eq!(decode_cursor("not-a-real-cursor"), None);
}
fn temp_store() -> (tempfile::TempDir, SessionHistoryStore) {
let dir = tempfile::tempdir().unwrap();
let store = SessionHistoryStore::new(dir.path());
(dir, store)
}
#[test]
fn append_assigns_monotonic_seq_and_persists_jsonl() {
let (_dir, store) = temp_store();
let a = store
.append(
"s1",
PendingEntry::Raw {
raw: "one".into(),
ts: 1,
},
)
.unwrap();
let b = store
.append(
"s1",
PendingEntry::Raw {
raw: "two".into(),
ts: 2,
},
)
.unwrap();
let HistoryEntry::Raw(a) = a else { panic!() };
let HistoryEntry::Raw(b) = b else { panic!() };
assert_eq!(a.seq, 1);
assert_eq!(b.seq, 2);
let contents = fs::read_to_string(store.file_path("s1")).unwrap();
assert_eq!(contents.lines().count(), 2);
}
#[test]
fn seq_counter_survives_process_restart_by_rehydrating_from_file() {
let dir = tempfile::tempdir().unwrap();
let store = SessionHistoryStore::new(dir.path());
for i in 0..3 {
store
.append(
"s1",
PendingEntry::Raw {
raw: format!("e{i}"),
ts: i,
},
)
.unwrap();
}
let store2 = SessionHistoryStore::new(dir.path());
let next = store2
.append(
"s1",
PendingEntry::Raw {
raw: "e3".into(),
ts: 3,
},
)
.unwrap();
let HistoryEntry::Raw(next) = next else {
panic!()
};
assert_eq!(next.seq, 4);
}
const TRIGGER: u64 = MAX_ENTRIES_PER_SESSION + TRIM_MARGIN + 1;
fn fill(store: &SessionHistoryStore, session: &str, count: u64) {
for i in 0..count {
store
.append(
session,
PendingEntry::Raw {
raw: format!("e{i}"),
ts: i as i64,
},
)
.unwrap();
}
}
#[test]
fn retention_trims_oldest_entries_once_past_cap_plus_margin() {
let (_dir, store) = temp_store();
fill(&store, "s1", TRIGGER);
let contents = fs::read_to_string(store.file_path("s1")).unwrap();
let line_count = contents.lines().count() as u64;
assert_eq!(line_count, MAX_ENTRIES_PER_SESSION);
let first_line: serde_json::Value =
serde_json::from_str(contents.lines().next().unwrap()).unwrap();
assert_eq!(first_line["seq"].as_u64().unwrap(), TRIM_MARGIN + 2);
fill(&store, "s1", 5);
let contents = fs::read_to_string(store.file_path("s1")).unwrap();
assert_eq!(contents.lines().count() as u64, MAX_ENTRIES_PER_SESSION + 5);
}
#[test]
fn cursor_stays_stable_across_a_trim() {
let (_dir, store) = temp_store();
fill(&store, "s1", TRIGGER);
let stale_cursor = 5u64;
let (page, next_cursor) = store.read_page("s1", Some(stale_cursor), 10).unwrap();
assert_eq!(page.len(), 10);
let expected_first_seq = TRIM_MARGIN + 2;
assert_eq!(page[0]["seq"].as_u64().unwrap(), expected_first_seq);
assert_eq!(next_cursor, expected_first_seq + 9);
}
#[test]
fn read_page_paginates_forward_with_limit() {
let (_dir, store) = temp_store();
for i in 0..5u64 {
store
.append(
"s1",
PendingEntry::Raw {
raw: format!("e{i}"),
ts: i as i64,
},
)
.unwrap();
}
let (page1, cursor1) = store.read_page("s1", None, 2).unwrap();
assert_eq!(page1.len(), 2);
assert_eq!(page1[0]["seq"], 1);
assert_eq!(page1[1]["seq"], 2);
assert_eq!(cursor1, 2);
let (page2, cursor2) = store.read_page("s1", Some(cursor1), 2).unwrap();
assert_eq!(page2[0]["seq"], 3);
assert_eq!(page2[1]["seq"], 4);
assert_eq!(cursor2, 4);
let (page3, cursor3) = store.read_page("s1", Some(cursor2), 2).unwrap();
assert_eq!(page3.len(), 1);
assert_eq!(page3[0]["seq"], 5);
assert_eq!(cursor3, 5);
let (page4, cursor4) = store.read_page("s1", Some(cursor3), 2).unwrap();
assert!(page4.is_empty());
assert_eq!(cursor4, cursor3);
}
#[test]
fn two_concurrent_attaches_only_one_acquires_the_feeder() {
let (_dir, store) = temp_store();
let store = Arc::new(store);
let first = store.try_acquire_feeder("s1");
assert!(first.is_some());
let second = store.try_acquire_feeder("s1");
assert!(second.is_none(), "a second attach must not double-feed");
drop(first);
let third = store.try_acquire_feeder("s1");
assert!(third.is_some(), "freed once the first attach disconnects");
}
}