1use 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
18pub 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 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 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 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 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 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 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 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 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 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 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 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 let _ = conn.execute_batch("PRAGMA wal_checkpoint(PASSIVE);");
581
582 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}