diff --git a/EvoScientist/channels/base.py b/EvoScientist/channels/base.py index 61938e2..1906a8f 100644 --- a/EvoScientist/channels/base.py +++ b/EvoScientist/channels/base.py @@ -865,7 +865,7 @@ class Channel(ChannelPlugin, ABC): Convenience method for subclass ``_on_message`` handlers. """ - msg = self._build_inbound(raw) + msg = await self._build_inbound_async(raw) if msg is None: return if raw.message_id: diff --git a/EvoScientist/channels/dingtalk/channel.py b/EvoScientist/channels/dingtalk/channel.py index 1b1bffb..d61c9e0 100644 --- a/EvoScientist/channels/dingtalk/channel.py +++ b/EvoScientist/channels/dingtalk/channel.py @@ -347,7 +347,7 @@ class DingTalkChannel(Channel, WebSocketMixin, TokenMixin): pass self._ws_task = None await self._stop_ws() - if self._http_client: + if hasattr(self, "_http_client") and self._http_client: await self._http_client.aclose() self._http_client = None self._access_token = None diff --git a/EvoScientist/channels/email/channel.py b/EvoScientist/channels/email/channel.py index 1336a07..b36eae0 100644 --- a/EvoScientist/channels/email/channel.py +++ b/EvoScientist/channels/email/channel.py @@ -81,7 +81,7 @@ class EmailChannel(Channel, PollingMixin): cfg = self.config if not cfg.imap_host or not cfg.imap_username: raise ChannelError("Email imap_host and imap_username are required") - loop = asyncio.get_event_loop() + loop = asyncio.get_running_loop() await loop.run_in_executor(None, self._connect_imap) self._running = True logger.info(f"Email channel started (IMAP: {cfg.imap_host}, poll {cfg.poll_interval}s)") @@ -120,7 +120,7 @@ class EmailChannel(Channel, PollingMixin): self._connect_imap() async def _poll_once(self) -> None: - loop = asyncio.get_event_loop() + loop = asyncio.get_running_loop() messages = await loop.run_in_executor(None, self._fetch_unseen) for m in messages: await self._process_email(m) @@ -270,7 +270,7 @@ class EmailChannel(Channel, PollingMixin): pass async def _send_chunk(self, chat_id, formatted_text, raw_text, reply_to, metadata): - loop = asyncio.get_event_loop() + loop = asyncio.get_running_loop() try: await loop.run_in_executor( None, self._smtp_send_html, chat_id, formatted_text, raw_text, metadata or {}, @@ -339,7 +339,7 @@ class EmailChannel(Channel, PollingMixin): metadata: dict | None = None, ) -> bool: """Send a file as an email attachment via SMTP.""" - loop = asyncio.get_event_loop() + loop = asyncio.get_running_loop() await loop.run_in_executor( None, self._smtp_send_attachment, recipient, file_path, caption, metadata or {}, ) diff --git a/EvoScientist/channels/qq/channel.py b/EvoScientist/channels/qq/channel.py index 35b47e0..5f6cd82 100644 --- a/EvoScientist/channels/qq/channel.py +++ b/EvoScientist/channels/qq/channel.py @@ -69,6 +69,7 @@ class QQChannel(Channel): self._processed_ids: deque = deque(maxlen=1000) self._msg_seq: dict[str, int] = {} # msg_id -> next seq number self._msg_seq_order: deque = deque(maxlen=500) + self._msg_seq_ids: set[str] = set() # companion set for O(1) lookup # ── Lifecycle ───────────────────────────────────────────────── @@ -167,10 +168,12 @@ class QQChannel(Channel): """Return the next msg_seq for *msg_id* and increment the counter.""" seq = self._msg_seq.get(msg_id, 1) self._msg_seq[msg_id] = seq + 1 - if msg_id not in set(self._msg_seq_order): + if msg_id not in self._msg_seq_ids: self._msg_seq_order.append(msg_id) + self._msg_seq_ids.add(msg_id) if len(self._msg_seq_order) > 500: oldest = self._msg_seq_order.popleft() + self._msg_seq_ids.discard(oldest) self._msg_seq.pop(oldest, None) return seq diff --git a/EvoScientist/channels/signal/channel.py b/EvoScientist/channels/signal/channel.py index e21986b..a918e55 100644 --- a/EvoScientist/channels/signal/channel.py +++ b/EvoScientist/channels/signal/channel.py @@ -45,6 +45,7 @@ class SignalChannel(Channel): # Cache message_id → sender for reaction targetAuthor (bounded) self._msg_senders: dict[str, str] = {} self._msg_senders_order: deque = deque(maxlen=200) + self._listen_task: asyncio.Task | None = None async def start(self) -> None: if not self.config.phone_number: @@ -69,7 +70,7 @@ class SignalChannel(Channel): self._listen_task = asyncio.create_task(self._listen_loop()) async def _cleanup(self) -> None: - if hasattr(self, "_listen_task") and self._listen_task: + if self._listen_task: self._listen_task.cancel() self._listen_task = None # Cancel any pending RPC futures @@ -180,7 +181,8 @@ class SignalChannel(Channel): try: await self._connect() except Exception: - pass + logger.warning("Signal reconnect failed, exiting listen loop") + break async def _handle_rpc(self, data: dict) -> None: """Handle a JSON RPC message from signal-cli."""