1use crate::{Session, SessionEntry, SessionError, SessionMetadata};
2use chrono::Utc;
3use std::fs::{self, OpenOptions};
4use std::io::{BufRead, BufReader, Read, Seek, SeekFrom, Write};
5use std::path::Path;
6use talos_core::message::{AgentEvent, Message};
7use uuid::Uuid;
8
9impl Session {
10 pub fn append(&self, message: &Message) -> Result<(), SessionError> {
11 self.append_with_metadata(message, SessionMetadata::default())
12 }
13
14 pub fn append_with_metadata(
15 &self,
16 message: &Message,
17 metadata: SessionMetadata,
18 ) -> Result<(), SessionError> {
19 let (role, content) = message_parts(message);
20 let entry = self.build_entry(&role, &content, metadata)?;
21 self.append_entry_locked(&entry)
22 }
23
24 pub fn append_event(&self, event: &AgentEvent) -> Result<(), SessionError> {
25 let content =
26 serde_json::to_string(event).map_err(|e| SessionError::InvalidJson(e.to_string()))?;
27 let entry = self.build_entry("system", &content, SessionMetadata::default())?;
28 self.append_entry_locked(&entry)
29 }
30
31 fn build_entry(
32 &self,
33 role: &str,
34 content: &str,
35 metadata: SessionMetadata,
36 ) -> Result<SessionEntry, SessionError> {
37 let parent_id = {
38 let guard = self
39 .last_entry_id
40 .lock()
41 .expect("last_entry_id mutex poisoned");
42 if guard.is_none() {
43 drop(guard);
44 let id = read_last_entry_id(&self.file_path);
45 *self
46 .last_entry_id
47 .lock()
48 .expect("last_entry_id mutex poisoned") = id.clone();
49 id
50 } else {
51 guard.clone()
52 }
53 };
54
55 Ok(SessionEntry {
56 id: Uuid::new_v4().to_string(),
57 parent_id,
58 timestamp: Utc::now(),
59 role: role.to_string(),
60 content: content.to_string(),
61 metadata,
62 })
63 }
64
65 fn append_entry_locked(&self, entry: &SessionEntry) -> Result<(), SessionError> {
66 let _lock = self.write_lock.lock().expect("write_lock mutex poisoned");
67 let line =
68 serde_json::to_string(entry).map_err(|e| SessionError::InvalidJson(e.to_string()))?;
69
70 if !self.file_path.exists()
71 && let Some(parent) = self.file_path.parent()
72 {
73 fs::create_dir_all(parent)?;
74 }
75
76 let mut file = OpenOptions::new()
77 .create(true)
78 .append(true)
79 .open(&self.file_path)?;
80 writeln!(file, "{line}")?;
81
82 *self
83 .last_entry_id
84 .lock()
85 .expect("last_entry_id mutex poisoned") = Some(entry.id.clone());
86 Ok(())
87 }
88
89 pub fn read_entries(&self) -> Result<Vec<SessionEntry>, SessionError> {
95 read_entries_from_path(&self.file_path)
96 }
97
98 pub fn read_messages(&self) -> Result<Vec<Message>, SessionError> {
103 let entries = self.read_entries()?;
104 let mut messages = Vec::new();
105
106 for entry in entries {
107 let msg = match entry.role.as_str() {
108 "user" => Some(Message::User {
109 content: entry.content,
110 }),
111 "assistant" => {
112 let tool_calls =
113 talos_core::message::extract_tool_calls_from_text(&entry.content);
114 let cleaned = talos_core::message::strip_tool_syntax(&entry.content);
115 Message::Assistant {
116 content: cleaned,
117 tool_calls,
118 }
119 .into()
120 }
121 "system" => {
122 if let Some(sys_content) = entry.content.strip_prefix("__SYSTEM__:") {
123 Some(Message::System {
124 content: sys_content.to_string(),
125 cache_markers: Vec::new(),
126 })
127 } else if serde_json::from_str::<AgentEvent>(&entry.content).is_ok() {
128 None
129 } else {
130 let (is_error, tool_use_id, content) = parse_tool_result(&entry.content);
131 Some(Message::Tool {
132 result: talos_core::message::MessageToolResult {
133 tool_use_id,
134 content,
135 is_error,
136 },
137 })
138 }
139 }
140 _ => None,
141 };
142
143 if let Some(msg) = msg {
144 messages.push(msg);
145 }
146 }
147
148 Ok(messages)
149 }
150
151 pub fn read_events(&self) -> Result<Vec<AgentEvent>, SessionError> {
153 let entries = self.read_entries()?;
154 let mut events = Vec::new();
155
156 for entry in entries {
157 if entry.role == "system"
158 && let Ok(event) = serde_json::from_str::<AgentEvent>(&entry.content)
159 {
160 events.push(event);
161 }
162 }
163
164 Ok(events)
165 }
166}
167
168fn read_last_entry_id(path: &Path) -> Option<String> {
169 let mut file = fs::File::open(path).ok()?;
170 let file_size = file.metadata().ok()?.len();
171 if file_size == 0 {
172 return None;
173 }
174 let read_size = std::cmp::min(file_size, 8192) as usize;
175 let seek_pos = file_size.saturating_sub(read_size as u64);
176 file.seek(SeekFrom::Start(seek_pos)).ok()?;
177 let mut buf = vec![0u8; read_size];
178 file.read_exact(&mut buf).ok()?;
179 let text = String::from_utf8_lossy(&buf);
180 let last_line = text.lines().rev().find(|l| !l.is_empty())?;
181 let entry: SessionEntry = serde_json::from_str(last_line).ok()?;
182 Some(entry.id)
183}
184
185pub(crate) fn scan_file(path: &Path) -> Result<(usize, String), SessionError> {
186 let file = fs::File::open(path)?;
187 let reader = BufReader::new(file);
188 let mut count = 0;
189 let mut last_preview = String::new();
190
191 for line in reader.lines() {
192 let line = line?;
193 if line.is_empty() {
194 continue;
195 }
196
197 if let Ok(entry) = serde_json::from_str::<SessionEntry>(&line) {
198 count += 1;
199 last_preview = preview_text(&entry.content);
200 continue;
201 }
202
203 if let Ok(value) = serde_json::from_str::<serde_json::Value>(&line)
204 && value.get("type").and_then(|t| t.as_str()) == Some("message")
205 {
206 count += 1;
207 if let Some(data) = value.get("data")
208 && let Ok(msg) = serde_json::from_value::<Message>(data.clone())
209 {
210 let (_, content) = message_parts(&msg);
211 last_preview = preview_text(&content);
212 }
213 }
214 }
215
216 Ok((count, last_preview))
217}
218
219fn read_entries_from_path(path: &Path) -> Result<Vec<SessionEntry>, SessionError> {
220 if !path.exists() {
221 return Ok(Vec::new());
222 }
223
224 let file = fs::File::open(path)?;
225 let reader = BufReader::new(file);
226 let mut entries = Vec::new();
227 let mut synthetic_counter: u64 = 0;
228
229 for line in reader.lines() {
230 let line = line?;
231 if line.is_empty() {
232 continue;
233 }
234
235 if let Ok(entry) = serde_json::from_str::<SessionEntry>(&line) {
236 entries.push(entry);
237 continue;
238 }
239
240 if let Ok(value) = serde_json::from_str::<serde_json::Value>(&line)
241 && value.get("type").and_then(|t| t.as_str()) == Some("message")
242 && let Some(data) = value.get("data")
243 && let Ok(msg) = serde_json::from_value::<Message>(data.clone())
244 {
245 let (role, content) = message_parts(&msg);
246 let id = format!("synthetic-{synthetic_counter}");
247 let parent_id = if synthetic_counter > 0 {
248 Some(format!("synthetic-{}", synthetic_counter - 1))
249 } else {
250 None
251 };
252
253 entries.push(SessionEntry {
254 id,
255 parent_id,
256 timestamp: Utc::now(),
257 role,
258 content,
259 metadata: SessionMetadata::default(),
260 });
261 synthetic_counter += 1;
262 }
263 }
265
266 Ok(entries)
267}
268
269fn parse_tool_result(content: &str) -> (bool, String, String) {
270 if let Some(rest) = content.strip_prefix("__ERROR__:")
271 && let Some((id, body)) = rest.split_once("__\n")
272 {
273 return (true, id.to_string(), body.to_string());
274 }
275 if let Some(rest) = content.strip_prefix("__OK__:")
276 && let Some((id, body)) = rest.split_once("__\n")
277 {
278 return (false, id.to_string(), body.to_string());
279 }
280 (false, "unknown".to_string(), content.to_string())
281}
282
283fn message_parts(message: &Message) -> (String, String) {
284 match message {
285 Message::User { content } => ("user".to_string(), content.clone()),
286 Message::Assistant {
287 content,
288 tool_calls,
289 } => {
290 if tool_calls.is_empty() {
291 return ("assistant".to_string(), content.clone());
292 }
293 let mut full = content.clone();
295 for tc in tool_calls {
296 let block = serde_json::json!({
297 "name": tc.name,
298 "args": tc.input,
299 });
300 full.push_str(&format!("\n```json-tool\n{block}\n```"));
301 }
302 ("assistant".to_string(), full)
303 }
304 Message::Tool { result } => {
305 let prefix = if result.is_error {
306 format!("__ERROR__:{}__\n", result.tool_use_id)
307 } else {
308 format!("__OK__:{}__\n", result.tool_use_id)
309 };
310 ("system".to_string(), format!("{prefix}{}", result.content))
311 }
312 Message::System { content, .. } => ("system".to_string(), format!("__SYSTEM__:{content}")),
313 Message::Context { content } => ("user".to_string(), content.clone()),
314 }
315}
316
317fn preview_text(content: &str) -> String {
318 const MAX_PREVIEW_CHARS: usize = 100;
319 let mut chars = content.chars();
320 let preview: String = chars.by_ref().take(MAX_PREVIEW_CHARS).collect();
321 if chars.next().is_some() {
322 format!("{preview}...")
323 } else {
324 preview
325 }
326}