run_scheduler.py 3.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112
  1. #!/usr/bin/env python
  2. """Run real-time control every ten minutes during delivery hours."""
  3. from __future__ import annotations
  4. import argparse
  5. import logging
  6. import signal
  7. import time
  8. from datetime import datetime, timedelta
  9. from zoneinfo import ZoneInfo
  10. from realtime_config import RealtimeControlConfig
  11. from run_once import load_environment, run_cycle, write_run_report
  12. SHANGHAI = ZoneInfo("Asia/Shanghai")
  13. logger = logging.getLogger("tencent_realtime_control.scheduler")
  14. stopping = False
  15. def request_stop(signum: int, _frame: object) -> None:
  16. global stopping
  17. logger.info("Received signal %s; scheduler will stop.", signum)
  18. stopping = True
  19. def next_wake(now: datetime, config: RealtimeControlConfig) -> datetime:
  20. if now.hour >= config.stop_hour:
  21. tomorrow = now.date() + timedelta(days=1)
  22. return datetime.combine(
  23. tomorrow,
  24. datetime.min.time().replace(hour=config.start_hour),
  25. SHANGHAI,
  26. )
  27. if now.hour < config.start_hour:
  28. return datetime.combine(
  29. now.date(),
  30. datetime.min.time().replace(hour=config.start_hour),
  31. SHANGHAI,
  32. )
  33. seconds = int(now.timestamp())
  34. next_boundary = (
  35. seconds // config.poll_seconds + 1
  36. ) * config.poll_seconds
  37. return datetime.fromtimestamp(next_boundary, SHANGHAI)
  38. def parse_args() -> argparse.Namespace:
  39. parser = argparse.ArgumentParser(description=__doc__)
  40. parser.add_argument(
  41. "--apply",
  42. action="store_true",
  43. help="Execute Tencent updates. Default scheduler is dry-run.",
  44. )
  45. return parser.parse_args()
  46. def run_forever(*, apply: bool, install_signal_handlers: bool = True) -> None:
  47. config = RealtimeControlConfig.from_env()
  48. if install_signal_handlers:
  49. signal.signal(signal.SIGINT, request_stop)
  50. signal.signal(signal.SIGTERM, request_stop)
  51. logger.info(
  52. "Scheduler started apply=%s hours=%02d:00-%02d:00 interval=%ss",
  53. apply,
  54. config.start_hour,
  55. config.stop_hour,
  56. config.poll_seconds,
  57. )
  58. while not stopping:
  59. try:
  60. payload = run_cycle(
  61. now=datetime.now(SHANGHAI),
  62. apply=apply,
  63. )
  64. report = write_run_report(payload)
  65. logger.info(
  66. "Cycle decision=%s cpm=%s updates=%s failures=%s notification=%s report=%s",
  67. payload.get("decision"),
  68. payload.get("observed_cpm"),
  69. (payload.get("summary") or {}).get("api_updates", 0),
  70. (payload.get("summary") or {}).get("failures", 0),
  71. ((payload.get("summary") or {}).get("notification") or {}).get(
  72. "status"
  73. ),
  74. report,
  75. )
  76. except Exception:
  77. logger.exception("Real-time control cycle failed")
  78. wake_at = next_wake(datetime.now(SHANGHAI), config)
  79. while not stopping:
  80. remaining = (wake_at - datetime.now(SHANGHAI)).total_seconds()
  81. if remaining <= 0:
  82. break
  83. time.sleep(min(remaining, 5))
  84. def main() -> None:
  85. load_environment()
  86. args = parse_args()
  87. run_forever(apply=args.apply)
  88. if __name__ == "__main__":
  89. logging.basicConfig(
  90. level=logging.INFO,
  91. format="%(asctime)s %(levelname)s %(name)s %(message)s",
  92. )
  93. main()