162 lines
6.0 KiB
Rust
162 lines
6.0 KiB
Rust
/// 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 ===");
|
|
}
|