73 lines
1.8 KiB
Python
73 lines
1.8 KiB
Python
"""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()
|