Skip to main content

magi_code/sessions/
read.rs

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        // Pin the byte window on the same validated handle. In particular, a file that was
334        // small at stat time must not fall back to an unbounded read while a writer appends.
335        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}