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| 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 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 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 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 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 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 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 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 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 let _ = conn.execute_batch("PRAGMA wal_checkpoint(PASSIVE);");
575
576 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}