Appearance
Job Queue Infrastructure
Redis/Valkey-based message bus between Orchestrator and Research Agents. Chosen for sub-millisecond enqueue/dequeue, built-in primitives (lists, streams, pub/sub, sorted sets), and Valkey as an open-source fork with no licensing concerns.
Epic: Queue Implementation
Plan: Build a Redis-based job queue with per-agent-type queues, structured job payloads, result delivery via pub/sub, and dead-letter handling for failed jobs. Use per-agent queues for simpler routing (each worker pool reads its own queue).
Architectural Context: The queue topology uses per-agent queues (jobs:dining, jobs:experiences, etc.) rather than a single stream with consumer groups. This is simpler and sufficient unless rebalancing across agent types is needed. Workers pull via BRPOP, write results to result:{job_id}, and publish completion events via PUBLISH. The Orchestrator subscribes to completion events.
Tasks
- Set up Redis/Valkey instance with persistence (AOF or RDB)
- Define job payload schema matching architecture spec:
json
{
"job_id": "uuid",
"session_id": "uuid",
"orchestrator_id": "orch-123",
"agent_type": "dining",
"task": "Find 3 dinner reservations in Tokyo for Feb 13-17, party of 2",
"constraints": {
"cuisines": ["sushi", "kaiseki"],
"max_price_per_person": 200,
"neighborhoods": ["Ginza", "Shibuya"],
"dietary": ["no_pork"]
},
"curated_shortlist": ["sukiyabashi-jiro", "den", "kyubey"],
"result_format": "RestaurantOption[]",
"created_at": "2026-01-15T10:30:00Z",
"ttl_seconds": 60
}- Implement job enqueue function with schema validation (reject invalid payloads)
- Build per-agent queues:
jobs:dining,jobs:experiences,jobs:transportation,jobs:flights,jobs:dmc - Implement result storage:
SET result:{job_id} {result} EX {ttl_seconds} - Build pub/sub channel for completion events:
PUBLISH agent:{type}:completed {job_id} - Implement Orchestrator subscriber:
SUBSCRIBE agent:*:completedto receive all completion events - Build dead-letter queue:
LPUSH dead:{agent_type} {payload}on max retries exceeded - Implement job retry logic: max 3 attempts, exponential backoff (1s, 2s, 4s)
- Add job TTL: auto-expire jobs that exceed TTL without completion
- Build queue cleanup: periodic job to remove expired jobs and stale results
- Implement Redis connection pooling: min 5, max 20 connections per service
- Add Redis health check: ping/pong on connection pool
Epic: Queue Monitoring
Plan: Expose queue metrics for autoscaling and alerting. Build visibility into queue depth, job latency, worker health, and dead-letter growth. Queue depth per agent type is the primary autoscaling signal.
Architectural Context: The Orchestrator monitors queue depth via LLEN jobs:{agent_type} to determine if workers need scaling. Job latency (enqueue → completion) is the primary SLO. Worker health is tracked via heartbeat or job completion rate. Dead-letter growth indicates systemic issues.
Tasks
- Implement queue depth monitoring:
LLENcheck per agent type, expose as Prometheus gauge - Build job latency tracking: record timestamp on enqueue, compute delta on completion, expose as histogram
- Implement worker health monitoring: track active workers per agent type via heartbeat or job completion rate
- Build metrics endpoint: Prometheus-format
/metricsendpoint on queue service - Implement metrics: queue_depth, jobs_enqueued_total, jobs_completed_total, jobs_failed_total, job_latency_seconds, dead_letter_depth
- Add alerting thresholds:
- Queue depth > 100 for > 5 minutes (scale up workers)
- Queue depth = 0 for > 10 minutes (scale down workers)
- Job p95 latency > SLO (investigate worker performance)
- Dead-letter depth > 10 (investigate failures)
- Build Grafana dashboards: queue depth over time, job latency distribution, worker throughput
- Implement alerting via Grafana unified alerting for threshold breaches
- Add runbook links in alert descriptions