Deploy Poimen Memory K8s cluster with ArgoCD tracking (M2.2, M3.5-M3.7)
This commit is contained in:
@@ -0,0 +1,161 @@
|
||||
/// 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 ===");
|
||||
}
|
||||
Reference in New Issue
Block a user