2026-08-28 13:14:39 -07:00
|
|
|
//! Integration tests for M8.2 Queue Worker + DualWriteIndexer
|
|
|
|
|
//!
|
|
|
|
|
//! Tests the full pipeline:
|
|
|
|
|
//! 1. IngestWorker → queue_chunk()
|
|
|
|
|
//! 2. QueueWorker → receive_chunks()
|
|
|
|
|
//! 3. process_queued_chunk() → pgvector + OpenSearch write
|
|
|
|
|
//! 4. Message deletion or retry
|
|
|
|
|
|
2026-08-28 15:33:04 -07:00
|
|
|
use mem_cli::queue_adapter::{InMemoryQueueAdapter, QueueAdapter};
|
2026-08-28 13:14:39 -07:00
|
|
|
use mem_cli::dual_write_indexer::{DualWriteIndexer, ChunkInput};
|
|
|
|
|
use mem_cli::queue_worker::{QueueWorker, QueueWorkerConfig};
|
|
|
|
|
use std::sync::Arc;
|
|
|
|
|
use uuid::Uuid;
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn test_queue_worker_config_default() {
|
|
|
|
|
let config = QueueWorkerConfig::default();
|
|
|
|
|
assert_eq!(config.max_messages_per_batch, 10);
|
|
|
|
|
assert_eq!(config.visibility_timeout_secs, 300);
|
|
|
|
|
assert_eq!(config.wait_time_secs, 20);
|
|
|
|
|
assert_eq!(config.max_retries, 3);
|
|
|
|
|
assert!(!config.enable_metrics);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn test_queue_worker_config_custom() {
|
|
|
|
|
let config = QueueWorkerConfig {
|
|
|
|
|
max_messages_per_batch: 5,
|
|
|
|
|
visibility_timeout_secs: 600,
|
|
|
|
|
wait_time_secs: 30,
|
|
|
|
|
project: Some("test-proj".to_string()),
|
|
|
|
|
max_retries: 5,
|
|
|
|
|
enable_metrics: true,
|
|
|
|
|
..Default::default()
|
|
|
|
|
};
|
|
|
|
|
|
|
|
|
|
assert_eq!(config.max_messages_per_batch, 5);
|
|
|
|
|
assert_eq!(config.max_retries, 5);
|
|
|
|
|
assert!(config.enable_metrics);
|
|
|
|
|
assert_eq!(config.project, Some("test-proj".to_string()));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn test_queue_message_roundtrip() {
|
|
|
|
|
let queue = InMemoryQueueAdapter::new();
|
|
|
|
|
|
|
|
|
|
// Queue a message
|
|
|
|
|
let chunk_id = Uuid::new_v4();
|
|
|
|
|
let body = serde_json::json!({
|
|
|
|
|
"chunk_id": chunk_id,
|
|
|
|
|
"content": "hello world",
|
|
|
|
|
"source": "test",
|
|
|
|
|
}).to_string();
|
|
|
|
|
|
|
|
|
|
let mut attrs = std::collections::HashMap::new();
|
|
|
|
|
attrs.insert("source".to_string(), "test".to_string());
|
|
|
|
|
attrs.insert("level".to_string(), "L0".to_string());
|
|
|
|
|
|
|
|
|
|
let msg_id = queue
|
|
|
|
|
.send_chunk(chunk_id, body.clone(), "test-proj".to_string(), attrs)
|
|
|
|
|
.await
|
|
|
|
|
.unwrap();
|
|
|
|
|
|
|
|
|
|
// Receive it back
|
|
|
|
|
let messages = queue
|
|
|
|
|
.receive_chunks(10, 30, Some("test-proj"))
|
|
|
|
|
.await
|
|
|
|
|
.unwrap();
|
|
|
|
|
|
|
|
|
|
assert_eq!(messages.len(), 1);
|
|
|
|
|
assert_eq!(messages[0].message_id, msg_id);
|
|
|
|
|
assert_eq!(messages[0].body, body);
|
|
|
|
|
assert_eq!(messages[0].project, "test-proj");
|
|
|
|
|
|
|
|
|
|
// Delete it
|
|
|
|
|
queue
|
|
|
|
|
.delete_chunk(&messages[0].message_id, &messages[0].receipt_handle)
|
|
|
|
|
.await
|
|
|
|
|
.unwrap();
|
|
|
|
|
|
|
|
|
|
// Queue should be empty
|
|
|
|
|
let messages = queue
|
|
|
|
|
.receive_chunks(10, 30, None)
|
|
|
|
|
.await
|
|
|
|
|
.unwrap();
|
|
|
|
|
|
|
|
|
|
assert_eq!(messages.len(), 0);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn test_queue_multiple_messages() {
|
|
|
|
|
let queue = InMemoryQueueAdapter::new();
|
|
|
|
|
|
|
|
|
|
// Queue multiple messages
|
|
|
|
|
for i in 0..5 {
|
|
|
|
|
let _ = queue
|
|
|
|
|
.send_chunk(
|
|
|
|
|
Uuid::new_v4(),
|
|
|
|
|
format!(r#"{{"content": "msg {}"}}"#, i),
|
|
|
|
|
"test".to_string(),
|
|
|
|
|
std::collections::HashMap::new(),
|
|
|
|
|
)
|
|
|
|
|
.await;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Receive batch of 3
|
|
|
|
|
let messages = queue
|
|
|
|
|
.receive_chunks(3, 30, None)
|
|
|
|
|
.await
|
|
|
|
|
.unwrap();
|
|
|
|
|
|
|
|
|
|
assert_eq!(messages.len(), 3);
|
|
|
|
|
|
|
|
|
|
// Delete all 3
|
|
|
|
|
for msg in messages {
|
|
|
|
|
queue
|
|
|
|
|
.delete_chunk(&msg.message_id, &msg.receipt_handle)
|
|
|
|
|
.await
|
|
|
|
|
.ok();
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Should have 2 left
|
|
|
|
|
let remaining = queue.receive_chunks(10, 30, None).await.unwrap();
|
|
|
|
|
assert_eq!(remaining.len(), 2);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn test_queue_dlq_transition() {
|
|
|
|
|
let queue = InMemoryQueueAdapter::new();
|
|
|
|
|
|
|
|
|
|
let msg_id = queue
|
|
|
|
|
.send_chunk(
|
|
|
|
|
Uuid::new_v4(),
|
|
|
|
|
"body".to_string(),
|
|
|
|
|
"test".to_string(),
|
|
|
|
|
std::collections::HashMap::new(),
|
|
|
|
|
)
|
|
|
|
|
.await
|
|
|
|
|
.unwrap();
|
|
|
|
|
|
|
|
|
|
// Simulate max retries exceeded
|
|
|
|
|
queue
|
|
|
|
|
.send_to_dlq(&msg_id, "handle-xyz", "max_retries_exceeded")
|
|
|
|
|
.await
|
|
|
|
|
.unwrap();
|
|
|
|
|
|
|
|
|
|
// Should not appear in normal queue anymore
|
|
|
|
|
let messages = queue.receive_chunks(10, 30, None).await.unwrap();
|
|
|
|
|
assert!(messages.is_empty());
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn test_chunk_input_structure() {
|
|
|
|
|
let chunk = ChunkInput {
|
|
|
|
|
content: "test content".to_string(),
|
|
|
|
|
source: "test-source".to_string(),
|
|
|
|
|
project: "test-proj".to_string(),
|
|
|
|
|
level: "L0".to_string(),
|
|
|
|
|
breadcrumb: vec!["root".to_string(), "section".to_string()],
|
|
|
|
|
};
|
|
|
|
|
|
|
|
|
|
assert_eq!(chunk.content, "test content");
|
|
|
|
|
assert_eq!(chunk.level, "L0");
|
|
|
|
|
assert_eq!(chunk.breadcrumb.len(), 2);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[test]
|
|
|
|
|
fn test_chunk_levels_valid() {
|
|
|
|
|
let levels = vec!["L0", "L1", "L2", "R"];
|
|
|
|
|
|
|
|
|
|
for level in levels {
|
|
|
|
|
let chunk = ChunkInput {
|
|
|
|
|
content: "test".to_string(),
|
|
|
|
|
source: "test".to_string(),
|
|
|
|
|
project: "test".to_string(),
|
|
|
|
|
level: level.to_string(),
|
|
|
|
|
breadcrumb: vec![],
|
|
|
|
|
};
|
|
|
|
|
assert_eq!(chunk.level, level);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn test_queue_stats_tracking() {
|
|
|
|
|
let queue = InMemoryQueueAdapter::new();
|
|
|
|
|
|
|
|
|
|
// Queue 3 messages
|
|
|
|
|
for i in 0..3 {
|
|
|
|
|
let _ = queue
|
|
|
|
|
.send_chunk(
|
|
|
|
|
Uuid::new_v4(),
|
|
|
|
|
format!("msg {}", i),
|
|
|
|
|
"test".to_string(),
|
|
|
|
|
std::collections::HashMap::new(),
|
|
|
|
|
)
|
|
|
|
|
.await;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
let stats = queue.get_stats(None).await.unwrap();
|
|
|
|
|
assert_eq!(stats.available_messages, 3);
|
|
|
|
|
assert_eq!(stats.total_processed, 0);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[tokio::test]
|
|
|
|
|
async fn test_message_attributes_preserved() {
|
|
|
|
|
let queue = InMemoryQueueAdapter::new();
|
|
|
|
|
|
|
|
|
|
let mut attrs = std::collections::HashMap::new();
|
|
|
|
|
attrs.insert("custom_key".to_string(), "custom_value".to_string());
|
|
|
|
|
attrs.insert("another".to_string(), "test".to_string());
|
|
|
|
|
|
|
|
|
|
let msg_id = queue
|
|
|
|
|
.send_chunk(
|
|
|
|
|
Uuid::new_v4(),
|
|
|
|
|
"body".to_string(),
|
|
|
|
|
"proj".to_string(),
|
|
|
|
|
attrs.clone(),
|
|
|
|
|
)
|
|
|
|
|
.await
|
|
|
|
|
.unwrap();
|
|
|
|
|
|
|
|
|
|
let messages = queue.receive_chunks(10, 30, None).await.unwrap();
|
|
|
|
|
assert_eq!(messages.len(), 1);
|
|
|
|
|
assert_eq!(messages[0].attributes.get("custom_key"), Some(&"custom_value".to_string()));
|
|
|
|
|
assert_eq!(messages[0].attributes.get("another"), Some(&"test".to_string()));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[test]
|
|
|
|
|
fn test_embedding_size_validation() {
|
|
|
|
|
// Standard embedding size
|
|
|
|
|
let embedding: Vec<f32> = (0..768).map(|i| i as f32).collect();
|
|
|
|
|
assert_eq!(embedding.len(), 768);
|
|
|
|
|
|
|
|
|
|
// Verify nomic-embed-text-v2-moe compatibility
|
|
|
|
|
assert!(embedding.len() > 0);
|
|
|
|
|
assert!(embedding.len() <= 1024);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[test]
|
|
|
|
|
fn test_queue_message_ordering() {
|
|
|
|
|
// Verify that message IDs are unique
|
|
|
|
|
let msg_id_1 = format!("msg-{}", Uuid::new_v4());
|
|
|
|
|
let msg_id_2 = format!("msg-{}", Uuid::new_v4());
|
|
|
|
|
let msg_id_3 = format!("msg-{}", Uuid::new_v4());
|
|
|
|
|
|
|
|
|
|
let ids = vec![msg_id_1, msg_id_2, msg_id_3];
|
|
|
|
|
let unique_ids: std::collections::HashSet<_> = ids.iter().cloned().collect();
|
|
|
|
|
|
|
|
|
|
assert_eq!(unique_ids.len(), 3);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[test]
|
|
|
|
|
fn test_breadcrumb_path_structure() {
|
|
|
|
|
let breadcrumbs = vec![
|
|
|
|
|
vec!["root".to_string()],
|
|
|
|
|
vec!["root".to_string(), "folder".to_string()],
|
|
|
|
|
vec!["root".to_string(), "folder".to_string(), "section".to_string()],
|
|
|
|
|
];
|
|
|
|
|
|
|
|
|
|
for crumb in breadcrumbs {
|
|
|
|
|
let chunk = ChunkInput {
|
|
|
|
|
content: "test".to_string(),
|
|
|
|
|
source: "test".to_string(),
|
|
|
|
|
project: "test".to_string(),
|
|
|
|
|
level: "L0".to_string(),
|
|
|
|
|
breadcrumb: crumb.clone(),
|
|
|
|
|
};
|
|
|
|
|
|
|
|
|
|
assert_eq!(chunk.breadcrumb.len(), crumb.len());
|
|
|
|
|
}
|
|
|
|
|
}
|