Skip to main content

sqlite_graphrag/commands/ingest_codex/
run.rs

1//! Main entry point for `ingest --mode codex`.
2
3use super::binary::{find_codex_binary, validate_codex_version};
4use super::extract::extract_with_codex;
5use super::queue::{collect_matching_files, open_queue_db};
6use super::types::*;
7use crate::commands::ingest::IngestArgs;
8use crate::commands::ingest_claude::ExtractionResult;
9use crate::entity_type::EntityType;
10use crate::errors::AppError;
11use crate::output::emit_json_line as emit_json;
12use crate::paths::AppPaths;
13use crate::storage::connection::{ensure_db_ready, open_rw};
14use crate::storage::entities::{self, NewEntity, NewRelationship};
15use crate::storage::memories::{self, NewMemory};
16use std::time::Instant;
17
18/// Run `ingest --mode codex` with queue resume and per-file Codex extraction.
19pub fn run_codex_ingest(args: &IngestArgs) -> Result<(), AppError> {
20    let started = Instant::now();
21
22    if !args.dir.exists() {
23        return Err(AppError::Validation(
24            crate::i18n::validation::directory_not_found(&args.dir.display().to_string()),
25        ));
26    }
27
28    // G28-B (v1.0.68) + G30 (v1.0.69): acquire singleton before doing real
29    // work so two parallel `ingest --mode codex` invocations cannot co-exist
30    // on the same database. Scope includes the database hash so concurrent
31    // ingest against different databases is allowed.
32    let early_ns = crate::namespace::resolve_namespace(args.namespace.as_deref())?;
33    let early_paths = AppPaths::resolve(args.db.as_deref())?;
34    let queue_path = match args.queue_db.as_deref() {
35        Some(p) => std::path::PathBuf::from(p),
36        None => crate::paths::sidecar_path(&early_paths.db, ".ingest-queue.sqlite"),
37    };
38    let _singleton = crate::lock::acquire_job_singleton(
39        crate::lock::JobType::IngestCodex,
40        &early_ns,
41        &early_paths.db,
42        args.wait_job_singleton,
43        args.force_job_singleton,
44    )?;
45
46    // Stage 1: Validate binary
47    let codex_binary = find_codex_binary(args.codex_binary.as_deref())?;
48    let version = validate_codex_version(&codex_binary)?;
49    tracing::info!(
50        target: "ingest",
51        binary = %codex_binary.display(),
52        version = %version,
53        "Codex CLI binary validated"
54    );
55
56    emit_json(&PhaseEvent {
57        phase: "validate",
58        codex_path: codex_binary.to_str(),
59        version: Some(&version),
60        dir: None,
61        files_total: None,
62        files_new: None,
63        files_existing: None,
64    });
65
66    // Stage 2: Scan files
67    let files = collect_matching_files(&args.dir, &args.pattern, args.recursive, args.max_files)?;
68
69    let queue_conn = open_queue_db(&queue_path)?;
70
71    if args.resume {
72        let reset = queue_conn
73            .execute(
74                "UPDATE queue SET status='pending' WHERE status='processing'",
75                [],
76            )
77            .map_err(|e| AppError::Validation(crate::i18n::validation::queue_resume_failed(&e)))?;
78        if reset > 0 {
79            tracing::info!(target: "ingest", count = reset, "reset stuck processing files to pending");
80        }
81    }
82
83    if args.retry_failed {
84        let count = queue_conn
85            .execute(
86                "UPDATE queue SET status='pending', attempt=0 WHERE status='failed'",
87                [],
88            )
89            .map_err(|e| {
90                AppError::Validation(crate::i18n::validation::queue_retry_failed_reset_failed(&e))
91            })?;
92        tracing::info!(target: "ingest", count, "retrying failed files");
93    }
94
95    if !args.resume && !args.retry_failed {
96        queue_conn
97            .execute("DELETE FROM queue", [])
98            .map_err(|e| AppError::Validation(crate::i18n::validation::queue_clear_failed(&e)))?;
99    }
100
101    let mut new_count = 0usize;
102    let mut existing_count = 0usize;
103
104    if !args.retry_failed {
105        for file in &files {
106            let file_str = file.to_string_lossy().into_owned();
107            let inserted = queue_conn
108                .execute(
109                    "INSERT OR IGNORE INTO queue (file_path, status) VALUES (?1, 'pending')",
110                    rusqlite::params![file_str],
111                )
112                .map_err(|e| {
113                    AppError::Validation(crate::i18n::validation::queue_insert_failed(&e))
114                })?;
115            if inserted > 0 {
116                new_count += 1;
117            } else {
118                existing_count += 1;
119            }
120        }
121    }
122
123    emit_json(&PhaseEvent {
124        phase: "scan",
125        codex_path: None,
126        version: None,
127        dir: args.dir.to_str(),
128        files_total: Some(files.len()),
129        files_new: Some(new_count),
130        files_existing: Some(existing_count),
131    });
132
133    if args.dry_run {
134        for (idx, file) in files.iter().enumerate() {
135            let (name, _truncated, _orig) =
136                crate::commands::ingest::derive_kebab_name(file, args.max_name_length);
137            emit_json(&FileEvent {
138                file: &file.to_string_lossy(),
139                name: &name,
140                status: "preview",
141                memory_id: None,
142                entities: None,
143                rels: None,
144                cost_usd: None,
145                input_tokens: None,
146                output_tokens: None,
147                elapsed_ms: None,
148                error: None,
149                index: idx,
150                total: files.len(),
151            });
152        }
153        emit_json(&Summary {
154            summary: true,
155            files_total: files.len(),
156            completed: 0,
157            failed: 0,
158            skipped: 0,
159            entities_total: 0,
160            rels_total: 0,
161            input_tokens_total: 0,
162            output_tokens_total: 0,
163            elapsed_ms: started.elapsed().as_millis() as u64,
164        });
165        if !args.keep_queue {
166            let _ = std::fs::remove_file(&queue_path);
167        }
168        return Ok(());
169    }
170
171    // Stage 3: Process files
172    let paths = AppPaths::resolve(args.db.as_deref())?;
173    ensure_db_ready(&paths)?;
174    let conn = open_rw(&paths.db)?;
175    let namespace = crate::namespace::resolve_namespace(args.namespace.as_deref())?;
176    let memory_type_str = args.r#type.as_str().to_string();
177
178    // Write schema to temp file once (reused across all files)
179    let schema_tempfile = super::extract::write_schema_tempfile()?;
180    let schema_path = schema_tempfile.path().to_path_buf();
181
182    let mut completed = 0usize;
183    let mut failed = 0usize;
184    let skipped_initial: usize = queue_conn
185        .query_row("SELECT COUNT(*) FROM queue WHERE status='done'", [], |r| {
186            r.get::<_, usize>(0)
187        })
188        .unwrap_or(0);
189    let mut skipped = skipped_initial;
190    let mut entities_total = 0usize;
191    let mut rels_total = 0usize;
192    let mut input_tokens_total = 0u64;
193    let mut output_tokens_total = 0u64;
194    let total = files.len();
195
196    let mut backoff_secs = args.rate_limit_wait;
197    let rate_limit_deadline = std::time::Instant::now() + std::time::Duration::from_secs(3600);
198
199    loop {
200        if crate::shutdown_requested() {
201            tracing::info!(target: "ingest", "shutdown requested, stopping before next file");
202            break;
203        }
204
205        let pending: Option<(i64, String)> = queue_conn
206            .query_row(
207                "UPDATE queue SET status='processing', attempt=attempt+1 \
208                 WHERE id = (SELECT id FROM queue WHERE status='pending' ORDER BY id LIMIT 1) \
209                 RETURNING id, file_path",
210                [],
211                |row| Ok((row.get(0)?, row.get(1)?)),
212            )
213            .ok();
214
215        let (queue_id, file_path) = match pending {
216            Some(p) => p,
217            None => break,
218        };
219
220        let file_started = Instant::now();
221
222        // Reject files that exceed the 10 MB stdin limit
223        const MAX_FILE_SIZE: u64 = 10 * 1024 * 1024;
224        if let Ok(meta) = std::fs::metadata(&file_path) {
225            if meta.len() > MAX_FILE_SIZE {
226                let err_msg = format!("file exceeds 10MB stdin limit ({} bytes)", meta.len());
227                let _ = queue_conn.execute(
228                    "UPDATE queue SET status='failed', error=?1, done_at=datetime('now') WHERE id=?2",
229                    rusqlite::params![err_msg, queue_id],
230                );
231                let current_index = completed + failed + skipped;
232                failed += 1;
233                emit_json(&FileEvent {
234                    file: &file_path,
235                    name: "",
236                    status: "failed",
237                    memory_id: None,
238                    entities: None,
239                    rels: None,
240                    cost_usd: None,
241                    input_tokens: None,
242                    output_tokens: None,
243                    elapsed_ms: Some(file_started.elapsed().as_millis() as u64),
244                    error: Some(&err_msg),
245                    index: current_index,
246                    total,
247                });
248                if args.fail_fast {
249                    break;
250                }
251                continue;
252            }
253        }
254
255        let file_content = match std::fs::read(&file_path) {
256            Ok(c) => c,
257            Err(e) => {
258                let err_msg = format!("IO error: {e}");
259                let _ = queue_conn.execute(
260                    "UPDATE queue SET status='failed', error=?1, done_at=datetime('now') WHERE id=?2",
261                    rusqlite::params![err_msg, queue_id],
262                );
263                let current_index = completed + failed + skipped;
264                failed += 1;
265                emit_json(&FileEvent {
266                    file: &file_path,
267                    name: "",
268                    status: "failed",
269                    memory_id: None,
270                    entities: None,
271                    rels: None,
272                    cost_usd: None,
273                    input_tokens: None,
274                    output_tokens: None,
275                    elapsed_ms: Some(file_started.elapsed().as_millis() as u64),
276                    error: Some(&err_msg),
277                    index: current_index,
278                    total,
279                });
280                if args.fail_fast {
281                    break;
282                }
283                continue;
284            }
285        };
286
287        // Skip files exceeding body cap BEFORE sending to LLM to avoid wasting tokens
288        if file_content.len() > crate::constants::MAX_MEMORY_BODY_LEN {
289            let err_msg = format!(
290                "file body exceeds {} byte limit ({} bytes) — skipping to avoid wasting LLM tokens",
291                crate::constants::MAX_MEMORY_BODY_LEN,
292                file_content.len()
293            );
294            tracing::warn!(target: "ingest", file = %file_path, size = file_content.len(), "body exceeds limit, skipping LLM extraction");
295            let _ = queue_conn.execute(
296                "UPDATE queue SET status='skipped', error=?1, done_at=datetime('now') WHERE id=?2",
297                rusqlite::params![err_msg, queue_id],
298            );
299            let current_index = completed + failed + skipped;
300            skipped += 1;
301            emit_json(&FileEvent {
302                file: &file_path,
303                name: "",
304                status: "skipped",
305                memory_id: None,
306                entities: None,
307                rels: None,
308                cost_usd: None,
309                input_tokens: None,
310                output_tokens: None,
311                elapsed_ms: Some(file_started.elapsed().as_millis() as u64),
312                error: Some(&err_msg),
313                index: current_index,
314                total,
315            });
316            continue;
317        }
318
319        // Retry once on cold-start failure
320        let max_extract_attempts: u32 = 2;
321        let mut extraction_result: Option<(ExtractionResult, Option<CodexUsage>)> = None;
322        let mut last_extract_err: Option<String> = None;
323        let mut last_was_rate_limited = false;
324
325        for attempt in 1..=max_extract_attempts {
326            match extract_with_codex(
327                &codex_binary,
328                &file_content,
329                args.codex_model.as_deref(),
330                args.codex_timeout,
331                &schema_path,
332            ) {
333                Ok(result) => {
334                    extraction_result = Some(result);
335                    break;
336                }
337                Err(ref e) if matches!(e, AppError::RateLimited { .. }) => {
338                    last_extract_err = Some(format!("{e}"));
339                    last_was_rate_limited = true;
340                    break;
341                }
342                Err(e) => {
343                    let msg = format!("{e}");
344                    if attempt < max_extract_attempts {
345                        let cold_start_delay = 2 * attempt as u64;
346                        tracing::warn!(
347                            target: "ingest",
348                            attempt,
349                            delay_secs = cold_start_delay,
350                            error = %msg,
351                            "codex extraction failed, retrying"
352                        );
353                        std::thread::sleep(std::time::Duration::from_secs(cold_start_delay));
354                    }
355                    last_extract_err = Some(msg);
356                }
357            }
358        }
359
360        if let Some((extraction, usage)) = extraction_result {
361            backoff_secs = args.rate_limit_wait;
362
363            let in_tok = usage.as_ref().map(|u| u.input_tokens).unwrap_or(0);
364            let out_tok = usage.as_ref().map(|u| u.output_tokens).unwrap_or(0);
365
366            let name = &extraction.name;
367            let ent_count = extraction.entities.len();
368            let rel_count = 0;
369
370            // GAP-SG-47: fold non-canonical labels onto the nearest canonical
371            // kind instead of discarding the entity (no silent data loss).
372            let new_entities: Vec<NewEntity> = extraction
373                .entities
374                .iter()
375                .map(|e| NewEntity {
376                    name: e.name.clone(),
377                    entity_type: EntityType::map_to_canonical(&e.entity_type),
378                    description: None,
379                })
380                .collect();
381
382            // GAP-SG-48: rewrite non-canonical relations to canonical instead
383            // of normalizing-and-accepting them raw.
384            let new_relationships: Vec<NewRelationship> = extraction
385                .relationships
386                .iter()
387                .map(|r| NewRelationship {
388                    source: r.source.clone(),
389                    target: r.target.clone(),
390                    relation: crate::parsers::map_to_canonical_relation(&r.relation),
391                    strength: r.strength,
392                    description: None,
393                })
394                .collect();
395
396            let body_str = String::from_utf8(file_content.clone())
397                .map_err(|e| AppError::Validation(crate::i18n::validation::file_not_utf8(&e)))?;
398            let body_hash = blake3::hash(body_str.as_bytes()).to_hex().to_string();
399            let new_memory = NewMemory {
400                name: name.clone(),
401                namespace: namespace.clone(),
402                memory_type: memory_type_str.clone(),
403                description: extraction.description.clone(),
404                body: body_str.to_string(),
405                body_hash,
406                session_id: None,
407                source: "agent".to_string(),
408                metadata: serde_json::Value::Object(serde_json::Map::new()),
409            };
410
411            // Deduplication: update existing memory instead of failing on UNIQUE
412            let memory_id = match memories::find_by_name_any_state(&conn, &namespace, name)? {
413                Some((existing_id, is_deleted)) => {
414                    if is_deleted {
415                        memories::clear_deleted_at(&conn, existing_id)?;
416                    }
417                    let (old_name, old_desc, old_body): (String, String, String) = conn.query_row(
418                        "SELECT name, COALESCE(description,''), COALESCE(body,'') FROM memories WHERE id=?1",
419                        rusqlite::params![existing_id],
420                        |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
421                    )?;
422                    memories::update(&conn, existing_id, &new_memory, None)?;
423                    memories::sync_fts_after_update(
424                        &conn,
425                        existing_id,
426                        &old_name,
427                        &old_desc,
428                        &old_body,
429                        &new_memory.name,
430                        &new_memory.description,
431                        &new_memory.body,
432                    )?;
433                    tracing::info!(target: "ingest", name, memory_id = existing_id, "updated existing memory (force-merge)");
434                    existing_id
435                }
436                None => match memories::insert(&conn, &new_memory) {
437                    Ok(id) => id,
438                    Err(e) => {
439                        let err_msg = format!("{e}");
440                        let _ = queue_conn.execute(
441                            "UPDATE queue SET status='failed', error=?1, done_at=datetime('now') WHERE id=?2",
442                            rusqlite::params![err_msg, queue_id],
443                        );
444                        let current_index = completed + failed + skipped;
445                        failed += 1;
446                        emit_json(&FileEvent {
447                            file: &file_path,
448                            name,
449                            status: "failed",
450                            memory_id: None,
451                            entities: None,
452                            rels: None,
453                            cost_usd: None,
454                            input_tokens: Some(in_tok),
455                            output_tokens: Some(out_tok),
456                            elapsed_ms: Some(file_started.elapsed().as_millis() as u64),
457                            error: Some(&err_msg),
458                            index: current_index,
459                            total,
460                        });
461                        input_tokens_total += in_tok;
462                        output_tokens_total += out_tok;
463                        if args.fail_fast {
464                            break;
465                        }
466                        continue;
467                    }
468                },
469            };
470
471            for ent in &new_entities {
472                if let Ok(eid) = entities::upsert_entity(&conn, &namespace, ent) {
473                    let _ = entities::link_memory_entity(&conn, memory_id, eid);
474                }
475            }
476            for rel in &new_relationships {
477                crate::parsers::warn_if_non_canonical(&rel.relation);
478                let src_id = entities::find_entity_id(&conn, &namespace, &rel.source);
479                let tgt_id = entities::find_entity_id(&conn, &namespace, &rel.target);
480                if let (Ok(Some(sid)), Ok(Some(tid))) = (src_id, tgt_id) {
481                    let _ = conn.execute(
482                        "INSERT OR IGNORE INTO relationships (namespace, source_id, target_id, relation, weight) VALUES (?1, ?2, ?3, ?4, ?5)",
483                        rusqlite::params![namespace, sid, tid, rel.relation, rel.strength],
484                    );
485                }
486            }
487
488            let _ = queue_conn.execute(
489                "UPDATE queue SET status='done', name=?1, memory_id=?2, entities=?3, rels=?4, \
490                 input_tokens=?5, output_tokens=?6, elapsed_ms=?7, done_at=datetime('now') WHERE id=?8",
491                rusqlite::params![
492                    name,
493                    memory_id,
494                    ent_count,
495                    rel_count,
496                    in_tok,
497                    out_tok,
498                    file_started.elapsed().as_millis() as i64,
499                    queue_id
500                ],
501            );
502
503            let current_index = completed + failed + skipped;
504            completed += 1;
505            entities_total += ent_count;
506            rels_total += rel_count;
507            input_tokens_total += in_tok;
508            output_tokens_total += out_tok;
509
510            emit_json(&FileEvent {
511                file: &file_path,
512                name,
513                status: "done",
514                memory_id: Some(memory_id),
515                entities: Some(ent_count),
516                rels: Some(rel_count),
517                cost_usd: None,
518                input_tokens: Some(in_tok),
519                output_tokens: Some(out_tok),
520                elapsed_ms: Some(file_started.elapsed().as_millis() as u64),
521                error: None,
522                index: current_index,
523                total,
524            });
525        } else if let Some(ref err_str) = last_extract_err {
526            if last_was_rate_limited {
527                if crate::retry::is_kill_switch_active() {
528                    tracing::warn!(target: "ingest", "retry.disable is set, skipping rate-limit retry");
529                } else if std::time::Instant::now() >= rate_limit_deadline {
530                    tracing::error!(target: "ingest", "rate-limit retry deadline (1h) exhausted");
531                } else {
532                    let half = backoff_secs / 2;
533                    let jitter = if half == 0 { 0 } else { fastrand::u64(0..half) };
534                    let actual_wait = half + jitter;
535                    tracing::warn!(target: "ingest", delay_secs = actual_wait, error_kind = "rate_limited", "Codex rate limited, backing off");
536                    let _ = queue_conn.execute(
537                        "UPDATE queue SET status='pending' WHERE id=?1",
538                        rusqlite::params![queue_id],
539                    );
540                    std::thread::sleep(std::time::Duration::from_secs(actual_wait));
541                    backoff_secs = (backoff_secs * 2).min(900);
542                    continue;
543                }
544            } else {
545                let _ = queue_conn.execute(
546                    "UPDATE queue SET status='failed', error=?1, done_at=datetime('now') WHERE id=?2",
547                    rusqlite::params![err_str, queue_id],
548                );
549                let current_index = completed + failed + skipped;
550                failed += 1;
551                emit_json(&FileEvent {
552                    file: &file_path,
553                    name: "",
554                    status: "failed",
555                    memory_id: None,
556                    entities: None,
557                    rels: None,
558                    cost_usd: None,
559                    input_tokens: None,
560                    output_tokens: None,
561                    elapsed_ms: Some(file_started.elapsed().as_millis() as u64),
562                    error: Some(err_str),
563                    index: current_index,
564                    total,
565                });
566                if args.fail_fast {
567                    break;
568                }
569            }
570        }
571    }
572
573    // WAL checkpoint before summary
574    let _ = conn.execute_batch("PRAGMA wal_checkpoint(PASSIVE);");
575
576    // Stage 4: Summary
577    emit_json(&Summary {
578        summary: true,
579        files_total: total,
580        completed,
581        failed,
582        skipped,
583        entities_total,
584        rels_total,
585        input_tokens_total,
586        output_tokens_total,
587        elapsed_ms: started.elapsed().as_millis() as u64,
588    });
589
590    if !args.keep_queue && failed == 0 {
591        let _ = std::fs::remove_file(&queue_path);
592    }
593
594    Ok(())
595}