Skip to main content

talos_session/
jsonl.rs

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    /// Read all entries from the session's JSONL file.
90    ///
91    /// Entries are reconstructed from the JSONL format. Entries without `id` or
92    /// `parent_id` (backward compatibility) are assigned synthetic IDs and treated
93    /// as a single linear branch.
94    pub fn read_entries(&self) -> Result<Vec<SessionEntry>, SessionError> {
95        read_entries_from_path(&self.file_path)
96    }
97
98    /// Read all messages from the session's JSONL file for the current branch.
99    ///
100    /// Only entries with role `"user"`, `"assistant"`, or `"system"` that contain
101    /// valid message data are returned.
102    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    /// Read all events from the session's JSONL file.
152    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        // Invalid lines are silently skipped (crash-safety guarantee)
264    }
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            // Embed tool calls as json-tool blocks so they survive JSONL round-trip.
294            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}