#!/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()