use super::args::IngestArgs;
use super::persist::persist_staged;
use super::plan::SlotMeta;
use super::report::{FileSuccess, IngestFileEvent, IngestSummary};
use super::stage_producer::StageMessage;
use crate::errors::AppError;
use crate::output;
use rusqlite::Connection;
use std::sync::mpsc;
#[derive(Debug, Default, Clone, Copy)]
pub(super) struct IngestTally {
pub(super) succeeded: usize,
pub(super) failed: usize,
pub(super) skipped: usize,
}
pub(super) struct PersistContext<'a> {
pub(super) args: &'a IngestArgs,
pub(super) namespace: &'a str,
pub(super) memory_type: &'a str,
pub(super) total: usize,
pub(super) started: std::time::Instant,
}
struct SlotEvent<'a> {
file: &'a str,
name: &'a str,
truncated: bool,
original_name: Option<String>,
original_filename: Option<&'a str>,
}
pub(super) fn drain_and_persist(
ctx: &PersistContext<'_>,
slots_meta: &[SlotMeta],
results: mpsc::Receiver<StageMessage>,
conn_or_err: &mut Result<Connection, String>,
) -> Result<IngestTally, AppError> {
let mut tally = IngestTally::default();
for meta in slots_meta {
if let SlotMeta::Skip {
file_str,
derived_base,
name_truncated,
original_name,
original_filename,
reason,
} = meta
{
emit_event(
&SlotEvent {
file: file_str,
name: derived_base,
truncated: *name_truncated,
original_name: original_name.clone(),
original_filename: original_filename.as_deref(),
},
"skipped",
Some(reason.clone()),
None,
None,
0,
)?;
tally.skipped += 1;
}
}
let meta_index: std::collections::HashMap<usize, &SlotMeta> = slots_meta
.iter()
.enumerate()
.filter(|(_, m)| matches!(m, SlotMeta::Process { .. }))
.collect();
tracing::info!(
target: "ingest",
phase = "persist_start",
files = meta_index.len(),
"phase B starting: persisting files incrementally as Phase A completes each one",
);
for (idx, stage_result) in results {
if crate::shutdown_requested() {
tracing::info!(target: "ingest", "shutdown requested, stopping persistence loop");
break;
}
let meta = meta_index.get(&idx).ok_or_else(|| {
AppError::Internal(anyhow::anyhow!(
"channel idx {idx} has no corresponding Process slot"
))
})?;
let SlotMeta::Process {
file_str,
derived_name,
name_truncated,
original_name,
original_filename,
} = meta
else {
unreachable!("channel only carries Process results")
};
let slot = SlotEvent {
file: file_str,
name: derived_name,
truncated: *name_truncated,
original_name: original_name.clone(),
original_filename: original_filename.as_deref(),
};
let conn = match conn_or_err.as_mut() {
Ok(c) => c,
Err(err_msg) => {
let err_msg = err_msg.clone();
report_failure(ctx, &mut tally, &slot, err_msg)?;
continue;
}
};
match stage_result {
Ok(parts) => {
for staged in parts {
let part_name = staged.name.clone();
let part = SlotEvent {
file: slot.file,
name: &part_name,
truncated: slot.truncated,
original_name: slot.original_name.clone(),
original_filename: slot.original_filename,
};
persist_one(ctx, &mut tally, conn, &part, staged)?;
}
}
Err(e) => report_failure(ctx, &mut tally, &slot, format!("{e}"))?,
}
}
Ok(tally)
}
fn persist_one(
ctx: &PersistContext<'_>,
tally: &mut IngestTally,
conn: &mut Connection,
slot: &SlotEvent<'_>,
staged: super::stage::StagedFile,
) -> Result<(), AppError> {
match persist_staged(
conn,
ctx.namespace,
ctx.memory_type,
staged,
ctx.args.force_merge,
) {
Ok(FileSuccess {
memory_id,
action,
body_length,
backend_invoked,
}) => {
output::emit_stream_record(&IngestFileEvent {
file: slot.file,
name: slot.name,
status: "indexed",
truncated: slot.truncated,
original_name: slot.original_name.clone(),
original_filename: slot.original_filename,
error: None,
memory_id: Some(memory_id),
action: Some(action),
body_length,
backend_invoked,
..Default::default()
})?;
tally.succeeded += 1;
Ok(())
}
Err(ref e) if matches!(e, AppError::Duplicate(_)) => {
emit_event(
slot,
"skipped",
Some(format!("{e}")),
None,
Some("duplicate".to_string()),
0,
)?;
tally.skipped += 1;
Ok(())
}
Err(e) => report_failure(ctx, tally, slot, format!("{e}")),
}
}
fn report_failure(
ctx: &PersistContext<'_>,
tally: &mut IngestTally,
slot: &SlotEvent<'_>,
error: String,
) -> Result<(), AppError> {
emit_event(slot, "failed", Some(error.clone()), None, None, 0)?;
tally.failed += 1;
if ctx.args.fail_fast {
emit_summary(ctx, *tally)?;
return Err(AppError::Validation(
crate::i18n::validation::ingest_aborted_on_first_failure(&error),
));
}
Ok(())
}
fn emit_event(
slot: &SlotEvent<'_>,
status: &'static str,
error: Option<String>,
memory_id: Option<i64>,
action: Option<String>,
body_length: usize,
) -> Result<(), AppError> {
output::emit_stream_record(&IngestFileEvent {
file: slot.file,
name: slot.name,
status,
truncated: slot.truncated,
original_name: slot.original_name.clone(),
original_filename: slot.original_filename,
error,
memory_id,
action,
body_length,
backend_invoked: None,
..Default::default()
})
}
pub(super) fn emit_summary(ctx: &PersistContext<'_>, tally: IngestTally) -> Result<(), AppError> {
output::emit_stream_trailer(&IngestSummary {
summary: true,
dir: ctx.args.dir.display().to_string(),
pattern: ctx.args.pattern.clone(),
recursive: ctx.args.recursive,
files_total: ctx.total,
files_succeeded: tally.succeeded,
files_failed: tally.failed,
files_skipped: tally.skipped,
elapsed_ms: ctx.started.elapsed().as_millis() as u64,
..Default::default()
})
}