/// End-to-end pipeline test: source -> chunker -> records consumed /// This proves the full system works, not just individual components use mem_ingest::PiSessionSource; use mem_chunk::{RecordSource, chunks, ChunkPolicy}; use mem_core::Role; use futures::stream::StreamExt; use std::path::PathBuf; #[tokio::test] async fn e2e_pi_session_full_pipeline() { let fixture = PathBuf::from("fixtures/pi-session-small.jsonl"); println!("\n=== E2E PIPELINE TEST ==="); println!("Fixture: {}", fixture.display()); // Step 1: Source reads project key println!("\nStep 1: Reading project key from source..."); let source = PiSessionSource::new(fixture.clone()); let project = source.read_project_key().await.expect("Failed to read project key"); println!("✓ Project key: {}", project); assert_eq!(project, "/tmp/my-project"); // Step 2: Source streams records println!("\nStep 2: Streaming records from source..."); let source = PiSessionSource::new(fixture.clone()); let mut records_stream = source.records(); let mut records = Vec::new(); let mut user_count = 0; let mut assistant_count = 0; let mut tool_count = 0; let mut system_count = 0; while let Some(result) = records_stream.next().await { match result { Ok(record) => { match record.role { Role::User => user_count += 1, Role::Assistant => assistant_count += 1, Role::ToolResult => tool_count += 1, Role::System => system_count += 1, } records.push(record); } Err(e) => panic!("Error reading record: {}", e), } } println!("✓ Records streamed: {}", records.len()); println!(" - User: {}", user_count); println!(" - Assistant: {}", assistant_count); println!(" - ToolResult: {}", tool_count); println!(" - System: {}", system_count); assert!(records.len() > 0, "Should have parsed records"); assert!(user_count > 0, "Should have user messages"); assert!(assistant_count > 0, "Should have assistant messages"); // Step 3: Chunker processes records println!("\nStep 3: Chunking records through policy..."); let source = PiSessionSource::new(fixture.clone()); let policy = ChunkPolicy::default(); let mut chunk_stream = chunks(source, policy); let mut chunks_vec = Vec::new(); let mut total_chunk_records = 0; let mut min_tokens = usize::MAX; let mut max_tokens = 0; while let Some(result) = chunk_stream.next().await { match result { Ok(chunk) => { let token_count = chunk.tokens; total_chunk_records += chunk.records.len(); min_tokens = min_tokens.min(token_count); max_tokens = max_tokens.max(token_count); chunks_vec.push(chunk); } Err(e) => panic!("Error chunking: {}", e), } } println!("✓ Chunks produced: {}", chunks_vec.len()); println!(" - Total records in chunks: {}", total_chunk_records); println!(" - Token range: {} to {} (budget: 5000)", min_tokens, max_tokens); assert!(chunks_vec.len() > 0, "Should produce at least one chunk"); // Step 4: Verify losslessness println!("\nStep 4: Verifying losslessness..."); assert_eq!(records.len(), total_chunk_records, "All records must flow into chunks without loss"); println!("✓ Lossless: {} records in == {} records out", records.len(), total_chunk_records); // Step 5: Verify chunk integrity println!("\nStep 5: Verifying chunk integrity..."); for chunk in &chunks_vec { assert!(!chunk.records.is_empty(), "Chunk must have records"); assert!(chunk.t > 0, "Turn index must be positive"); // Verify no record was split for record in &chunk.records { assert!(!record.text.is_empty(), "Record must have content"); } } println!("✓ All chunks have valid turn indices and records"); // Verify turn indices are contiguous let mut expected_t = 1u32; for chunk in &chunks_vec { assert_eq!(chunk.t, expected_t, "Turn indices must be contiguous"); expected_t += 1; } println!("✓ Turn indices are contiguous (1..{})", chunks_vec.len()); println!("\n=== E2E PIPELINE SUCCESS ==="); println!("Project: {}", project); println!("Records: {} (user: {}, asst: {}, tool: {}, sys: {})", records.len(), user_count, assistant_count, tool_count, system_count); println!("Chunks: {}", chunks_vec.len()); println!("Lossless: ✓"); println!("Integrity: ✓"); } #[tokio::test] async fn e2e_claude_transcript_full_pipeline() { let fixture = PathBuf::from("fixtures/claude-transcript-small.jsonl"); println!("\n=== E2E CLAUDE PIPELINE TEST ==="); // Same full pipeline but with Claude source use mem_ingest::ClaudeTranscriptSource; let source = ClaudeTranscriptSource::new(fixture.clone()); let project = source.read_project_key().await.expect("Failed to read project"); println!("✓ Project: {}", project); let source = ClaudeTranscriptSource::new(fixture.clone()); let mut records_stream = source.records(); let mut record_count = 0; while let Some(result) = records_stream.next().await { if result.is_ok() { record_count += 1; } } println!("✓ Records: {}", record_count); let source = ClaudeTranscriptSource::new(fixture); let policy = ChunkPolicy::default(); let mut chunk_stream = chunks(source, policy); let mut chunk_count = 0; let mut total = 0; while let Some(Ok(chunk)) = chunk_stream.next().await { chunk_count += 1; total += chunk.records.len(); } println!("✓ Chunks: {}", chunk_count); assert_eq!(record_count, total, "Claude pipeline must also be lossless"); println!("✓ Lossless: {} == {}", record_count, total); println!("\n=== E2E CLAUDE PIPELINE SUCCESS ==="); }