1use crate::events::{AgentEvent, ContextEvent, LlmCallOutcome, MessageEvent, ModelEvent, ToolEvent, TurnEvent};
2use serde::{Deserialize, Serialize};
3use std::fs::File;
4use std::io::{BufRead, BufReader};
5use std::path::{Path, PathBuf};
6
7#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
8#[serde(rename_all = "camelCase")]
9pub struct SessionMeta {
10 pub session_id: String,
11 pub cwd: PathBuf,
12 pub model: String,
13 #[serde(default)]
14 pub selected_mode: Option<String>,
15 pub created_at: String,
16}
17
18#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
19#[serde(tag = "type", rename_all = "camelCase")]
20pub enum UserEvent {
21 Message { content: Vec<llm::ContentBlock> },
22 ClearContext,
23}
24
25#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
26#[serde(tag = "type", rename_all = "camelCase")]
27pub enum SessionControlEvent {
28 AgentSwitched { from: Option<String>, to: Option<String> },
29}
30
31#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
32#[serde(tag = "kind", content = "data", rename_all = "camelCase")]
33#[allow(clippy::large_enum_variant)]
34pub enum SessionEvent {
35 User(UserEvent),
36 Agent(AgentEvent),
37 Control(SessionControlEvent),
38}
39
40impl SessionEvent {
41 pub fn content(&self) -> Option<String> {
42 self.user_content().or_else(|| match self {
43 Self::Agent(event) => event.content(),
44 Self::User(_) | Self::Control(_) => None,
45 })
46 }
47
48 pub fn user_content(&self) -> Option<String> {
49 match self {
50 Self::User(UserEvent::Message { content }) => {
51 let text = llm::ContentBlock::join_text(content);
52 (!text.is_empty()).then_some(text)
53 }
54 Self::User(UserEvent::ClearContext) | Self::Agent(_) | Self::Control(_) => None,
55 }
56 }
57
58 pub fn is_persisted(&self) -> bool {
59 match self {
60 Self::User(_) | Self::Control(_) => true,
61 Self::Agent(event) => match event {
62 AgentEvent::Message(
63 MessageEvent::Text { is_complete, .. } | MessageEvent::Thought { is_complete, .. },
64 ) => *is_complete,
65 AgentEvent::Tool(
66 ToolEvent::Call { .. }
67 | ToolEvent::Result { .. }
68 | ToolEvent::Error { .. }
69 | ToolEvent::TaskCreated { .. }
70 | ToolEvent::TaskCompleted { .. }
71 | ToolEvent::TaskFailed { .. }
72 | ToolEvent::TaskCancelled { .. },
73 )
74 | AgentEvent::Turn(
75 TurnEvent::RetryScheduled { .. }
76 | TurnEvent::AutoContinue { .. }
77 | TurnEvent::Ended { .. }
78 | TurnEvent::LlmCallEnded { outcome: LlmCallOutcome::Failed { .. }, .. },
79 )
80 | AgentEvent::Context(
81 ContextEvent::CompactionStarted { .. }
82 | ContextEvent::CompactionEnded { .. }
83 | ContextEvent::CompactionResult { .. }
84 | ContextEvent::UsageUpdated { .. }
85 | ContextEvent::Cleared,
86 )
87 | AgentEvent::Model(ModelEvent::Switched { .. }) => true,
88 AgentEvent::Tool(
89 ToolEvent::CallUpdate { .. }
90 | ToolEvent::ExecutionStarted { .. }
91 | ToolEvent::Progress { .. }
92 | ToolEvent::TaskStatus { .. }
93 | ToolEvent::DefinitionsUpdated { .. },
94 )
95 | AgentEvent::Turn(
96 TurnEvent::Started { .. }
97 | TurnEvent::LlmCallStarted { .. }
98 | TurnEvent::LlmCallEnded {
99 outcome: LlmCallOutcome::Completed { .. } | LlmCallOutcome::Cancelled,
100 ..
101 },
102 ) => false,
103 },
104 }
105 }
106}
107
108#[derive(Debug, Clone)]
109pub struct SessionLine {
110 pub line_number: usize,
111 pub bytes_read: usize,
112 pub raw: String,
113}
114
115#[derive(Debug)]
116pub enum SessionLogEntry {
117 Persisted { line: SessionLine, event: Box<SessionEvent> },
118 Transient { line: SessionLine },
119 Malformed { line: SessionLine, error: serde_json::Error },
120}
121
122impl SessionLogEntry {
123 pub fn line(&self) -> &SessionLine {
124 match self {
125 Self::Persisted { line, .. } | Self::Transient { line } | Self::Malformed { line, .. } => line,
126 }
127 }
128}
129
130#[derive(Debug, thiserror::Error)]
131pub enum SessionLogError {
132 #[error(transparent)]
133 Io(#[from] std::io::Error),
134 #[error("missing session metadata line")]
135 MissingMetadata,
136 #[error("invalid session metadata on line {line_number}: {source}")]
137 InvalidMetadata { line_number: usize, source: serde_json::Error },
138}
139
140pub struct SessionLog<T: BufRead> {
141 reader: T,
142 pub meta: SessionMeta,
143 line_number: usize,
144}
145
146impl SessionLog<BufReader<File>> {
147 pub fn open(path: impl AsRef<Path>) -> Result<Self, SessionLogError> {
148 Self::from_reader(BufReader::new(File::open(path.as_ref())?))
149 }
150}
151
152impl<T: BufRead> SessionLog<T> {
153 pub fn from_reader(mut reader: T) -> Result<Self, SessionLogError> {
154 let mut line = String::new();
155 let mut line_number = 0;
156 loop {
157 line.clear();
158 if reader.read_line(&mut line)? == 0 {
159 return Err(SessionLogError::MissingMetadata);
160 }
161 line_number += 1;
162 if !line.trim().is_empty() {
163 break;
164 }
165 }
166 let meta = serde_json::from_str(line.trim())
167 .map_err(|source| SessionLogError::InvalidMetadata { line_number, source })?;
168 Ok(Self { reader, meta, line_number })
169 }
170
171 pub fn next_entry(&mut self) -> std::io::Result<Option<SessionLogEntry>> {
172 let Some(line) = self.next_line()? else {
173 return Ok(None);
174 };
175 let entry = match serde_json::from_str::<SessionEvent>(&line.raw) {
176 Ok(event) if event.is_persisted() => SessionLogEntry::Persisted { line, event: Box::new(event) },
177 Ok(_) => SessionLogEntry::Transient { line },
178 Err(error) => SessionLogEntry::Malformed { line, error },
179 };
180 Ok(Some(entry))
181 }
182
183 fn next_line(&mut self) -> std::io::Result<Option<SessionLine>> {
184 let mut line = String::new();
185 loop {
186 line.clear();
187 let bytes_read = self.reader.read_line(&mut line)?;
188 if bytes_read == 0 {
189 return Ok(None);
190 }
191 self.line_number += 1;
192 let trimmed = line.trim();
193 if !trimmed.is_empty() {
194 return Ok(Some(SessionLine { line_number: self.line_number, bytes_read, raw: trimmed.to_string() }));
195 }
196 }
197 }
198}
199
200pub fn last_agent_from_events(initial: Option<String>, events: &[SessionEvent]) -> Option<String> {
201 events
202 .iter()
203 .rev()
204 .find_map(|event| match event {
205 SessionEvent::Control(SessionControlEvent::AgentSwitched { to, .. }) => Some(to.clone()),
206 _ => None,
207 })
208 .unwrap_or(initial)
209}
210
211#[cfg(test)]
212mod tests {
213 use super::*;
214 use crate::events::{LlmCallPurpose, StreamState, TurnOutcome};
215
216 fn agent(event: AgentEvent) -> SessionEvent {
217 SessionEvent::Agent(event)
218 }
219
220 #[test]
221 fn persistence_policy_covers_every_event_variant() {
222 let retry = agent(AgentEvent::Turn(TurnEvent::RetryScheduled {
223 purpose: LlmCallPurpose::Chat,
224 attempt: 1,
225 max_attempts: 3,
226 delay_ms: 10,
227 }));
228 let cancelled = agent(AgentEvent::Turn(TurnEvent::Ended { outcome: TurnOutcome::Cancelled }));
229 let partial = agent(AgentEvent::text("m", "partial", StreamState::Partial));
230 let compaction_ended = agent(AgentEvent::Context(ContextEvent::CompactionEnded {
231 outcome: crate::events::CompactionOutcome::Completed,
232 }));
233
234 assert!(retry.is_persisted());
235 assert!(cancelled.is_persisted());
236 assert!(compaction_ended.is_persisted());
237 assert!(!partial.is_persisted());
238 }
239}