use std::collections::{HashMap, HashSet};
use std::fs::{self, File, OpenOptions};
use std::io::{BufRead, BufReader, Read, Seek, SeekFrom, 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;
pub const MAX_WIRE_OUTPUT_BYTES: usize = 16 * 1024;
pub const MAX_PAGE_BYTES: usize = 512 * 1024;
const BACKWARD_CHUNK: usize = 64 * 1024;
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<PageCursor>,
limit: usize,
) -> anyhow::Result<Page> {
let path = self.file_path(session);
let floor = cursor.map(|c| c.seq).unwrap_or(0);
let zero = Page {
entries: Vec::new(),
next_seq: 0,
next_offset: 0,
prev_seq: 0,
prev_offset: 0,
has_older: false,
};
let Ok(file) = File::open(&path) else {
return Ok(zero);
};
let file_len = file.metadata()?.len();
if file_len == 0 {
return Ok(zero);
}
let mut reader = BufReader::new(file);
let from_offset = match cursor.and_then(|c| c.offset) {
Some(offset) if offset == file_len => {
return Ok(Page {
entries: Vec::new(),
next_seq: floor,
next_offset: offset,
prev_seq: floor,
prev_offset: offset,
has_older: offset > 0,
})
}
Some(offset) if offset < file_len && resumes_at_seq(&mut reader, offset, floor)? => {
offset
}
_ => 0,
};
scan_forward(&mut reader, from_offset, floor, limit)
}
pub fn read_tail(&self, session: &str, count: usize) -> anyhow::Result<Page> {
let path = self.file_path(session);
self.collect_backwards(BackwardLines::open(&path)?, u64::MAX, count, 0)
}
pub fn read_before(
&self,
session: &str,
cursor: PageCursor,
limit: usize,
) -> anyhow::Result<Page> {
let path = self.file_path(session);
let trusted = match cursor.offset {
Some(offset) => offset_resumes_at_seq(&path, offset, cursor.seq)?,
None => None,
};
let lines = match trusted {
Some(offset) => BackwardLines::open_at(&path, offset)?,
None => BackwardLines::open(&path)?,
};
let fallback_offset = trusted.unwrap_or(0);
self.collect_backwards(lines, cursor.seq, limit, fallback_offset)
}
fn collect_backwards(
&self,
lines: Option<BackwardLines>,
ceiling: u64,
count: usize,
empty_offset: u64,
) -> anyhow::Result<Page> {
let mut entries: Vec<serde_json::Value> = Vec::new();
let mut next_seq = 0u64;
let mut next_offset = 0u64;
let mut oldest_seq = 0u64;
let mut oldest_start = empty_offset;
let mut has_older = false;
let mut used = 0usize;
if let Some(mut lines) = lines {
loop {
let Some((start, bytes)) = lines.next_line()? else {
break;
};
let text = String::from_utf8_lossy(&bytes);
if text.trim().is_empty() {
continue;
}
let Ok(mut value) = serde_json::from_str::<serde_json::Value>(&text) else {
continue;
};
let Some(seq) = value.get("seq").and_then(|s| s.as_u64()) else {
continue;
};
if seq == 0 || seq > ceiling {
continue;
}
if entries.len() >= count {
has_older = true;
break;
}
cap_output_for_wire(&mut value);
let size = serde_json::to_string(&value)?.len();
if !entries.is_empty() && used + size > MAX_PAGE_BYTES {
has_older = true;
break;
}
used += size;
if entries.is_empty() {
next_seq = seq;
next_offset = start + bytes.len() as u64 + 1;
}
oldest_seq = seq;
oldest_start = start;
entries.push(value);
}
}
if entries.is_empty() {
return Ok(Page {
entries,
next_seq,
next_offset,
prev_seq: 0,
prev_offset: empty_offset,
has_older: false,
});
}
entries.reverse();
Ok(Page {
entries,
next_seq,
next_offset,
prev_seq: oldest_seq.saturating_sub(1),
prev_offset: oldest_start,
has_older,
})
}
}
pub struct Page {
pub entries: Vec<serde_json::Value>,
pub next_seq: u64,
pub next_offset: u64,
pub prev_seq: u64,
pub prev_offset: u64,
pub has_older: bool,
}
fn cap_output_for_wire(value: &mut serde_json::Value) {
let Some(object) = value.as_object_mut() else {
return;
};
let Some(output) = object.get("output").and_then(|o| o.as_str()) else {
return;
};
if output.len() <= MAX_WIRE_OUTPUT_BYTES {
return;
}
let mut dropped = output.len() - MAX_WIRE_OUTPUT_BYTES;
while !output.is_char_boundary(dropped) {
dropped += 1;
}
let kept = output[dropped..].to_string();
object.insert("output".to_string(), serde_json::Value::String(kept));
object.insert(
"outputTruncatedBytes".to_string(),
serde_json::Value::from(dropped),
);
}
fn offset_resumes_at_seq(path: &Path, offset: u64, cursor_seq: u64) -> anyhow::Result<Option<u64>> {
let Ok(file) = File::open(path) else {
return Ok(None);
};
if offset >= file.metadata()?.len() {
return Ok(None);
}
let mut reader = BufReader::new(file);
if resumes_at_seq(&mut reader, offset, cursor_seq)? {
return Ok(Some(offset));
}
Ok(None)
}
fn resumes_at_seq(
reader: &mut BufReader<File>,
offset: u64,
cursor_seq: u64,
) -> anyhow::Result<bool> {
reader.seek(SeekFrom::Start(offset))?;
let mut line = Vec::new();
reader.read_until(b'\n', &mut line)?;
if !line.ends_with(b"\n") {
return Ok(false);
}
let Ok(value) = serde_json::from_slice::<serde_json::Value>(&line) else {
return Ok(false);
};
Ok(value.get("seq").and_then(|s| s.as_u64()) == Some(cursor_seq + 1))
}
fn scan_forward(
reader: &mut BufReader<File>,
from_offset: u64,
floor: u64,
limit: usize,
) -> anyhow::Result<Page> {
let mut entries: Vec<serde_json::Value> = Vec::new();
let mut next_seq = floor;
let mut next_offset = from_offset;
let mut prev_seq = floor;
let mut prev_offset = from_offset;
let mut used = 0usize;
reader.seek(SeekFrom::Start(from_offset))?;
let mut position = from_offset;
let mut line = Vec::new();
loop {
line.clear();
let read = reader.read_until(b'\n', &mut line)?;
if read == 0 || !line.ends_with(b"\n") {
break;
}
let line_start = position;
position += read as u64;
let Ok(mut value) = serde_json::from_slice::<serde_json::Value>(&line) else {
next_offset = position;
continue;
};
let seq = value.get("seq").and_then(|s| s.as_u64()).unwrap_or(0);
if seq <= floor {
next_offset = position;
continue;
}
cap_output_for_wire(&mut value);
let size = serde_json::to_string(&value)?.len();
if !entries.is_empty() && used + size > MAX_PAGE_BYTES {
break;
}
used += size;
next_seq = seq;
next_offset = position;
if entries.is_empty() {
prev_seq = seq.saturating_sub(1);
prev_offset = line_start;
}
entries.push(value);
if entries.len() >= limit {
break;
}
}
Ok(Page {
entries,
next_seq,
next_offset,
prev_seq,
prev_offset,
has_older: prev_offset > 0,
})
}
struct BackwardLines {
file: File,
remaining: u64,
buf: Vec<u8>,
}
impl BackwardLines {
fn open(path: &Path) -> anyhow::Result<Option<Self>> {
let Ok(file) = File::open(path) else {
return Ok(None);
};
let end = file.metadata()?.len();
Self::from_end(file, end)
}
fn open_at(path: &Path, end: u64) -> anyhow::Result<Option<Self>> {
let Ok(file) = File::open(path) else {
return Ok(None);
};
let len = file.metadata()?.len();
Self::from_end(file, end.min(len))
}
fn from_end(mut file: File, end: u64) -> anyhow::Result<Option<Self>> {
let mut scan_end = end;
while scan_end > 0 {
let start = scan_end.saturating_sub(BACKWARD_CHUNK as u64);
let mut chunk = vec![0u8; (scan_end - start) as usize];
file.seek(SeekFrom::Start(start))?;
file.read_exact(&mut chunk)?;
if let Some(i) = chunk.iter().rposition(|&b| b == b'\n') {
return Ok(Some(Self {
file,
remaining: start + i as u64,
buf: Vec::new(),
}));
}
scan_end = start;
}
Ok(None)
}
fn next_line(&mut self) -> anyhow::Result<Option<(u64, Vec<u8>)>> {
loop {
if let Some(i) = self.buf.iter().rposition(|&b| b == b'\n') {
let line = self.buf.split_off(i + 1);
self.buf.truncate(i);
return Ok(Some((self.remaining + i as u64 + 1, line)));
}
if self.remaining == 0 {
if self.buf.is_empty() {
return Ok(None);
}
return Ok(Some((0, std::mem::take(&mut self.buf))));
}
let take = BACKWARD_CHUNK.min(self.remaining as usize);
self.remaining -= take as u64;
let mut chunk = vec![0u8; take];
self.file.seek(SeekFrom::Start(self.remaining))?;
self.file.read_exact(&mut chunk)?;
chunk.append(&mut self.buf);
self.buf = chunk;
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct PageCursor {
pub seq: u64,
pub offset: Option<u64>,
}
pub fn encode_cursor(seq: u64, offset: u64) -> String {
BASE64URL.encode(format!("v2:{seq}:{offset}"))
}
pub fn decode_cursor(cursor: &str) -> Option<PageCursor> {
let bytes = BASE64URL.decode(cursor).ok()?;
let text = String::from_utf8(bytes).ok()?;
if let Some(rest) = text.strip_prefix("v2:") {
let (seq, offset) = rest.split_once(':')?;
return Some(PageCursor {
seq: seq.parse().ok()?,
offset: Some(offset.parse().ok()?),
});
}
Some(PageCursor {
seq: text.strip_prefix("v1:")?.parse().ok()?,
offset: None,
})
}
#[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, 900);
assert_ne!(c, "42");
assert_eq!(
decode_cursor(&c),
Some(PageCursor {
seq: 42,
offset: Some(900)
})
);
assert_eq!(decode_cursor("not-a-real-cursor"), None);
}
#[test]
fn v1_cursors_still_decode_with_an_unknown_offset() {
let legacy = BASE64URL.encode("v1:7");
assert_eq!(
decode_cursor(&legacy),
Some(PageCursor {
seq: 7,
offset: 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 = PageCursor {
seq: 5,
offset: None,
};
let page = store.read_page("s1", Some(stale_cursor), 10).unwrap();
assert_eq!(page.entries.len(), 10);
let expected_first_seq = TRIM_MARGIN + 2;
assert_eq!(page.entries[0]["seq"].as_u64().unwrap(), expected_first_seq);
assert_eq!(page.next_seq, expected_first_seq + 9);
}
#[test]
fn an_offset_invalidated_by_a_trim_falls_back_to_a_full_scan() {
let (_dir, store) = temp_store();
fill(&store, "s1", 20);
let before = store.read_page("s1", None, 10).unwrap();
assert_eq!(before.next_seq, 10);
fill(&store, "s1", TRIGGER);
let resumed = store
.read_page(
"s1",
Some(PageCursor {
seq: before.next_seq,
offset: Some(before.next_offset),
}),
10,
)
.unwrap();
let oldest_retained = TRIM_MARGIN + 2;
assert_eq!(
resumed.entries[0]["seq"].as_u64().unwrap(),
oldest_retained,
"a stale offset must fall back to a scan, never return a wrong page"
);
}
#[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 = store.read_page("s1", None, 2).unwrap();
assert_eq!(page1.entries.len(), 2);
assert_eq!(page1.entries[0]["seq"], 1);
assert_eq!(page1.entries[1]["seq"], 2);
assert_eq!(page1.next_seq, 2);
let page2 = store.read_page("s1", Some(cursor_of(&page1)), 2).unwrap();
assert_eq!(page2.entries[0]["seq"], 3);
assert_eq!(page2.entries[1]["seq"], 4);
assert_eq!(page2.next_seq, 4);
let page3 = store.read_page("s1", Some(cursor_of(&page2)), 2).unwrap();
assert_eq!(page3.entries.len(), 1);
assert_eq!(page3.entries[0]["seq"], 5);
assert_eq!(page3.next_seq, 5);
let page4 = store.read_page("s1", Some(cursor_of(&page3)), 2).unwrap();
assert!(page4.entries.is_empty());
assert_eq!(page4.next_seq, page3.next_seq);
assert_eq!(page4.next_offset, page3.next_offset);
}
fn cursor_of(page: &Page) -> PageCursor {
PageCursor {
seq: page.next_seq,
offset: Some(page.next_offset),
}
}
fn append_output(store: &SessionHistoryStore, session: &str, output: String) {
store
.append(
session,
PendingEntry::Command {
command: "cmd".into(),
output,
exit_code: Some(0),
started_at: 1,
ended_at: 2,
},
)
.unwrap();
}
#[test]
fn a_caught_up_offset_cursor_lands_exactly_at_the_end_of_the_file() {
let (_dir, store) = temp_store();
fill(&store, "s1", 5);
let page = store.read_page("s1", None, 50).unwrap();
let file_len = fs::metadata(store.file_path("s1")).unwrap().len();
assert_eq!(page.next_offset, file_len);
assert_eq!(page.next_seq, 5);
}
#[test]
fn a_wrong_offset_still_returns_the_correct_page() {
let (_dir, store) = temp_store();
fill(&store, "s1", 5);
let contents = fs::read_to_string(store.file_path("s1")).unwrap();
let past_seq_three = contents.lines().take(3).map(|l| l.len() + 1).sum::<usize>() as u64;
for offset in [3, past_seq_three] {
let page = store
.read_page(
"s1",
Some(PageCursor {
seq: 2,
offset: Some(offset),
}),
50,
)
.unwrap();
let seqs: Vec<u64> = page
.entries
.iter()
.map(|e| e["seq"].as_u64().unwrap())
.collect();
assert_eq!(seqs, vec![3, 4, 5], "offset {offset}");
}
}
#[test]
fn an_offset_past_the_end_of_the_file_falls_back_to_a_full_scan() {
let (_dir, store) = temp_store();
fill(&store, "s1", 3);
let page = store
.read_page(
"s1",
Some(PageCursor {
seq: 1,
offset: Some(u64::MAX),
}),
50,
)
.unwrap();
assert_eq!(page.entries.len(), 2);
assert_eq!(page.entries[0]["seq"], 2);
}
#[test]
fn a_line_written_without_its_newline_yet_is_not_served() {
let (_dir, store) = temp_store();
fill(&store, "s1", 2);
let caught_up = store.read_page("s1", None, 50).unwrap();
let mut f = OpenOptions::new()
.append(true)
.open(store.file_path("s1"))
.unwrap();
write!(
f,
"{}",
serde_json::json!({"seq": 3, "raw": "torn", "ts": 0})
)
.unwrap();
drop(f);
let page = store
.read_page("s1", Some(cursor_of(&caught_up)), 50)
.unwrap();
assert!(page.entries.is_empty());
assert_eq!(page.next_offset, caught_up.next_offset);
}
#[test]
fn a_line_torn_inside_a_multibyte_character_is_skipped_not_an_error() {
let (_dir, store) = temp_store();
fill(&store, "s1", 2);
let caught_up = store.read_page("s1", None, 50).unwrap();
let line = serde_json::json!({"seq": 3, "raw": "€€€", "ts": 0}).to_string();
let bytes = line.as_bytes();
let cut = bytes.iter().position(|&b| b == 0xe2).unwrap() + 1;
let mut f = OpenOptions::new()
.append(true)
.open(store.file_path("s1"))
.unwrap();
f.write_all(&bytes[..cut]).unwrap();
drop(f);
for cursor in [
None,
Some(cursor_of(&caught_up)),
Some(PageCursor {
seq: 1,
offset: None,
}),
] {
let page = store.read_page("s1", cursor, 50).unwrap();
let seqs: Vec<u64> = page
.entries
.iter()
.map(|e| e["seq"].as_u64().unwrap())
.collect();
let expected: Vec<u64> = match cursor.map(|c| c.seq).unwrap_or(0) {
0 => vec![1, 2],
1 => vec![2],
_ => vec![],
};
assert_eq!(seqs, expected, "cursor {cursor:?}");
}
let mut f = OpenOptions::new()
.append(true)
.open(store.file_path("s1"))
.unwrap();
f.write_all(&bytes[cut..]).unwrap();
f.write_all(b"\n").unwrap();
drop(f);
let page = store.read_page("s1", None, 50).unwrap();
assert_eq!(page.entries.len(), 3);
assert_eq!(page.entries[2]["raw"], "€€€");
}
fn decoy_line(seq: u64, byte_length: usize) -> String {
let base = serde_json::json!({"seq": seq, "raw": "", "ts": 0})
.to_string()
.len();
let padded = "p".repeat(byte_length - 1 - base);
format!(
"{}\n",
serde_json::json!({"seq": seq, "raw": padded, "ts": 0})
)
}
#[test]
fn a_caught_up_offset_answers_without_reading_the_files_contents() {
let (_dir, store) = temp_store();
fill(&store, "s1", 5);
let caught_up = store.read_page("s1", None, 50).unwrap();
assert_eq!(caught_up.next_seq, 5);
let path = store.file_path("s1");
let file_len = fs::metadata(&path).unwrap().len() as usize;
fs::write(&path, decoy_line(90, file_len)).unwrap();
assert_eq!(fs::metadata(&path).unwrap().len() as usize, file_len);
let scanned = store
.read_page(
"s1",
Some(PageCursor {
seq: 5,
offset: None,
}),
50,
)
.unwrap();
assert_eq!(
scanned.entries[0]["seq"].as_u64().unwrap(),
90,
"the decoy must be something a full scan returns"
);
let polled = store
.read_page("s1", Some(cursor_of(&caught_up)), 50)
.unwrap();
assert!(
polled.entries.is_empty(),
"a cursor at the end of the file must not read its contents"
);
assert_eq!(polled.next_seq, caught_up.next_seq);
assert_eq!(polled.next_offset, caught_up.next_offset);
}
#[test]
fn a_trusted_offset_seeks_past_content_a_full_scan_would_return() {
let (_dir, store) = temp_store();
fill(&store, "s1", 5);
let first = store.read_page("s1", None, 3).unwrap();
assert_eq!(first.next_seq, 3);
let path = store.file_path("s1");
let bytes = fs::read(&path).unwrap();
let offset = first.next_offset as usize;
let mut rewritten = decoy_line(90, offset).into_bytes();
rewritten.extend_from_slice(&bytes[offset..]);
fs::write(&path, &rewritten).unwrap();
let scanned = store
.read_page(
"s1",
Some(PageCursor {
seq: 3,
offset: None,
}),
50,
)
.unwrap();
let scanned_seqs: Vec<u64> = scanned
.entries
.iter()
.map(|e| e["seq"].as_u64().unwrap())
.collect();
assert_eq!(
scanned_seqs,
vec![90, 4, 5],
"the decoy must be something a full scan returns"
);
let seeked = store.read_page("s1", Some(cursor_of(&first)), 50).unwrap();
let seeked_seqs: Vec<u64> = seeked
.entries
.iter()
.map(|e| e["seq"].as_u64().unwrap())
.collect();
assert_eq!(
seeked_seqs,
vec![4, 5],
"a trusted offset must seek, never rescan from byte 0"
);
}
#[test]
fn output_over_the_wire_cap_keeps_its_end_and_reports_what_it_dropped() {
let (_dir, store) = temp_store();
let output = format!("{}TAIL", "x".repeat(MAX_WIRE_OUTPUT_BYTES));
append_output(&store, "s1", output);
let page = store.read_page("s1", None, 50).unwrap();
let entry = &page.entries[0];
assert_eq!(
entry["output"].as_str().unwrap().len(),
MAX_WIRE_OUTPUT_BYTES
);
assert!(entry["output"].as_str().unwrap().ends_with("TAIL"));
assert_eq!(entry["outputTruncatedBytes"].as_u64().unwrap(), 4);
}
#[test]
fn output_under_the_wire_cap_carries_no_truncation_field() {
let (_dir, store) = temp_store();
append_output(&store, "s1", "short".into());
let page = store.read_page("s1", None, 50).unwrap();
assert!(page.entries[0].get("outputTruncatedBytes").is_none());
}
#[test]
fn a_forward_page_stops_at_the_byte_budget_without_skipping_a_seq() {
let (_dir, store) = temp_store();
for _ in 0..64 {
append_output(&store, "s1", "x".repeat(MAX_WIRE_OUTPUT_BYTES));
}
let page = store.read_page("s1", None, 500).unwrap();
assert!(page.entries.len() < 64);
assert_eq!(
page.next_seq,
page.entries.last().unwrap()["seq"].as_u64().unwrap()
);
let next = store.read_page("s1", Some(cursor_of(&page)), 500).unwrap();
assert_eq!(next.entries[0]["seq"].as_u64().unwrap(), page.next_seq + 1);
}
#[test]
fn a_page_returns_its_first_entry_even_when_that_entry_alone_blows_the_budget() {
let (_dir, store) = temp_store();
store
.append(
"s1",
PendingEntry::Raw {
raw: "x".repeat(MAX_PAGE_BYTES + 1),
ts: 1,
},
)
.unwrap();
store
.append(
"s1",
PendingEntry::Raw {
raw: "second".into(),
ts: 2,
},
)
.unwrap();
let page = store.read_page("s1", None, 50).unwrap();
assert_eq!(page.entries.len(), 1);
assert_eq!(page.entries[0]["seq"], 1);
let tail = store.read_tail("s1", 50).unwrap();
assert_eq!(tail.entries.len(), 1);
assert_eq!(tail.entries[0]["seq"], 2);
}
#[test]
fn tail_returns_the_newest_entries_oldest_first() {
let (_dir, store) = temp_store();
fill(&store, "s1", 10);
let page = store.read_tail("s1", 3).unwrap();
let seqs: Vec<u64> = page
.entries
.iter()
.map(|e| e["seq"].as_u64().unwrap())
.collect();
assert_eq!(seqs, vec![8, 9, 10]);
assert_eq!(page.next_seq, 10);
assert_eq!(
page.next_offset,
fs::metadata(store.file_path("s1")).unwrap().len()
);
}
#[test]
fn tail_over_a_chunk_boundary_still_walks_the_file_correctly() {
let (_dir, store) = temp_store();
for _ in 0..12 {
append_output(&store, "s1", "x".repeat(BACKWARD_CHUNK));
}
let page = store.read_tail("s1", 4).unwrap();
let seqs: Vec<u64> = page
.entries
.iter()
.map(|e| e["seq"].as_u64().unwrap())
.collect();
assert_eq!(seqs, vec![9, 10, 11, 12]);
}
#[test]
fn tail_drops_from_its_oldest_end_when_the_budget_bites() {
let (_dir, store) = temp_store();
for _ in 0..64 {
append_output(&store, "s1", "x".repeat(MAX_WIRE_OUTPUT_BYTES));
}
let page = store.read_tail("s1", 64).unwrap();
assert!(page.entries.len() < 64);
assert_eq!(page.entries.last().unwrap()["seq"], 64);
assert_eq!(page.next_seq, 64);
}
#[test]
fn a_cursor_against_a_session_with_no_history_gets_the_zero_cursor_back() {
let (_dir, store) = temp_store();
for cursor in [
PageCursor {
seq: 7,
offset: None,
},
PageCursor {
seq: 7,
offset: Some(0),
},
PageCursor {
seq: 7,
offset: Some(400),
},
] {
let page = store.read_page("never-written", Some(cursor), 50).unwrap();
assert!(page.entries.is_empty());
assert_eq!(page.next_seq, 0, "cursor {cursor:?}");
assert_eq!(page.next_offset, 0, "cursor {cursor:?}");
}
}
#[test]
fn tail_of_a_session_with_no_history_is_the_zero_cursor() {
let (_dir, store) = temp_store();
let page = store.read_tail("never-written", 50).unwrap();
assert!(page.entries.is_empty());
assert_eq!(page.next_seq, 0);
assert_eq!(page.next_offset, 0);
}
#[test]
fn a_cursor_from_a_tail_page_resumes_forward_with_no_gap() {
let (_dir, store) = temp_store();
fill(&store, "s1", 10);
let tail = store.read_tail("s1", 3).unwrap();
let resumed = store.read_page("s1", Some(cursor_of(&tail)), 50).unwrap();
assert!(resumed.entries.is_empty());
fill(&store, "s1", 2);
let resumed = store.read_page("s1", Some(cursor_of(&tail)), 50).unwrap();
let seqs: Vec<u64> = resumed
.entries
.iter()
.map(|e| e["seq"].as_u64().unwrap())
.collect();
assert_eq!(seqs, vec![11, 12]);
}
#[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");
}
fn back_cursor(page: &Page) -> PageCursor {
PageCursor {
seq: page.prev_seq,
offset: Some(page.prev_offset),
}
}
fn seqs_of(page: &Page) -> Vec<u64> {
page.entries
.iter()
.map(|e| e["seq"].as_u64().unwrap_or(0))
.collect()
}
#[test]
fn read_before_walks_back_over_the_whole_record() {
let (_dir, store) = temp_store();
fill(&store, "s1", 500);
let tail = store.read_tail("s1", 100).unwrap();
assert_eq!(seqs_of(&tail).first().copied(), Some(401));
assert!(tail.has_older);
let mut seen: Vec<u64> = seqs_of(&tail);
let mut page = tail;
while page.has_older {
page = store.read_before("s1", back_cursor(&page), 100).unwrap();
let mut batch = seqs_of(&page);
assert!(
!batch.is_empty(),
"a page claiming older entries must carry some"
);
batch.extend(seen);
seen = batch;
}
assert_eq!(seen, (1..=500).collect::<Vec<u64>>());
}
#[test]
fn read_before_stops_at_the_oldest_retained_entry() {
let (_dir, store) = temp_store();
fill(&store, "s1", 5);
let tail = store.read_tail("s1", 2).unwrap();
assert!(tail.has_older);
let older = store.read_before("s1", back_cursor(&tail), 10).unwrap();
assert_eq!(seqs_of(&older), vec![1, 2, 3]);
assert!(!older.has_older);
let past_the_start = store.read_before("s1", back_cursor(&older), 10).unwrap();
assert!(past_the_start.entries.is_empty());
assert!(!past_the_start.has_older);
}
#[test]
fn read_before_falls_back_to_a_scan_on_a_stale_offset() {
let (_dir, store) = temp_store();
fill(&store, "s1", 20);
let tail = store.read_tail("s1", 5).unwrap();
let honest = store.read_before("s1", back_cursor(&tail), 5).unwrap();
let stale = PageCursor {
seq: tail.prev_seq,
offset: Some(tail.prev_offset + 3),
};
let recovered = store.read_before("s1", stale, 5).unwrap();
assert_eq!(seqs_of(&recovered), seqs_of(&honest));
assert_eq!(seqs_of(&recovered), vec![11, 12, 13, 14, 15]);
let offsetless = PageCursor {
seq: tail.prev_seq,
offset: None,
};
let from_v1 = store.read_before("s1", offsetless, 5).unwrap();
assert_eq!(seqs_of(&from_v1), seqs_of(&honest));
}
#[test]
fn read_before_reaches_the_start_of_a_trimmed_record() {
let (_dir, store) = temp_store();
fill(&store, "s1", TRIGGER);
let oldest_retained = TRIM_MARGIN + 2;
let mut page = store.read_tail("s1", MAX_LIMIT).unwrap();
let mut oldest_seen = *seqs_of(&page).first().unwrap();
while page.has_older {
page = store
.read_before("s1", back_cursor(&page), MAX_LIMIT)
.unwrap();
let first = *seqs_of(&page).first().unwrap();
assert!(first < oldest_seen);
oldest_seen = first;
}
assert_eq!(oldest_seen, oldest_retained);
}
#[test]
fn read_before_drops_from_its_oldest_end_when_the_budget_bites() {
let (_dir, store) = temp_store();
for _ in 0..64 {
append_output(&store, "s1", "x".repeat(MAX_WIRE_OUTPUT_BYTES));
}
let tail = store.read_tail("s1", 1).unwrap();
let page = store.read_before("s1", back_cursor(&tail), 64).unwrap();
assert!(page.entries.len() < 63);
assert_eq!(page.entries.last().unwrap()["seq"], 63);
assert!(page.has_older);
}
#[test]
fn tail_reports_whether_the_record_continues_past_it() {
let (_dir, store) = temp_store();
fill(&store, "s1", 10);
let partial = store.read_tail("s1", 3).unwrap();
assert!(partial.has_older);
let older = store.read_before("s1", back_cursor(&partial), 3).unwrap();
assert_eq!(seqs_of(&older), vec![5, 6, 7]);
let whole = store.read_tail("s1", 50).unwrap();
assert!(!whole.has_older);
}
#[test]
fn a_line_without_a_seq_does_not_end_the_walk() {
let (_dir, store) = temp_store();
fill(&store, "s1", 10);
let path = store.file_path("s1");
let mut lines: Vec<String> = fs::read_to_string(&path)
.unwrap()
.lines()
.map(|l| l.to_string())
.collect();
lines.insert(5, r#"{"kind":"raw","raw":"legacy"}"#.to_string());
fs::write(&path, lines.join("\n") + "\n").unwrap();
let mut page = store.read_tail("s1", 3).unwrap();
let mut seen = seqs_of(&page);
let mut guard = 0;
while page.has_older {
assert!(guard < 20, "the walk must terminate");
guard += 1;
page = store.read_before("s1", back_cursor(&page), 3).unwrap();
let mut batch = seqs_of(&page);
assert!(!batch.is_empty());
batch.extend(seen);
seen = batch;
}
assert_eq!(seen, (1..=10).collect::<Vec<u64>>());
}
#[test]
fn an_empty_record_has_nothing_older() {
let (_dir, store) = temp_store();
let tail = store.read_tail("s1", 50).unwrap();
assert!(tail.entries.is_empty());
assert!(!tail.has_older);
let older = store.read_before("s1", back_cursor(&tail), 50).unwrap();
assert!(older.entries.is_empty());
assert!(!older.has_older);
}
}