""" 定时任务调度器 使用 asyncio 循环驱动 Ping 引擎,协调 Pinger 和 Alerter。 """ import asyncio import logging from datetime import datetime from sqlalchemy.ext.asyncio import AsyncSession from app.config import settings from app.services.pinger import Pinger from app.services.alerter import Alerter from app.core.deps import async_session logger = logging.getLogger("pingwatch.scheduler") class PingScheduler: """ 调度器职责: 1. 按间隔驱动 Ping 引擎 2. 每轮结束后触发 Alerter 处理待发送告警 3. 控制并发和清理 """ def __init__(self): self._pinger = Pinger() self._alerter = Alerter() self._running = False self._task: asyncio.Task | None = None # 注册状态变化回调 self._pinger.on_state_change(self._on_state_change) async def _on_state_change(self, change): """收到设备状态变化,转给 alerter""" async with async_session() as db: try: await self._alerter.on_state_change(change, db) except Exception as e: logger.error(f"告警处理异常: {e}", exc_info=True) async def _run_loop(self): """主循环""" logger.info("Ping 调度器已启动") self._running = True while self._running: try: async with async_session() as db: # 执行一轮 ping results = await self._pinger.run_one_round(db) if results: # 直接用启用的设备数 total = len(results) # 处理待发送告警 await self._alerter.flush_pending(db, total) except asyncio.CancelledError: break except Exception as e: logger.error(f"调度器异常: {e}", exc_info=True) # 等待下一轮 await asyncio.sleep(settings.PING_INTERVAL_SECONDS) logger.info("Ping 调度器已停止") def start(self): """启动调度器(后台任务)""" if self._running: logger.warning("调度器已在运行") return self._task = asyncio.create_task(self._run_loop()) async def stop(self): """停止调度器""" self._running = False if self._task: self._task.cancel() try: await self._task except asyncio.CancelledError: pass self._task = None # 全局调度器实例 scheduler = PingScheduler()