Skip to main content

aether_core/
session.rs

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}