feat(feishu): add WebSocket long connection subscription mode (#87)
* feat(feishu): add WebSocket long connection subscription mode Add WebSocket (长连接) mode as an alternative to webhook for Feishu event subscription. This allows running without a public IP, port forwarding, or tunnel — ideal for local dev and NAT/firewall setups. - New `feishu_subscription_mode` config: "webhook" (default) or "websocket" - WebSocket mode uses official `lark-oapi` SDK with thread-safe queue bridge - Onboard wizard: mode selection, SDK install prompt for websocket - CLI: `--mode webhook|websocket` for standalone serve - `pip install evoscientist[feishu]` optional dependency - 5 new tests covering config, SDK missing error, message bridge, cleanup - Docs: subscription mode comparison table, prerequisites per mode * Fix: Ruff * Fix: small fix
This commit is contained in:
@@ -201,7 +201,7 @@ Mention detection is platform-specific:
|
||||
| Telegram | HTML | 4000 | S/R | R | R | R | | 4s | emoji | | G | @ | yes | | yes | yes |
|
||||
| Discord | Discord | 2000 | S/R | | | | | 8s | emoji | yes | G | @ | yes | | yes | yes |
|
||||
| Slack | Mrkdwn | 4000 | S/R | | | | | post | emoji | yes | G | @ | yes | | yes | yes |
|
||||
| Feishu | Post | 4096 | S/R | R | R | | | | emoji | | G | @ | no | 2h | yes | yes |
|
||||
| Feishu | Post | 4096 | S/R | R | R | | | | emoji | | G | @ | ws mode | 2h | yes | yes |
|
||||
| WeChat | MD | 4096 | S/R | R | | R | R | recall | | | G | @ | no | 2h | yes | yes |
|
||||
| DingTalk | MD | 4096 | S/R | R | | | R | | | | G | @ | yes | 2h | yes | yes |
|
||||
| QQ | Plain | 4096 | S/R | | | | | | | | G | @ | yes | | | yes |
|
||||
@@ -218,7 +218,7 @@ Legend: **S** = send, **R** = receive, **G** = group chat, **@** = @mention dete
|
||||
| Telegram | HTTPS | Long polling (`getUpdates`) | -- |
|
||||
| Discord | WebSocket | Gateway events (`discord.py`) | -- |
|
||||
| Slack | WebSocket | Socket Mode (`slack-sdk`) | -- |
|
||||
| Feishu | HTTP | Webhook `POST /webhook/event` | 9000 |
|
||||
| Feishu | HTTP/WS | Webhook or WebSocket long connection | 9000/-- |
|
||||
| WeChat | HTTP | Webhook `POST /wechat/callback` | 9001 |
|
||||
| DingTalk | WebSocket | Stream Mode (DingTalk gateway) | -- |
|
||||
| QQ | WebSocket | Bot Gateway (`qq-botpy`) | -- |
|
||||
@@ -456,12 +456,40 @@ slack_proxy: ""
|
||||
|
||||
1. Go to [Feishu Open Platform](https://open.feishu.cn/app) (international: [Lark Developer](https://open.larksuite.com/app)) -> create a custom app.
|
||||
2. Copy the **App ID** and **App Secret**.
|
||||
3. Left menu **Event Subscriptions** -> set request URL to `http://your-host:9000/webhook/event` -> copy **Verification Token** and **Encrypt Key**.
|
||||
3. Left menu **Event Subscriptions**:
|
||||
- **Webhook mode**: set request URL to `http://your-host:9000/webhook/event` -> copy **Verification Token** and **Encrypt Key**.
|
||||
- **WebSocket mode**: select **长连接** (Long Connection) as the subscription method. No URL needed.
|
||||
4. Add event: `im.message.receive_v1` (receive messages).
|
||||
5. Left menu **Permissions** -> enable `im:message:send_as_bot`.
|
||||
6. Create a version and publish.
|
||||
|
||||
> Webhook must be publicly reachable. For local dev, use `ngrok http 9000`.
|
||||
> **Webhook mode** requires a publicly reachable URL. For local dev, use `ngrok http 9000` or [natapp](https://natapp.cn/) (recommended for China).
|
||||
|
||||
#### Subscription Modes
|
||||
|
||||
Feishu supports two subscription modes:
|
||||
|
||||
| Mode | Transport | Public IP? | Best For |
|
||||
|------|-----------|:----------:|----------|
|
||||
| `webhook` (default) | HTTP POST callback | Yes | Servers with public IP / cloud deployment |
|
||||
| `websocket` | WebSocket long connection | **No** | Local dev, behind NAT/firewall, China (no ngrok needed) |
|
||||
|
||||
**WebSocket mode** uses the official `lark-oapi` SDK to maintain an outbound WebSocket connection to Feishu servers. No public IP, port forwarding, or tunnel is required.
|
||||
|
||||
To use WebSocket mode:
|
||||
|
||||
```bash
|
||||
# Install the SDK
|
||||
pip install 'evoscientist[feishu]'
|
||||
|
||||
# Via config file
|
||||
feishu_subscription_mode: "websocket"
|
||||
|
||||
# Via CLI (standalone)
|
||||
python -m EvoScientist.channels.feishu.serve --app-id ID --app-secret SECRET --mode websocket
|
||||
```
|
||||
|
||||
> **Note:** In WebSocket mode, `feishu_verification_token`, `feishu_encrypt_key`, and `feishu_webhook_port` are not used — the SDK handles authentication and encryption internally.
|
||||
|
||||
**Configuration:**
|
||||
|
||||
@@ -469,6 +497,7 @@ slack_proxy: ""
|
||||
channel_enabled: "feishu"
|
||||
feishu_app_id: "cli_xxxxxxx"
|
||||
feishu_app_secret: "xxxxxxxxxxxxxxxxxx"
|
||||
feishu_subscription_mode: "webhook" # or "websocket"
|
||||
feishu_webhook_port: 9000
|
||||
feishu_allowed_senders: "" # open_id
|
||||
feishu_domain: "https://open.feishu.cn"
|
||||
@@ -483,6 +512,7 @@ feishu_proxy: ""
|
||||
| `feishu_allowed_senders` | `str` | `""` | Comma-separated open_ids |
|
||||
| `feishu_domain` | `str` | `"https://open.feishu.cn"` | API domain (use `https://open.larksuite.com` for Lark) |
|
||||
| `feishu_proxy` | `str` | `""` | HTTPS proxy URL |
|
||||
| `feishu_subscription_mode` | `str` | `"webhook"` | `"webhook"` or `"websocket"` (WebSocket long connection, no public IP) |
|
||||
|
||||
**Env vars:** `EVOSCIENTIST_FEISHU_APP_ID`, `EVOSCIENTIST_FEISHU_APP_SECRET`, `EVOSCIENTIST_FEISHU_WEBHOOK_PORT`, `EVOSCIENTIST_FEISHU_DOMAIN`
|
||||
|
||||
@@ -860,7 +890,7 @@ pip install evoscientist[telegram] # or discord, slack, feishu, etc.
|
||||
**Webhook channels (Feishu, WeChat) not receiving messages**
|
||||
|
||||
1. Ensure the webhook URL is publicly reachable (not behind NAT without port forwarding)
|
||||
2. For local development, use a tunnel: `ngrok http 9000`
|
||||
2. For local development, use a tunnel: `ngrok http 9000` or [natapp](https://natapp.cn/) for China
|
||||
3. Verify the callback URL matches exactly (including path: `/webhook/event` for Feishu, `/wechat/callback` for WeChat)
|
||||
4. Check that signature verification tokens match between the platform config and your local config
|
||||
|
||||
|
||||
@@ -17,6 +17,7 @@ def create_from_config(config) -> FeishuChannel:
|
||||
allowed_senders=allowed,
|
||||
feishu_domain=config.feishu_domain,
|
||||
proxy=proxy,
|
||||
subscription_mode=getattr(config, "feishu_subscription_mode", "webhook"),
|
||||
)
|
||||
)
|
||||
|
||||
|
||||
@@ -1,16 +1,20 @@
|
||||
"""Feishu (飞书/Lark) channel implementation.
|
||||
|
||||
Receives messages via an HTTP event subscription webhook (aiohttp),
|
||||
sends replies via Feishu Open API REST endpoints.
|
||||
Receives messages via HTTP event subscription webhook (aiohttp) or
|
||||
WebSocket long connection (lark-oapi SDK), 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
|
||||
Event subscription (two modes):
|
||||
- **Webhook**: URL verification challenge on first request,
|
||||
``im.message.receive_v1`` events via HTTP POST callback.
|
||||
Requires a publicly reachable URL.
|
||||
- **WebSocket**: outbound long connection via ``lark-oapi`` SDK.
|
||||
No public IP required.
|
||||
|
||||
Send API:
|
||||
- ``POST /open-apis/im/v1/messages?receive_id_type=chat_id``
|
||||
@@ -18,11 +22,14 @@ Send API:
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import base64
|
||||
import hashlib
|
||||
import json
|
||||
import logging
|
||||
import queue
|
||||
import re
|
||||
import threading
|
||||
from dataclasses import dataclass
|
||||
from datetime import datetime
|
||||
from pathlib import Path
|
||||
@@ -240,6 +247,7 @@ class FeishuConfig(BaseChannelConfig):
|
||||
webhook_port: int = 9000
|
||||
text_chunk_limit: int = 4096
|
||||
feishu_domain: str = "https://open.feishu.cn"
|
||||
subscription_mode: str = "webhook" # "webhook" | "websocket"
|
||||
|
||||
|
||||
class FeishuChannel(Channel, WebhookMixin, TokenMixin):
|
||||
@@ -262,6 +270,10 @@ class FeishuChannel(Channel, WebhookMixin, TokenMixin):
|
||||
def __init__(self, config: FeishuConfig):
|
||||
super().__init__(config)
|
||||
self._mention_names: list[str] = [] # bot mention keys from events
|
||||
self._main_loop: asyncio.AbstractEventLoop | None = None
|
||||
self._lark_ws_thread: threading.Thread | None = None
|
||||
self._ws_event_queue: queue.Queue | None = None
|
||||
self._ws_consumer_task: asyncio.Task | None = None
|
||||
|
||||
# ── WebhookMixin overrides ────────────────────────────────────
|
||||
|
||||
@@ -292,7 +304,25 @@ class FeishuChannel(Channel, WebhookMixin, TokenMixin):
|
||||
|
||||
# ── Lifecycle ─────────────────────────────────────────────────
|
||||
|
||||
_VALID_SUBSCRIPTION_MODES = ("webhook", "websocket")
|
||||
|
||||
async def start(self) -> None:
|
||||
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")
|
||||
if self.config.subscription_mode not in self._VALID_SUBSCRIPTION_MODES:
|
||||
raise ChannelError(
|
||||
f"Invalid feishu_subscription_mode: {self.config.subscription_mode!r}. "
|
||||
f"Must be one of {self._VALID_SUBSCRIPTION_MODES}"
|
||||
)
|
||||
|
||||
if self.config.subscription_mode == "websocket":
|
||||
await self._start_websocket_mode()
|
||||
else:
|
||||
await self._start_webhook_mode()
|
||||
|
||||
async def _start_webhook_mode(self) -> None:
|
||||
try:
|
||||
import httpx # noqa: F401
|
||||
from aiohttp import web # noqa: F401
|
||||
@@ -302,11 +332,6 @@ class FeishuChannel(Channel, WebhookMixin, TokenMixin):
|
||||
"Install with: pip install aiohttp httpx"
|
||||
) from None
|
||||
|
||||
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()
|
||||
|
||||
@@ -318,8 +343,161 @@ class FeishuChannel(Channel, WebhookMixin, TokenMixin):
|
||||
f"Feishu channel started (webhook on port {self.config.webhook_port})"
|
||||
)
|
||||
|
||||
async def _start_websocket_mode(self) -> None:
|
||||
try:
|
||||
import lark_oapi as lark
|
||||
except ImportError:
|
||||
raise ChannelError(
|
||||
"lark-oapi not installed. Install with: pip install 'lark-oapi>=1.4.0'"
|
||||
) from None
|
||||
|
||||
import httpx
|
||||
|
||||
proxy = getattr(self.config, "proxy", None) or None
|
||||
self._http_client = httpx.AsyncClient(timeout=15, proxy=proxy)
|
||||
|
||||
# Verify credentials by fetching initial token
|
||||
await self._refresh_token()
|
||||
|
||||
self._main_loop = asyncio.get_running_loop()
|
||||
|
||||
# Thread-safe queue: SDK thread puts events, main loop consumes
|
||||
self._ws_event_queue = queue.Queue()
|
||||
|
||||
# Set _running BEFORE creating consumer task — the task checks
|
||||
# `while self._running` and would exit immediately otherwise.
|
||||
self._running = True
|
||||
self._ws_consumer_task = asyncio.create_task(self._consume_ws_events())
|
||||
|
||||
# Build SDK event handler
|
||||
handler = (
|
||||
lark.EventDispatcherHandler.builder("", "")
|
||||
.register_p2_im_message_receive_v1(self._on_lark_sdk_message)
|
||||
.build()
|
||||
)
|
||||
|
||||
ws_client = lark.ws.Client(
|
||||
self.config.app_id,
|
||||
self.config.app_secret,
|
||||
event_handler=handler,
|
||||
log_level=lark.LogLevel.WARNING,
|
||||
)
|
||||
|
||||
def _run_ws():
|
||||
# Workaround: nest_asyncio (used by the main process) patches
|
||||
# the event loop policy globally but does NOT patch Handle._run
|
||||
# for context re-entry (versions < 1.7). The SDK creates its
|
||||
# own event loop in this thread; its transport callbacks hit
|
||||
# RuntimeError: cannot enter context: … is already entered
|
||||
# Fix: patch Handle._run to retry with a context copy.
|
||||
_orig_handle_run = asyncio.Handle._run
|
||||
|
||||
def _safe_handle_run(self):
|
||||
try:
|
||||
self._context.run(self._callback, *self._args)
|
||||
except RuntimeError as exc:
|
||||
if "cannot enter context" not in str(exc):
|
||||
raise
|
||||
ctx = self._context.copy()
|
||||
ctx.run(self._callback, *self._args)
|
||||
|
||||
asyncio.Handle._run = _safe_handle_run
|
||||
try:
|
||||
ws_client.start()
|
||||
except Exception:
|
||||
logger.exception(
|
||||
"Feishu WebSocket SDK thread exited unexpectedly. "
|
||||
"The channel will no longer receive messages. "
|
||||
"Check app_id/app_secret and connection limits."
|
||||
)
|
||||
finally:
|
||||
asyncio.Handle._run = _orig_handle_run
|
||||
|
||||
self._lark_ws_thread = threading.Thread(target=_run_ws, daemon=True)
|
||||
self._lark_ws_thread.start()
|
||||
|
||||
logger.info("Feishu channel started (WebSocket long connection mode)")
|
||||
|
||||
def _on_lark_sdk_message(self, data) -> None:
|
||||
"""Sync callback invoked in lark-oapi SDK thread.
|
||||
|
||||
Converts the SDK event object to a dict and puts it on a
|
||||
thread-safe queue. The ``_consume_ws_events`` task on the main
|
||||
asyncio loop picks it up — no asyncio cross-thread calls needed,
|
||||
avoiding the nest_asyncio + Python 3.11 contextvars conflict.
|
||||
"""
|
||||
try:
|
||||
event = data.event
|
||||
msg = event.message
|
||||
sender = event.sender
|
||||
|
||||
# Rebuild mentions list from SDK objects
|
||||
mentions_list = []
|
||||
if msg.mentions:
|
||||
for m in msg.mentions:
|
||||
mention_dict: dict[str, Any] = {"key": m.key, "id": {}}
|
||||
if m.id:
|
||||
mention_dict["id"] = {
|
||||
"open_id": getattr(m.id, "open_id", ""),
|
||||
"user_id": getattr(m.id, "user_id", ""),
|
||||
}
|
||||
mentions_list.append(mention_dict)
|
||||
|
||||
event_dict = {
|
||||
"sender": {
|
||||
"sender_id": {
|
||||
"open_id": sender.sender_id.open_id if sender.sender_id else "",
|
||||
"user_id": getattr(sender.sender_id, "user_id", "")
|
||||
if sender.sender_id
|
||||
else "",
|
||||
},
|
||||
"sender_type": sender.sender_type or "",
|
||||
},
|
||||
"message": {
|
||||
"chat_id": msg.chat_id or "",
|
||||
"message_type": msg.message_type or "",
|
||||
"message_id": msg.message_id or "",
|
||||
"chat_type": msg.chat_type or "",
|
||||
"content": msg.content or "{}",
|
||||
"create_time": msg.create_time or "",
|
||||
"mentions": mentions_list,
|
||||
},
|
||||
}
|
||||
|
||||
self._ws_event_queue.put(event_dict)
|
||||
except Exception:
|
||||
logger.exception("Feishu SDK message handler error")
|
||||
|
||||
async def _consume_ws_events(self) -> None:
|
||||
"""Main-loop task that drains the thread-safe event queue."""
|
||||
while self._running:
|
||||
try:
|
||||
event_dict = self._ws_event_queue.get_nowait()
|
||||
try:
|
||||
await self._on_message(event_dict)
|
||||
except Exception:
|
||||
logger.exception("Feishu WS event processing error")
|
||||
except queue.Empty:
|
||||
await asyncio.sleep(0.05)
|
||||
|
||||
async def _cleanup(self) -> None:
|
||||
await self._stop_webhook_server()
|
||||
if self.config.subscription_mode == "websocket":
|
||||
if self._ws_consumer_task:
|
||||
self._ws_consumer_task.cancel()
|
||||
try:
|
||||
await self._ws_consumer_task
|
||||
except asyncio.CancelledError:
|
||||
pass
|
||||
self._ws_consumer_task = None
|
||||
if self._http_client:
|
||||
await self._http_client.aclose()
|
||||
self._http_client = None
|
||||
# Daemon thread exits with the process; no explicit stop needed
|
||||
self._lark_ws_thread = None
|
||||
self._main_loop = None
|
||||
self._ws_event_queue = None
|
||||
else:
|
||||
await self._stop_webhook_server()
|
||||
self._access_token = None
|
||||
logger.info("Feishu channel stopped")
|
||||
|
||||
|
||||
@@ -75,6 +75,12 @@ def parse_args():
|
||||
dest="allowed_senders",
|
||||
help="Allowed sender (Feishu open_id). Can be used multiple times.",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--mode",
|
||||
choices=["webhook", "websocket"],
|
||||
default="webhook",
|
||||
help="Subscription mode: webhook (default) or websocket (long connection, no public IP needed)",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--agent",
|
||||
action="store_true",
|
||||
@@ -100,6 +106,7 @@ def main():
|
||||
webhook_port=args.webhook_port,
|
||||
feishu_domain=args.domain,
|
||||
allowed_senders=set(args.allowed_senders) if args.allowed_senders else None,
|
||||
subscription_mode=args.mode,
|
||||
)
|
||||
|
||||
send_thinking = args.thinking and args.agent
|
||||
|
||||
@@ -2140,25 +2140,76 @@ def _step_channels(config: EvoScientistConfig) -> dict[str, object]:
|
||||
raise KeyboardInterrupt()
|
||||
updates[field_name] = value.strip()
|
||||
|
||||
# Feishu optional fields (verification_token & encrypt_key)
|
||||
# Feishu: subscription mode + optional fields
|
||||
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()
|
||||
mode_choices = [
|
||||
Choice(
|
||||
title="Webhook (requires public IP / port forwarding)",
|
||||
value="webhook",
|
||||
),
|
||||
Choice(
|
||||
title="WebSocket long connection (no public IP needed)",
|
||||
value="websocket",
|
||||
),
|
||||
]
|
||||
sub_mode = questionary.select(
|
||||
"Subscription mode:",
|
||||
choices=mode_choices,
|
||||
default="webhook",
|
||||
style=WIZARD_STYLE,
|
||||
qmark=f" {QMARK}",
|
||||
use_indicator=True,
|
||||
).ask()
|
||||
if sub_mode is None:
|
||||
raise KeyboardInterrupt()
|
||||
updates["feishu_subscription_mode"] = sub_mode
|
||||
|
||||
if sub_mode == "websocket":
|
||||
# WebSocket mode needs lark-oapi SDK
|
||||
try:
|
||||
__import__("lark_oapi")
|
||||
except ImportError:
|
||||
console.print(
|
||||
' [yellow]✗ WebSocket mode requires "lark-oapi".[/yellow]'
|
||||
)
|
||||
install_sdk = questionary.confirm(
|
||||
'Install "lark-oapi>=1.4.0" now?',
|
||||
default=True,
|
||||
style=WIZARD_STYLE,
|
||||
qmark=f" {QMARK}",
|
||||
).ask()
|
||||
if install_sdk is None:
|
||||
raise KeyboardInterrupt() from None
|
||||
if install_sdk:
|
||||
console.print(' [dim]Installing "lark-oapi"...[/dim]')
|
||||
if install_pip_package("lark-oapi>=1.4.0"):
|
||||
console.print(" [green]✓ Installed successfully.[/green]")
|
||||
else:
|
||||
console.print(" [red]✗ Installation failed.[/red]")
|
||||
console.print(
|
||||
f" [dim]Run manually:[/dim] {pip_install_hint()} "
|
||||
'"lark-oapi>=1.4.0"'
|
||||
)
|
||||
else:
|
||||
# Webhook mode: prompt optional verification/encryption fields
|
||||
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"
|
||||
|
||||
@@ -130,6 +130,7 @@ class EvoScientistConfig:
|
||||
feishu_allowed_senders: str = ""
|
||||
feishu_domain: str = "https://open.feishu.cn"
|
||||
feishu_proxy: str = ""
|
||||
feishu_subscription_mode: str = "webhook" # "webhook" | "websocket"
|
||||
|
||||
# WeChat Settings
|
||||
wechat_backend: str = "wecom"
|
||||
|
||||
@@ -62,6 +62,7 @@ telegram = ["python-telegram-bot>=21.0"]
|
||||
discord = ["discord.py>=2.3"]
|
||||
slack = ["slack-sdk>=3.27", "aiohttp>=3.9"]
|
||||
wechat = ["pycryptodome>=3.20"]
|
||||
feishu = ["lark-oapi>=1.4.0"]
|
||||
qq = ["qq-botpy>=1.0"]
|
||||
stt = ["faster-whisper>=1.0"]
|
||||
oauth = ["ccproxy-api>=0.2.4"]
|
||||
@@ -71,6 +72,7 @@ all-channels = [
|
||||
"aiohttp>=3.9",
|
||||
"slack-sdk>=3.27",
|
||||
"pycryptodome>=3.20",
|
||||
"lark-oapi>=1.4.0",
|
||||
"qq-botpy>=1.0",
|
||||
]
|
||||
|
||||
|
||||
@@ -1,7 +1,8 @@
|
||||
"""Tests for Feishu channel implementation."""
|
||||
|
||||
import json
|
||||
from unittest.mock import AsyncMock, MagicMock
|
||||
import sys
|
||||
from unittest.mock import AsyncMock, MagicMock, patch
|
||||
|
||||
import pytest
|
||||
|
||||
@@ -27,6 +28,7 @@ class TestFeishuConfig:
|
||||
assert config.text_chunk_limit == 4096
|
||||
assert config.feishu_domain == "https://open.feishu.cn"
|
||||
assert config.allowed_senders is None
|
||||
assert config.subscription_mode == "webhook"
|
||||
|
||||
def test_custom_values(self):
|
||||
config = FeishuConfig(
|
||||
@@ -495,3 +497,131 @@ class TestFeishuProbe:
|
||||
ok, msg = _run(validate_feishu_credentials("id", ""))
|
||||
assert ok is False
|
||||
assert "app_secret" in msg
|
||||
|
||||
|
||||
class TestFeishuWebSocketMode:
|
||||
"""Tests for WebSocket (长连接) subscription mode."""
|
||||
|
||||
def test_config_subscription_mode_websocket(self):
|
||||
config = FeishuConfig(
|
||||
app_id="test-id",
|
||||
app_secret="test-secret",
|
||||
subscription_mode="websocket",
|
||||
)
|
||||
assert config.subscription_mode == "websocket"
|
||||
|
||||
def test_start_websocket_raises_without_lark_oapi(self):
|
||||
config = FeishuConfig(
|
||||
app_id="test-id",
|
||||
app_secret="test-secret",
|
||||
subscription_mode="websocket",
|
||||
)
|
||||
channel = FeishuChannel(config)
|
||||
# Temporarily hide lark_oapi if it's installed
|
||||
with patch.dict(sys.modules, {"lark_oapi": None}):
|
||||
with pytest.raises(ChannelError, match="lark-oapi"):
|
||||
_run(channel.start())
|
||||
|
||||
def test_start_webhook_mode_still_works(self):
|
||||
"""Ensure subscription_mode='webhook' still validates as before."""
|
||||
config = FeishuConfig(
|
||||
app_id="",
|
||||
app_secret="test-secret",
|
||||
subscription_mode="webhook",
|
||||
)
|
||||
channel = FeishuChannel(config)
|
||||
with pytest.raises(ChannelError, match="app_id"):
|
||||
_run(channel.start())
|
||||
|
||||
def test_invalid_subscription_mode_raises(self):
|
||||
config = FeishuConfig(
|
||||
app_id="test-id",
|
||||
app_secret="test-secret",
|
||||
subscription_mode="websockeet",
|
||||
)
|
||||
channel = FeishuChannel(config)
|
||||
with pytest.raises(ChannelError, match="Invalid feishu_subscription_mode"):
|
||||
_run(channel.start())
|
||||
|
||||
def test_on_lark_sdk_message_bridges_to_on_message(self):
|
||||
"""Test that _on_lark_sdk_message enqueues event dict via queue."""
|
||||
import queue as queue_mod
|
||||
|
||||
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()
|
||||
channel._ws_event_queue = queue_mod.Queue()
|
||||
|
||||
# Build a mock SDK event object matching lark_oapi structure
|
||||
mock_sender_id = MagicMock()
|
||||
mock_sender_id.open_id = "ou_test_ws"
|
||||
mock_sender_id.user_id = "user_ws"
|
||||
|
||||
mock_sender = MagicMock()
|
||||
mock_sender.sender_id = mock_sender_id
|
||||
mock_sender.sender_type = "user"
|
||||
|
||||
mock_msg = MagicMock()
|
||||
mock_msg.chat_id = "oc_ws_chat"
|
||||
mock_msg.message_type = "text"
|
||||
mock_msg.message_id = "msg_ws_1"
|
||||
mock_msg.chat_type = "p2p"
|
||||
mock_msg.content = json.dumps({"text": "hello from websocket"})
|
||||
mock_msg.create_time = "1700000000000"
|
||||
mock_msg.mentions = None
|
||||
|
||||
mock_event = MagicMock()
|
||||
mock_event.message = mock_msg
|
||||
mock_event.sender = mock_sender
|
||||
|
||||
mock_data = MagicMock()
|
||||
mock_data.event = mock_event
|
||||
|
||||
# Call the SDK callback (sync, puts on queue)
|
||||
channel._on_lark_sdk_message(mock_data)
|
||||
|
||||
# Verify event was queued
|
||||
assert not channel._ws_event_queue.empty()
|
||||
event_dict = channel._ws_event_queue.get_nowait()
|
||||
assert event_dict["sender"]["sender_id"]["open_id"] == "ou_test_ws"
|
||||
assert event_dict["message"]["chat_id"] == "oc_ws_chat"
|
||||
assert event_dict["message"]["content"] == json.dumps(
|
||||
{"text": "hello from websocket"}
|
||||
)
|
||||
|
||||
# Verify the consumer processes it correctly
|
||||
_run(channel._on_message(event_dict))
|
||||
channel._enqueue_raw.assert_called_once()
|
||||
raw = channel._enqueue_raw.call_args[0][0]
|
||||
assert raw.text == "hello from websocket"
|
||||
assert raw.sender_id == "ou_test_ws"
|
||||
assert raw.is_group is False
|
||||
|
||||
def test_cleanup_websocket_mode(self):
|
||||
config = FeishuConfig(
|
||||
app_id="test-id",
|
||||
app_secret="test-secret",
|
||||
subscription_mode="websocket",
|
||||
)
|
||||
channel = FeishuChannel(config)
|
||||
mock_client = MagicMock()
|
||||
mock_client.aclose = AsyncMock()
|
||||
channel._http_client = mock_client
|
||||
channel._lark_ws_thread = MagicMock()
|
||||
channel._main_loop = MagicMock()
|
||||
channel._ws_event_queue = MagicMock()
|
||||
channel._ws_consumer_task = None
|
||||
channel._access_token = "fake-token"
|
||||
|
||||
_run(channel._cleanup())
|
||||
|
||||
mock_client.aclose.assert_called_once()
|
||||
assert channel._http_client is None
|
||||
assert channel._lark_ws_thread is None
|
||||
assert channel._main_loop is None
|
||||
assert channel._ws_event_queue is None
|
||||
assert channel._access_token is None
|
||||
|
||||
@@ -902,6 +902,7 @@ dependencies = [
|
||||
all-channels = [
|
||||
{ name = "aiohttp" },
|
||||
{ name = "discord-py" },
|
||||
{ name = "lark-oapi" },
|
||||
{ name = "pycryptodome" },
|
||||
{ name = "python-telegram-bot" },
|
||||
{ name = "qq-botpy" },
|
||||
@@ -918,6 +919,9 @@ dev = [
|
||||
discord = [
|
||||
{ name = "discord-py" },
|
||||
]
|
||||
feishu = [
|
||||
{ name = "lark-oapi" },
|
||||
]
|
||||
oauth = [
|
||||
{ name = "ccproxy-api" },
|
||||
]
|
||||
@@ -968,6 +972,8 @@ requires-dist = [
|
||||
{ name = "langchain-openai", specifier = ">=1.1" },
|
||||
{ name = "langgraph-checkpoint-sqlite", specifier = ">=3.0.0" },
|
||||
{ name = "langgraph-cli", extras = ["inmem"], specifier = ">=0.4" },
|
||||
{ name = "lark-oapi", marker = "extra == 'all-channels'", specifier = ">=1.4.0" },
|
||||
{ name = "lark-oapi", marker = "extra == 'feishu'", specifier = ">=1.4.0" },
|
||||
{ name = "markdownify", specifier = ">=0.14" },
|
||||
{ name = "nest-asyncio", specifier = ">=1.6" },
|
||||
{ name = "pre-commit", marker = "extra == 'dev'", specifier = ">=3.5.0" },
|
||||
@@ -992,7 +998,7 @@ requires-dist = [
|
||||
{ name = "textual", specifier = ">=0.80" },
|
||||
{ name = "typer", specifier = ">=0.12" },
|
||||
]
|
||||
provides-extras = ["dev", "telegram", "discord", "slack", "wechat", "qq", "stt", "oauth", "all-channels"]
|
||||
provides-extras = ["dev", "telegram", "discord", "slack", "wechat", "feishu", "qq", "stt", "oauth", "all-channels"]
|
||||
|
||||
[package.metadata.requires-dev]
|
||||
dev = [
|
||||
@@ -2183,6 +2189,21 @@ otel = [
|
||||
{ name = "opentelemetry-sdk" },
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "lark-oapi"
|
||||
version = "1.5.3"
|
||||
source = { registry = "https://pypi.org/simple" }
|
||||
dependencies = [
|
||||
{ name = "httpx" },
|
||||
{ name = "pycryptodome" },
|
||||
{ name = "requests" },
|
||||
{ name = "requests-toolbelt" },
|
||||
{ name = "websockets" },
|
||||
]
|
||||
wheels = [
|
||||
{ url = "https://files.pythonhosted.org/packages/bf/ff/2ece5d735ebfa2af600a53176f2636ae47af2bf934e08effab64f0d1e047/lark_oapi-1.5.3-py3-none-any.whl", hash = "sha256:fda6b32bb38d21b6bdaae94979c600b94c7c521e985adade63a54e4b3e20cc36", size = 6993016, upload-time = "2026-01-27T08:21:49.307Z" },
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "linkify-it-py"
|
||||
version = "2.1.0"
|
||||
|
||||
Reference in New Issue
Block a user