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