from __future__ import annotations import logging import os import signal import socket import threading from supply_infra.config import get_infra_settings from supply_infra.db import dispose_engine from supply_infra.pipeline.health import touch_scheduler_heartbeat from supply_infra.scheduler.app import start_scheduler, stop_scheduler def main() -> None: logging.basicConfig( level=logging.INFO, format="%(asctime)s [%(levelname)s] %(name)s: %(message)s", ) stopped = threading.Event() def _stop(*_args: object) -> None: stopped.set() signal.signal(signal.SIGTERM, _stop) signal.signal(signal.SIGINT, _stop) scheduler = start_scheduler() if scheduler is None: logging.getLogger(__name__).warning("Scheduler process exits because it is disabled") return settings = get_infra_settings() owner = f"{socket.gethostname()}:{os.getpid()}" try: while not stopped.is_set(): try: touch_scheduler_heartbeat(owner=owner, settings=settings) except Exception: logging.getLogger(__name__).exception("Scheduler heartbeat failed") stopped.wait(settings.pipeline_heartbeat_seconds) finally: stop_scheduler() dispose_engine() if __name__ == "__main__": main()