848f804169
- FastAPI 后端 + Vue 3 前端 - Docker Compose 一键部署 - Casdoor OAuth 认证集成 - LogHive 集中式日志 - 设备批量 CSV 导入/导出 - WebSocket 实时状态推送 - 企业微信告警通知 - fping 高性能并发 Ping 检测 Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
95 lines
2.6 KiB
Python
95 lines
2.6 KiB
Python
"""
|
|
定时任务调度器
|
|
|
|
使用 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()
|