run_scheduler.py 3.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107
  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 main() -> None:
  47. load_environment()
  48. args = parse_args()
  49. config = RealtimeControlConfig.from_env()
  50. signal.signal(signal.SIGINT, request_stop)
  51. signal.signal(signal.SIGTERM, request_stop)
  52. logger.info(
  53. "Scheduler started apply=%s hours=%02d:00-%02d:00 interval=%ss",
  54. args.apply,
  55. config.start_hour,
  56. config.stop_hour,
  57. config.poll_seconds,
  58. )
  59. while not stopping:
  60. try:
  61. payload = run_cycle(
  62. now=datetime.now(SHANGHAI),
  63. apply=args.apply,
  64. )
  65. report = write_run_report(payload)
  66. logger.info(
  67. "Cycle decision=%s cpm=%s updates=%s failures=%s notification=%s report=%s",
  68. payload.get("decision"),
  69. payload.get("observed_cpm"),
  70. (payload.get("summary") or {}).get("api_updates", 0),
  71. (payload.get("summary") or {}).get("failures", 0),
  72. ((payload.get("summary") or {}).get("notification") or {}).get(
  73. "status"
  74. ),
  75. report,
  76. )
  77. except Exception:
  78. logger.exception("Real-time control cycle failed")
  79. wake_at = next_wake(datetime.now(SHANGHAI), config)
  80. while not stopping:
  81. remaining = (wake_at - datetime.now(SHANGHAI)).total_seconds()
  82. if remaining <= 0:
  83. break
  84. time.sleep(min(remaining, 5))
  85. if __name__ == "__main__":
  86. logging.basicConfig(
  87. level=logging.INFO,
  88. format="%(asctime)s %(levelname)s %(name)s %(message)s",
  89. )
  90. main()