Spaces:
Sleeping
Sleeping
| #!/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()) | |