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(ToolEvent::Call { .. } | ToolEvent::Result { .. } | ToolEvent::Error { .. })
66                | AgentEvent::Turn(
67                    TurnEvent::RetryScheduled { .. }
68                    | TurnEvent::AutoContinue { .. }
69                    | TurnEvent::Ended { .. }
70                    | TurnEvent::LlmCallEnded { outcome: LlmCallOutcome::Failed { .. }, .. },
71                )
72                | AgentEvent::Context(
73                    ContextEvent::CompactionStarted { .. }
74                    | ContextEvent::CompactionEnded { .. }
75                    | ContextEvent::CompactionResult { .. }
76                    | ContextEvent::UsageUpdated { .. }
77                    | ContextEvent::Cleared,
78                )
79                | AgentEvent::Model(ModelEvent::Switched { .. }) => true,
80                AgentEvent::Tool(
81                    ToolEvent::CallUpdate { .. }
82                    | ToolEvent::ExecutionStarted { .. }
83                    | ToolEvent::Progress { .. }
84                    | ToolEvent::DefinitionsUpdated { .. },
85                )
86                | AgentEvent::Turn(
87                    TurnEvent::Started { .. }
88                    | TurnEvent::LlmCallStarted { .. }
89                    | TurnEvent::LlmCallEnded {
90                        outcome: LlmCallOutcome::Completed { .. } | LlmCallOutcome::Cancelled,
91                        ..
92                    },
93                ) => false,
94            },
95        }
96    }
97}
98
99#[derive(Debug, Clone)]
100pub struct SessionLine {
101    pub line_number: usize,
102    pub bytes_read: usize,
103    pub raw: String,
104}
105
106#[derive(Debug)]
107pub enum SessionLogEntry {
108    Persisted { line: SessionLine, event: Box<SessionEvent> },
109    Transient { line: SessionLine },
110    Malformed { line: SessionLine, error: serde_json::Error },
111}
112
113impl SessionLogEntry {
114    pub fn line(&self) -> &SessionLine {
115        match self {
116            Self::Persisted { line, .. } | Self::Transient { line } | Self::Malformed { line, .. } => line,
117        }
118    }
119}
120
121#[derive(Debug, thiserror::Error)]
122pub enum SessionLogError {
123    #[error(transparent)]
124    Io(#[from] std::io::Error),
125    #[error("missing session metadata line")]
126    MissingMetadata,
127    #[error("invalid session metadata on line {line_number}: {source}")]
128    InvalidMetadata { line_number: usize, source: serde_json::Error },
129}
130
131pub struct SessionLog<T: BufRead> {
132    reader: T,
133    pub meta: SessionMeta,
134    line_number: usize,
135}
136
137impl SessionLog<BufReader<File>> {
138    pub fn open(path: impl AsRef<Path>) -> Result<Self, SessionLogError> {
139        Self::from_reader(BufReader::new(File::open(path.as_ref())?))
140    }
141}
142
143impl<T: BufRead> SessionLog<T> {
144    pub fn from_reader(mut reader: T) -> Result<Self, SessionLogError> {
145        let mut line = String::new();
146        let mut line_number = 0;
147        loop {
148            line.clear();
149            if reader.read_line(&mut line)? == 0 {
150                return Err(SessionLogError::MissingMetadata);
151            }
152            line_number += 1;
153            if !line.trim().is_empty() {
154                break;
155            }
156        }
157        let meta = serde_json::from_str(line.trim())
158            .map_err(|source| SessionLogError::InvalidMetadata { line_number, source })?;
159        Ok(Self { reader, meta, line_number })
160    }
161
162    pub fn next_entry(&mut self) -> std::io::Result<Option<SessionLogEntry>> {
163        let Some(line) = self.next_line()? else {
164            return Ok(None);
165        };
166        let entry = match serde_json::from_str::<SessionEvent>(&line.raw) {
167            Ok(event) if event.is_persisted() => SessionLogEntry::Persisted { line, event: Box::new(event) },
168            Ok(_) => SessionLogEntry::Transient { line },
169            Err(error) => SessionLogEntry::Malformed { line, error },
170        };
171        Ok(Some(entry))
172    }
173
174    fn next_line(&mut self) -> std::io::Result<Option<SessionLine>> {
175        let mut line = String::new();
176        loop {
177            line.clear();
178            let bytes_read = self.reader.read_line(&mut line)?;
179            if bytes_read == 0 {
180                return Ok(None);
181            }
182            self.line_number += 1;
183            let trimmed = line.trim();
184            if !trimmed.is_empty() {
185                return Ok(Some(SessionLine { line_number: self.line_number, bytes_read, raw: trimmed.to_string() }));
186            }
187        }
188    }
189}
190
191pub fn last_agent_from_events(initial: Option<String>, events: &[SessionEvent]) -> Option<String> {
192    events
193        .iter()
194        .rev()
195        .find_map(|event| match event {
196            SessionEvent::Control(SessionControlEvent::AgentSwitched { to, .. }) => Some(to.clone()),
197            _ => None,
198        })
199        .unwrap_or(initial)
200}
201
202#[cfg(test)]
203mod tests {
204    use super::*;
205    use crate::events::{LlmCallPurpose, TurnOutcome};
206
207    fn agent(event: AgentEvent) -> SessionEvent {
208        SessionEvent::Agent(event)
209    }
210
211    #[test]
212    fn persistence_policy_covers_every_event_variant() {
213        let retry = agent(AgentEvent::Turn(TurnEvent::RetryScheduled {
214            purpose: LlmCallPurpose::Chat,
215            attempt: 1,
216            max_attempts: 3,
217            delay_ms: 10,
218        }));
219        let cancelled = agent(AgentEvent::Turn(TurnEvent::Ended { outcome: TurnOutcome::Cancelled }));
220        let partial = agent(AgentEvent::text("m", "partial", false));
221        let compaction_ended = agent(AgentEvent::Context(ContextEvent::CompactionEnded {
222            outcome: crate::events::CompactionOutcome::Completed,
223        }));
224
225        assert!(retry.is_persisted());
226        assert!(cancelled.is_persisted());
227        assert!(compaction_ended.is_persisted());
228        assert!(!partial.is_persisted());
229    }
230}