RandomZ / jobs_worker.py
StormShadow308's picture
feat: async pipeline, job queue, generation hardening, and docs
732b14f
Raw
History Blame Contribute Delete
1.24 kB
#!/usr/bin/env python3
"""Redis jobs worker for report generation (Phase 3).
Usage (from repo root, Redis running, ENABLE_JOB_QUEUE=true, REDIS_URL set):
python jobs_worker.py
Run multiple replicas for throughput; each honours ``JOB_QUEUE_MAX_CONCURRENT``.
"""
from __future__ import annotations
import asyncio
import logging
from app.config import settings
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)
async def main() -> None:
if not settings.enable_job_queue:
raise SystemExit("ENABLE_JOB_QUEUE must be true for jobs_worker.py")
if not (settings.redis_url or "").strip():
raise SystemExit("REDIS_URL must be set for jobs_worker.py")
from app.db.database import init_db
from app.jobs.queue import process_jobs_forever
from app.redis_client import get_redis
from app.services.generation_stale import generation_stale_sweeper_loop
await init_db()
await get_redis()
if int(settings.generation_stale_sweep_seconds) > 0:
asyncio.create_task(generation_stale_sweeper_loop())
logger.info("Jobs worker started queue=%s", settings.job_queue_key)
await process_jobs_forever()
if __name__ == "__main__":
asyncio.run(main())