import logging from typing import Optional, Dict from apscheduler.schedulers.asyncio import AsyncIOScheduler from apscheduler.triggers.interval import IntervalTrigger from banban.service.pending_voice_message import pending_voice_message_service from services.connection_manager import connection_manager from handlers.mqtt_handler import TalkingQMQTTService from services.offline_audio_cache import offline_audio_cache from config import settings from utils.logger import session_logger as logger class TaskScheduler: _instance = None def __init__(self): self._scheduler = AsyncIOScheduler(timezone="Asia/Shanghai") @classmethod def get_instance(cls) -> "TaskScheduler": if cls._instance is None: cls._instance = cls() return cls._instance @classmethod def reset_instance(cls): if cls._instance is not None: cls._instance.shutdown() cls._instance = None def start(self): if not self._scheduler.running: self._scheduler.start() logger.system_info("", "定时任务调度器已启动") def shutdown(self): if self._scheduler.running: self._scheduler.shutdown(wait=False) logger.system_info("", "定时任务调度器已关闭") async def _execute_task(self): try: service = await TalkingQMQTTService.get_instance() pending_devices = await pending_voice_message_service.list_devices_with_pending() audio_cache = await offline_audio_cache.get_all_audio_cache() fallback_devices = list(audio_cache.keys()) if audio_cache else [] device_ids = list(dict.fromkeys([*pending_devices, *fallback_devices])) if not device_ids: logger.warning("", "", "无离线音频缓存,跳过本次执行") return for device_id in device_ids: websocket = await connection_manager.get_connection(device_id) if websocket is not None: audio_url = f"http://{settings.server_host}:{settings.server_port}/assets/audio/new_message.mp3" await service.send_nfc_notice(device_id, audio_url) logger.info(device_id, "", f"[定时任务] 发送音频成功: device={device_id}, audio_url={audio_url}") else: logger.info(device_id, "", f"[定时任务] 跳过音频发送: 设备 {device_id} 离线") except Exception as e: logger.error("", "", f"[定时任务] 执行失败: {e}") def add_interval_task(self, interval_seconds: int = 600): task_id = "Voice_Message_Task" task_name = "Voice_Message_Task" job = self._scheduler.add_job( self._execute_task, trigger=IntervalTrigger(seconds=interval_seconds), id=task_id, name=task_name, replace_existing=True, ) logger.system_info("", f"[定时任务] 已添加: {task_id} | 间隔={interval_seconds}s | 下次执行={job.next_run_time}")