From b85faf75164f7a9d58a9bd387ff67e6bcb62c055 Mon Sep 17 00:00:00 2001 From: MuXinCG <202322130196@mail.sdu.edu.cn> Date: Tue, 17 Feb 2026 10:54:15 +0800 Subject: [PATCH 1/3] Add Dingtalk Feishu --- EvoScientist/channels/__init__.py | 2 +- EvoScientist/channels/dingtalk/__init__.py | 29 + EvoScientist/channels/dingtalk/channel.py | 354 ++++++++++ EvoScientist/channels/dingtalk/probe.py | 33 + EvoScientist/channels/dingtalk/serve.py | 92 +++ EvoScientist/channels/feishu/__init__.py | 21 + EvoScientist/channels/feishu/channel.py | 773 +++++++++++++++++++++ EvoScientist/channels/feishu/probe.py | 39 ++ EvoScientist/channels/feishu/serve.py | 113 +++ EvoScientist/config/onboard.py | 16 + EvoScientist/config/settings.py | 2 +- tests/test_dingtalk_channel.py | 122 ++++ tests/test_feishu_channel.py | 204 ++++++ 13 files changed, 1798 insertions(+), 2 deletions(-) create mode 100644 EvoScientist/channels/dingtalk/__init__.py create mode 100644 EvoScientist/channels/dingtalk/channel.py create mode 100644 EvoScientist/channels/dingtalk/probe.py create mode 100644 EvoScientist/channels/dingtalk/serve.py create mode 100644 EvoScientist/channels/feishu/__init__.py create mode 100644 EvoScientist/channels/feishu/channel.py create mode 100644 EvoScientist/channels/feishu/probe.py create mode 100644 EvoScientist/channels/feishu/serve.py create mode 100644 tests/test_dingtalk_channel.py create mode 100644 tests/test_feishu_channel.py diff --git a/EvoScientist/channels/__init__.py b/EvoScientist/channels/__init__.py index c910120..abec61f 100644 --- a/EvoScientist/channels/__init__.py +++ b/EvoScientist/channels/__init__.py @@ -1,7 +1,7 @@ """Communication channels for EvoScientist. This module provides an extensible interface for different messaging channels -(iMessage, Telegram, Discord, Slack, WeChat) to communicate with the EvoScientist agent. +(iMessage, Telegram, Discord, Slack, WeChat, DingTalk, Feishu) to communicate with the EvoScientist agent. """ from .base import Channel, RawIncoming, IncomingMessage, OutgoingMessage, chunk_text diff --git a/EvoScientist/channels/dingtalk/__init__.py b/EvoScientist/channels/dingtalk/__init__.py new file mode 100644 index 0000000..2026d82 --- /dev/null +++ b/EvoScientist/channels/dingtalk/__init__.py @@ -0,0 +1,29 @@ +"""DingTalk (钉钉) channel for EvoScientist. + +Uses Stream Mode (WebSocket) for receiving messages — no public IP needed. +Sends replies via HTTP API. + +Usage in config: + channel_enabled = "dingtalk" + dingtalk_client_id = "your_app_key" + dingtalk_client_secret = "your_app_secret" +""" + +from .channel import DingTalkChannel, DingTalkConfig +from ..channel_manager import register_channel, _parse_csv + +__all__ = ["DingTalkChannel", "DingTalkConfig"] + + +def create_from_config(config) -> DingTalkChannel: + allowed = _parse_csv(getattr(config, "dingtalk_allowed_senders", "")) + proxy = getattr(config, "dingtalk_proxy", "") or None + return DingTalkChannel(DingTalkConfig( + client_id=getattr(config, "dingtalk_client_id", ""), + client_secret=getattr(config, "dingtalk_client_secret", ""), + allowed_senders=allowed, + proxy=proxy, + )) + + +register_channel("dingtalk", create_from_config) diff --git a/EvoScientist/channels/dingtalk/channel.py b/EvoScientist/channels/dingtalk/channel.py new file mode 100644 index 0000000..1b1bffb --- /dev/null +++ b/EvoScientist/channels/dingtalk/channel.py @@ -0,0 +1,354 @@ +"""DingTalk channel — refactored with WebSocketMixin + TokenMixin.""" + +import asyncio +import json +import logging +from urllib.parse import quote_plus +from dataclasses import dataclass +from datetime import datetime +from pathlib import Path + +from ..base import Channel, RawIncoming, ChannelError +from ..capabilities import DINGTALK as DINGTALK_CAPS +from ..mixins import WebSocketMixin, TokenMixin +from ..config import BaseChannelConfig + +logger = logging.getLogger(__name__) + +GATEWAY_URL = "https://api.dingtalk.com/v1.0/gateway/connections/open" +TOKEN_URL = "https://api.dingtalk.com/v1.0/oauth2/accessToken" +SEND_URL = "https://api.dingtalk.com/v1.0/robot/oToMessages/batchSend" +MEDIA_SEND_URL = "https://api.dingtalk.com/v1.0/robot/oToMessages/batchSend" +MEDIA_UPLOAD_URL = "https://oapi.dingtalk.com/media/upload" +FILE_DOWNLOAD_URL = "https://api.dingtalk.com/v1.0/robot/messageFiles/download" + + +@dataclass +class DingTalkConfig(BaseChannelConfig): + client_id: str = "" + client_secret: str = "" + text_chunk_limit: int = 4096 + + +class DingTalkChannel(Channel, WebSocketMixin, TokenMixin): + capabilities = DINGTALK_CAPS + name = "dingtalk" + _ready_attrs = ("_http_client", "_access_token") + _non_retryable_patterns = ("invalidauthentication", "forbidden", "40014") + _mention_pattern = r"@\S+\s*" + _mention_strip_count = 1 + + def __init__(self, config: DingTalkConfig): + super().__init__(config) + + async def start(self) -> None: + import httpx + if not self.config.client_id or not self.config.client_secret: + raise ChannelError("DingTalk client_id and client_secret are required") + self._http_client = httpx.AsyncClient(timeout=15, proxy=self.config.proxy) + await self._refresh_token() + self._running = True + logger.info("DingTalk channel starting (Stream Mode)...") + self._ws_task = asyncio.create_task(self._ws_loop()) + + # ── TokenMixin ──────────────────────────────────────────────── + + async def _fetch_token(self) -> tuple[str, int]: + data = await self._api_post(TOKEN_URL, { + "appKey": self.config.client_id, + "appSecret": self.config.client_secret, + }) + token = data.get("accessToken") + if not token: + raise ChannelError(f"DingTalk auth error: {data}") + return token, int(data.get("expireIn", 7200)) + + async def _api_post(self, url, body, headers=None): + resp = await self._http_client.post(url, json=body, headers=headers) + return resp.json() + + async def _resolve_download_code(self, download_code: str) -> str | None: + """Exchange a DingTalk downloadCode for a real download URL.""" + try: + token = await self._ensure_token() + data = await self._api_post( + FILE_DOWNLOAD_URL, + {"downloadCode": download_code, "robotCode": self.config.client_id}, + headers={"x-acs-dingtalk-access-token": token}, + ) + url = data.get("downloadUrl") or "" + if url: + return url + logger.warning(f"DingTalk downloadCode resolve failed: {data}") + except Exception as e: + logger.warning(f"DingTalk downloadCode resolve error: {e}") + return None + + # ── WebSocketMixin ──────────────────────────────────────────── + + async def _get_ws_url(self) -> str: + resp = await self._http_client.post(GATEWAY_URL, json={ + "clientId": self.config.client_id, + "clientSecret": self.config.client_secret, + "subscriptions": [{"type": "CALLBACK", "topic": "/v1.0/im/bot/messages/get"}], + "ua": "dingtalk-sdk-python/v0.24.3-union", + }) + data = resp.json() + endpoint, ticket = data.get("endpoint"), data.get("ticket") + if not endpoint or not ticket: + raise ChannelError(f"DingTalk gateway failed: {data}") + return f"{endpoint}?ticket={quote_plus(ticket)}" + + async def _on_ws_message(self, data) -> None: + if not isinstance(data, dict): + return + headers = data.get("headers", {}) + msg_id = headers.get("messageId", "") + + # System ping + if data.get("type") == "SYSTEM" and headers.get("topic") == "ping": + await self._ws_send_json({"code": 200, "headers": headers, "message": "OK", "data": data.get("data", "")}) + return + + # ACK + await self._ws_send_json({"code": 200, "headers": {"contentType": "application/json", "messageId": msg_id}, "message": "OK", "data": "{}"}) + + if data.get("type") != "CALLBACK": + return + + payload = data.get("data", "{}") + payload = json.loads(payload) if isinstance(payload, str) else payload + text_obj = payload.get("text", {}) + content = (text_obj.get("content", "") if isinstance(text_obj, dict) else str(text_obj)).strip() + if not content: + raw_content = payload.get("content", "") + content = raw_content.strip() if isinstance(raw_content, str) else "" + + # Download attachments if present + annotations: list[str] = [] + media_paths: list[str] = [] + + # DingTalk file/image messages may put download info in + # payload["content"] (as a dict) instead of in a dedicated + # "fileContent"/"imageContent" key. + raw_content_obj = payload.get("content") + if isinstance(raw_content_obj, dict) and raw_content_obj not in [ + payload.get(k) for k in ("imageContent", "fileContent", "videoContent", "audioContent") + ]: + msg_type = payload.get("msgtype") or payload.get("msgType") or "" + media_label = msg_type or "file" + file_size = raw_content_obj.get("fileSize") or raw_content_obj.get("downloadSize") or 0 + file_name = raw_content_obj.get("fileName") or raw_content_obj.get("name") or f"dingtalk_{msg_type}" + download_code = raw_content_obj.get("downloadCode") or "" + download_url = raw_content_obj.get("downloadUrl") or "" + # downloadCode is NOT a URL — resolve it via DingTalk API first + if download_code and not download_code.startswith("http"): + resolved = await self._resolve_download_code(download_code) + if resolved: + download_url = resolved + elif download_code: + download_url = download_code + if download_url: + try: + dl_token = await self._ensure_token() + dl_headers = {"x-acs-dingtalk-access-token": dl_token} + except Exception: + dl_headers = None + local, ann = await self._download_attachment( + download_url, f"dingtalk_{file_name}", + headers=dl_headers, + file_size=int(file_size) if file_size else None, + ) + if local: + media_paths.append(local) + if ann: + ann = ann.replace("[attachment:", f"[{media_label}:") + annotations.append(ann) + elif file_name: + annotations.append(f"[{media_label}: {file_name}]") + + for att_key in ("imageContent", "fileContent", "videoContent", "audioContent"): + att = payload.get(att_key) + if att and isinstance(att, dict): + file_size = att.get("fileSize") or att.get("downloadSize") or 0 + file_name = att.get("fileName", att_key) + download_code = att.get("downloadCode") or "" + download_url = att.get("downloadUrl") or "" + # Resolve downloadCode via API if it's not a URL + if download_code and not download_code.startswith("http"): + resolved = await self._resolve_download_code(download_code) + if resolved: + download_url = resolved + elif download_code: + download_url = download_code + # DingTalk audioContent is voice messages + media_label = "voice" if att_key == "audioContent" else att_key + if download_url and (self.config.include_attachments if hasattr(self.config, 'include_attachments') else True): + # DingTalk download URLs require access token + try: + dl_token = await self._ensure_token() + dl_headers = {"x-acs-dingtalk-access-token": dl_token} + except Exception: + dl_headers = None + local, ann = await self._download_attachment( + download_url, f"dingtalk_{file_name}", + headers=dl_headers, + file_size=int(file_size) if file_size else None, + ) + if local: + media_paths.append(local) + if ann: + ann = ann.replace("[attachment:", f"[{media_label}:") + annotations.append(ann) + elif file_size: + too_large = self._check_attachment_size(int(file_size), file_name) + if too_large: + annotations.append(too_large) + else: + annotations.append(f"[{media_label}: {file_name}]") + + if not content and not media_paths and not annotations: + return + + sender_id = payload.get("senderStaffId") or payload.get("senderId", "") + is_group = payload.get("conversationType") == "2" + # For send API (oToMessages/batchSend), userIds needs staffId, not conversationId + chat_id = sender_id + create_time = payload.get("createAt") or payload.get("createTime", "") + + # Mention gating: DMs always pass; groups require @bot + was_mentioned = not is_group + if is_group: + # isInAtList is set by DingTalk when bot is @mentioned + if payload.get("isInAtList"): + was_mentioned = True + else: + # Fallback: check atUsers array + at_users = payload.get("atUsers") or [] + for u in at_users: + if u.get("dingtalkId") == self.config.client_id: + was_mentioned = True + break + + try: + ts = datetime.fromtimestamp(int(create_time) / 1000) if create_time else datetime.now() + except (ValueError, TypeError, OSError): + ts = datetime.now() + + await self._enqueue_raw(RawIncoming( + sender_id=sender_id, chat_id=chat_id, text=content, timestamp=ts, + message_id=msg_id, is_group=is_group, was_mentioned=was_mentioned, + media_files=media_paths, + content_annotations=annotations, + metadata={"chat_id": chat_id, "sender_nick": payload.get("senderNick", ""), "backend": "dingtalk"}, + )) + + # _send_typing_action: inherited no-op (DingTalk has no typing API) + # _format_chunk: inherited from base (UnifiedFormatter) + + # ── Send ────────────────────────────────────────────────────── + + async def _send_chunk(self, chat_id, formatted_text, raw_text, reply_to, metadata): + token = await self._ensure_token() + data = await self._api_post(SEND_URL, { + "robotCode": self.config.client_id, + "userIds": [chat_id], + "msgKey": "sampleMarkdown", + "msgParam": json.dumps({"text": raw_text, "title": "EvoScientist"}), + }, headers={"x-acs-dingtalk-access-token": token}) + return data + + # ── Media send ──────────────────────────────────────────────── + + _IMAGE_EXTS = {".jpg", ".jpeg", ".png", ".gif", ".bmp", ".webp"} + + async def _send_media_impl( + self, + recipient: str, + file_path: str, + caption: str = "", + metadata: dict | None = None, + ) -> bool: + """Send a media file through DingTalk. + + For images: uploads via /media/upload to get media_id, then sends + as sampleImageMsg. Non-image files are sent as markdown links + (DingTalk robot API does not support arbitrary file uploads). + """ + token = await self._ensure_token() + chat_id = self._resolve_media_chat_id(recipient, metadata) + headers = {"x-acs-dingtalk-access-token": token} + ext = Path(file_path).suffix.lower() + + if ext in self._IMAGE_EXTS: + # Try uploading image to get media_id for native image message + media_id = await self._upload_dingtalk_media(token, file_path, "image") + if media_id: + await self._api_post(MEDIA_SEND_URL, { + "robotCode": self.config.client_id, + "userIds": [chat_id], + "msgKey": "sampleImageMsg", + "msgParam": json.dumps({"photoURL": media_id}), + }, headers=headers) + else: + # Fallback to markdown with file path + await self._api_post(MEDIA_SEND_URL, { + "robotCode": self.config.client_id, + "userIds": [chat_id], + "msgKey": "sampleMarkdown", + "msgParam": json.dumps({ + "text": f"![image]({file_path})" + (f"\n{caption}" if caption else ""), + "title": caption or "Image", + }), + }, headers=headers) + else: + # Non-image: send as markdown with filename + name = Path(file_path).name + text = f"[文件] {name}" + (f"\n{caption}" if caption else "") + await self._api_post(MEDIA_SEND_URL, { + "robotCode": self.config.client_id, + "userIds": [chat_id], + "msgKey": "sampleMarkdown", + "msgParam": json.dumps({"text": text, "title": name}), + }, headers=headers) + + if caption and ext in self._IMAGE_EXTS: + # Send caption separately for image messages + await self._api_post(MEDIA_SEND_URL, { + "robotCode": self.config.client_id, + "userIds": [chat_id], + "msgKey": "sampleMarkdown", + "msgParam": json.dumps({"text": caption, "title": "Caption"}), + }, headers=headers) + return True + + async def _upload_dingtalk_media( + self, token: str, file_path: str, media_type: str = "image", + ) -> str | None: + """Upload a file to DingTalk media API and return the media_id.""" + try: + url = f"{MEDIA_UPLOAD_URL}?access_token={token}&type={media_type}" + with open(file_path, "rb") as f: + resp = await self._http_client.post( + url, files={"media": (Path(file_path).name, f)}, + ) + data = resp.json() + return data.get("media_id") + except Exception as e: + logger.warning(f"DingTalk media upload failed: {e}") + return None + + async def _cleanup(self) -> None: + if hasattr(self, "_ws_task") and self._ws_task: + self._ws_task.cancel() + try: + await self._ws_task + except (asyncio.CancelledError, Exception): + pass + self._ws_task = None + await self._stop_ws() + if self._http_client: + await self._http_client.aclose() + self._http_client = None + self._access_token = None + logger.info("DingTalk channel stopped") diff --git a/EvoScientist/channels/dingtalk/probe.py b/EvoScientist/channels/dingtalk/probe.py new file mode 100644 index 0000000..eb0eeb2 --- /dev/null +++ b/EvoScientist/channels/dingtalk/probe.py @@ -0,0 +1,33 @@ +"""DingTalk credential validation.""" + +import logging + +logger = logging.getLogger(__name__) + + +async def validate_dingtalk( + client_id: str, + client_secret: str, + proxy: str | None = None, +) -> tuple[bool, str]: + """Validate DingTalk credentials by fetching an access token.""" + if not client_id or not client_secret: + return False, "client_id and client_secret are required" + + try: + import httpx + except ImportError: + return False, "httpx not installed" + + url = "https://api.dingtalk.com/v1.0/oauth2/accessToken" + body = {"appKey": client_id, "appSecret": client_secret} + + try: + async with httpx.AsyncClient(proxy=proxy) as client: + resp = await client.post(url, json=body, timeout=10) + data = resp.json() + if data.get("accessToken"): + return True, "DingTalk credentials valid" + return False, f"Error: {data.get('message', data)}" + except Exception as e: + return False, f"Error: {e}" diff --git a/EvoScientist/channels/dingtalk/serve.py b/EvoScientist/channels/dingtalk/serve.py new file mode 100644 index 0000000..8fb4abb --- /dev/null +++ b/EvoScientist/channels/dingtalk/serve.py @@ -0,0 +1,92 @@ +"""DingTalk channel server. + +Standalone script to run the DingTalk channel with CLI options. + +Usage: + python -m EvoScientist.channels.dingtalk.serve --client-id ID --client-secret SECRET [OPTIONS] + +Examples: + # Basic usage + python -m EvoScientist.channels.dingtalk.serve --client-id ID --client-secret SECRET + + # With proxy and allowed senders + python -m EvoScientist.channels.dingtalk.serve --client-id ID --client-secret SECRET --proxy http://proxy:8080 --allow user123 + + # With agent and thinking + python -m EvoScientist.channels.dingtalk.serve --client-id ID --client-secret SECRET --agent --thinking +""" + +import argparse +import logging + +from .channel import DingTalkChannel, DingTalkConfig +from ..bus import MessageBus +from ..standalone import run_standalone + +logging.basicConfig( + level=logging.DEBUG, + format="%(asctime)s [%(levelname)s] %(name)s: %(message)s", + datefmt="%H:%M:%S", +) +logger = logging.getLogger(__name__) + + +def parse_args(): + """Parse command line arguments.""" + parser = argparse.ArgumentParser( + description="DingTalk channel server", + formatter_class=argparse.RawDescriptionHelpFormatter, + ) + parser.add_argument( + "--client-id", + required=True, + help="DingTalk app client ID", + ) + parser.add_argument( + "--client-secret", + required=True, + help="DingTalk app client secret", + ) + parser.add_argument( + "--allow", + action="append", + dest="allowed_senders", + help="Allowed sender (DingTalk user ID). Can be used multiple times.", + ) + parser.add_argument( + "--proxy", + help="HTTP proxy URL", + ) + parser.add_argument( + "--agent", + action="store_true", + help="Use EvoScientist agent as handler (default: echo)", + ) + parser.add_argument( + "--thinking", + action="store_true", + help="Send thinking content as intermediate messages (requires --agent)", + ) + return parser.parse_args() + + +def main(): + """Entry point.""" + args = parse_args() + + config = DingTalkConfig( + client_id=args.client_id, + client_secret=args.client_secret, + allowed_senders=set(args.allowed_senders) if args.allowed_senders else None, + proxy=args.proxy, + ) + + send_thinking = args.thinking and args.agent + bus = MessageBus() + channel = DingTalkChannel(config) + + run_standalone(channel, bus, use_agent=args.agent, send_thinking=send_thinking) + + +if __name__ == "__main__": + main() diff --git a/EvoScientist/channels/feishu/__init__.py b/EvoScientist/channels/feishu/__init__.py new file mode 100644 index 0000000..4cc68e5 --- /dev/null +++ b/EvoScientist/channels/feishu/__init__.py @@ -0,0 +1,21 @@ +from .channel import FeishuChannel, FeishuConfig +from ..channel_manager import register_channel, _parse_csv + +__all__ = ["FeishuChannel", "FeishuConfig"] + + +def create_from_config(config) -> FeishuChannel: + allowed = _parse_csv(config.feishu_allowed_senders) + return FeishuChannel(FeishuConfig( + app_id=config.feishu_app_id, + app_secret=config.feishu_app_secret, + verification_token=config.feishu_verification_token, + encrypt_key=config.feishu_encrypt_key, + webhook_port=config.feishu_webhook_port, + allowed_senders=allowed, + feishu_domain=config.feishu_domain, + proxy=getattr(config, 'feishu_proxy', '') or None, + )) + + +register_channel("feishu", create_from_config) diff --git a/EvoScientist/channels/feishu/channel.py b/EvoScientist/channels/feishu/channel.py new file mode 100644 index 0000000..8edebd1 --- /dev/null +++ b/EvoScientist/channels/feishu/channel.py @@ -0,0 +1,773 @@ +"""Feishu (飞书/Lark) channel implementation. + +Receives messages via an HTTP event subscription webhook (aiohttp), +sends replies via Feishu Open API REST endpoints. + +Feishu Open API docs: https://open.feishu.cn/document + +Authentication: + - App ID + App Secret → tenant_access_token (2-hour TTL, auto-refreshed) + +Event subscription: + - URL verification challenge on first request + - ``im.message.receive_v1`` events for incoming messages + +Send API: + - ``POST /open-apis/im/v1/messages?receive_id_type=chat_id`` +""" + +from __future__ import annotations + +import json +import logging +import re +from typing import Any, TYPE_CHECKING +from dataclasses import dataclass +from datetime import datetime +from pathlib import Path + +if TYPE_CHECKING: + from aiohttp import web + +from ..base import Channel, RawIncoming, ChannelError +from ..capabilities import FEISHU as FEISHU_CAPS +from ..mixins import WebhookMixin, TokenMixin +from ..config import BaseChannelConfig + +logger = logging.getLogger(__name__) + + +# ── Markdown → Feishu Post conversion ──────────────────────────── + + +def _parse_inline_text(text: str) -> list[dict]: + """Parse inline Markdown elements into Feishu post tag dicts. + + Handles: `code`, **bold**, ~~strikethrough~~, [link](url), _italic_. + """ + elements: list[dict] = [] + # Pattern order matters: code first (protect content), then bold, strikethrough, link, italic + pattern = re.compile( + r"`([^`]+)`" # inline code + r"|\*\*(.+?)\*\*" # bold + r"|~~(.+?)~~" # strikethrough + r"|\[([^\]]+)\]\(([^)]+)\)" # link + r"|_(.+?)_" # italic + ) + pos = 0 + for m in pattern.finditer(text): + # Plain text before this match + if m.start() > pos: + elements.append({"tag": "text", "text": text[pos:m.start()]}) + + if m.group(1) is not None: + # inline code → code_block would be block-level; use text with style + elements.append({ + "tag": "text", "text": m.group(1), + "style": ["code_block"], + }) + elif m.group(2) is not None: + elements.append({ + "tag": "text", "text": m.group(2), + "style": ["bold"], + }) + elif m.group(3) is not None: + elements.append({ + "tag": "text", "text": m.group(3), + "style": ["strikethrough"], + }) + elif m.group(4) is not None: + elements.append({ + "tag": "a", "text": m.group(4), "href": m.group(5), + }) + elif m.group(6) is not None: + elements.append({ + "tag": "text", "text": m.group(6), + "style": ["italic"], + }) + pos = m.end() + + # Remaining plain text + if pos < len(text): + elements.append({"tag": "text", "text": text[pos:]}) + return elements + + +def _parse_inline_elements(line: str) -> list[dict]: + """Parse a single Markdown line into a list of Feishu post elements. + + Handles headings (→ bold), blockquotes (→ italic with prefix), + list items (→ bullet prefix), and plain lines. + """ + # Heading: # Title → bold text + heading_match = re.match(r"^(#{1,6})\s+(.+)$", line) + if heading_match: + return [{"tag": "text", "text": heading_match.group(2), "style": ["bold"]}] + + # Blockquote: > text → italic with "▎" prefix + quote_match = re.match(r"^>\s*(.*)$", line) + if quote_match: + inner = quote_match.group(1) + elements = [{"tag": "text", "text": "▎", "style": ["italic"]}] + elements.extend(_parse_inline_text(inner)) + return elements + + # Unordered list: - item or * item → "• " prefix + list_match = re.match(r"^[\-\*]\s+(.+)$", line) + if list_match: + elements = [{"tag": "text", "text": "• "}] + elements.extend(_parse_inline_text(list_match.group(1))) + return elements + + # Ordered list: 1. item → keep number prefix + ol_match = re.match(r"^(\d+)\.\s+(.+)$", line) + if ol_match: + elements = [{"tag": "text", "text": f"{ol_match.group(1)}. "}] + elements.extend(_parse_inline_text(ol_match.group(2))) + return elements + + # Plain line + return _parse_inline_text(line) + + +def _markdown_to_feishu_post(text: str) -> dict | None: + """Convert Markdown text to Feishu post (rich text) JSON structure. + + Returns a dict like {"zh_cn": {"content": [[...]]}} suitable for + Feishu msg_type="post", or None if the text is empty. + """ + if not text or not text.strip(): + return None + + paragraphs: list[list[dict]] = [] + current_paragraph: list[dict] = [] + in_code_block = False + code_lines: list[str] = [] + code_lang = "" + + for line in text.split("\n"): + # Code block fences + if line.startswith("```"): + if not in_code_block: + # Flush any pending paragraph + if current_paragraph: + paragraphs.append(current_paragraph) + current_paragraph = [] + in_code_block = True + code_lang = line[3:].strip() + code_lines = [] + else: + # End of code block + code_text = "\n".join(code_lines) + paragraphs.append([{ + "tag": "code_block", + "language": code_lang or "plain", + "text": code_text, + }]) + in_code_block = False + code_lines = [] + code_lang = "" + continue + + if in_code_block: + code_lines.append(line) + continue + + # Empty line → new paragraph + if not line.strip(): + if current_paragraph: + paragraphs.append(current_paragraph) + current_paragraph = [] + continue + + # Non-empty line + elements = _parse_inline_elements(line) + if elements: + # Each visual line becomes its own paragraph in Feishu post + if current_paragraph: + paragraphs.append(current_paragraph) + current_paragraph = elements + + # Flush remaining + if in_code_block and code_lines: + code_text = "\n".join(code_lines) + paragraphs.append([{ + "tag": "code_block", + "language": code_lang or "plain", + "text": code_text, + }]) + elif current_paragraph: + paragraphs.append(current_paragraph) + + if not paragraphs: + return None + + return {"zh_cn": {"content": paragraphs}} + + +@dataclass +class FeishuConfig(BaseChannelConfig): + app_id: str = "" + app_secret: str = "" + verification_token: str = "" + encrypt_key: str = "" + webhook_port: int = 9000 + text_chunk_limit: int = 4096 + feishu_domain: str = "https://open.feishu.cn" + + +class FeishuChannel(Channel, WebhookMixin, TokenMixin): + capabilities = FEISHU_CAPS + """Feishu channel using Open API + event subscription webhook.""" + + name = "feishu" + _ready_attrs = ("_http_client", "_access_token") + _rate_limit_patterns = ("99991400", "rate limit", "频率限制") + _rate_limit_delay = 2.0 + + def __init__(self, config: FeishuConfig): + super().__init__(config) + self._mention_names: list[str] = [] # bot mention keys from events + + # ── WebhookMixin overrides ──────────────────────────────────── + + def _get_webhook_port(self) -> int: + return self.config.webhook_port + + def _webhook_routes(self) -> list[tuple[str, Any]]: + return [("POST", "/webhook/event", self._handle_event)] + + # ── TokenMixin overrides ────────────────────────────────────── + + async def _fetch_token(self) -> tuple[str, int]: + """Fetch Feishu tenant_access_token.""" + url = f"{self.config.feishu_domain}/open-apis/auth/v3/tenant_access_token/internal" + body = { + "app_id": self.config.app_id, + "app_secret": self.config.app_secret, + } + try: + resp = await self._http_client.post(url, json=body) + data = resp.json() + except Exception as e: + raise ChannelError(f"Failed to get Feishu access token: {e}") + + if data.get("code") != 0: + raise ChannelError( + f"Feishu auth error: {data.get('msg', 'unknown')}" + ) + return data["tenant_access_token"], data.get("expire", 7200) + + # ── Lifecycle ───────────────────────────────────────────────── + + async def start(self) -> None: + try: + from aiohttp import web # noqa: F401 + import httpx # noqa: F401 + except ImportError: + raise ChannelError( + "aiohttp or httpx not installed. " + "Install with: pip install aiohttp httpx" + ) + + if not self.config.app_id: + raise ChannelError("Feishu app_id is required") + if not self.config.app_secret: + raise ChannelError("Feishu app_secret is required") + + # Start webhook server (sets up self._http_client) + await self._start_webhook_server() + + # Verify credentials by fetching initial token + await self._refresh_token() + + self._running = True + logger.info( + f"Feishu channel started " + f"(webhook on port {self.config.webhook_port})" + ) + + async def _cleanup(self) -> None: + await self._stop_webhook_server() + self._access_token = None + logger.info("Feishu channel stopped") + + # ── Token helpers (adapt old API to mixin) ──────────────────── + + async def _ensure_token(self) -> str: + """Return a valid access token, refreshing if needed.""" + return await TokenMixin._ensure_token(self) + + # ── Send (template method overrides) ────────────────────────── + + async def _feishu_send(self, url: str, body: dict, headers: dict) -> bool: + """POST to Feishu API and return True if code==0.""" + try: + resp = await self._http_client.post(url, json=body, headers=headers) + return resp.json().get("code") == 0 + except Exception as e: + logger.warning(f"Feishu send error: {e}") + return False + + async def _send_chunk( + self, chat_id, formatted_text, raw_text, reply_to, metadata, + ): + token = await self._ensure_token() + headers = {"Authorization": f"Bearer {token}"} + post_content = _markdown_to_feishu_post(raw_text) + + # If reply_to is set, try the reply API first + if reply_to: + reply_url = ( + f"{self.config.feishu_domain}" + f"/open-apis/im/v1/messages/{reply_to}/reply" + ) + if post_content is not None: + body = {"msg_type": "post", "content": json.dumps(post_content)} + else: + body = {"msg_type": "text", "content": json.dumps({"text": formatted_text})} + if await self._feishu_send(reply_url, body, headers): + return + + # Normal send (non-reply or reply fallback) + url = ( + f"{self.config.feishu_domain}" + f"/open-apis/im/v1/messages?receive_id_type=chat_id" + ) + + # Try post format first + if post_content is not None: + body = { + "receive_id": chat_id, + "msg_type": "post", + "content": json.dumps(post_content), + } + if await self._feishu_send(url, body, headers): + return + + # Fallback: plain text + body = { + "receive_id": chat_id, + "msg_type": "text", + "content": json.dumps({"text": formatted_text}), + } + if not await self._feishu_send(url, body, headers): + raise RuntimeError("Feishu send failed") + + # ── Media helpers ────────────────────────────────────────────── + + _IMAGE_EXTENSIONS = {".jpg", ".jpeg", ".png", ".gif", ".bmp", ".webp"} + + async def _download_media( + self, message_id: str, file_key: str, msg_type: str, + ) -> str | None: + """Download an image or file attachment from Feishu. + + Returns the local file path on success, or None on failure. + """ + token = await self._ensure_token() + resource_type = "image" if msg_type == "image" else "file" + url = ( + f"{self.config.feishu_domain}" + f"/open-apis/im/v1/messages/{message_id}" + f"/resources/{file_key}?type={resource_type}" + ) + headers = {"Authorization": f"Bearer {token}"} + try: + resp = await self._http_client.get(url, headers=headers, timeout=30) + if resp.status_code != 200: + logger.warning( + f"Feishu media download failed: HTTP {resp.status_code}" + ) + return None + + # Check attachment size before writing to disk + cl = resp.headers.get("content-length") + if cl: + try: + too_large = self._check_attachment_size(int(cl), file_key) + if too_large: + logger.warning(too_large) + return None + except (ValueError, TypeError): + pass + from ..base import MAX_ATTACHMENT_BYTES + if len(resp.content) > MAX_ATTACHMENT_BYTES: + logger.warning( + f"Feishu media too large: {len(resp.content)} bytes" + ) + return None + + # Determine extension from Content-Type or default + content_type = resp.headers.get("content-type", "") + ext_map = { + "image/jpeg": ".jpg", + "image/png": ".png", + "image/gif": ".gif", + "image/webp": ".webp", + "image/bmp": ".bmp", + } + ext = ext_map.get(content_type, ".bin") + local_path = self._media_path(f"feishu_{message_id}_{file_key}{ext}") + local_path.write_bytes(resp.content) + return str(local_path) + except Exception as e: + logger.warning(f"Failed to download Feishu media: {e}") + return None + + async def _upload_feishu_resource( + self, url: str, headers: dict, file_path: str, + field_name: str, extra_data: dict, + ) -> dict | None: + """Upload a file to Feishu API. Returns response data or None on failure.""" + with open(file_path, "rb") as f: + resp = await self._http_client.post( + url, headers=headers, data=extra_data, + files={field_name: (Path(file_path).name, f)}, + ) + data = resp.json() + if data.get("code") != 0: + logger.error(f"Feishu upload failed: {data.get('msg')}") + return None + return data["data"] + + async def _send_media_impl( + self, + recipient: str, + file_path: str, + caption: str = "", + metadata: dict | None = None, + ) -> bool: + """Send a media file through Feishu.""" + token = await self._ensure_token() + headers = {"Authorization": f"Bearer {token}"} + chat_id = self._resolve_media_chat_id(recipient, metadata) + + path = Path(file_path) + ext = path.suffix.lower() + is_image = ext in self._IMAGE_EXTENSIONS + + send_url = ( + f"{self.config.feishu_domain}" + f"/open-apis/im/v1/messages?receive_id_type=chat_id" + ) + + if is_image: + upload_url = f"{self.config.feishu_domain}/open-apis/im/v1/images" + data = await self._upload_feishu_resource( + upload_url, headers, file_path, "image", {"image_type": "message"}, + ) + if not data: + return False + body = { + "receive_id": chat_id, + "msg_type": "image", + "content": json.dumps({"image_key": data["image_key"]}), + } + else: + upload_url = f"{self.config.feishu_domain}/open-apis/im/v1/files" + data = await self._upload_feishu_resource( + upload_url, headers, file_path, "file", + {"file_type": "stream", "file_name": path.name}, + ) + if not data: + return False + body = { + "receive_id": chat_id, + "msg_type": "file", + "content": json.dumps({"file_key": data["file_key"]}), + } + + if not await self._feishu_send(send_url, body, headers): + return False + + # Send caption as a separate text message if provided + if caption: + cap_body = { + "receive_id": chat_id, + "msg_type": "text", + "content": json.dumps({"text": caption}), + } + await self._feishu_send(send_url, cap_body, headers) + + return True + + # ── ACK reaction ─────────────────────────────────────────────── + + async def _send_ack_reaction(self, chat_id: str, message_id: str, emoji: str = "THUMBSUP") -> None: + """Send an acknowledgment reaction via Feishu Open API.""" + try: + token = await self._ensure_token() + url = f"{self.config.feishu_domain}/open-apis/im/v1/messages/{message_id}/reactions" + await self._http_client.post( + url, + json={"reaction_type": {"emoji_type": emoji}}, + headers={"Authorization": f"Bearer {token}"}, + ) + except Exception as e: + logger.debug(f"Feishu ack reaction failed: {e}") + + async def _remove_ack_reaction(self, chat_id: str, message_id: str, emoji: str = "THUMBSUP") -> None: + """Remove ACK reaction via Feishu Open API. + + Feishu's DELETE /reactions endpoint requires the reaction_id, which + we don't track. No-op for now. + """ + pass + + # ── Mention stripping ───────────────────────────────────────── + + def _strip_mention(self, text: str) -> str: + """Strip bot @mention placeholders from Feishu text. + + In Feishu v2 events the text contains placeholders like ``@_user_1`` + for each mention. ``_mention_names`` caches the placeholder keys that + belong to the bot (identified during ``_on_message``). + """ + result = text + for key in self._mention_names: + result = result.replace(key, "") + # Clean up extra whitespace left behind + return re.sub(r" +", " ", result).strip() + + # ── Webhook event handler ───────────────────────────────────── + + async def _handle_event(self, request) -> "web.Response": + """Handle POST /webhook/event from Feishu.""" + from aiohttp import web + + try: + body = await request.json() + except Exception: + return web.Response(status=400) + + # ── URL verification challenge ── + if body.get("type") == "url_verification": + challenge = body.get("challenge", "") + return web.json_response({"challenge": challenge}) + + # ── v2 event schema ── + schema = body.get("schema") + if schema == "2.0": + header = body.get("header", {}) + + # Verify token if configured + if self.config.verification_token: + token = header.get("token", "") + if token != self.config.verification_token: + logger.warning("Feishu event token mismatch") + return web.Response(status=403) + + event_type = header.get("event_type", "") + if event_type == "im.message.receive_v1": + await self._on_message(body.get("event", {})) + + # ── v1 event schema (legacy) ── + elif "event" in body: + if self.config.verification_token: + token = body.get("token", "") + if token != self.config.verification_token: + logger.warning("Feishu event token mismatch (v1)") + return web.Response(status=403) + + event = body["event"] + msg_type = event.get("type", "") + if msg_type == "message": + await self._on_message_v1(event) + + return web.Response(status=200) + + async def _on_message(self, event: dict) -> None: + """Handle im.message.receive_v1 event (v2 schema).""" + sender_info = event.get("sender", {}) + sender_id_info = sender_info.get("sender_id", {}) + sender_id = ( + sender_id_info.get("open_id") + or sender_id_info.get("user_id") + or "" + ) + sender_type = sender_info.get("sender_type", "") + + # Skip bot's own messages + if sender_type == "app": + return + + message = event.get("message", {}) + chat_id = message.get("chat_id", "") + msg_type = message.get("message_type", "") + message_id = message.get("message_id", "") + + # In group chats, detect mention status for centralized gating + chat_type = message.get("chat_type", "") + is_group = chat_type == "group" + was_mentioned = True + if is_group: + mentions = message.get("mentions", []) + was_mentioned = bool(mentions) + # Cache bot mention keys — bot mentions have empty user IDs + bot_keys = [] + for m in mentions: + m_id = m.get("id", {}) + # Bot/app mentions have no open_id / user_id + if not m_id.get("open_id") and not m_id.get("user_id"): + key = m.get("key", "") + if key: + bot_keys.append(key) + if bot_keys: + self._mention_names = bot_keys + + # Parse content JSON + content_str = message.get("content", "{}") + try: + content_data = json.loads(content_str) + except json.JSONDecodeError: + content_data = {} + + text = "" + annotations: list[str] = [] + media_paths: list[str] = [] + + if msg_type == "text": + text = content_data.get("text", "") + elif msg_type == "post": + text = self._extract_post_text(content_data) + elif msg_type == "image" and self.config.include_attachments: + image_key = content_data.get("image_key", "") + if image_key: + local = await self._download_media(message_id, image_key, "image") + if local: + media_paths.append(local) + annotations.append(f"[attachment: {local}]") + else: + annotations.append("[image message - download failed]") + else: + annotations.append("[image message]") + elif msg_type == "file" and self.config.include_attachments: + file_key = content_data.get("file_key", "") + file_name = content_data.get("file_name", "unknown") + if file_key: + local = await self._download_media(message_id, file_key, "file") + if local: + media_paths.append(local) + annotations.append(f"[attachment: {local}]") + else: + annotations.append(f"[file: {file_name} - download failed]") + else: + annotations.append(f"[file message: {file_name}]") + elif msg_type in ("audio", "media") and self.config.include_attachments: + # Feishu audio messages are voice recordings + media_label = "voice" if msg_type == "audio" else msg_type + file_key = content_data.get("file_key", "") + if file_key: + local = await self._download_media(message_id, file_key, "file") + if local: + media_paths.append(local) + annotations.append(f"[{media_label}: {local}]") + else: + annotations.append(f"[{media_label} message - download failed]") + else: + annotations.append(f"[{media_label} message]") + elif msg_type == "sticker": + sticker_key = content_data.get("file_key", "") + if sticker_key and self.config.include_attachments: + local = await self._download_media(message_id, sticker_key, "image") + if local: + media_paths.append(local) + annotations.append(f"[sticker: {local}]") + else: + annotations.append("[sticker message]") + else: + annotations.append("[sticker message]") + else: + text = f"[{msg_type} message]" + + if not text and not media_paths and not annotations: + return + + # Parse timestamp (milliseconds) + create_time = message.get("create_time", "") + try: + timestamp = datetime.fromtimestamp( + int(create_time) / 1000 + ) if create_time else datetime.now() + except (ValueError, TypeError, OSError): + timestamp = datetime.now() + + await self._enqueue_raw(RawIncoming( + sender_id=sender_id, + chat_id=chat_id, + text=text, + media_files=media_paths, + content_annotations=annotations, + timestamp=timestamp, + message_id=message_id, + metadata={ + "chat_id": chat_id, + "chat_type": message.get("chat_type", ""), + }, + is_group=is_group, + was_mentioned=was_mentioned, + )) + + async def _on_message_v1(self, event: dict) -> None: + """Handle v1 schema message event (legacy).""" + sender_id = event.get("open_id", "") + if not sender_id: + return + + # Detect group and mention status for centralized gating + chat_type = event.get("chat_type", "") + is_group = chat_type == "group" + was_mentioned = True + if is_group: + text_without_at = event.get("text_without_at_bot", "") + was_mentioned = bool(text_without_at) + + text = event.get("text_without_at_bot", "") or event.get("text", "") + if not text: + return + + chat_id = event.get("open_chat_id", "") + message_id = event.get("open_message_id", "") + + await self._enqueue_raw(RawIncoming( + sender_id=sender_id, + chat_id=chat_id, + text=text, + timestamp=datetime.now(), + message_id=message_id, + metadata={ + "chat_id": chat_id, + "chat_type": event.get("chat_type", ""), + }, + is_group=is_group, + was_mentioned=was_mentioned, + )) + + @staticmethod + def _extract_post_text(content: dict) -> str: + """Extract plain text from Feishu post (rich text) content.""" + parts: list[str] = [] + # Post content has locale keys like "zh_cn", "en_us" + for locale_key in ("zh_cn", "en_us", "ja_jp"): + locale_content = content.get(locale_key) + if locale_content: + title = locale_content.get("title", "") + if title: + parts.append(title) + for paragraph in locale_content.get("content", []): + line_parts: list[str] = [] + for element in paragraph: + tag = element.get("tag", "") + if tag == "text": + line_parts.append(element.get("text", "")) + elif tag == "a": + line_parts.append(element.get("text", "")) + elif tag == "at": + # Skip @mentions of the bot + pass + line = "".join(line_parts).strip() + if line: + parts.append(line) + break # Use first available locale + return "\n".join(parts) diff --git a/EvoScientist/channels/feishu/probe.py b/EvoScientist/channels/feishu/probe.py new file mode 100644 index 0000000..46d962e --- /dev/null +++ b/EvoScientist/channels/feishu/probe.py @@ -0,0 +1,39 @@ +"""Feishu (飞书/Lark) app credential validation.""" + +import logging + +logger = logging.getLogger(__name__) + + +async def validate_feishu_credentials( + app_id: str, + app_secret: str, + domain: str = "https://open.feishu.cn", +) -> tuple[bool, str]: + """Validate Feishu app credentials by requesting a tenant_access_token. + + Returns: + Tuple of (is_valid, message). + """ + if not app_id: + return False, "No app_id provided" + if not app_secret: + return False, "No app_secret provided" + + try: + import httpx + except ImportError: + return False, "httpx not installed" + + url = f"{domain}/open-apis/auth/v3/tenant_access_token/internal" + body = {"app_id": app_id, "app_secret": app_secret} + try: + async with httpx.AsyncClient() as client: + resp = await client.post(url, json=body, timeout=10) + data = resp.json() + if data.get("code") == 0: + return True, f"App: {app_id}" + msg = data.get("msg", "unknown error") + return False, f"Auth failed: {msg}" + except Exception as e: + return False, f"Error: {e}" diff --git a/EvoScientist/channels/feishu/serve.py b/EvoScientist/channels/feishu/serve.py new file mode 100644 index 0000000..c9d0709 --- /dev/null +++ b/EvoScientist/channels/feishu/serve.py @@ -0,0 +1,113 @@ +"""Feishu (飞书/Lark) channel server. + +Standalone script to run the Feishu channel with CLI options. + +Usage: + python -m EvoScientist.channels.feishu.serve --app-id ID --app-secret SECRET [OPTIONS] + +Examples: + # Basic setup + python -m EvoScientist.channels.feishu.serve --app-id ID --app-secret SECRET + + # With verification token and custom port + python -m EvoScientist.channels.feishu.serve --app-id ID --app-secret SECRET \\ + --verification-token TOKEN --webhook-port 9000 + + # With agent and thinking + python -m EvoScientist.channels.feishu.serve --app-id ID --app-secret SECRET --agent --thinking +""" + +import argparse +import logging + +from .channel import FeishuChannel, FeishuConfig +from ..bus import MessageBus +from ..standalone import run_standalone + +logging.basicConfig( + level=logging.DEBUG, + format="%(asctime)s [%(levelname)s] %(name)s: %(message)s", + datefmt="%H:%M:%S", +) +logger = logging.getLogger(__name__) + + +def parse_args(): + """Parse command line arguments.""" + parser = argparse.ArgumentParser( + description="Feishu (飞书/Lark) channel server", + formatter_class=argparse.RawDescriptionHelpFormatter, + ) + parser.add_argument( + "--app-id", + required=True, + help="Feishu App ID", + ) + parser.add_argument( + "--app-secret", + required=True, + help="Feishu App Secret", + ) + parser.add_argument( + "--verification-token", + default="", + help="Feishu event verification token", + ) + parser.add_argument( + "--encrypt-key", + default="", + help="Feishu event encrypt key", + ) + parser.add_argument( + "--webhook-port", + type=int, + default=9000, + help="Port for webhook HTTP server (default: 9000)", + ) + parser.add_argument( + "--domain", + default="https://open.feishu.cn", + help="Feishu API domain (use https://open.larksuite.com for Lark)", + ) + parser.add_argument( + "--allow", + action="append", + dest="allowed_senders", + help="Allowed sender (Feishu open_id). Can be used multiple times.", + ) + parser.add_argument( + "--agent", + action="store_true", + help="Use EvoScientist agent as handler (default: echo)", + ) + parser.add_argument( + "--thinking", + action="store_true", + help="Send thinking content as intermediate messages (requires --agent)", + ) + return parser.parse_args() + + +def main(): + """Entry point.""" + args = parse_args() + + config = FeishuConfig( + app_id=args.app_id, + app_secret=args.app_secret, + verification_token=args.verification_token, + encrypt_key=args.encrypt_key, + webhook_port=args.webhook_port, + feishu_domain=args.domain, + allowed_senders=set(args.allowed_senders) if args.allowed_senders else None, + ) + + send_thinking = args.thinking and args.agent + bus = MessageBus() + channel = FeishuChannel(config) + + run_standalone(channel, bus, use_agent=args.agent, send_thinking=send_thinking) + + +if __name__ == "__main__": + main() diff --git a/EvoScientist/config/onboard.py b/EvoScientist/config/onboard.py index b8b680c..6c9584e 100644 --- a/EvoScientist/config/onboard.py +++ b/EvoScientist/config/onboard.py @@ -1319,6 +1319,8 @@ def _step_channels(config: EvoScientistConfig) -> dict[str, object]: ("telegram", "Telegram", [("telegram_bot_token", "Bot token (from @BotFather)")], "telegram", "telegram"), ("discord", "Discord", [("discord_bot_token", "Bot token")], "discord", "discord"), ("slack", "Slack", [("slack_bot_token", "Bot token (xoxb-...)"), ("slack_app_token", "App token for Socket Mode (xapp-...)")], "slack_sdk", "slack"), + ("feishu", "Feishu", [("feishu_app_id", "App ID"), ("feishu_app_secret", "App Secret"), ("feishu_verification_token", "Verification Token"), ("feishu_encrypt_key", "Encrypt Key")], "aiohttp", "feishu"), + ("dingtalk", "DingTalk", [("dingtalk_client_id", "Client ID (AppKey)"), ("dingtalk_client_secret", "Client Secret (AppSecret)")], "aiohttp", "dingtalk"), ("wechat", "WeChat", [("wechat_wecom_corp_id", "WeCom Corp ID"), ("wechat_wecom_agent_id", "WeCom Agent ID"), ("wechat_wecom_secret", "WeCom Secret")], "aiohttp", "wechat"), ("imessage", "iMessage", [], None, None), # handled via _setup_imessage() ] @@ -1512,6 +1514,20 @@ def _probe_channel( _val("wechat_wecom_secret"), _val("wechat_proxy") or None, ) + elif ch_name == "feishu": + from ..channels.feishu.probe import validate_feishu_credentials + return await validate_feishu_credentials( + _val("feishu_app_id"), + _val("feishu_app_secret"), + _val("feishu_domain", "https://open.feishu.cn"), + ) + elif ch_name == "dingtalk": + from ..channels.dingtalk.probe import validate_dingtalk + return await validate_dingtalk( + _val("dingtalk_client_id"), + _val("dingtalk_client_secret"), + _val("dingtalk_proxy") or None, + ) else: return True, "No probe available" diff --git a/EvoScientist/config/settings.py b/EvoScientist/config/settings.py index 4a40e56..d94006d 100644 --- a/EvoScientist/config/settings.py +++ b/EvoScientist/config/settings.py @@ -85,7 +85,7 @@ class EvoScientistConfig: show_thinking: bool = True # Channel Settings - channel_enabled: str = "" # "imessage" | "telegram" | "discord" | "slack" | "wechat" | "" (comma-separated for multiple) + channel_enabled: str = "" # "imessage" | "telegram" | "discord" | "slack" | "wechat" | "dingtalk" | "feishu" | "" (comma-separated for multiple) channel_send_thinking: bool = True # forward thinking to any channel require_mention: str = "group" # "always" | "group" | "off" text_chunk_limit: int = 0 # 0 = use capability default diff --git a/tests/test_dingtalk_channel.py b/tests/test_dingtalk_channel.py new file mode 100644 index 0000000..91870b6 --- /dev/null +++ b/tests/test_dingtalk_channel.py @@ -0,0 +1,122 @@ +"""Tests for DingTalk channel implementation.""" + +import asyncio + +import pytest + +from EvoScientist.channels.dingtalk.channel import DingTalkChannel, DingTalkConfig +from EvoScientist.channels.base import ChannelError + + +def _run(coro): + """Run an async coroutine safely, creating a fresh event loop.""" + loop = asyncio.new_event_loop() + try: + return loop.run_until_complete(coro) + finally: + loop.close() + + +class TestDingTalkConfig: + def test_default_values(self): + config = DingTalkConfig() + assert config.client_id == "" + assert config.client_secret == "" + assert config.allowed_senders is None + assert config.text_chunk_limit == 4096 + + def test_custom_values(self): + config = DingTalkConfig( + client_id="test-id", + client_secret="test-secret", + allowed_senders={"user1"}, + text_chunk_limit=2000, + proxy="http://proxy:8080", + ) + assert config.client_id == "test-id" + assert config.client_secret == "test-secret" + assert config.allowed_senders == {"user1"} + assert config.text_chunk_limit == 2000 + assert config.proxy == "http://proxy:8080" + + +class TestDingTalkChannel: + def test_init(self): + config = DingTalkConfig(client_id="test-id", client_secret="test-secret") + channel = DingTalkChannel(config) + assert channel.config is config + assert channel._running is False + assert channel.name == "dingtalk" + + def test_start_raises_without_credentials(self): + config = DingTalkConfig(client_id="", client_secret="") + channel = DingTalkChannel(config) + with pytest.raises(ChannelError, match="client_id and client_secret"): + _run(channel.start()) + + def test_start_raises_without_client_id(self): + config = DingTalkConfig(client_id="", client_secret="secret") + channel = DingTalkChannel(config) + with pytest.raises(ChannelError, match="client_id and client_secret"): + _run(channel.start()) + + def test_start_raises_without_client_secret(self): + config = DingTalkConfig(client_id="id", client_secret="") + channel = DingTalkChannel(config) + with pytest.raises(ChannelError, match="client_id and client_secret"): + _run(channel.start()) + + def test_stop_when_not_running(self): + config = DingTalkConfig(client_id="test-id", client_secret="test-secret") + channel = DingTalkChannel(config) + _run(channel.stop()) + + def test_send_returns_false_without_client(self): + from EvoScientist.channels.base import OutboundMessage + + config = DingTalkConfig(client_id="test-id", client_secret="test-secret") + channel = DingTalkChannel(config) + msg = OutboundMessage( + channel="dingtalk", + chat_id="user123", + content="hello", + metadata={"chat_id": "user123"}, + ) + result = _run(channel.send(msg)) + assert result is False + + def test_capabilities(self): + from EvoScientist.channels.capabilities import DINGTALK + config = DingTalkConfig() + channel = DingTalkChannel(config) + assert channel.capabilities is DINGTALK + assert channel.capabilities.format_type == "markdown" + assert channel.capabilities.groups is True + assert channel.capabilities.mentions is True + assert channel.capabilities.media_send is True + assert channel.capabilities.media_receive is True + + +class TestDingTalkChannelRegistration: + def test_dingtalk_registered(self): + from EvoScientist.channels.channel_manager import available_channels + channels = available_channels() + assert "dingtalk" in channels + + +class TestDingTalkProbe: + def test_missing_credentials(self): + from EvoScientist.channels.dingtalk.probe import validate_dingtalk + ok, msg = _run(validate_dingtalk("", "")) + assert ok is False + assert "required" in msg + + def test_missing_client_id(self): + from EvoScientist.channels.dingtalk.probe import validate_dingtalk + ok, msg = _run(validate_dingtalk("", "secret")) + assert ok is False + + def test_missing_client_secret(self): + from EvoScientist.channels.dingtalk.probe import validate_dingtalk + ok, msg = _run(validate_dingtalk("id", "")) + assert ok is False diff --git a/tests/test_feishu_channel.py b/tests/test_feishu_channel.py new file mode 100644 index 0000000..1b779ed --- /dev/null +++ b/tests/test_feishu_channel.py @@ -0,0 +1,204 @@ +"""Tests for Feishu channel implementation.""" + +import asyncio +import json + +import pytest + +from EvoScientist.channels.feishu.channel import ( + FeishuChannel, + FeishuConfig, + _markdown_to_feishu_post, + _parse_inline_text, +) +from EvoScientist.channels.base import ChannelError + + +def _run(coro): + """Run an async coroutine safely, creating a fresh event loop.""" + loop = asyncio.new_event_loop() + try: + return loop.run_until_complete(coro) + finally: + loop.close() + + +class TestFeishuConfig: + def test_default_values(self): + config = FeishuConfig() + assert config.app_id == "" + assert config.app_secret == "" + assert config.verification_token == "" + assert config.encrypt_key == "" + assert config.webhook_port == 9000 + assert config.text_chunk_limit == 4096 + assert config.feishu_domain == "https://open.feishu.cn" + assert config.allowed_senders is None + + def test_custom_values(self): + config = FeishuConfig( + app_id="test-id", + app_secret="test-secret", + verification_token="token123", + encrypt_key="key123", + webhook_port=8080, + allowed_senders={"user1"}, + feishu_domain="https://open.larksuite.com", + ) + assert config.app_id == "test-id" + assert config.app_secret == "test-secret" + assert config.verification_token == "token123" + assert config.encrypt_key == "key123" + assert config.webhook_port == 8080 + assert config.allowed_senders == {"user1"} + assert config.feishu_domain == "https://open.larksuite.com" + + +class TestFeishuChannel: + def test_init(self): + config = FeishuConfig(app_id="test-id", app_secret="test-secret") + channel = FeishuChannel(config) + assert channel.config is config + assert channel._running is False + assert channel.name == "feishu" + + def test_start_raises_without_app_id(self): + config = FeishuConfig(app_id="", app_secret="test-secret") + channel = FeishuChannel(config) + with pytest.raises(ChannelError, match="app_id"): + _run(channel.start()) + + def test_start_raises_without_app_secret(self): + config = FeishuConfig(app_id="test-id", app_secret="") + channel = FeishuChannel(config) + with pytest.raises(ChannelError, match="app_secret"): + _run(channel.start()) + + def test_stop_when_not_running(self): + config = FeishuConfig(app_id="test-id", app_secret="test-secret") + channel = FeishuChannel(config) + _run(channel.stop()) + + def test_send_returns_false_without_client(self): + from EvoScientist.channels.base import OutboundMessage + + config = FeishuConfig(app_id="test-id", app_secret="test-secret") + channel = FeishuChannel(config) + msg = OutboundMessage( + channel="feishu", + chat_id="oc_test", + content="hello", + metadata={"chat_id": "oc_test"}, + ) + result = _run(channel.send(msg)) + assert result is False + + def test_capabilities(self): + from EvoScientist.channels.capabilities import FEISHU + config = FeishuConfig() + channel = FeishuChannel(config) + assert channel.capabilities is FEISHU + assert channel.capabilities.format_type == "markdown" + assert channel.capabilities.groups is True + assert channel.capabilities.mentions is True + assert channel.capabilities.media_send is True + assert channel.capabilities.media_receive is True + assert channel.capabilities.reactions is True + assert channel.capabilities.voice is True + assert channel.capabilities.stickers is True + + def test_extract_post_text(self): + content = { + "zh_cn": { + "title": "Test Title", + "content": [ + [{"tag": "text", "text": "Hello "}, {"tag": "a", "text": "world", "href": "http://example.com"}], + [{"tag": "text", "text": "Second line"}], + ], + } + } + result = FeishuChannel._extract_post_text(content) + assert "Test Title" in result + assert "Hello world" in result + assert "Second line" in result + + def test_extract_post_text_empty(self): + result = FeishuChannel._extract_post_text({}) + assert result == "" + + def test_strip_mention(self): + config = FeishuConfig() + channel = FeishuChannel(config) + channel._mention_names = ["@_user_1"] + result = channel._strip_mention("@_user_1 hello world") + assert result == "hello world" + + +class TestFeishuMarkdownConversion: + def test_empty_text(self): + assert _markdown_to_feishu_post("") is None + assert _markdown_to_feishu_post(" ") is None + + def test_plain_text(self): + result = _markdown_to_feishu_post("Hello world") + assert result is not None + assert "zh_cn" in result + content = result["zh_cn"]["content"] + assert len(content) >= 1 + + def test_code_block(self): + md = "```python\nprint('hello')\n```" + result = _markdown_to_feishu_post(md) + assert result is not None + content = result["zh_cn"]["content"] + # Should have a code_block element + found = False + for para in content: + for elem in para: + if elem.get("tag") == "code_block": + found = True + assert elem["language"] == "python" + assert "print" in elem["text"] + assert found + + def test_bold_text(self): + elements = _parse_inline_text("**bold text**") + assert any( + e.get("style") == ["bold"] and e["text"] == "bold text" + for e in elements + ) + + def test_inline_code(self): + elements = _parse_inline_text("`code`") + assert any( + "code_block" in (e.get("style") or []) and e["text"] == "code" + for e in elements + ) + + def test_link(self): + elements = _parse_inline_text("[click](http://example.com)") + assert any( + e.get("tag") == "a" and e["text"] == "click" + for e in elements + ) + + +class TestFeishuChannelRegistration: + def test_feishu_registered(self): + from EvoScientist.channels.channel_manager import available_channels + channels = available_channels() + assert "feishu" in channels + + +class TestFeishuProbe: + def test_missing_app_id(self): + from EvoScientist.channels.feishu.probe import validate_feishu_credentials + ok, msg = _run(validate_feishu_credentials("", "secret")) + assert ok is False + assert "app_id" in msg + + def test_missing_app_secret(self): + from EvoScientist.channels.feishu.probe import validate_feishu_credentials + ok, msg = _run(validate_feishu_credentials("id", "")) + assert ok is False + assert "app_secret" in msg From bec3fc0b675c3a2c4853d51cff023b5bac689814 Mon Sep 17 00:00:00 2001 From: MuXinCG <202322130196@mail.sdu.edu.cn> Date: Tue, 17 Feb 2026 15:40:57 +0800 Subject: [PATCH 2/3] pass linter --- EvoScientist/channels/dingtalk/__init__.py | 8 +- EvoScientist/channels/feishu/__init__.py | 3 +- EvoScientist/channels/feishu/channel.py | 8 + EvoScientist/config/onboard.py | 20 ++- tests/test_dingtalk_channel.py | 181 ++++++++++++++++++++- 5 files changed, 211 insertions(+), 9 deletions(-) diff --git a/EvoScientist/channels/dingtalk/__init__.py b/EvoScientist/channels/dingtalk/__init__.py index 2026d82..df73768 100644 --- a/EvoScientist/channels/dingtalk/__init__.py +++ b/EvoScientist/channels/dingtalk/__init__.py @@ -16,11 +16,11 @@ __all__ = ["DingTalkChannel", "DingTalkConfig"] def create_from_config(config) -> DingTalkChannel: - allowed = _parse_csv(getattr(config, "dingtalk_allowed_senders", "")) - proxy = getattr(config, "dingtalk_proxy", "") or None + allowed = _parse_csv(config.dingtalk_allowed_senders) + proxy = config.dingtalk_proxy if config.dingtalk_proxy else None return DingTalkChannel(DingTalkConfig( - client_id=getattr(config, "dingtalk_client_id", ""), - client_secret=getattr(config, "dingtalk_client_secret", ""), + client_id=config.dingtalk_client_id, + client_secret=config.dingtalk_client_secret, allowed_senders=allowed, proxy=proxy, )) diff --git a/EvoScientist/channels/feishu/__init__.py b/EvoScientist/channels/feishu/__init__.py index 4cc68e5..d24ce16 100644 --- a/EvoScientist/channels/feishu/__init__.py +++ b/EvoScientist/channels/feishu/__init__.py @@ -6,6 +6,7 @@ __all__ = ["FeishuChannel", "FeishuConfig"] def create_from_config(config) -> FeishuChannel: allowed = _parse_csv(config.feishu_allowed_senders) + proxy = config.feishu_proxy if config.feishu_proxy else None return FeishuChannel(FeishuConfig( app_id=config.feishu_app_id, app_secret=config.feishu_app_secret, @@ -14,7 +15,7 @@ def create_from_config(config) -> FeishuChannel: webhook_port=config.feishu_webhook_port, allowed_senders=allowed, feishu_domain=config.feishu_domain, - proxy=getattr(config, 'feishu_proxy', '') or None, + proxy=proxy, )) diff --git a/EvoScientist/channels/feishu/channel.py b/EvoScientist/channels/feishu/channel.py index 8edebd1..9010605 100644 --- a/EvoScientist/channels/feishu/channel.py +++ b/EvoScientist/channels/feishu/channel.py @@ -222,6 +222,14 @@ class FeishuChannel(Channel, WebhookMixin, TokenMixin): name = "feishu" _ready_attrs = ("_http_client", "_access_token") + _non_retryable_patterns = ( + "app_access_token is empty", # invalid credentials + "10003", # invalid app_id + "10014", # invalid app_secret + "99991401", # permission denied + "99991663", # no permission + "99991672", # feature not enabled + ) _rate_limit_patterns = ("99991400", "rate limit", "频率限制") _rate_limit_delay = 2.0 diff --git a/EvoScientist/config/onboard.py b/EvoScientist/config/onboard.py index 6c9584e..d99e3e1 100644 --- a/EvoScientist/config/onboard.py +++ b/EvoScientist/config/onboard.py @@ -1319,7 +1319,7 @@ def _step_channels(config: EvoScientistConfig) -> dict[str, object]: ("telegram", "Telegram", [("telegram_bot_token", "Bot token (from @BotFather)")], "telegram", "telegram"), ("discord", "Discord", [("discord_bot_token", "Bot token")], "discord", "discord"), ("slack", "Slack", [("slack_bot_token", "Bot token (xoxb-...)"), ("slack_app_token", "App token for Socket Mode (xapp-...)")], "slack_sdk", "slack"), - ("feishu", "Feishu", [("feishu_app_id", "App ID"), ("feishu_app_secret", "App Secret"), ("feishu_verification_token", "Verification Token"), ("feishu_encrypt_key", "Encrypt Key")], "aiohttp", "feishu"), + ("feishu", "Feishu", [("feishu_app_id", "App ID"), ("feishu_app_secret", "App Secret")], "aiohttp", "feishu"), ("dingtalk", "DingTalk", [("dingtalk_client_id", "Client ID (AppKey)"), ("dingtalk_client_secret", "Client Secret (AppSecret)")], "aiohttp", "dingtalk"), ("wechat", "WeChat", [("wechat_wecom_corp_id", "WeCom Corp ID"), ("wechat_wecom_agent_id", "WeCom Agent ID"), ("wechat_wecom_secret", "WeCom Secret")], "aiohttp", "wechat"), ("imessage", "iMessage", [], None, None), # handled via _setup_imessage() @@ -1412,6 +1412,24 @@ def _step_channels(config: EvoScientistConfig) -> dict[str, object]: raise KeyboardInterrupt() updates[field_name] = value.strip() + # Feishu optional fields (verification_token & encrypt_key) + if ch_name == "feishu": + console.print(" [dim]The following fields are optional (press Enter to skip):[/dim]") + for field_name, prompt_label in [ + ("feishu_verification_token", "Verification Token (optional)"), + ("feishu_encrypt_key", "Encrypt Key (optional)"), + ]: + current = getattr(config, field_name, "") + value = questionary.text( + f"{prompt_label}:", + default=current, + style=WIZARD_STYLE, + qmark=f" {QMARK}", + ).ask() + if value is None: + raise KeyboardInterrupt() + updates[field_name] = value.strip() + # Allowed senders (common for all channels) senders_field = f"{ch_name}_allowed_senders" if hasattr(config, senders_field): diff --git a/tests/test_dingtalk_channel.py b/tests/test_dingtalk_channel.py index 91870b6..e3deba6 100644 --- a/tests/test_dingtalk_channel.py +++ b/tests/test_dingtalk_channel.py @@ -1,11 +1,13 @@ """Tests for DingTalk channel implementation.""" import asyncio +import json +from unittest.mock import AsyncMock, MagicMock, patch import pytest from EvoScientist.channels.dingtalk.channel import DingTalkChannel, DingTalkConfig -from EvoScientist.channels.base import ChannelError +from EvoScientist.channels.base import ChannelError, OutboundMessage def _run(coro): @@ -72,8 +74,6 @@ class TestDingTalkChannel: _run(channel.stop()) def test_send_returns_false_without_client(self): - from EvoScientist.channels.base import OutboundMessage - config = DingTalkConfig(client_id="test-id", client_secret="test-secret") channel = DingTalkChannel(config) msg = OutboundMessage( @@ -97,6 +97,181 @@ class TestDingTalkChannel: assert channel.capabilities.media_receive is True +class TestDingTalkErrorPatterns: + """Test non-retryable and rate-limit pattern detection.""" + + def test_non_retryable_patterns_defined(self): + config = DingTalkConfig() + channel = DingTalkChannel(config) + assert "invalidauthentication" in channel._non_retryable_patterns + assert "forbidden" in channel._non_retryable_patterns + assert "40014" in channel._non_retryable_patterns + + def test_non_retryable_returns_none(self): + config = DingTalkConfig() + channel = DingTalkChannel(config) + exc = Exception("invalidauthentication: bad credentials") + result = channel._extract_retry_after(exc) + assert result is None + + def test_rate_limit_returns_delay(self): + config = DingTalkConfig() + channel = DingTalkChannel(config) + # Base class default includes "429" and "ratelimit" + exc = Exception("HTTP 429 ratelimit exceeded") + result = channel._extract_retry_after(exc) + assert result is not None + assert result > 0 + + +class TestDingTalkWsMessageParsing: + """Test _on_ws_message parsing logic with mocked bus.""" + + def _make_channel(self): + config = DingTalkConfig(client_id="test-app", client_secret="test-secret") + channel = DingTalkChannel(config) + channel._running = True + channel._ws_session = MagicMock() + channel._ws_session.send_str = AsyncMock() + channel._http_client = MagicMock() + channel._access_token = "fake-token" + channel._token_expires = 9999999999 + return channel + + def test_system_ping_ack(self): + channel = self._make_channel() + data = { + "type": "SYSTEM", + "headers": {"topic": "ping", "messageId": "ping-1"}, + "data": "pong-data", + } + _run(channel._on_ws_message(data)) + channel._ws_session.send_str.assert_called_once() + sent = json.loads(channel._ws_session.send_str.call_args[0][0]) + assert sent["code"] == 200 + assert sent["data"] == "pong-data" + + def test_callback_text_message(self): + channel = self._make_channel() + channel._enqueue_raw = AsyncMock() + + payload = { + "text": {"content": "hello bot"}, + "senderStaffId": "staff123", + "conversationType": "1", + "createAt": "1700000000000", + } + data = { + "type": "CALLBACK", + "headers": {"messageId": "msg-1", "contentType": "application/json"}, + "data": json.dumps(payload), + } + _run(channel._on_ws_message(data)) + channel._enqueue_raw.assert_called_once() + raw = channel._enqueue_raw.call_args[0][0] + assert raw.text == "hello bot" + assert raw.sender_id == "staff123" + assert raw.is_group is False + + def test_callback_group_message_mention(self): + channel = self._make_channel() + channel._enqueue_raw = AsyncMock() + + payload = { + "text": {"content": "@bot hello"}, + "senderStaffId": "staff456", + "conversationType": "2", + "isInAtList": True, + "createAt": "1700000000000", + } + data = { + "type": "CALLBACK", + "headers": {"messageId": "msg-2"}, + "data": json.dumps(payload), + } + _run(channel._on_ws_message(data)) + raw = channel._enqueue_raw.call_args[0][0] + assert raw.is_group is True + assert raw.was_mentioned is True + + def test_callback_group_no_mention(self): + channel = self._make_channel() + channel._enqueue_raw = AsyncMock() + + payload = { + "text": {"content": "just chatting"}, + "senderStaffId": "staff789", + "conversationType": "2", + "createAt": "1700000000000", + } + data = { + "type": "CALLBACK", + "headers": {"messageId": "msg-3"}, + "data": json.dumps(payload), + } + _run(channel._on_ws_message(data)) + raw = channel._enqueue_raw.call_args[0][0] + assert raw.is_group is True + assert raw.was_mentioned is False + + def test_ignores_non_callback(self): + channel = self._make_channel() + channel._enqueue_raw = AsyncMock() + + data = { + "type": "EVENT", + "headers": {"messageId": "msg-x"}, + "data": "{}", + } + _run(channel._on_ws_message(data)) + channel._enqueue_raw.assert_not_called() + + def test_ignores_empty_content(self): + channel = self._make_channel() + channel._enqueue_raw = AsyncMock() + + payload = { + "text": {"content": ""}, + "senderStaffId": "staff0", + "conversationType": "1", + } + data = { + "type": "CALLBACK", + "headers": {"messageId": "msg-e"}, + "data": json.dumps(payload), + } + _run(channel._on_ws_message(data)) + channel._enqueue_raw.assert_not_called() + + def test_non_dict_data_ignored(self): + channel = self._make_channel() + channel._enqueue_raw = AsyncMock() + _run(channel._on_ws_message("not a dict")) + channel._enqueue_raw.assert_not_called() + + +class TestDingTalkSendChunk: + """Test _send_chunk with mocked HTTP client.""" + + def test_send_chunk_calls_api(self): + config = DingTalkConfig(client_id="test-app", client_secret="test-secret") + channel = DingTalkChannel(config) + channel._access_token = "fake-token" + channel._token_expires = 9999999999 + + mock_response = MagicMock() + mock_response.json.return_value = {"processQueryKey": "ok"} + channel._http_client = MagicMock() + channel._http_client.post = AsyncMock(return_value=mock_response) + + _run(channel._send_chunk("user1", "formatted", "raw text", None, {})) + channel._http_client.post.assert_called_once() + call_args = channel._http_client.post.call_args + body = call_args.kwargs.get("json") or call_args[1].get("json") + assert body["robotCode"] == "test-app" + assert body["userIds"] == ["user1"] + + class TestDingTalkChannelRegistration: def test_dingtalk_registered(self): from EvoScientist.channels.channel_manager import available_channels From 5fd8e1a8ade7b19ba1b5175a21ab49427402054b Mon Sep 17 00:00:00 2001 From: MuXinCG <202322130196@mail.sdu.edu.cn> Date: Tue, 17 Feb 2026 16:08:41 +0800 Subject: [PATCH 3/3] fix feishu bugs --- EvoScientist/channels/feishu/channel.py | 48 +++- tests/test_dingtalk_channel.py | 2 +- tests/test_feishu_channel.py | 308 +++++++++++++++++++++++- 3 files changed, 351 insertions(+), 7 deletions(-) diff --git a/EvoScientist/channels/feishu/channel.py b/EvoScientist/channels/feishu/channel.py index 9010605..16e6fbd 100644 --- a/EvoScientist/channels/feishu/channel.py +++ b/EvoScientist/channels/feishu/channel.py @@ -18,6 +18,8 @@ Send API: from __future__ import annotations +import base64 +import hashlib import json import logging import re @@ -538,6 +540,30 @@ class FeishuChannel(Channel, WebhookMixin, TokenMixin): # Clean up extra whitespace left behind return re.sub(r" +", " ", result).strip() + # ── Event decryption ───────────────────────────────────────── + + def _decrypt_event(self, encrypted: str) -> dict: + """Decrypt a Feishu encrypted event payload (AES-256-CBC). + + Feishu encryption spec: + key = SHA256(encrypt_key) + data = base64_decode(encrypted) + iv = data[:16] + plain = AES_CBC_decrypt(data[16:], key, iv) # PKCS7 padded + """ + from cryptography.hazmat.primitives.ciphers import Cipher, algorithms, modes + + key = hashlib.sha256(self.config.encrypt_key.encode()).digest() + data = base64.b64decode(encrypted) + iv, ciphertext = data[:16], data[16:] + cipher = Cipher(algorithms.AES(key), modes.CBC(iv)) + decryptor = cipher.decryptor() + padded = decryptor.update(ciphertext) + decryptor.finalize() + # Remove PKCS7 padding + pad_len = padded[-1] + plaintext = padded[:-pad_len].decode() + return json.loads(plaintext) + # ── Webhook event handler ───────────────────────────────────── async def _handle_event(self, request) -> "web.Response": @@ -549,6 +575,14 @@ class FeishuChannel(Channel, WebhookMixin, TokenMixin): except Exception: return web.Response(status=400) + # ── Decrypt if encrypt_key is configured ── + if self.config.encrypt_key and "encrypt" in body: + try: + body = self._decrypt_event(body["encrypt"]) + except Exception: + logger.exception("Feishu event decryption failed") + return web.Response(status=400) + # ── URL verification challenge ── if body.get("type") == "url_verification": challenge = body.get("challenge", "") @@ -567,8 +601,12 @@ class FeishuChannel(Channel, WebhookMixin, TokenMixin): return web.Response(status=403) event_type = header.get("event_type", "") + logger.info(f"Feishu v2 event received: {event_type}") if event_type == "im.message.receive_v1": - await self._on_message(body.get("event", {})) + try: + await self._on_message(body.get("event", {})) + except Exception: + logger.exception("Feishu _on_message failed") # ── v1 event schema (legacy) ── elif "event" in body: @@ -580,8 +618,14 @@ class FeishuChannel(Channel, WebhookMixin, TokenMixin): event = body["event"] msg_type = event.get("type", "") + logger.info(f"Feishu v1 event received: type={msg_type}") if msg_type == "message": - await self._on_message_v1(event) + try: + await self._on_message_v1(event) + except Exception: + logger.exception("Feishu _on_message_v1 failed") + else: + logger.info(f"Feishu event ignored: schema={schema}") return web.Response(status=200) diff --git a/tests/test_dingtalk_channel.py b/tests/test_dingtalk_channel.py index e3deba6..3f40f32 100644 --- a/tests/test_dingtalk_channel.py +++ b/tests/test_dingtalk_channel.py @@ -2,7 +2,7 @@ import asyncio import json -from unittest.mock import AsyncMock, MagicMock, patch +from unittest.mock import AsyncMock, MagicMock import pytest diff --git a/tests/test_feishu_channel.py b/tests/test_feishu_channel.py index 1b779ed..dad38f1 100644 --- a/tests/test_feishu_channel.py +++ b/tests/test_feishu_channel.py @@ -2,6 +2,7 @@ import asyncio import json +from unittest.mock import AsyncMock, MagicMock import pytest @@ -10,8 +11,9 @@ from EvoScientist.channels.feishu.channel import ( FeishuConfig, _markdown_to_feishu_post, _parse_inline_text, + _parse_inline_elements, ) -from EvoScientist.channels.base import ChannelError +from EvoScientist.channels.base import ChannelError, OutboundMessage def _run(coro): @@ -80,8 +82,6 @@ class TestFeishuChannel: _run(channel.stop()) def test_send_returns_false_without_client(self): - from EvoScientist.channels.base import OutboundMessage - config = FeishuConfig(app_id="test-id", app_secret="test-secret") channel = FeishuChannel(config) msg = OutboundMessage( @@ -126,6 +126,17 @@ class TestFeishuChannel: result = FeishuChannel._extract_post_text({}) assert result == "" + def test_extract_post_text_skips_at_mentions(self): + content = { + "zh_cn": { + "content": [ + [{"tag": "text", "text": "Hi "}, {"tag": "at", "user_id": "bot"}], + ], + } + } + result = FeishuChannel._extract_post_text(content) + assert result == "Hi" + def test_strip_mention(self): config = FeishuConfig() channel = FeishuChannel(config) @@ -133,6 +144,236 @@ class TestFeishuChannel: result = channel._strip_mention("@_user_1 hello world") assert result == "hello world" + def test_strip_mention_multiple(self): + config = FeishuConfig() + channel = FeishuChannel(config) + channel._mention_names = ["@_user_1", "@_user_2"] + result = channel._strip_mention("@_user_1 @_user_2 hello") + assert result == "hello" + + def test_strip_mention_no_match(self): + config = FeishuConfig() + channel = FeishuChannel(config) + channel._mention_names = [] + result = channel._strip_mention("hello world") + assert result == "hello world" + + +class TestFeishuErrorPatterns: + """Test non-retryable and rate-limit pattern detection.""" + + def test_non_retryable_patterns_defined(self): + config = FeishuConfig() + channel = FeishuChannel(config) + assert len(channel._non_retryable_patterns) > 0 + assert "10003" in channel._non_retryable_patterns + assert "99991401" in channel._non_retryable_patterns + + def test_non_retryable_returns_none(self): + config = FeishuConfig() + channel = FeishuChannel(config) + exc = Exception("error code 10003: invalid app_id") + result = channel._extract_retry_after(exc) + assert result is None + + def test_rate_limit_patterns_defined(self): + config = FeishuConfig() + channel = FeishuChannel(config) + assert "99991400" in channel._rate_limit_patterns + assert "频率限制" in channel._rate_limit_patterns + + def test_rate_limit_returns_delay(self): + config = FeishuConfig() + channel = FeishuChannel(config) + exc = Exception("99991400 频率限制") + result = channel._extract_retry_after(exc) + assert result is not None + assert result == channel._rate_limit_delay + + +class TestFeishuWebhookEvent: + """Test _on_message parsing with mocked bus.""" + + def _make_channel(self): + config = FeishuConfig(app_id="test-app", app_secret="test-secret") + channel = FeishuChannel(config) + channel._running = True + channel._http_client = MagicMock() + channel._access_token = "fake-token" + channel._token_expires = 9999999999 + channel._enqueue_raw = AsyncMock() + return channel + + def test_text_message_v2(self): + channel = self._make_channel() + event = { + "sender": { + "sender_id": {"open_id": "ou_test123"}, + "sender_type": "user", + }, + "message": { + "chat_id": "oc_chat1", + "message_type": "text", + "message_id": "msg_1", + "chat_type": "p2p", + "content": json.dumps({"text": "hello feishu"}), + "create_time": "1700000000000", + }, + } + _run(channel._on_message(event)) + channel._enqueue_raw.assert_called_once() + raw = channel._enqueue_raw.call_args[0][0] + assert raw.text == "hello feishu" + assert raw.sender_id == "ou_test123" + assert raw.chat_id == "oc_chat1" + assert raw.is_group is False + + def test_group_message_with_mention(self): + channel = self._make_channel() + event = { + "sender": { + "sender_id": {"open_id": "ou_sender"}, + "sender_type": "user", + }, + "message": { + "chat_id": "oc_group1", + "message_type": "text", + "message_id": "msg_2", + "chat_type": "group", + "content": json.dumps({"text": "@_user_1 do something"}), + "create_time": "1700000000000", + "mentions": [{"key": "@_user_1", "id": {}}], + }, + } + _run(channel._on_message(event)) + raw = channel._enqueue_raw.call_args[0][0] + assert raw.is_group is True + assert raw.was_mentioned is True + assert channel._mention_names == ["@_user_1"] + + def test_group_message_no_mention(self): + channel = self._make_channel() + event = { + "sender": { + "sender_id": {"open_id": "ou_sender"}, + "sender_type": "user", + }, + "message": { + "chat_id": "oc_group2", + "message_type": "text", + "message_id": "msg_3", + "chat_type": "group", + "content": json.dumps({"text": "just talking"}), + "create_time": "1700000000000", + }, + } + _run(channel._on_message(event)) + raw = channel._enqueue_raw.call_args[0][0] + assert raw.is_group is True + assert raw.was_mentioned is False + + def test_skips_bot_messages(self): + channel = self._make_channel() + event = { + "sender": { + "sender_id": {"open_id": "ou_bot"}, + "sender_type": "app", + }, + "message": { + "chat_id": "oc_chat", + "message_type": "text", + "message_id": "msg_bot", + "content": json.dumps({"text": "bot reply"}), + }, + } + _run(channel._on_message(event)) + channel._enqueue_raw.assert_not_called() + + def test_post_message(self): + channel = self._make_channel() + post_content = { + "zh_cn": { + "title": "Test", + "content": [[{"tag": "text", "text": "Post body"}]], + } + } + event = { + "sender": { + "sender_id": {"open_id": "ou_test"}, + "sender_type": "user", + }, + "message": { + "chat_id": "oc_chat", + "message_type": "post", + "message_id": "msg_post", + "chat_type": "p2p", + "content": json.dumps(post_content), + "create_time": "1700000000000", + }, + } + _run(channel._on_message(event)) + raw = channel._enqueue_raw.call_args[0][0] + assert "Test" in raw.text + assert "Post body" in raw.text + + def test_unsupported_msg_type_annotation(self): + channel = self._make_channel() + event = { + "sender": { + "sender_id": {"open_id": "ou_test"}, + "sender_type": "user", + }, + "message": { + "chat_id": "oc_chat", + "message_type": "share_chat", + "message_id": "msg_share", + "chat_type": "p2p", + "content": "{}", + "create_time": "1700000000000", + }, + } + _run(channel._on_message(event)) + raw = channel._enqueue_raw.call_args[0][0] + assert "share_chat" in raw.text + + +class TestFeishuSendChunk: + """Test _send_chunk with mocked HTTP client.""" + + def test_send_chunk_post_format(self): + config = FeishuConfig(app_id="test-app", app_secret="test-secret") + channel = FeishuChannel(config) + channel._access_token = "fake-token" + channel._token_expires = 9999999999 + + mock_response = MagicMock() + mock_response.json.return_value = {"code": 0} + channel._http_client = MagicMock() + channel._http_client.post = AsyncMock(return_value=mock_response) + + _run(channel._send_chunk("oc_chat1", "formatted", "raw **text**", None, {})) + channel._http_client.post.assert_called() + # Should try post format first + call_args = channel._http_client.post.call_args + body = call_args.kwargs.get("json") or call_args[1].get("json") + assert body["receive_id"] == "oc_chat1" + + def test_send_chunk_with_reply(self): + config = FeishuConfig(app_id="test-app", app_secret="test-secret") + channel = FeishuChannel(config) + channel._access_token = "fake-token" + channel._token_expires = 9999999999 + + mock_response = MagicMock() + mock_response.json.return_value = {"code": 0} + channel._http_client = MagicMock() + channel._http_client.post = AsyncMock(return_value=mock_response) + + _run(channel._send_chunk("oc_chat1", "reply", "reply text", "om_reply_id", {})) + # Should call the reply API endpoint + first_call_url = channel._http_client.post.call_args_list[0][0][0] + assert "reply" in first_call_url + class TestFeishuMarkdownConversion: def test_empty_text(self): @@ -151,7 +392,6 @@ class TestFeishuMarkdownConversion: result = _markdown_to_feishu_post(md) assert result is not None content = result["zh_cn"]["content"] - # Should have a code_block element found = False for para in content: for elem in para: @@ -161,6 +401,16 @@ class TestFeishuMarkdownConversion: assert "print" in elem["text"] assert found + def test_code_block_no_language(self): + md = "```\nsome code\n```" + result = _markdown_to_feishu_post(md) + assert result is not None + content = result["zh_cn"]["content"] + for para in content: + for elem in para: + if elem.get("tag") == "code_block": + assert elem["language"] == "plain" + def test_bold_text(self): elements = _parse_inline_text("**bold text**") assert any( @@ -182,6 +432,56 @@ class TestFeishuMarkdownConversion: for e in elements ) + def test_strikethrough(self): + elements = _parse_inline_text("~~deleted~~") + assert any( + e.get("style") == ["strikethrough"] and e["text"] == "deleted" + for e in elements + ) + + def test_italic(self): + elements = _parse_inline_text("_italic text_") + assert any( + e.get("style") == ["italic"] and e["text"] == "italic text" + for e in elements + ) + + def test_heading(self): + elements = _parse_inline_elements("## My Heading") + assert any( + e.get("style") == ["bold"] and e["text"] == "My Heading" + for e in elements + ) + + def test_blockquote(self): + elements = _parse_inline_elements("> quoted text") + # Should have ▎ prefix with italic style + assert any(e.get("text") == "▎" for e in elements) + assert any(e.get("text") == "quoted text" for e in elements) + + def test_unordered_list(self): + elements = _parse_inline_elements("- list item") + assert any(e.get("text") == "• " for e in elements) + assert any(e.get("text") == "list item" for e in elements) + + def test_ordered_list(self): + elements = _parse_inline_elements("3. third item") + assert any(e.get("text") == "3. " for e in elements) + assert any(e.get("text") == "third item" for e in elements) + + def test_multi_paragraph(self): + md = "First paragraph\n\nSecond paragraph" + result = _markdown_to_feishu_post(md) + content = result["zh_cn"]["content"] + assert len(content) == 2 + + def test_mixed_content(self): + md = "# Title\n\nSome **bold** text\n\n```python\ncode\n```" + result = _markdown_to_feishu_post(md) + assert result is not None + content = result["zh_cn"]["content"] + assert len(content) >= 3 + class TestFeishuChannelRegistration: def test_feishu_registered(self):