use std::io::{Read, Write};
use std::os::unix::fs::OpenOptionsExt;
use std::path::PathBuf;
use serde::{Deserialize, Serialize};
use crate::log::fields;
use crate::log::{Logger, now_ms};
use crate::session::event::SessionEvent;
use crate::session::views::Recorder;
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct Journaled {
pub at: i64,
#[serde(skip_serializing_if = "Option::is_none")]
pub turn: Option<u32>,
pub entry: SessionEvent,
}
pub const TRANSCRIPT_FILENAME: &str = "transcript.jsonl";
pub const MAX_REPLAYED: usize = 2_000;
pub const OPENING_SCAN_BYTES: usize = 64 * 1024;
#[derive(Debug, Clone, PartialEq)]
pub struct StoredTranscript {
pub entries: Vec<Journaled>,
pub dropped: usize,
}
fn journaled(text: &str) -> Vec<Journaled> {
text.split('\n')
.filter(|line| !line.trim().is_empty())
.filter_map(|line| serde_json::from_str(line).ok())
.collect()
}
pub struct Transcript {
path: PathBuf,
log: Option<Logger>,
}
impl Transcript {
pub fn new(path: impl Into<PathBuf>, log: Option<Logger>) -> Self {
Self {
path: path.into(),
log,
}
}
pub fn append_at(&self, entry: &SessionEvent, turn: Option<u32>, at: i64) {
let line = serde_json::to_string(&Journaled {
at,
turn,
entry: entry.clone(),
})
.expect("a transcript entry serializes");
let written = std::fs::OpenOptions::new()
.create(true)
.append(true)
.mode(0o600)
.open(&self.path)
.and_then(|mut file| file.write_all(format!("{line}\n").as_bytes()));
if let Err(error) = written
&& let Some(log) = &self.log
{
log.warn(
"a transcript entry could not be written",
&fields([("detail", error.to_string().into())]),
);
}
}
pub fn opening(&self) -> Option<String> {
let head = self.read_head(OPENING_SCAN_BYTES).ok()?;
journaled(&head)
.into_iter()
.find_map(|held| match held.entry {
SessionEvent::Prompt { text, .. } => Some(text.trim().to_owned()),
_ => None,
})
}
pub fn read(&self) -> StoredTranscript {
self.read_up_to(MAX_REPLAYED)
}
pub fn read_up_to(&self, limit: usize) -> StoredTranscript {
let Ok(text) = std::fs::read_to_string(&self.path) else {
return StoredTranscript {
entries: Vec::new(),
dropped: 0,
};
};
let mut parsed = journaled(&text);
if parsed.len() <= limit {
return StoredTranscript {
entries: parsed,
dropped: 0,
};
}
let dropped = parsed.len() - limit;
StoredTranscript {
entries: parsed.split_off(dropped),
dropped,
}
}
fn read_head(&self, bytes: usize) -> std::io::Result<String> {
let mut file = std::fs::File::open(&self.path)?;
let mut buffer = vec![0; bytes];
let read = file.read(&mut buffer)?;
buffer.truncate(read);
let text = String::from_utf8_lossy(&buffer);
let cut = text.rfind('\n').unwrap_or(text.len());
Ok(text[..cut].to_owned())
}
}
impl Recorder for Transcript {
fn append(&self, entry: &SessionEvent, turn: u32) {
self.append_at(entry, Some(turn), now_ms());
}
}
#[cfg(test)]
mod tests;