From 65cec644452ec60f74410afbf108d4b5a6c4cb94 Mon Sep 17 00:00:00 2001 From: Ziheng Zhang <142805986+MuXinCG2004@users.noreply.github.com> Date: Sun, 22 Mar 2026 22:45:28 +0800 Subject: [PATCH] feat(feishu): add WebSocket long connection subscription mode (#87) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * 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 --- EvoScientist/channels/README.md | 40 ++++- EvoScientist/channels/feishu/__init__.py | 1 + EvoScientist/channels/feishu/channel.py | 200 +++++++++++++++++++++-- EvoScientist/channels/feishu/serve.py | 7 + EvoScientist/config/onboard.py | 87 ++++++++-- EvoScientist/config/settings.py | 1 + pyproject.toml | 2 + tests/test_feishu_channel.py | 132 ++++++++++++++- uv.lock | 23 ++- 9 files changed, 457 insertions(+), 36 deletions(-) diff --git a/EvoScientist/channels/README.md b/EvoScientist/channels/README.md index cb2ebb4..09b7e1c 100644 --- a/EvoScientist/channels/README.md +++ b/EvoScientist/channels/README.md @@ -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 diff --git a/EvoScientist/channels/feishu/__init__.py b/EvoScientist/channels/feishu/__init__.py index 5b8c553..55c6ed4 100644 --- a/EvoScientist/channels/feishu/__init__.py +++ b/EvoScientist/channels/feishu/__init__.py @@ -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"), ) ) diff --git a/EvoScientist/channels/feishu/channel.py b/EvoScientist/channels/feishu/channel.py index f5339ce..cbc0320 100644 --- a/EvoScientist/channels/feishu/channel.py +++ b/EvoScientist/channels/feishu/channel.py @@ -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") diff --git a/EvoScientist/channels/feishu/serve.py b/EvoScientist/channels/feishu/serve.py index 945e90e..4888c21 100644 --- a/EvoScientist/channels/feishu/serve.py +++ b/EvoScientist/channels/feishu/serve.py @@ -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 diff --git a/EvoScientist/config/onboard.py b/EvoScientist/config/onboard.py index 2d677d1..842dc71 100644 --- a/EvoScientist/config/onboard.py +++ b/EvoScientist/config/onboard.py @@ -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" diff --git a/EvoScientist/config/settings.py b/EvoScientist/config/settings.py index de31d97..b612242 100644 --- a/EvoScientist/config/settings.py +++ b/EvoScientist/config/settings.py @@ -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" diff --git a/pyproject.toml b/pyproject.toml index 7288064..396c29b 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -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", ] diff --git a/tests/test_feishu_channel.py b/tests/test_feishu_channel.py index f322dd5..623bc06 100644 --- a/tests/test_feishu_channel.py +++ b/tests/test_feishu_channel.py @@ -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 diff --git a/uv.lock b/uv.lock index 6d60e08..492a6d3 100644 --- a/uv.lock +++ b/uv.lock @@ -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"