| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112 |
- #!/usr/bin/env python
- """Run real-time control every ten minutes during delivery hours."""
- from __future__ import annotations
- import argparse
- import logging
- import signal
- import time
- from datetime import datetime, timedelta
- from zoneinfo import ZoneInfo
- from realtime_config import RealtimeControlConfig
- from run_once import load_environment, run_cycle, write_run_report
- SHANGHAI = ZoneInfo("Asia/Shanghai")
- logger = logging.getLogger("tencent_realtime_control.scheduler")
- stopping = False
- def request_stop(signum: int, _frame: object) -> None:
- global stopping
- logger.info("Received signal %s; scheduler will stop.", signum)
- stopping = True
- def next_wake(now: datetime, config: RealtimeControlConfig) -> datetime:
- if now.hour >= config.stop_hour:
- tomorrow = now.date() + timedelta(days=1)
- return datetime.combine(
- tomorrow,
- datetime.min.time().replace(hour=config.start_hour),
- SHANGHAI,
- )
- if now.hour < config.start_hour:
- return datetime.combine(
- now.date(),
- datetime.min.time().replace(hour=config.start_hour),
- SHANGHAI,
- )
- seconds = int(now.timestamp())
- next_boundary = (
- seconds // config.poll_seconds + 1
- ) * config.poll_seconds
- return datetime.fromtimestamp(next_boundary, SHANGHAI)
- def parse_args() -> argparse.Namespace:
- parser = argparse.ArgumentParser(description=__doc__)
- parser.add_argument(
- "--apply",
- action="store_true",
- help="Execute Tencent updates. Default scheduler is dry-run.",
- )
- return parser.parse_args()
- def run_forever(*, apply: bool, install_signal_handlers: bool = True) -> None:
- config = RealtimeControlConfig.from_env()
- if install_signal_handlers:
- signal.signal(signal.SIGINT, request_stop)
- signal.signal(signal.SIGTERM, request_stop)
- logger.info(
- "Scheduler started apply=%s hours=%02d:00-%02d:00 interval=%ss",
- apply,
- config.start_hour,
- config.stop_hour,
- config.poll_seconds,
- )
- while not stopping:
- try:
- payload = run_cycle(
- now=datetime.now(SHANGHAI),
- apply=apply,
- )
- report = write_run_report(payload)
- logger.info(
- "Cycle decision=%s cpm=%s updates=%s failures=%s notification=%s report=%s",
- payload.get("decision"),
- payload.get("observed_cpm"),
- (payload.get("summary") or {}).get("api_updates", 0),
- (payload.get("summary") or {}).get("failures", 0),
- ((payload.get("summary") or {}).get("notification") or {}).get(
- "status"
- ),
- report,
- )
- except Exception:
- logger.exception("Real-time control cycle failed")
- wake_at = next_wake(datetime.now(SHANGHAI), config)
- while not stopping:
- remaining = (wake_at - datetime.now(SHANGHAI)).total_seconds()
- if remaining <= 0:
- break
- time.sleep(min(remaining, 5))
- def main() -> None:
- load_environment()
- args = parse_args()
- run_forever(apply=args.apply)
- if __name__ == "__main__":
- logging.basicConfig(
- level=logging.INFO,
- format="%(asctime)s %(levelname)s %(name)s %(message)s",
- )
- main()
|