Merge pull request #13 from EvoScientist/feature/channel-dingtalk-feishu

Add Dingtalk Feishu
This commit is contained in:
Xi Zhang
2026-02-17 13:31:40 +00:00
committed by GitHub
13 changed files with 2344 additions and 2 deletions
+1 -1
View File
@@ -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
@@ -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(config.dingtalk_allowed_senders)
proxy = config.dingtalk_proxy if config.dingtalk_proxy else None
return DingTalkChannel(DingTalkConfig(
client_id=config.dingtalk_client_id,
client_secret=config.dingtalk_client_secret,
allowed_senders=allowed,
proxy=proxy,
))
register_channel("dingtalk", create_from_config)
+354
View File
@@ -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")
+33
View File
@@ -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}"
+92
View File
@@ -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()
+22
View File
@@ -0,0 +1,22 @@
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)
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,
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=proxy,
))
register_channel("feishu", create_from_config)
+825
View File
@@ -0,0 +1,825 @@
"""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 base64
import hashlib
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")
_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
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()
# ── 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":
"""Handle POST /webhook/event from Feishu."""
from aiohttp import web
try:
body = await request.json()
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", "")
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", "")
logger.info(f"Feishu v2 event received: {event_type}")
if event_type == "im.message.receive_v1":
try:
await self._on_message(body.get("event", {}))
except Exception:
logger.exception("Feishu _on_message failed")
# ── 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", "")
logger.info(f"Feishu v1 event received: type={msg_type}")
if msg_type == "message":
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)
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)
+39
View File
@@ -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}"
+113
View File
@@ -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()
+34
View File
@@ -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")], "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()
]
@@ -1410,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):
@@ -1512,6 +1532,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"
+1 -1
View File
@@ -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
+297
View File
@@ -0,0 +1,297 @@
"""Tests for DingTalk channel implementation."""
import asyncio
import json
from unittest.mock import AsyncMock, MagicMock
import pytest
from EvoScientist.channels.dingtalk.channel import DingTalkChannel, DingTalkConfig
from EvoScientist.channels.base import ChannelError, OutboundMessage
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):
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 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
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
+504
View File
@@ -0,0 +1,504 @@
"""Tests for Feishu channel implementation."""
import asyncio
import json
from unittest.mock import AsyncMock, MagicMock
import pytest
from EvoScientist.channels.feishu.channel import (
FeishuChannel,
FeishuConfig,
_markdown_to_feishu_post,
_parse_inline_text,
_parse_inline_elements,
)
from EvoScientist.channels.base import ChannelError, OutboundMessage
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):
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_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)
channel._mention_names = ["@_user_1"]
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):
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"]
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_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(
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
)
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):
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