//! Append-only JSONL event log, one per session, with monotonically //! increasing sequence numbers -- the phone's resume cursor. //! //! One line per event: `{"seq":N,"ts":...,"type":...,...}`. The writer //! assigns sequence numbers; readers replay everything after a cursor. //! Reopening an existing file continues the numbering, which is what makes //! a backend restart invisible to a phone holding a cursor. use std::fs::{File, OpenOptions}; use std::io::{BufRead, BufReader, Write}; use std::os::unix::fs::OpenOptionsExt; use std::path::Path; use anyhow::{Context, Result}; use serde::{Deserialize, Serialize}; use super::driver::{Event, SessionStatus}; /// One transcript line: an [`Event`] plus its position and time. The event /// is flattened so the wire shape stays one flat object. #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] pub struct SeqEvent { pub seq: u64, /// Epoch seconds. pub ts: f64, #[serde(flatten)] pub event: Event, } pub struct Transcript { file: File, next_seq: u64, last_status: Option, } impl Transcript { /// Opens (or creates) the log at `path`, continuing the sequence from /// the last line if one exists. pub fn open(path: &Path) -> Result { // One pass for both answers. They are wanted at the same moment by // the same caller, and reading the file twice to get them doubled // the cost of starting every session -- which is paid per session, // at the point a restart is trying to be quick. let existing = read_after(path, 0)?; let last_seq = existing.last().map(|entry| entry.seq).unwrap_or(0); let last_status = existing.iter().rev().find_map(|entry| match entry.event { Event::Status { state } => Some(state), _ => None, }); // Owner-only: a transcript is the whole conversation, including // whatever the session read, wrote, or was told. let file = OpenOptions::new() .create(true) .append(true) .mode(0o600) .open(path) .with_context(|| format!("open transcript {}", path.display()))?; Ok(Self { file, next_seq: last_seq + 1, last_status, }) } /// The state the session was last reported to be in, as of opening. /// /// Read from the file rather than assumed, because a server that has /// just restarted has been told nothing yet and this is the only thing /// it knows. Assuming idle claimed a session was waiting for you when /// it had exited hours earlier, and would now also claim it of one /// whose process is still mid-turn. /// /// `None` for a transcript that never carried a status, which is a new /// session and genuinely has no prior state. pub fn last_status(&self) -> Option { self.last_status } /// Appends `event`, assigning it the next sequence number. Flushed per /// event: each line is tiny, and the transcript is the source of truth /// a crash must not lose the tail of. pub fn append(&mut self, event: Event, ts: f64) -> Result { let entry = SeqEvent { seq: self.next_seq, ts, event, }; let mut line = serde_json::to_string(&entry).context("serialize event")?; line.push('\n'); self.file .write_all(line.as_bytes()) .context("append to transcript")?; self.next_seq += 1; Ok(entry) } } /// A window of the transcript ending just before `before`, newest-biased. /// /// The screen opens on the end of a conversation, not the start of it, and /// the end is all it can show at once. Replaying the whole file to get /// there costs one network frame per event -- on an 863-event import that /// was several seconds of messages arriving oldest-first, which reads as /// the app loading top-down because that is exactly what it was doing. /// /// `before` pages backwards for history somebody actually scrolls to; the /// file is read whole each time because a transcript is small and a /// seek-backwards reader would be a lot of machinery for a list that fits /// in memory anyway. pub fn read_window(path: &Path, before: Option, limit: usize) -> Result> { let mut all = read_after(path, 0)?; if let Some(before) = before { all.retain(|entry| entry.seq < before); } if all.len() > limit { all.drain(..all.len() - limit); } Ok(all) } /// How far behind a reconnecting subscriber can be and still be handed the /// backlog one event at a time. /// /// Past this it is served better by rebuilding its view from the newest /// window than by receiving everything it missed. The events are the same /// either way; what differs is that one arrives as a single window and the /// other as thousands of frames a screen renders one by one. Set well /// above a screenful (`transcript`'s page is 80) so an ordinary blip -- a /// phone asleep, a tunnel reconnecting, a backend restart -- still streams /// continuously, and only a genuine backlog changes mode. pub const CATCH_UP_LIMIT: usize = 200; /// What a subscriber asking for "everything after my cursor" gets back. /// /// Two answers rather than one list, because they mean different things to /// the screen holding the cursor: one continues what it already has, the /// other replaces it. Collapsing them into a list would leave the client /// splicing a window onto rows it has no way to know are no longer /// adjacent to it -- a seam that looks exactly like ordinary output. #[derive(Debug, Clone, PartialEq)] pub enum CatchUp { /// The events after the cursor, continuing what the subscriber holds. Continue(Vec), /// The subscriber was further behind than [`CATCH_UP_LIMIT`]: the /// newest window, replacing whatever it holds. Earlier history is /// still there to be paged backwards through, exactly as it is when a /// session is first opened. Restart(Vec), } /// Everything after `after`, or the newest `limit` when that is more than /// `limit` events. pub fn catch_up(path: &Path, after: u64, limit: usize) -> Result { let mut events = read_after(path, after)?; if events.len() > limit { events.drain(..events.len() - limit); return Ok(CatchUp::Restart(events)); } Ok(CatchUp::Continue(events)) } /// Replays every event with `seq > after`, oldest first. A missing file is /// an empty transcript, not an error -- the session just hasn't produced an /// event yet. pub fn read_after(path: &Path, after: u64) -> Result> { let file = match File::open(path) { Ok(file) => file, Err(err) if err.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()), Err(err) => return Err(err).with_context(|| format!("read transcript {}", path.display())), }; let mut events = Vec::new(); for line in BufReader::new(file).lines() { let line = line.context("read transcript line")?; if line.trim().is_empty() { continue; } let entry: SeqEvent = serde_json::from_str(&line) .with_context(|| format!("bad transcript line in {}", path.display()))?; if entry.seq > after { events.push(entry); } } Ok(events) } #[cfg(test)] mod tests { use super::*; use crate::session::driver::{QuestionOption, SessionStatus}; fn text(delta: &str) -> Event { Event::AssistantText { delta: delta.to_string(), } } #[test] fn assigns_increasing_seqs_and_replays_after_a_cursor() { let dir = tempfile::tempdir().expect("tempdir"); let path = dir.path().join("transcript.jsonl"); let mut transcript = Transcript::open(&path).expect("open"); assert_eq!(transcript.append(text("a"), 1.0).expect("append").seq, 1); assert_eq!(transcript.append(text("b"), 2.0).expect("append").seq, 2); assert_eq!(transcript.append(text("c"), 3.0).expect("append").seq, 3); let replay = read_after(&path, 1).expect("read"); assert_eq!(replay.len(), 2); assert_eq!(replay[0].seq, 2); assert_eq!(replay[0].event, text("b")); assert_eq!(replay[1].seq, 3); // A cursor at or past the end replays nothing. assert!(read_after(&path, 3).expect("read").is_empty()); } #[test] fn reopening_continues_the_numbering() { let dir = tempfile::tempdir().expect("tempdir"); let path = dir.path().join("transcript.jsonl"); let mut transcript = Transcript::open(&path).expect("open"); transcript.append(text("a"), 1.0).expect("append"); transcript.append(text("b"), 2.0).expect("append"); drop(transcript); let mut reopened = Transcript::open(&path).expect("reopen"); assert_eq!(reopened.append(text("c"), 3.0).expect("append").seq, 3); } #[test] fn a_short_backlog_continues_and_a_long_one_restarts() { let dir = tempfile::tempdir().expect("tempdir"); let path = dir.path().join("transcript.jsonl"); let mut transcript = Transcript::open(&path).expect("open"); for n in 0..10 { transcript .append(text(&n.to_string()), 0.0) .expect("append"); } // Within the limit the subscriber keeps what it has. let CatchUp::Continue(events) = catch_up(&path, 7, 5).expect("catch up") else { panic!("a backlog of 3 should continue"); }; assert_eq!(events.len(), 3); assert_eq!(events[0].seq, 8); // Past it, the newest window replaces what it has -- and it is the // newest, not the oldest, that survives the trim. let CatchUp::Restart(events) = catch_up(&path, 0, 5).expect("catch up") else { panic!("a backlog of 10 should restart"); }; assert_eq!(events.len(), 5); assert_eq!(events[0].seq, 6); assert_eq!(events[4].seq, 10); // Exactly at the limit is still a continuation: the boundary // belongs to the cheaper answer, so a client is not reset for // being one event behind the threshold. assert!(matches!( catch_up(&path, 5, 5).expect("catch up"), CatchUp::Continue(_) )); } #[test] fn reopening_reports_the_state_it_was_last_left_in() { let dir = tempfile::tempdir().expect("tempdir"); let path = dir.path().join("transcript.jsonl"); // Nothing recorded yet: no prior state to report, which is not the // same as reporting idle. assert_eq!(Transcript::open(&path).expect("open").last_status(), None); let mut transcript = Transcript::open(&path).expect("open"); transcript .append( Event::Status { state: SessionStatus::Running, }, 1.0, ) .expect("append"); transcript .append( Event::Status { state: SessionStatus::Exited, }, 2.0, ) .expect("append"); // Events after the last status must not hide it. transcript.append(text("trailing"), 3.0).expect("append"); drop(transcript); let reopened = Transcript::open(&path).expect("reopen"); assert_eq!(reopened.last_status(), Some(SessionStatus::Exited)); // And the same pass still continues the numbering. assert_eq!(reopened.next_seq, 4); } #[test] fn a_missing_file_reads_as_empty() { let dir = tempfile::tempdir().expect("tempdir"); assert!( read_after(&dir.path().join("nope.jsonl"), 0) .expect("read") .is_empty() ); } #[test] fn round_trips_every_event_shape() { let dir = tempfile::tempdir().expect("tempdir"); let path = dir.path().join("transcript.jsonl"); let events = vec![ Event::UserMessage { text: "hi".into() }, text("hello"), Event::ToolStart { id: "t1".into(), tool: "bash".into(), input: serde_json::json!({"command": "ls"}), }, Event::ToolUpdate { id: "t1".into(), output: "partial".into(), }, Event::ToolEnd { id: "t1".into(), output: "done".into(), }, Event::Image { image: "img1".into(), about: None, }, Event::Question { id: "q1".into(), prompt: "Allow?".into(), header: None, options: vec![QuestionOption::plain("Yes"), QuestionOption::plain("No")], multi_select: false, about: None, }, Event::Answered { id: "q1".into(), answers: vec!["Yes".into()], }, Event::Status { state: SessionStatus::Idle, }, Event::UsageDelta { tokens: 42 }, Event::Error { message: "boom".into(), }, ]; let mut transcript = Transcript::open(&path).expect("open"); for event in &events { transcript.append(event.clone(), 0.0).expect("append"); } let replayed: Vec = read_after(&path, 0) .expect("read") .into_iter() .map(|entry| entry.event) .collect(); assert_eq!(replayed, events); } }