use super::event::SessionEvent;
use super::manager::Session;
use sha2::{Digest, Sha256};
use std::{
collections::VecDeque,
fs,
io::{BufRead, BufReader, Read, Seek, SeekFrom},
path::PathBuf,
};
pub fn validate_session_id(id: String) -> anyhow::Result<String> {
if id.is_empty() {
anyhow::bail!("session id must not be empty");
}
if id == "." || id == ".." || id.contains("..") {
anyhow::bail!("session id must not contain '..'");
}
if id.contains('/') || id.contains('\\') {
anyhow::bail!("session id must not contain path separators");
}
if PathBuf::from(&id).is_absolute() {
anyhow::bail!("session id must not be an absolute path");
}
if !id
.chars()
.all(|character| character.is_ascii_alphanumeric() || character == '_' || character == '-')
{
anyhow::bail!("session id must match [A-Za-z0-9_-]+");
}
Ok(id)
}
fn open_session_file(session: &Session) -> anyhow::Result<fs::File> {
let root = session
.path
.parent()
.ok_or_else(|| anyhow::anyhow!("session file has no parent"))?;
super::store::open_existing_primary(root, &session.id)?
.ok_or_else(|| anyhow::anyhow!("session JSONL is missing"))
}
#[derive(Debug)]
pub(crate) enum BoundedReadError {
BudgetExceeded(String),
}
impl std::fmt::Display for BoundedReadError {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::BudgetExceeded(message) => formatter.write_str(message),
}
}
}
impl std::error::Error for BoundedReadError {}
fn budget_error(message: impl Into<String>) -> anyhow::Error {
anyhow::Error::new(BoundedReadError::BudgetExceeded(message.into()))
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct SessionReadDiagnostic {
pub(crate) line: usize,
pub(crate) message: String,
}
#[derive(Debug, Clone, PartialEq)]
pub(crate) struct TolerantSessionEvents {
pub(crate) events: Vec<SessionEvent>,
pub(crate) diagnostics: Vec<SessionReadDiagnostic>,
pub(crate) cutoff_bytes: u64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) struct SessionReadStats {
pub lines_read: usize,
pub bytes_read: usize,
pub content_digest: [u8; 32],
}
const MAX_TOLERANT_READ_DIAGNOSTICS: usize = 64;
fn push_tolerant_read_diagnostic(
diagnostics: &mut Vec<SessionReadDiagnostic>,
omitted_count: &mut usize,
diagnostic: SessionReadDiagnostic,
) {
if diagnostics.len() < MAX_TOLERANT_READ_DIAGNOSTICS {
diagnostics.push(diagnostic);
} else {
*omitted_count = omitted_count.saturating_add(1);
}
}
fn finalize_tolerant_read_diagnostics(
mut diagnostics: Vec<SessionReadDiagnostic>,
omitted_count: usize,
) -> Vec<SessionReadDiagnostic> {
if omitted_count == 0 {
return diagnostics;
}
if diagnostics.len() == MAX_TOLERANT_READ_DIAGNOSTICS {
diagnostics.pop();
}
diagnostics.push(SessionReadDiagnostic {
line: 0,
message: format!(
"omitted {omitted_count} additional session JSONL diagnostics after cap of {MAX_TOLERANT_READ_DIAGNOSTICS}"
),
});
diagnostics
}
pub(crate) const MAX_METADATA_VISIT_LINES: usize = 100_000;
pub(crate) const MAX_METADATA_VISIT_BYTES: usize = 64 * 1024 * 1024;
impl Session {
pub(crate) fn read_events_tolerant_bounded(
&self,
max_lines: usize,
max_bytes: usize,
) -> anyhow::Result<TolerantSessionEvents> {
Ok(self
.read_events_tolerant_bounded_with_stats(max_lines, max_bytes)?
.0)
}
pub(crate) fn read_events_tolerant_bounded_with_stats(
&self,
max_lines: usize,
max_bytes: usize,
) -> anyhow::Result<(TolerantSessionEvents, SessionReadStats)> {
let mut events = Vec::new();
let (diagnostics, cutoff_bytes, stats) =
self.visit_events_tolerant_bounded_with_stats(max_lines, max_bytes, |event| {
events.push(event)
})?;
Ok((
TolerantSessionEvents {
events,
diagnostics,
cutoff_bytes: cutoff_bytes as u64,
},
stats,
))
}
pub(crate) fn history_bytes(&self) -> anyhow::Result<u64> {
let root = self
.path
.parent()
.ok_or_else(|| anyhow::anyhow!("session file has no parent"))?;
let Some(file) = super::store::open_existing_primary(root, &self.id)? else {
return Ok(0);
};
Ok(file.metadata()?.len())
}
pub(crate) fn content_fingerprint_bounded(
&self,
max_bytes: usize,
) -> anyhow::Result<([u8; 32], usize)> {
validate_session_id(self.id.clone())?;
if !self.path.exists() {
return Ok((Sha256::digest([]).into(), 0));
}
let file = open_session_file(self)?;
let mut reader = file.take(
u64::try_from(max_bytes)
.unwrap_or(u64::MAX)
.saturating_add(1),
);
let mut digest = Sha256::new();
let mut total_bytes = 0usize;
let mut buffer = [0u8; 8192];
loop {
let read = reader.read(&mut buffer)?;
if read == 0 {
break;
}
total_bytes = total_bytes.saturating_add(read);
if total_bytes > max_bytes {
return Err(budget_error(format!(
"session JSONL tolerant read limit exceeded: {max_bytes} bytes"
)));
}
digest.update(&buffer[..read]);
}
Ok((digest.finalize().into(), total_bytes))
}
pub(crate) fn visit_events_tolerant_bounded(
&self,
max_lines: usize,
max_bytes: usize,
visit: impl FnMut(SessionEvent),
) -> anyhow::Result<(Vec<SessionReadDiagnostic>, usize)> {
let (diagnostics, bytes, _) =
self.visit_events_tolerant_bounded_with_stats(max_lines, max_bytes, visit)?;
Ok((diagnostics, bytes))
}
fn visit_events_tolerant_bounded_with_stats(
&self,
max_lines: usize,
max_bytes: usize,
mut visit: impl FnMut(SessionEvent),
) -> anyhow::Result<(Vec<SessionReadDiagnostic>, usize, SessionReadStats)> {
validate_session_id(self.id.clone())?;
if !self.path.exists() {
return Ok((
Vec::new(),
0,
SessionReadStats {
lines_read: 0,
bytes_read: 0,
content_digest: Sha256::digest([]).into(),
},
));
}
let file = open_session_file(self)?;
let mut reader = BufReader::new(file);
let mut diagnostics = Vec::new();
let mut omitted_diagnostics = 0usize;
let mut total_bytes = 0usize;
let mut lines_read = 0usize;
let mut digest = Sha256::new();
let mut line = Vec::new();
for line_number in 1..=max_lines {
line.clear();
let remaining = max_bytes.saturating_sub(total_bytes);
if remaining == 0 {
if !reader.fill_buf()?.is_empty() {
return Err(budget_error(format!(
"session JSONL tolerant read limit exceeded: {max_bytes} bytes"
)));
}
break;
}
let read = (&mut reader)
.take(
u64::try_from(remaining)
.unwrap_or(u64::MAX)
.saturating_add(1),
)
.read_until(b'\n', &mut line)?;
if read == 0 {
break;
}
if read > remaining {
return Err(budget_error(format!(
"session JSONL tolerant read limit exceeded at line {line_number}: {max_bytes} bytes"
)));
}
total_bytes = total_bytes.saturating_add(read);
lines_read = lines_read.saturating_add(1);
digest.update(&line);
if line.last() == Some(&b'\n') {
line.pop();
if line.last() == Some(&b'\r') {
line.pop();
}
}
match serde_json::from_slice::<SessionEvent>(&line) {
Ok(event) => visit(event),
Err(_) => push_tolerant_read_diagnostic(
&mut diagnostics,
&mut omitted_diagnostics,
SessionReadDiagnostic {
line: line_number,
message: format!(
"operation=replay category=session_jsonl failed to parse session JSONL at line {line_number}"
),
},
),
}
}
if !reader.fill_buf()?.is_empty() {
return Err(budget_error(format!(
"session JSONL tolerant read limit exceeded: more than {max_lines} lines or {max_bytes} bytes"
)));
}
let diagnostics = finalize_tolerant_read_diagnostics(diagnostics, omitted_diagnostics);
Ok((
diagnostics,
total_bytes,
SessionReadStats {
lines_read,
bytes_read: total_bytes,
content_digest: digest.finalize().into(),
},
))
}
pub(crate) fn read_recent_events_tolerant(
&self,
max_events: usize,
max_bytes: usize,
) -> anyhow::Result<TolerantSessionEvents> {
validate_session_id(self.id.clone())?;
if !self.path.exists() || max_events == 0 || max_bytes == 0 {
return Ok(TolerantSessionEvents {
events: Vec::new(),
diagnostics: Vec::new(),
cutoff_bytes: 0,
});
}
let file = open_session_file(self)?;
let lines = BufReader::new(file)
.split(b'\n')
.enumerate()
.map(|(index, line)| (index + 1, line));
self.collect_recent_events_tolerant_lines(lines, max_events, max_bytes)
}
pub(crate) fn read_recent_events_tolerant_tail(
&self,
max_events: usize,
max_retained_bytes: usize,
max_read_bytes: usize,
) -> anyhow::Result<TolerantSessionEvents> {
validate_session_id(self.id.clone())?;
if max_events == 0 || max_retained_bytes == 0 || max_read_bytes == 0 {
anyhow::bail!(
"session JSONL tolerant tail read limits must be non-zero: max_events={max_events}, max_retained_bytes={max_retained_bytes}, max_read_bytes={max_read_bytes}"
);
}
if !self.path.exists() {
return Ok(TolerantSessionEvents {
events: Vec::new(),
diagnostics: Vec::new(),
cutoff_bytes: 0,
});
}
let mut file = open_session_file(self)?;
let file_len = file.metadata()?.len();
let read_bytes = file_len.min(u64::try_from(max_read_bytes).unwrap_or(u64::MAX));
let start = file_len - read_bytes;
file.seek(SeekFrom::Start(start))?;
let mut tail = Vec::with_capacity(read_bytes as usize);
file.take(read_bytes).read_to_end(&mut tail)?;
let first_complete_line = if start == 0 {
0
} else {
tail.iter()
.position(|byte| *byte == b'\n')
.map_or(tail.len(), |index| index + 1)
};
let tail = &tail[first_complete_line..];
let mut line_count = 0;
let lines = tail
.split_inclusive(|byte| *byte == b'\n')
.enumerate()
.inspect(|_| line_count += 1)
.map(|(index, line)| {
let line = line.strip_suffix(b"\n").unwrap_or(line);
(index + 1, Ok(line.to_vec()))
});
let mut result =
self.collect_recent_events_tolerant_lines(lines, max_events, max_retained_bytes)?;
if start > 0 || (line_count > result.events.len() && result.diagnostics.is_empty()) {
result.diagnostics.insert(
0,
SessionReadDiagnostic {
line: 0,
message: format!(
"operation=recent_context category=session_jsonl bounded tail window; omitted older lines; read final {max_read_bytes} of {file_len} bytes"
),
},
);
}
Ok(result)
}
fn collect_recent_events_tolerant_lines(
&self,
lines: impl IntoIterator<Item = (usize, Result<Vec<u8>, std::io::Error>)>,
max_events: usize,
max_bytes: usize,
) -> anyhow::Result<TolerantSessionEvents> {
let mut retained = VecDeque::new();
let mut retained_bytes = 0usize;
let mut malformed_bytes = 0usize;
let mut diagnostics = Vec::new();
let mut omitted_diagnostics = 0usize;
for (line_number, line) in lines {
match line {
Ok(line) => {
let line_bytes = line.len() + 1;
match serde_json::from_slice::<SessionEvent>(&line) {
Ok(event) => {
retained_bytes = retained_bytes.saturating_add(line_bytes);
retained.push_back((event, line_bytes));
while retained.len() > max_events || retained_bytes > max_bytes {
if let Some((_, bytes)) = retained.pop_front() {
retained_bytes = retained_bytes.saturating_sub(bytes);
} else {
break;
}
}
}
Err(_) => {
malformed_bytes = malformed_bytes.saturating_add(line_bytes);
if malformed_bytes > max_bytes {
anyhow::bail!(
"operation=recent_context category=session_jsonl tolerant recent read limit exceeded: malformed bytes > {max_bytes}"
);
}
push_tolerant_read_diagnostic(
&mut diagnostics,
&mut omitted_diagnostics,
SessionReadDiagnostic {
line: line_number,
message: format!(
"operation=recent_context category=session_jsonl failed to parse session JSONL at line {line_number}"
),
},
);
}
}
}
Err(_) => push_tolerant_read_diagnostic(
&mut diagnostics,
&mut omitted_diagnostics,
SessionReadDiagnostic {
line: line_number,
message: format!(
"operation=recent_context category=session_jsonl read failure at line {line_number}"
),
},
),
}
}
Ok(TolerantSessionEvents {
events: retained.into_iter().map(|(event, _)| event).collect(),
diagnostics: finalize_tolerant_read_diagnostics(diagnostics, omitted_diagnostics),
cutoff_bytes: 0,
})
}
}