1use super::event::SessionEvent;
2use super::manager::Session;
3use sha2::{Digest, Sha256};
4use std::{
5 collections::VecDeque,
6 fs,
7 io::{BufRead, BufReader, Read, Seek, SeekFrom},
8 path::PathBuf,
9};
10pub fn validate_session_id(id: String) -> anyhow::Result<String> {
11 if id.is_empty() {
12 anyhow::bail!("session id must not be empty");
13 }
14 if id == "." || id == ".." || id.contains("..") {
15 anyhow::bail!("session id must not contain '..'");
16 }
17 if id.contains('/') || id.contains('\\') {
18 anyhow::bail!("session id must not contain path separators");
19 }
20 if PathBuf::from(&id).is_absolute() {
21 anyhow::bail!("session id must not be an absolute path");
22 }
23 if !id
24 .chars()
25 .all(|character| character.is_ascii_alphanumeric() || character == '_' || character == '-')
26 {
27 anyhow::bail!("session id must match [A-Za-z0-9_-]+");
28 }
29 Ok(id)
30}
31
32fn open_session_file(session: &Session) -> anyhow::Result<fs::File> {
33 let root = session
34 .path
35 .parent()
36 .ok_or_else(|| anyhow::anyhow!("session file has no parent"))?;
37 super::store::open_existing_primary(root, &session.id)?
38 .ok_or_else(|| anyhow::anyhow!("session JSONL is missing"))
39}
40
41#[derive(Debug)]
42pub(crate) enum BoundedReadError {
43 BudgetExceeded(String),
44}
45
46impl std::fmt::Display for BoundedReadError {
47 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
48 match self {
49 Self::BudgetExceeded(message) => formatter.write_str(message),
50 }
51 }
52}
53
54impl std::error::Error for BoundedReadError {}
55
56fn budget_error(message: impl Into<String>) -> anyhow::Error {
57 anyhow::Error::new(BoundedReadError::BudgetExceeded(message.into()))
58}
59#[derive(Debug, Clone, PartialEq, Eq)]
60pub(crate) struct SessionReadDiagnostic {
61 pub(crate) line: usize,
62 pub(crate) message: String,
63}
64
65#[derive(Debug, Clone, PartialEq)]
66pub(crate) struct TolerantSessionEvents {
67 pub(crate) events: Vec<SessionEvent>,
68 pub(crate) diagnostics: Vec<SessionReadDiagnostic>,
69 pub(crate) cutoff_bytes: u64,
70}
71
72#[derive(Debug, Clone, Copy, PartialEq, Eq)]
73pub(crate) struct SessionReadStats {
74 pub lines_read: usize,
75 pub bytes_read: usize,
76 pub content_digest: [u8; 32],
77}
78
79const MAX_TOLERANT_READ_DIAGNOSTICS: usize = 64;
80
81fn push_tolerant_read_diagnostic(
82 diagnostics: &mut Vec<SessionReadDiagnostic>,
83 omitted_count: &mut usize,
84 diagnostic: SessionReadDiagnostic,
85) {
86 if diagnostics.len() < MAX_TOLERANT_READ_DIAGNOSTICS {
87 diagnostics.push(diagnostic);
88 } else {
89 *omitted_count = omitted_count.saturating_add(1);
90 }
91}
92
93fn finalize_tolerant_read_diagnostics(
94 mut diagnostics: Vec<SessionReadDiagnostic>,
95 omitted_count: usize,
96) -> Vec<SessionReadDiagnostic> {
97 if omitted_count == 0 {
98 return diagnostics;
99 }
100 if diagnostics.len() == MAX_TOLERANT_READ_DIAGNOSTICS {
101 diagnostics.pop();
102 }
103 diagnostics.push(SessionReadDiagnostic {
104 line: 0,
105 message: format!(
106 "omitted {omitted_count} additional session JSONL diagnostics after cap of {MAX_TOLERANT_READ_DIAGNOSTICS}"
107 ),
108 });
109 diagnostics
110}
111
112pub(crate) const MAX_METADATA_VISIT_LINES: usize = 100_000;
113pub(crate) const MAX_METADATA_VISIT_BYTES: usize = 64 * 1024 * 1024;
114
115impl Session {
116 pub(crate) fn read_events_tolerant_bounded(
117 &self,
118 max_lines: usize,
119 max_bytes: usize,
120 ) -> anyhow::Result<TolerantSessionEvents> {
121 Ok(self
122 .read_events_tolerant_bounded_with_stats(max_lines, max_bytes)?
123 .0)
124 }
125
126 pub(crate) fn read_events_tolerant_bounded_with_stats(
127 &self,
128 max_lines: usize,
129 max_bytes: usize,
130 ) -> anyhow::Result<(TolerantSessionEvents, SessionReadStats)> {
131 let mut events = Vec::new();
132 let (diagnostics, cutoff_bytes, stats) =
133 self.visit_events_tolerant_bounded_with_stats(max_lines, max_bytes, |event| {
134 events.push(event)
135 })?;
136 Ok((
137 TolerantSessionEvents {
138 events,
139 diagnostics,
140 cutoff_bytes: cutoff_bytes as u64,
141 },
142 stats,
143 ))
144 }
145
146 pub(crate) fn history_bytes(&self) -> anyhow::Result<u64> {
147 let root = self
148 .path
149 .parent()
150 .ok_or_else(|| anyhow::anyhow!("session file has no parent"))?;
151 let Some(file) = super::store::open_existing_primary(root, &self.id)? else {
152 return Ok(0);
153 };
154 Ok(file.metadata()?.len())
155 }
156
157 pub(crate) fn content_fingerprint_bounded(
158 &self,
159 max_bytes: usize,
160 ) -> anyhow::Result<([u8; 32], usize)> {
161 validate_session_id(self.id.clone())?;
162 if !self.path.exists() {
163 return Ok((Sha256::digest([]).into(), 0));
164 }
165 let file = open_session_file(self)?;
166 let mut reader = file.take(
167 u64::try_from(max_bytes)
168 .unwrap_or(u64::MAX)
169 .saturating_add(1),
170 );
171 let mut digest = Sha256::new();
172 let mut total_bytes = 0usize;
173 let mut buffer = [0u8; 8192];
174 loop {
175 let read = reader.read(&mut buffer)?;
176 if read == 0 {
177 break;
178 }
179 total_bytes = total_bytes.saturating_add(read);
180 if total_bytes > max_bytes {
181 return Err(budget_error(format!(
182 "session JSONL tolerant read limit exceeded: {max_bytes} bytes"
183 )));
184 }
185 digest.update(&buffer[..read]);
186 }
187 Ok((digest.finalize().into(), total_bytes))
188 }
189
190 pub(crate) fn visit_events_tolerant_bounded(
191 &self,
192 max_lines: usize,
193 max_bytes: usize,
194 visit: impl FnMut(SessionEvent),
195 ) -> anyhow::Result<(Vec<SessionReadDiagnostic>, usize)> {
196 let (diagnostics, bytes, _) =
197 self.visit_events_tolerant_bounded_with_stats(max_lines, max_bytes, visit)?;
198 Ok((diagnostics, bytes))
199 }
200
201 fn visit_events_tolerant_bounded_with_stats(
202 &self,
203 max_lines: usize,
204 max_bytes: usize,
205 mut visit: impl FnMut(SessionEvent),
206 ) -> anyhow::Result<(Vec<SessionReadDiagnostic>, usize, SessionReadStats)> {
207 validate_session_id(self.id.clone())?;
208 if !self.path.exists() {
209 return Ok((
210 Vec::new(),
211 0,
212 SessionReadStats {
213 lines_read: 0,
214 bytes_read: 0,
215 content_digest: Sha256::digest([]).into(),
216 },
217 ));
218 }
219 let file = open_session_file(self)?;
220 let mut reader = BufReader::new(file);
221 let mut diagnostics = Vec::new();
222 let mut omitted_diagnostics = 0usize;
223 let mut total_bytes = 0usize;
224 let mut lines_read = 0usize;
225 let mut digest = Sha256::new();
226 let mut line = Vec::new();
227 for line_number in 1..=max_lines {
228 line.clear();
229 let remaining = max_bytes.saturating_sub(total_bytes);
230 if remaining == 0 {
231 if !reader.fill_buf()?.is_empty() {
232 return Err(budget_error(format!(
233 "session JSONL tolerant read limit exceeded: {max_bytes} bytes"
234 )));
235 }
236 break;
237 }
238 let read = (&mut reader)
239 .take(
240 u64::try_from(remaining)
241 .unwrap_or(u64::MAX)
242 .saturating_add(1),
243 )
244 .read_until(b'\n', &mut line)?;
245 if read == 0 {
246 break;
247 }
248 if read > remaining {
249 return Err(budget_error(format!(
250 "session JSONL tolerant read limit exceeded at line {line_number}: {max_bytes} bytes"
251 )));
252 }
253 total_bytes = total_bytes.saturating_add(read);
254 lines_read = lines_read.saturating_add(1);
255 digest.update(&line);
256 if line.last() == Some(&b'\n') {
257 line.pop();
258 if line.last() == Some(&b'\r') {
259 line.pop();
260 }
261 }
262 match serde_json::from_slice::<SessionEvent>(&line) {
263 Ok(event) => visit(event),
264 Err(_) => push_tolerant_read_diagnostic(
265 &mut diagnostics,
266 &mut omitted_diagnostics,
267 SessionReadDiagnostic {
268 line: line_number,
269 message: format!(
270 "operation=replay category=session_jsonl failed to parse session JSONL at line {line_number}"
271 ),
272 },
273 ),
274 }
275 }
276 if !reader.fill_buf()?.is_empty() {
277 return Err(budget_error(format!(
278 "session JSONL tolerant read limit exceeded: more than {max_lines} lines or {max_bytes} bytes"
279 )));
280 }
281 let diagnostics = finalize_tolerant_read_diagnostics(diagnostics, omitted_diagnostics);
282 Ok((
283 diagnostics,
284 total_bytes,
285 SessionReadStats {
286 lines_read,
287 bytes_read: total_bytes,
288 content_digest: digest.finalize().into(),
289 },
290 ))
291 }
292
293 pub(crate) fn read_recent_events_tolerant(
294 &self,
295 max_events: usize,
296 max_bytes: usize,
297 ) -> anyhow::Result<TolerantSessionEvents> {
298 validate_session_id(self.id.clone())?;
299 if !self.path.exists() || max_events == 0 || max_bytes == 0 {
300 return Ok(TolerantSessionEvents {
301 events: Vec::new(),
302 diagnostics: Vec::new(),
303 cutoff_bytes: 0,
304 });
305 }
306 let file = open_session_file(self)?;
307 let lines = BufReader::new(file)
308 .split(b'\n')
309 .enumerate()
310 .map(|(index, line)| (index + 1, line));
311 self.collect_recent_events_tolerant_lines(lines, max_events, max_bytes)
312 }
313
314 pub(crate) fn read_recent_events_tolerant_tail(
315 &self,
316 max_events: usize,
317 max_retained_bytes: usize,
318 max_read_bytes: usize,
319 ) -> anyhow::Result<TolerantSessionEvents> {
320 validate_session_id(self.id.clone())?;
321 if max_events == 0 || max_retained_bytes == 0 || max_read_bytes == 0 {
322 anyhow::bail!(
323 "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}"
324 );
325 }
326 if !self.path.exists() {
327 return Ok(TolerantSessionEvents {
328 events: Vec::new(),
329 diagnostics: Vec::new(),
330 cutoff_bytes: 0,
331 });
332 }
333 let mut file = open_session_file(self)?;
336 let file_len = file.metadata()?.len();
337 let read_bytes = file_len.min(u64::try_from(max_read_bytes).unwrap_or(u64::MAX));
338 let start = file_len - read_bytes;
339 file.seek(SeekFrom::Start(start))?;
340 let mut tail = Vec::with_capacity(read_bytes as usize);
341 file.take(read_bytes).read_to_end(&mut tail)?;
342
343 let first_complete_line = if start == 0 {
344 0
345 } else {
346 tail.iter()
347 .position(|byte| *byte == b'\n')
348 .map_or(tail.len(), |index| index + 1)
349 };
350 let tail = &tail[first_complete_line..];
351 let mut line_count = 0;
352 let lines = tail
353 .split_inclusive(|byte| *byte == b'\n')
354 .enumerate()
355 .inspect(|_| line_count += 1)
356 .map(|(index, line)| {
357 let line = line.strip_suffix(b"\n").unwrap_or(line);
358 (index + 1, Ok(line.to_vec()))
359 });
360 let mut result =
361 self.collect_recent_events_tolerant_lines(lines, max_events, max_retained_bytes)?;
362 if start > 0 || (line_count > result.events.len() && result.diagnostics.is_empty()) {
363 result.diagnostics.insert(
364 0,
365 SessionReadDiagnostic {
366 line: 0,
367 message: format!(
368 "operation=recent_context category=session_jsonl bounded tail window; omitted older lines; read final {max_read_bytes} of {file_len} bytes"
369 ),
370 },
371 );
372 }
373 Ok(result)
374 }
375
376 fn collect_recent_events_tolerant_lines(
377 &self,
378 lines: impl IntoIterator<Item = (usize, Result<Vec<u8>, std::io::Error>)>,
379 max_events: usize,
380 max_bytes: usize,
381 ) -> anyhow::Result<TolerantSessionEvents> {
382 let mut retained = VecDeque::new();
383 let mut retained_bytes = 0usize;
384 let mut malformed_bytes = 0usize;
385 let mut diagnostics = Vec::new();
386 let mut omitted_diagnostics = 0usize;
387 for (line_number, line) in lines {
388 match line {
389 Ok(line) => {
390 let line_bytes = line.len() + 1;
391 match serde_json::from_slice::<SessionEvent>(&line) {
392 Ok(event) => {
393 retained_bytes = retained_bytes.saturating_add(line_bytes);
394 retained.push_back((event, line_bytes));
395 while retained.len() > max_events || retained_bytes > max_bytes {
396 if let Some((_, bytes)) = retained.pop_front() {
397 retained_bytes = retained_bytes.saturating_sub(bytes);
398 } else {
399 break;
400 }
401 }
402 }
403 Err(_) => {
404 malformed_bytes = malformed_bytes.saturating_add(line_bytes);
405 if malformed_bytes > max_bytes {
406 anyhow::bail!(
407 "operation=recent_context category=session_jsonl tolerant recent read limit exceeded: malformed bytes > {max_bytes}"
408 );
409 }
410 push_tolerant_read_diagnostic(
411 &mut diagnostics,
412 &mut omitted_diagnostics,
413 SessionReadDiagnostic {
414 line: line_number,
415 message: format!(
416 "operation=recent_context category=session_jsonl failed to parse session JSONL at line {line_number}"
417 ),
418 },
419 );
420 }
421 }
422 }
423 Err(_) => push_tolerant_read_diagnostic(
424 &mut diagnostics,
425 &mut omitted_diagnostics,
426 SessionReadDiagnostic {
427 line: line_number,
428 message: format!(
429 "operation=recent_context category=session_jsonl read failure at line {line_number}"
430 ),
431 },
432 ),
433 }
434 }
435 Ok(TolerantSessionEvents {
436 events: retained.into_iter().map(|(event, _)| event).collect(),
437 diagnostics: finalize_tolerant_read_diagnostics(diagnostics, omitted_diagnostics),
438 cutoff_bytes: 0,
439 })
440 }
441}