"""SQS long-poll consumer. One process per task, not per gunicorn worker.""" from __future__ import annotations import json import logging import os import signal import sys import time import boto3 from server.jobs import run_job logger = logging.getLogger(__name__) logging.basicConfig(level=logging.INFO, stream=sys.stderr) _running = True def _stop(_signum, _frame) -> None: global _running _running = False def main() -> None: signal.signal(signal.SIGTERM, _stop) signal.signal(signal.SIGINT, _stop) queue_url = os.environ.get("JOBS_QUEUE_URL", "").strip() if not queue_url: logger.info("JOBS_QUEUE_URL unset; worker idle") while _running: time.sleep(1) return sqs = boto3.client("sqs") logger.info("Polling jobs queue") while _running: resp = sqs.receive_message( QueueUrl=queue_url, MaxNumberOfMessages=1, WaitTimeSeconds=20, VisibilityTimeout=180, ) for msg in resp.get("Messages", []): receipt = msg["ReceiptHandle"] try: payload = json.loads(msg["Body"]) result = run_job(payload) logger.info("job result %s", result) sqs.delete_message(QueueUrl=queue_url, ReceiptHandle=receipt) except Exception: logger.exception("job failed; leaving message for retry") if __name__ == "__main__": main()