"""PostgreSQL-backed worker for registered feature jobs.""" import asyncio import inspect import logging import os from pathlib import Path import socket import time from dotenv import load_dotenv from core import jobs from core.registry import discover_modules load_dotenv(Path(__file__).resolve().parents[1] / ".env", override=False) logging.basicConfig(level=os.getenv("LOG_LEVEL", "INFO")) logger = logging.getLogger(__name__) POLL_INTERVAL = float(os.getenv("JOB_POLL_INTERVAL", 5)) JOB_BATCH_SIZE = int(os.getenv("JOB_BATCH_SIZE", 20)) JOB_LEASE_SECONDS = int(os.getenv("JOB_LEASE_SECONDS", 300)) WORKER_ID = f"scheduler:{socket.gethostname()}:{os.getpid()}" module_registry = discover_modules() def runJob(job): handler = module_registry.get_job_handler(job["job_type"]) if not handler: jobs.fail_job(job["id"], WORKER_ID, f"unknown job type: {job['job_type']}") return try: result = handler(job, WORKER_ID) if inspect.isawaitable(result): asyncio.run(result) current = jobs.get_job(job["id"]) if current and current["status"] == "running": jobs.complete_job(job["id"], WORKER_ID) except Exception as error: logger.exception("Job failed: %s", job["id"]) jobs.fail_job(job["id"], WORKER_ID, error) def pollJobs(): claimed = jobs.claim_due_jobs( WORKER_ID, limit=JOB_BATCH_SIZE, lease_seconds=JOB_LEASE_SECONDS, ) for job in claimed: runJob(job) return len(claimed) def daemonLoop(): logger.info("Scheduler starting as %s", WORKER_ID) while True: try: claimed = pollJobs() except Exception: logger.exception("Scheduler poll failed") claimed = 0 if claimed == 0: time.sleep(POLL_INTERVAL) if __name__ == "__main__": daemonLoop()