File size: 1,236 Bytes
732b14f
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
#!/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())