__main__.py 1.3 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546
  1. from __future__ import annotations
  2. import logging
  3. import os
  4. import signal
  5. import socket
  6. import threading
  7. from supply_infra.config import get_infra_settings
  8. from supply_infra.db import dispose_engine
  9. from supply_infra.pipeline.health import touch_scheduler_heartbeat
  10. from supply_infra.scheduler.app import start_scheduler, stop_scheduler
  11. def main() -> None:
  12. logging.basicConfig(
  13. level=logging.INFO,
  14. format="%(asctime)s [%(levelname)s] %(name)s: %(message)s",
  15. )
  16. stopped = threading.Event()
  17. def _stop(*_args: object) -> None:
  18. stopped.set()
  19. signal.signal(signal.SIGTERM, _stop)
  20. signal.signal(signal.SIGINT, _stop)
  21. scheduler = start_scheduler()
  22. if scheduler is None:
  23. logging.getLogger(__name__).warning("Scheduler process exits because it is disabled")
  24. return
  25. settings = get_infra_settings()
  26. owner = f"{socket.gethostname()}:{os.getpid()}"
  27. try:
  28. while not stopped.is_set():
  29. try:
  30. touch_scheduler_heartbeat(owner=owner, settings=settings)
  31. except Exception:
  32. logging.getLogger(__name__).exception("Scheduler heartbeat failed")
  33. stopped.wait(settings.pipeline_heartbeat_seconds)
  34. finally:
  35. stop_scheduler()
  36. dispose_engine()
  37. if __name__ == "__main__":
  38. main()