# M8.2 — Gateway Queue Adapter for SQS/kmsvc **Status**: Implementation complete **Version**: 1.0 **Architecture**: Unified queue API via api.riotpiao.com gateway --- ## Overview The Gateway Queue Adapter provides a unified interface for enqueueing chunk dual-write operations via the `api.riotpiao.com` gateway. Rather than connecting directly to kmsvc gRPC, this adapter uses standard HTTP/REST with JWT bearer tokens. ### Design Rationale ``` ┌─────────────────────────────────────────────────────────────────┐ │ Traditional Direct gRPC Approach (NOT used) │ ├─────────────────────────────────────────────────────────────────┤ │ │ │ mem-cli kmsvc (gRPC) │ │ │ │ │ │ │──── gRPC stub ───────>│ (complex connection mgmt) │ │ │ └─────────────────────────────────────────────────────────────────┘ ┌─────────────────────────────────────────────────────────────────┐ │ NEW: Gateway-based approach (THIS IMPLEMENTATION) │ ├─────────────────────────────────────────────────────────────────┤ │ │ │ mem-cli api.riotpiao.com kmsvc │ │ │ │ │ │ │ │─ HTTP + JWT ──────>│──RoutedBy──────>│ │ │ │ (Bearer token) │ X-Service: sqs │ │ │ │ │ │ │ │ (gateway validates JWT before routing)│ │ │ └─────────────────────────────────────────────────────────────────┘ ``` **Benefits**: - ✅ JWT tokens handled by Authentik (same as HTTP API) - ✅ Standard HTTP/REST interface (easier debugging via curl) - ✅ Leverage existing API gateway infrastructure - ✅ No direct gRPC connection management - ✅ Unified authentication across all services --- ## API Reference ### Trait: `QueueAdapter` ```rust #[async_trait] pub trait QueueAdapter: Send + Sync { async fn send_chunk( &self, chunk_id: Uuid, body: String, project: String, attributes: HashMap, ) -> Result; async fn receive_chunks( &self, max_messages: i32, visibility_timeout_secs: i32, project: Option<&str>, ) -> Result>; async fn delete_chunk( &self, message_id: &str, receipt_handle: &str, ) -> Result<()>; async fn change_visibility( &self, message_id: &str, receipt_handle: &str, visibility_timeout_secs: i32, ) -> Result<()>; async fn send_to_dlq( &self, message_id: &str, receipt_handle: &str, reason: &str, ) -> Result<()>; async fn get_stats(&self, project: Option<&str>) -> Result; async fn purge(&self, project: Option<&str>) -> Result; async fn health_check(&self) -> Result<()>; } ``` ### Queue Message Format ```rust pub struct QueueMessage { pub message_id: String, // From SQS pub chunk_id: Uuid, // Original chunk ID pub body: String, // Serialized chunk data pub receive_count: i32, // Number of receives pub receipt_handle: String, // For delete/visibility ops pub project: String, // Project context pub attributes: HashMap, // Metadata } ``` ### Token Provider Trait ```rust #[async_trait] pub trait TokenProvider: Send + Sync { async fn token(&self) -> Result; } ``` Implementations: - `StaticTokenProvider` — Fixed token (testing) - `AuthentikTokenProvider` — OAuth2 client credentials flow (production) --- ## Usage Examples ### Setup: Static Token (Testing) ```rust use mem_cli::gateway_queue_adapter::GatewayQueueAdapter; use mem_cli::queue_adapter::QueueAdapter; use uuid::Uuid; use std::collections::HashMap; #[tokio::main] async fn main() -> anyhow::Result<()> { // Create adapter with static token let adapter = GatewayQueueAdapter::with_static_token( "https://api.riotpiao.com".to_string(), "eyJ...my-jwt-token".to_string(), ); // Queue a chunk let msg_id = adapter.send_chunk( Uuid::new_v4(), r#"{"content": "hello world"}"#.to_string(), "myproject".to_string(), HashMap::new(), ).await?; println!("Queued: {}", msg_id); Ok(()) } ``` ### Setup: Authentik Token (Production) ```rust let adapter = GatewayQueueAdapter::with_authentik( "https://api.riotpiao.com".to_string(), "https://authentik.riotpiao.com/application/o/poimen-memory/".to_string(), "poimen-memory".to_string(), // client_id "your-client-secret".to_string(), ); // Token is automatically refreshed when expired ``` ### Queue a Chunk ```rust let mut attrs = std::collections::HashMap::new(); attrs.insert("source".to_string(), "obsidian".to_string()); attrs.insert("level".to_string(), "L0".to_string()); attrs.insert("breadcrumb".to_string(), serde_json::to_string(&vec!["root", "section"])?, ); let message_id = adapter.send_chunk( Uuid::new_v4(), serde_json::json!({ "content": "chunk text", "metadata": "...", }).to_string(), "myproject".to_string(), attrs, ).await?; tracing::info!("Chunk queued: {}", message_id); ``` ### Receive Messages (Long-poll) ```rust // Receive up to 10 messages, wait up to 20 seconds for availability let messages = adapter.receive_chunks( 10, // max_messages (1-10) 30, // visibility_timeout_secs Some("myproject"), // optional project filter ).await?; for msg in messages { println!("Message ID: {}", msg.message_id); println!("Receive count: {}", msg.receive_count); println!("Receipt handle: {}", msg.receipt_handle); // Process the message... match process_chunk(&msg).await { Ok(_) => { // Delete on success adapter.delete_chunk(&msg.message_id, &msg.receipt_handle).await?; } Err(e) if msg.receive_count < 3 => { // Retry: extend visibility for 5 minutes adapter.change_visibility( &msg.message_id, &msg.receipt_handle, 300, ).await?; } Err(e) => { // Max retries: send to DLQ adapter.send_to_dlq( &msg.message_id, &msg.receipt_handle, &e.to_string(), ).await?; } } } ``` ### Health Check ```rust if let Err(e) = adapter.health_check().await { eprintln!("Gateway unavailable: {}", e); } ``` ### Monitor Queue ```rust let stats = adapter.get_stats(Some("myproject")).await?; println!("Available: {}", stats.available_messages); println!("In-flight: {}", stats.in_flight_messages); println!("DLQ: {}", stats.dead_letter_messages); println!("Processed: {}", stats.total_processed); println!("Avg delay: {}s", stats.average_delay_secs); ``` --- ## HTTP Message Flow ### 1. Send Chunk (POST) **Request**: ```bash POST https://api.riotpiao.com/ X-Service: sqs Authorization: Bearer eyJ... Content-Type: application/json { "messageBody": "aGVsbG8gd29ybGQ=", # Base64-encoded chunk data "messageAttributes": { "values": { "chunk_id": "550e8400-e29b-41d4-a716-446655440000", "project": "myproject", "source": "obsidian", "level": "L0", "breadcrumb": "[\"root\", \"section\"]" } }, "delaySeconds": 0 } ``` **Response** (200 OK): ```json { "messageId": "d9f94e63-b2c1-4e9f-8c5f-8d5e3c1b7a0f" } ``` ### 2. Receive Messages (GET) **Request**: ```bash GET https://api.riotpiao.com/?X-Service=sqs&queue=poimen-chunks-myproject&maxNumberOfMessages=10&waitTimeSeconds=20&visibilityTimeoutSeconds=30 Authorization: Bearer eyJ... ``` **Response** (200 OK): ```json { "messages": [ { "messageId": "d9f94e63-b2c1-4e9f-8c5f-8d5e3c1b7a0f", "receiptHandle": "AQEBxxxx...", "body": "aGVsbG8gd29ybGQ=", # Base64-encoded "attributes": { "values": { "chunk_id": "550e8400-e29b-41d4-a716-446655440000", "project": "myproject", "source": "obsidian" } }, "receiveCount": 1 } ] } ``` ### 3. Delete Message (DELETE) **Request**: ```bash DELETE https://api.riotpiao.com/ X-Service: sqs Authorization: Bearer eyJ... Content-Type: application/json { "receiptHandle": "AQEBxxxx..." } ``` **Response** (204 No Content) --- ## Integration with DualWriteIndexer The `DualWriteIndexer` uses the queue adapter for concurrent dual-write processing: ```rust use mem_cli::dual_write_indexer::DualWriteIndexer; use mem_cli::gateway_queue_adapter::GatewayQueueAdapter; use std::sync::Arc; // Create queue adapter let queue = Arc::new( GatewayQueueAdapter::with_authentik( "https://api.riotpiao.com".to_string(), issuer, client_id, client_secret, ) ); // Create dual-write indexer with queue let indexer = DualWriteIndexer::new( pg_pool, opensearch_client, queue, ); // Queue chunk for processing let message_id = indexer.queue_chunk(&chunk_input, &embedding).await?; // Concurrent workers receive and process let messages = queue.receive_chunks(10, 30, None).await?; for msg in messages { match indexer.process_queued_chunk(&msg, &embedding).await { Ok(result) => { queue.delete_chunk(&msg.message_id, &msg.receipt_handle).await?; } Err(e) => { queue.change_visibility(&msg.message_id, &msg.receipt_handle, 300).await?; } } } ``` --- ## Error Handling ### Common Errors | Status | Meaning | Recovery | |--------|---------|----------| | `401 Unauthorized` | Missing/expired token | Refresh token via TokenProvider | | `403 Forbidden` | Token valid but no permission | Check JWT claims in Authentik | | `404 Not Found` | Queue doesn't exist | Create queue via Queue CRD | | `429 Too Many Requests` | Rate limited | Implement backoff | | `502 Bad Gateway` | kmsvc unreachable | Retry with exponential backoff | | `503 Service Unavailable` | Gateway overloaded | Circuit breaker pattern | ### Retry Strategy ```rust use std::time::Duration; let mut retries = 0; const MAX_RETRIES: usize = 3; loop { match adapter.send_chunk(...).await { Ok(msg_id) => { tracing::info!("Sent: {}", msg_id); break; } Err(e) if retries < MAX_RETRIES => { retries += 1; let backoff = Duration::from_millis(100 * 2_u64.pow(retries as u32)); tracing::warn!("Retry {} in {:?}: {}", retries, backoff, e); tokio::time::sleep(backoff).await; } Err(e) => { tracing::error!("Max retries exceeded: {}", e); return Err(e); } } } ``` --- ## Configuration (Environment Variables) ```bash # Gateway endpoint export GATEWAY_URL=https://api.riotpiao.com # Authentik (for JWT) export AUTHENTIK_ISSUER=https://authentik.riotpiao.com/application/o/poimen-memory/ export AUTHENTIK_CLIENT_ID=poimen-memory export AUTHENTIK_CLIENT_SECRET=your-secret # Optional: Static token (for testing) export STATIC_JWT_TOKEN=eyJ... ``` --- ## Testing ### Unit Tests ```bash cargo test --lib gateway_queue_adapter ``` ### Integration Tests ```bash # Requires running api.riotpiao.com cargo test --test it_gateway_queue_adapter -- --ignored ``` ### Manual Testing with curl ```bash # Get token TOKEN=$(curl -s -X POST https://authentik.riotpiao.com/application/o/token/ \ -d "grant_type=client_credentials&client_id=poimen-memory&client_secret=secret" \ | jq -r '.access_token') # Send message curl -X POST https://api.riotpiao.com/ \ -H "X-Service: sqs" \ -H "Authorization: Bearer $TOKEN" \ -H "Content-Type: application/json" \ -d '{"messageBody": "aGVsbG8gd29ybGQ=", "messageAttributes": {"values": {}}}' # Receive messages curl -X GET "https://api.riotpiao.com/?X-Service=sqs&queue=poimen-chunks-test&maxNumberOfMessages=10&waitTimeSeconds=20" \ -H "Authorization: Bearer $TOKEN" | jq . ``` --- ## Security Notes ✅ **JWT Validation**: Gateway validates token signature and claims before routing ✅ **Bearer Token Format**: Strictly requires `Authorization: Bearer ` ✅ **Token Expiry**: Automatic refresh via TokenProvider ✅ **HTTPS Only**: All calls to api.riotpiao.com are encrypted ✅ **Header Validation**: X-Service header validated by gateway --- ## Future Enhancements - [ ] ChangeMessageVisibility support in gateway - [ ] GetQueueAttributes for monitoring - [ ] Batch operations (SendMessageBatch, DeleteMessageBatch) - [ ] Circuit breaker pattern for fault tolerance - [ ] Metrics export (Prometheus) - [ ] Tracing integration (OpenTelemetry) --- ## References - [SERVICE-USAGE.md](../homelab-frontend/docs/SERVICE-USAGE.md) — Gateway usage guide - [kmsvc-SDK README](../kmsvc-SDK/README.md) — Underlying SQS implementation - [Authentik Docs](https://goauthentik.io/) — JWT token provider