1use crate::cursor::TranscriptCursor;
2use crate::filter::{classify_record, update_progress_last, FilterAction};
3use serde::{Deserialize, Serialize};
4use std::collections::HashMap;
5use std::io::{Read, Seek, SeekFrom, Write};
6use std::path::Path;
7
8const DEFAULT_MAX_BYTES: u64 = 4 * 1024 * 1024; pub type IndexWriterFn = dyn Fn(&str, u64, u64, &serde_json::Value) -> anyhow::Result<()>;
12
13#[derive(Debug, Serialize, Deserialize, Clone)]
14pub struct IngestStats {
15 pub records_read: usize,
16 pub records_kept: usize,
17 pub records_dropped: usize,
18 pub bytes_read: u64,
19 pub kept_by_type: HashMap<String, usize>,
20 pub dropped_by_type: HashMap<String, usize>,
21 pub from_offset: u64,
22 pub to_offset: u64,
23}
24
25#[allow(clippy::too_many_lines)] pub fn ingest_transcript_delta(
35 project_dir: &Path,
36 session_id: &str,
37 transcript_path: &Path,
38 index_writer: Option<&IndexWriterFn>,
39) -> anyhow::Result<IngestStats> {
40 let state_dir = project_dir.join("state");
41 std::fs::create_dir_all(&state_dir)?;
42
43 let lock_path = state_dir.join(format!("ingest.{session_id}.lock"));
45 let _lock = edda_store::lock_file(&lock_path)?;
46
47 let mut cursor = TranscriptCursor::load(&state_dir, session_id)?.unwrap_or(TranscriptCursor {
49 offset: 0,
50 file_size: 0,
51 mtime_unix: 0,
52 updated_at_unix: 0,
53 });
54
55 let meta = std::fs::metadata(transcript_path)?;
57 let file_size = meta.len();
58
59 cursor.detect_truncation(file_size);
61
62 if cursor.offset >= file_size {
63 return Ok(IngestStats {
65 records_read: 0,
66 records_kept: 0,
67 records_dropped: 0,
68 bytes_read: 0,
69 kept_by_type: HashMap::new(),
70 dropped_by_type: HashMap::new(),
71 from_offset: cursor.offset,
72 to_offset: cursor.offset,
73 });
74 }
75
76 let max_bytes: u64 = std::env::var("EDDA_TRANSCRIPT_MAX_BYTES")
77 .ok()
78 .and_then(|v| v.parse().ok())
79 .unwrap_or(DEFAULT_MAX_BYTES);
80
81 let mut file = std::fs::File::open(transcript_path)?;
83 file.seek(SeekFrom::Start(cursor.offset))?;
84
85 let bytes_to_read = (file_size - cursor.offset).min(max_bytes);
86 let mut buf = vec![0u8; bytes_to_read as usize];
87 let actually_read = file.read(&mut buf)?;
88 buf.truncate(actually_read);
89
90 let consumable_len = match buf.iter().rposition(|&b| b == b'\n') {
92 Some(pos) => pos + 1,
93 None => 0, };
95
96 if consumable_len == 0 {
97 return Ok(IngestStats {
98 records_read: 0,
99 records_kept: 0,
100 records_dropped: 0,
101 bytes_read: 0,
102 kept_by_type: HashMap::new(),
103 dropped_by_type: HashMap::new(),
104 from_offset: cursor.offset,
105 to_offset: cursor.offset,
106 });
107 }
108
109 let from_offset = cursor.offset;
110 let data = &buf[..consumable_len];
111
112 let transcripts_dir = project_dir.join("transcripts");
114 std::fs::create_dir_all(&transcripts_dir)?;
115 let store_path = transcripts_dir.join(format!("{session_id}.jsonl"));
116 let mut store_file = std::fs::OpenOptions::new()
117 .create(true)
118 .append(true)
119 .open(&store_path)?;
120
121 let progress_path = state_dir.join(format!("progress_last.{session_id}.json"));
123 let mut progress_map: HashMap<String, serde_json::Value> = if progress_path.exists() {
124 let content = std::fs::read_to_string(&progress_path)?;
125 serde_json::from_str(&content).unwrap_or_default()
126 } else {
127 HashMap::new()
128 };
129
130 let mut stats = IngestStats {
131 records_read: 0,
132 records_kept: 0,
133 records_dropped: 0,
134 bytes_read: consumable_len as u64,
135 kept_by_type: HashMap::new(),
136 dropped_by_type: HashMap::new(),
137 from_offset,
138 to_offset: from_offset + consumable_len as u64,
139 };
140
141 for raw_line in data.split(|&b| b == b'\n') {
143 if raw_line.is_empty() {
144 continue;
145 }
146
147 stats.records_read += 1;
148
149 let parsed: serde_json::Value = match serde_json::from_slice(raw_line) {
150 Ok(v) => v,
151 Err(_) => {
152 stats.records_dropped += 1;
153 *stats
154 .dropped_by_type
155 .entry("parse_error".into())
156 .or_insert(0) += 1;
157 continue;
158 }
159 };
160
161 let record_type = parsed
162 .get("type")
163 .and_then(|v| v.as_str())
164 .unwrap_or("unknown")
165 .to_string();
166
167 match classify_record(&parsed) {
168 FilterAction::Keep => {
169 let store_offset = store_file.seek(SeekFrom::End(0)).unwrap_or(0);
171
172 store_file.write_all(raw_line)?;
174 store_file.write_all(b"\n")?;
175
176 let store_len = raw_line.len() as u64 + 1; if let Some(writer) = index_writer {
180 let raw_str = std::str::from_utf8(raw_line).unwrap_or("");
181 writer(raw_str, store_offset, store_len, &parsed)?;
182 }
183
184 stats.records_kept += 1;
185 *stats.kept_by_type.entry(record_type).or_insert(0) += 1;
186 }
187 FilterAction::Progress => {
188 update_progress_last(&mut progress_map, &parsed);
189 stats.records_dropped += 1;
190 *stats.dropped_by_type.entry(record_type).or_insert(0) += 1;
191 }
192 FilterAction::Drop => {
193 stats.records_dropped += 1;
194 *stats.dropped_by_type.entry(record_type).or_insert(0) += 1;
195 }
196 }
197 }
198
199 let progress_json = serde_json::to_string_pretty(&progress_map)?;
201 edda_store::write_atomic(&progress_path, progress_json.as_bytes())?;
202
203 cursor.offset = stats.to_offset;
205 cursor.file_size = file_size;
206 cursor.updated_at_unix = std::time::SystemTime::now()
207 .duration_since(std::time::UNIX_EPOCH)
208 .map(|d| d.as_secs() as i64)
209 .unwrap_or(0);
210 cursor.save(&state_dir, session_id)?;
211
212 Ok(stats)
213}
214
215#[cfg(test)]
216mod tests {
217 use super::*;
218 use std::io::Write;
219
220 fn write_transcript(dir: &Path, lines: &[&str]) -> std::path::PathBuf {
221 let path = dir.join("transcript.jsonl");
222 let mut f = std::fs::File::create(&path).unwrap();
223 for line in lines {
224 writeln!(f, "{line}").unwrap();
225 }
226 path
227 }
228
229 #[test]
230 fn ingest_basic_keep_and_drop() {
231 let tmp = tempfile::tempdir().unwrap();
232 let project_dir = tmp.path().join("project");
233 std::fs::create_dir_all(&project_dir).unwrap();
234
235 let transcript = write_transcript(
236 tmp.path(),
237 &[
238 r#"{"type":"user","uuid":"u1","message":{"content":"hello"}}"#,
239 r#"{"type":"assistant","uuid":"a1","parentUuid":"u1","message":{"content":[{"type":"text","text":"hi"}]}}"#,
240 r#"{"type":"progress","toolUseID":"t1","data":{"output":"running"}}"#,
241 r#"{"type":"system","subtype":"turn_duration","duration_ms":100}"#,
242 ],
243 );
244
245 let stats = ingest_transcript_delta(&project_dir, "sess1", &transcript, None).unwrap();
246
247 assert_eq!(stats.records_read, 4);
248 assert_eq!(stats.records_kept, 2); assert_eq!(stats.records_dropped, 2); let store = project_dir.join("transcripts").join("sess1.jsonl");
253 let content = std::fs::read_to_string(&store).unwrap();
254 let lines: Vec<&str> = content.lines().collect();
255 assert_eq!(lines.len(), 2);
256 assert!(lines[0].contains("\"type\":\"user\""));
257 assert!(lines[1].contains("\"type\":\"assistant\""));
258 }
259
260 #[test]
261 fn ingest_cursor_based_delta() {
262 let tmp = tempfile::tempdir().unwrap();
263 let project_dir = tmp.path().join("project");
264 std::fs::create_dir_all(&project_dir).unwrap();
265
266 let transcript_path = tmp.path().join("transcript.jsonl");
267
268 {
270 let mut f = std::fs::File::create(&transcript_path).unwrap();
271 writeln!(
272 f,
273 r#"{{"type":"user","uuid":"u1","message":{{"content":"first"}}}}"#
274 )
275 .unwrap();
276 }
277 let stats1 =
278 ingest_transcript_delta(&project_dir, "sess1", &transcript_path, None).unwrap();
279 assert_eq!(stats1.records_kept, 1);
280
281 {
283 let mut f = std::fs::OpenOptions::new()
284 .append(true)
285 .open(&transcript_path)
286 .unwrap();
287 writeln!(
288 f,
289 r#"{{"type":"user","uuid":"u2","message":{{"content":"second"}}}}"#
290 )
291 .unwrap();
292 }
293 let stats2 =
294 ingest_transcript_delta(&project_dir, "sess1", &transcript_path, None).unwrap();
295 assert_eq!(stats2.records_kept, 1); assert_eq!(stats2.from_offset, stats1.to_offset);
297
298 let store = project_dir.join("transcripts").join("sess1.jsonl");
300 let content = std::fs::read_to_string(&store).unwrap();
301 assert_eq!(content.lines().count(), 2);
302 }
303
304 #[test]
305 fn ingest_with_index_writer() {
306 let tmp = tempfile::tempdir().unwrap();
307 let project_dir = tmp.path().join("project");
308 std::fs::create_dir_all(&project_dir).unwrap();
309
310 let transcript = write_transcript(
311 tmp.path(),
312 &[r#"{"type":"user","uuid":"u1","message":{"content":"hello"}}"#],
313 );
314
315 let called = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
316 let called_clone = called.clone();
317
318 let writer = move |_raw: &str,
319 _offset: u64,
320 _len: u64,
321 _json: &serde_json::Value|
322 -> anyhow::Result<()> {
323 called_clone.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
324 Ok(())
325 };
326
327 ingest_transcript_delta(&project_dir, "sess1", &transcript, Some(&writer)).unwrap();
328
329 assert_eq!(called.load(std::sync::atomic::Ordering::SeqCst), 1);
330 }
331}