use super::args::IngestArgs;
use super::persist::init_storage;
use super::persist_loop::{self, PersistContext};
use super::plan::build_plan;
use super::scan_fs::collect_files;
use super::validate::validate_mode_conditional_flags_ingest;
use super::{dry_run, enrich_after, stage_producer};
use crate::errors::AppError;
use crate::paths::AppPaths;
use std::path::PathBuf;
pub fn run(args: IngestArgs, backends: crate::cli::BackendChoice) -> Result<(), AppError> {
validate_mode_conditional_flags_ingest(&args)?;
crate::agent_surface::stream::open(crate::agent_surface::get(), &[], 0)?;
tracing::debug!(target: "ingest", dir = %args.dir.display(), mode = ?args.mode, "starting ingest");
let started = std::time::Instant::now();
let files = scan(&args)?;
let total = files.len();
let namespace = crate::namespace::resolve_namespace(args.namespace.as_deref())?;
let memory_type_str = args.r#type.as_str().to_string();
let paths = AppPaths::resolve(args.db.as_deref())?;
let mut conn_or_err = init_storage(&paths).map_err(|e| format!("{e}"));
let plan = build_plan(&args, &files)?;
if args.dry_run {
return dry_run::emit_preview(&args, &plan.slots_meta, total, started);
}
let parallelism = stage_producer::resolve_worker_count(&args)?;
stage_producer::validate_extraction_flags(&args)?;
let total_to_process = plan.process_items.len();
tracing::info!(
target: "ingest",
phase = "pipeline_start",
files = total_to_process,
ingest_parallelism = parallelism,
"incremental pipeline starting: Phase A (rayon) → channel → Phase B (main thread)",
);
let producer = stage_producer::spawn(&args, plan.process_items, &paths, parallelism, backends)?;
let ctx = PersistContext {
args: &args,
namespace: &namespace,
memory_type: &memory_type_str,
total,
started,
};
let tally = persist_loop::drain_and_persist(
&ctx,
&plan.slots_meta,
producer.results,
&mut conn_or_err,
)?;
producer
.handle
.join()
.map_err(|_| AppError::Internal(anyhow::anyhow!("ingest producer thread panicked")))?;
if let Ok(ref conn) = conn_or_err {
if tally.succeeded > 0 {
let _ = conn.execute_batch("PRAGMA wal_checkpoint(TRUNCATE);");
}
}
persist_loop::emit_summary(&ctx, tally)?;
if args.enrich_after && tally.succeeded > 0 {
enrich_after::run(&args, backends)?;
}
Ok(())
}
fn scan(args: &IngestArgs) -> Result<Vec<PathBuf>, AppError> {
if !args.dir.exists() {
return Err(AppError::Validation(
crate::i18n::validation::directory_not_found(&args.dir.display().to_string()),
));
}
if !args.dir.is_dir() {
return Err(AppError::Validation(
crate::i18n::validation::not_a_directory(&args.dir.display().to_string()),
));
}
let mut files: Vec<PathBuf> = Vec::with_capacity(128);
collect_files(&args.dir, &args.pattern, args.recursive, &mut files)?;
files.sort_unstable();
if files.len() > args.max_files {
return Err(AppError::Validation(
crate::i18n::validation::max_files_exceeded_matching(files.len(), args.max_files),
));
}
Ok(files)
}