Skip to main content

edda_transcript/
ingest.rs

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; // 4MB
9
10/// Callback type for index generation during ingest.
11pub 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/// Perform cursor-based delta ingest from a Claude transcript JSONL file.
26///
27/// Reads from `transcript_path` starting at the cursor offset (or 0 if new),
28/// classifies records, writes kept records verbatim to the store,
29/// and returns ingest statistics.
30///
31/// If `index_writer` is Some, calls it for each kept record with
32/// (raw_line, store_offset, store_len, parsed_json) for index generation.
33#[allow(clippy::too_many_lines)] // 186 lines at #779; split tracked in none
34pub 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    // Session-level lock
44    let lock_path = state_dir.join(format!("ingest.{session_id}.lock"));
45    let _lock = edda_store::lock_file(&lock_path)?;
46
47    // Load or create cursor
48    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    // Check file metadata
56    let meta = std::fs::metadata(transcript_path)?;
57    let file_size = meta.len();
58
59    // Truncation detection
60    cursor.detect_truncation(file_size);
61
62    if cursor.offset >= file_size {
63        // Nothing new to read
64        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    // Open and seek
82    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    // Partial line protection: only consume up to the last newline
91    let consumable_len = match buf.iter().rposition(|&b| b == b'\n') {
92        Some(pos) => pos + 1,
93        None => 0, // no complete line
94    };
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    // Prepare store path (verbatim append)
113    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    // Load progress_last map
122    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    // Process line by line
142    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                // Record store_offset before write
170                let store_offset = store_file.seek(SeekFrom::End(0)).unwrap_or(0);
171
172                // Write raw line verbatim (CONTRACT BRIDGE-03)
173                store_file.write_all(raw_line)?;
174                store_file.write_all(b"\n")?;
175
176                let store_len = raw_line.len() as u64 + 1; // +1 for newline
177
178                // Call index writer if provided
179                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    // Save progress_last map
200    let progress_json = serde_json::to_string_pretty(&progress_map)?;
201    edda_store::write_atomic(&progress_path, progress_json.as_bytes())?;
202
203    // Update and save cursor
204    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); // user + assistant
249        assert_eq!(stats.records_dropped, 2); // progress + turn_duration
250
251        // Verify verbatim store
252        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        // First write
269        {
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        // Append more
282        {
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); // only the new line
296        assert_eq!(stats2.from_offset, stats1.to_offset);
297
298        // Store should have 2 lines total
299        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}