You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.

87 lines
2.8 KiB

"""消息轮询引擎"""
from __future__ import annotations
import logging
import time
from collections.abc import Callable
from ..config import Config
from ..models.message import Message
from ..ui.message import MessageList
from .event import MessageEvent
logger = logging.getLogger(__name__)
# 回调签名
MessageCallback = Callable[[MessageEvent], None]
class PollingEngine:
"""以轮询方式持续从聊天面板提取新消息"""
def __init__(self, message_list: MessageList, config: Config) -> None:
self._message_list = message_list
self.config = config
self._callbacks: list[MessageCallback] = []
self._seen_ids: set[str] = set()
self._last_poll: float = 0
self._running = False
# ------------------------------------------------------------------
# 回调注册
# ------------------------------------------------------------------
def on_message(self, callback: MessageCallback) -> None:
"""注册新消息回调"""
self._callbacks.append(callback)
# ------------------------------------------------------------------
# 轮询控制
# ------------------------------------------------------------------
def poll_once(self, conversation: str = "") -> list[Message]:
"""执行一次消息提取,返回去重后的新消息"""
try:
messages = self._message_list.extract_messages(conversation)
except Exception:
logger.exception("提取消息失败")
return []
new_messages: list[Message] = []
for msg in messages:
if msg.msg_id not in self._seen_ids:
self._seen_ids.add(msg.msg_id)
new_messages.append(msg)
# 清理过期的去重 ID(防止内存泄漏)
if len(self._seen_ids) > 10_000:
self._seen_ids = set(list(self._seen_ids)[-5_000:])
# 触发回调
for msg in new_messages:
event = MessageEvent(message=msg)
for cb in self._callbacks:
try:
cb(event)
except Exception:
logger.exception("消息回调执行失败")
self._last_poll = time.time()
return new_messages
def start_loop(self, conversation: str = "", stop_event=None) -> None:
"""持续轮询,直到 stop_event 被 set 或调用 stop()"""
self._running = True
logger.info("轮询引擎启动 (间隔 %.1fs)", self.config.poll_interval)
while self._running:
if stop_event is not None and stop_event.is_set():
break
self.poll_once(conversation)
time.sleep(self.config.poll_interval)
logger.info("轮询引擎已停止")
def stop(self) -> None:
self._running = False