/// Ingest pipeline with DB persistence (Phase 2.6 integration) /// /// Orchestrates: /// 1. Run extraction pipeline /// 2. Save entities to DB /// 3. Save edges to DB /// 4. Return extraction result + DB IDs use anyhow::{Result, anyhow}; use mem_core::entity::Entity; use mem_core::edge::Edge; use mem_ingest::ingest_pipeline::{IngestPipeline, Episode, ExtractionResult}; use mem_store::db_repo::{PersistentEntityRepo, PersistentEdgeRepo, ReviewQueueRepo}; use sqlx::Pool; use sqlx::postgres::Postgres; use std::sync::Arc; use tracing::{debug, error, info}; /// Ingest result with DB persistence #[derive(Debug, Clone)] pub struct IngestWithDbResult { pub episode_id: String, pub entity_count: usize, pub entity_ids: Vec, pub edge_count: usize, pub edge_ids: Vec, pub contradiction_count: usize, pub extraction_errors: Vec, } /// Execute ingest pipeline with DB persistence pub async fn ingest_with_db_persistence( pool: &Pool, pipeline: &IngestPipeline, episode: &Episode, ) -> Result { debug!("Starting ingest with DB persistence for episode: {}", episode.id); // 1. Run extraction pipeline let extraction = pipeline.ingest(episode).await?; info!("Extraction complete: {} entities, {} edges, {} contradictions", extraction.entities.len(), extraction.edges.len(), extraction.reviews.len() ); // 2. Create repositories let entity_repo = PersistentEntityRepo::new(pool.clone()); let edge_repo = PersistentEdgeRepo::new(pool.clone()); let review_queue_repo = ReviewQueueRepo::new(pool.clone()); let mut entity_ids = Vec::new(); let mut edge_ids = Vec::new(); let mut errors = Vec::new(); // 3. Save entities for entity in &extraction.entities { match entity_repo.save(entity).await { Ok(id) => { debug!("Saved entity: {} → {}", entity.name, id); entity_ids.push(id); } Err(e) => { error!("Failed to save entity {}: {}", entity.name, e); errors.push(format!("Entity save failed: {}", e)); } } } // 4. Save edges for edge in &extraction.edges { match edge_repo.save(edge).await { Ok(id) => { debug!("Saved edge: {} → {} ({})", edge.source_id, edge.target_id, id); edge_ids.push(id); } Err(e) => { error!("Failed to save edge: {}", e); errors.push(format!("Edge save failed: {}", e)); } } } // 5. Queue contradictions for review (only high-confidence) for review_id in &extraction.reviews { match review_queue_repo.enqueue( &episode.project_id, review_id, "contradiction", 0.9, ).await { Ok(_) => { debug!("Queued contradiction for review: {}", review_id); } Err(e) => { error!("Failed to queue contradiction: {}", e); errors.push(format!("Review queue failed: {}", e)); } } } info!("Ingest complete: saved {} entities, {} edges, {} contradictions, {} errors", entity_ids.len(), edge_ids.len(), extraction.reviews.len(), errors.len() ); Ok(IngestWithDbResult { episode_id: episode.id.clone(), entity_count: entity_ids.len(), entity_ids, edge_count: edge_ids.len(), edge_ids, contradiction_count: extraction.reviews.len(), extraction_errors: errors, }) } #[cfg(test)] mod tests { use super::*; #[tokio::test] async fn test_ingest_with_db_result_creation() { let result = IngestWithDbResult { episode_id: "ep-1".to_string(), entity_count: 2, entity_ids: vec!["e1".to_string(), "e2".to_string()], edge_count: 1, edge_ids: vec!["edge-1".to_string()], contradiction_count: 0, extraction_errors: vec![], }; assert_eq!(result.entity_count, 2); assert_eq!(result.edge_count, 1); assert!(result.extraction_errors.is_empty()); } #[test] fn test_ingest_with_db_result_errors() { let result = IngestWithDbResult { episode_id: "ep-1".to_string(), entity_count: 1, entity_ids: vec!["e1".to_string()], edge_count: 0, edge_ids: vec![], contradiction_count: 0, extraction_errors: vec!["DB connection failed".to_string()], }; assert_eq!(result.extraction_errors.len(), 1); assert!(result.extraction_errors[0].contains("connection")); } }