| 12345678910111213141516171819202122232425262728293031323334353637383940414243444546 |
- 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()
|