Skip to content

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:*:completed to 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: LLEN check 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 /metrics endpoint 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

Marchay Platform Documentation