ソースを参照

feat: add Feishu alerts for business messages

AI-Co-Authored-By: Codex
chendeben 1 ヶ月 前
親
コミット
5b11806623

+ 124 - 1
admin-web/src/pages/BusinessAssistant.tsx

@@ -23,6 +23,7 @@ import {
 import type { ColumnsType } from 'antd/es/table';
 import dayjs from 'dayjs';
 import {
+  BellRing,
   BotMessageSquare,
   CircleCheck,
   CirclePlus,
@@ -113,9 +114,11 @@ export default function BusinessAssistantPage() {
   const [detailLoading, setDetailLoading] = useState(false);
   const [testQuestion, setTestQuestion] = useState('');
   const [testResult, setTestResult] = useState<Record<string, unknown> | null>(null);
+  const [feishuTesting, setFeishuTesting] = useState(false);
   const [providerForm] = Form.useForm<ProviderForm>();
   const [settingsForm] = Form.useForm<AssistantAccountSettings>();
   const [knowledgeForm] = Form.useForm<KnowledgeForm>();
+  const feishuWebhookDraft = Form.useWatch('feishu_webhook_url', settingsForm);
 
   const provider = bot?.business_assistant;
 
@@ -237,9 +240,14 @@ export default function BusinessAssistantPage() {
       okText: '确认保存',
       cancelText: '取消',
       onOk: async () => {
+        const payload = { ...values };
+        if (!payload.feishu_webhook_url?.trim()) delete payload.feishu_webhook_url;
+        if (!payload.feishu_webhook_signing_secret?.trim()) {
+          delete payload.feishu_webhook_signing_secret;
+        }
         const saved = await apiRequest<AssistantAccountSettings>(
           '/business-assistant/settings',
-          jsonOptions('PUT', { ...values, connection_id: selectedConnection, confirm: true }),
+          jsonOptions('PUT', { ...payload, connection_id: selectedConnection, confirm: true }),
         );
         setSettings(saved);
         settingsForm.setFieldsValue(saved);
@@ -249,6 +257,41 @@ export default function BusinessAssistantPage() {
     });
   };
 
+  const testFeishuWebhook = () => {
+    const values = settingsForm.getFieldsValue([
+      'feishu_webhook_url',
+      'feishu_webhook_signing_secret',
+    ]);
+    const payload: Record<string, unknown> = {
+      connection_id: selectedConnection,
+      confirm: true,
+    };
+    if (values.feishu_webhook_url?.trim()) {
+      payload.feishu_webhook_url = values.feishu_webhook_url.trim();
+    }
+    if (values.feishu_webhook_signing_secret?.trim()) {
+      payload.feishu_webhook_signing_secret = values.feishu_webhook_signing_secret.trim();
+    }
+    Modal.confirm({
+      title: '发送飞书测试提醒?',
+      content: '系统会向当前填写或已保存的群机器人 Webhook 发送一条测试消息。',
+      okText: '发送测试',
+      cancelText: '取消',
+      onOk: async () => {
+        setFeishuTesting(true);
+        try {
+          await apiRequest(
+            '/business-assistant/feishu-webhook/test',
+            jsonOptions('POST', payload),
+          );
+          message.success('飞书测试提醒发送成功');
+        } finally {
+          setFeishuTesting(false);
+        }
+      },
+    });
+  };
+
   const openKnowledge = (entry?: AssistantKnowledgeEntry) => {
     setKnowledgeTarget(entry || null);
     knowledgeForm.setFieldsValue(
@@ -633,6 +676,86 @@ export default function BusinessAssistantPage() {
                         <Col xs={24} md={8}>
                           <Form.Item name="digest_time" label="简报时间" extra="24 小时制 HH:MM"><Input placeholder="09:00" /></Form.Item>
                         </Col>
+                        <Col xs={24}>
+                          <Typography.Title level={5} style={{ marginTop: 8 }}>飞书新消息提醒</Typography.Title>
+                          <Alert
+                            type="info"
+                            showIcon
+                            message="第一期按每条客户消息提醒"
+                            description="启用后,每一条新的客户 Business 消息都会推送到飞书群;账号本人回复、机器人消息及 Telegram 原生欢迎/离线自动消息不会触发。"
+                            style={{ marginBottom: 16 }}
+                          />
+                        </Col>
+                        <Col xs={24} md={8}>
+                          <Form.Item name="feishu_webhook_enabled" label="飞书群提醒" valuePropName="checked">
+                            <Switch checkedChildren="开启" unCheckedChildren="关闭" />
+                          </Form.Item>
+                        </Col>
+                        <Col xs={24} md={8}>
+                          <Form.Item name="feishu_message_preview_enabled" label="包含消息内容" valuePropName="checked" extra="关闭时只发送账号、客户、消息类型和标识。">
+                            <Switch checkedChildren="包含" unCheckedChildren="隐藏" />
+                          </Form.Item>
+                        </Col>
+                        <Col xs={24} md={8}>
+                          <Form.Item label="配置状态">
+                            <Space wrap>
+                              <Tag color={settings.feishu_webhook_configured ? 'success' : 'warning'}>
+                                {settings.feishu_webhook_configured ? 'Webhook 已配置' : 'Webhook 未配置'}
+                              </Tag>
+                              {settings.feishu_signing_secret_configured ? <Tag color="success">签名已配置</Tag> : null}
+                            </Space>
+                          </Form.Item>
+                        </Col>
+                        <Col xs={24} lg={14}>
+                          <Form.Item
+                            name="feishu_webhook_url"
+                            label="飞书群机器人 Webhook"
+                            extra={settings.feishu_webhook_configured ? '已保存的地址不会回显;留空会保留,填写新地址会替换。' : '仅支持飞书或 Lark 官方群机器人 Webhook。'}
+                            dependencies={['feishu_webhook_enabled']}
+                            rules={[
+                              ({ getFieldValue }) => ({
+                                validator: (_, value) => {
+                                  if (getFieldValue('feishu_webhook_enabled') && !value && !settings.feishu_webhook_configured) {
+                                    return Promise.reject(new Error('开启提醒前请填写 Webhook 地址'));
+                                  }
+                                  return Promise.resolve();
+                                },
+                              }),
+                            ]}
+                          >
+                            <Input.Password
+                              autoComplete="new-password"
+                              placeholder="https://open.feishu.cn/open-apis/bot/v2/hook/..."
+                              style={{ minHeight: 44 }}
+                            />
+                          </Form.Item>
+                        </Col>
+                        <Col xs={24} lg={10}>
+                          <Form.Item
+                            name="feishu_webhook_signing_secret"
+                            label="签名密钥(可选)"
+                            extra={settings.feishu_signing_secret_configured ? '签名已配置;留空保留,填写新密钥会替换。' : '仅在飞书群机器人启用了“签名校验”时填写。'}
+                          >
+                            <Input.Password
+                              autoComplete="new-password"
+                              placeholder="飞书群机器人签名密钥"
+                              style={{ minHeight: 44 }}
+                            />
+                          </Form.Item>
+                        </Col>
+                        <Col xs={24}>
+                          <Form.Item>
+                            <Button
+                              icon={<BellRing size={16} />}
+                              loading={feishuTesting}
+                              disabled={!settings.feishu_webhook_configured && !feishuWebhookDraft?.trim()}
+                              onClick={testFeishuWebhook}
+                              style={{ minHeight: 44 }}
+                            >
+                              发送测试提醒
+                            </Button>
+                          </Form.Item>
+                        </Col>
                         <Col xs={24}>
                           <Form.Item name="handoff_message" label="转人工提示"><Input.TextArea rows={2} maxLength={1000} /></Form.Item>
                         </Col>

+ 6 - 0
admin-web/src/types.ts

@@ -99,6 +99,12 @@ export interface AssistantAccountSettings {
   digest_time: string;
   handoff_message: string;
   unsupported_message: string;
+  feishu_webhook_enabled: boolean;
+  feishu_webhook_url?: string;
+  feishu_webhook_signing_secret?: string;
+  feishu_webhook_configured: boolean;
+  feishu_signing_secret_configured: boolean;
+  feishu_message_preview_enabled: boolean;
 }
 
 export interface AssistantConnection {

+ 29 - 0
admin-web/tests/unit/BusinessAssistant.test.tsx

@@ -58,6 +58,12 @@ const settings = {
   digest_time: '09:00',
   handoff_message: '请稍候,已转人工。',
   unsupported_message: '该消息需要人工处理。',
+  feishu_webhook_enabled: true,
+  feishu_webhook_url: '',
+  feishu_webhook_signing_secret: '',
+  feishu_webhook_configured: true,
+  feishu_signing_secret_configured: true,
+  feishu_message_preview_enabled: false,
 };
 
 beforeEach(() => {
@@ -92,6 +98,9 @@ beforeEach(() => {
       };
     }
     if (path.startsWith('/business-assistant/settings')) return settings;
+    if (path === '/business-assistant/feishu-webhook/test') {
+      return { ok: true, account: '店主' };
+    }
     if (path.startsWith('/business-assistant/knowledge-sources')) {
       return {
         items: [
@@ -231,3 +240,23 @@ test('展示多源知识与待审核候选', async () => {
   expect(screen.getByText('仅支持市区配送。')).toBeInTheDocument();
   expect(screen.getByRole('button', { name: '批准' })).toBeInTheDocument();
 });
+
+test('展示脱敏的飞书提醒配置并可发送测试消息', async () => {
+  render(<BusinessAssistantPage />);
+  await screen.findAllByText('店主');
+
+  expect(screen.getByText('第一期按每条客户消息提醒')).toBeInTheDocument();
+  expect(screen.getByText('Webhook 已配置')).toBeInTheDocument();
+  expect(screen.getByText('签名已配置')).toBeInTheDocument();
+
+  fireEvent.click(screen.getByRole('button', { name: '发送测试提醒' }));
+  fireEvent.click(await screen.findByRole('button', { name: '发送测试' }));
+
+  await waitFor(() => {
+    expect(
+      vi.mocked(apiRequest).mock.calls.some(
+        ([path]) => path === '/business-assistant/feishu-webhook/test',
+      ),
+    ).toBe(true);
+  });
+});

+ 206 - 2
tests/test_business_assistant.py

@@ -115,6 +115,20 @@ class FakeBusinessApi:
         self.read.append((connection_id, chat_id, message_id))
 
 
+class FakeFeishuWebhook:
+    def __init__(self) -> None:
+        self.sent: list[dict] = []
+
+    async def send(self, webhook_url: str, text: str, *, signing_secret: str = "") -> None:
+        self.sent.append(
+            {
+                "webhook_url": webhook_url,
+                "text": text,
+                "signing_secret": signing_secret,
+            }
+        )
+
+
 async def prepare_runtime(app_modules, *, provider: FakeProvider | None = None):
     dbassistant = app_modules.load("wbb.utils.dbassistant")
     service = app_modules.load("wbb.services.business_assistant")
@@ -218,6 +232,41 @@ async def test_connection_settings_knowledge_quota_and_bot_scope(app_modules):
     assert error.value.code == "invalid_setting"
 
 
+async def test_feishu_settings_are_validated_and_never_exposed(app_modules):
+    dbassistant = app_modules.load("wbb.utils.dbassistant")
+    await dbassistant.upsert_business_connection(connection_payload())
+    webhook_url = "https://open.feishu.cn/open-apis/bot/v2/hook/12345678-abcd-4321-abcd-123456789abc"
+    stored = await dbassistant.update_account_settings(
+        "conn-a",
+        {
+            "feishu_webhook_enabled": True,
+            "feishu_webhook_url": webhook_url,
+            "feishu_webhook_signing_secret": "signing-secret",
+            "feishu_message_preview_enabled": True,
+        },
+    )
+    assert stored["feishu_webhook_url"] == webhook_url
+    assert stored["feishu_webhook_signing_secret"] == "signing-secret"
+
+    public = dbassistant.public_account_settings(stored)
+    assert public["feishu_webhook_url"] == ""
+    assert public["feishu_webhook_signing_secret"] == ""
+    assert public["feishu_webhook_configured"] is True
+    assert public["feishu_signing_secret_configured"] is True
+
+    connections, _total = await dbassistant.list_business_connections()
+    assert connections[0]["settings"]["feishu_webhook_url"] == ""
+    assert connections[0]["settings"]["feishu_webhook_configured"] is True
+
+    with pytest.raises(dbassistant.AssistantDataError):
+        dbassistant.normalize_account_settings(
+            {
+                "feishu_webhook_enabled": True,
+                "feishu_webhook_url": "https://example.com/internal-hook",
+            }
+        )
+
+
 def test_ai_contract_and_reply_window_validation(app_modules):
     service = app_modules.load("wbb.services.business_assistant")
     parsed = service.parse_ai_decision(
@@ -550,6 +599,93 @@ async def test_business_updates_are_idempotent_and_manual_reply_pauses(app_modul
     assert len(runtime.api.sent) == 1
 
 
+async def test_every_customer_message_notifies_feishu_once(app_modules):
+    dbassistant, _service, runtime, _knowledge = await prepare_runtime(app_modules)
+    webhook_url = "https://open.feishu.cn/open-apis/bot/v2/hook/12345678-abcd-4321-abcd-123456789abc"
+    await dbassistant.update_account_settings(
+        "conn-a",
+        {
+            "assistant_enabled": False,
+            "feishu_webhook_enabled": True,
+            "feishu_webhook_url": webhook_url,
+            "feishu_webhook_signing_secret": "secret",
+            "feishu_message_preview_enabled": False,
+        },
+    )
+    feishu = FakeFeishuWebhook()
+    runtime.feishu = feishu
+
+    await runtime.process_update(
+        {"update_id": 101, "business_message": customer_message(message_id=1)}
+    )
+    await runtime.process_update(
+        {"update_id": 102, "business_message": customer_message(message_id=1)}
+    )
+    assert len(feishu.sent) == 1
+    assert feishu.sent[0]["webhook_url"] == webhook_url
+    assert feishu.sent[0]["signing_secret"] == "secret"
+    assert "营业时间是什么" not in feishu.sent[0]["text"]
+
+    await dbassistant.update_account_settings(
+        "conn-a", {"feishu_message_preview_enabled": True}
+    )
+    await runtime.process_update(
+        {
+            "update_id": 103,
+            "business_message": customer_message(
+                message_id=2, text="第二条客户消息"
+            ),
+        }
+    )
+    assert len(feishu.sent) == 2
+    assert "第二条客户消息" in feishu.sent[1]["text"]
+
+    await runtime.process_update(
+        {
+            "update_id": 1030,
+            "business_message": customer_message(
+                message_id=20,
+                text=None,
+                caption="图片说明",
+                photo=[{"file_id": "photo"}],
+            ),
+        }
+    )
+    assert len(feishu.sent) == 3
+    assert "类型:非文本" in feishu.sent[2]["text"]
+    assert "内容:图片说明" in feishu.sent[2]["text"]
+
+    await runtime.process_update(
+        {
+            "update_id": 104,
+            "business_message": customer_message(
+                message_id=3, text="账号本人回复", sender_id=900
+            ),
+        }
+    )
+    await runtime.process_update(
+        {
+            "update_id": 105,
+            "business_message": customer_message(message_id=4, is_from_offline=True),
+        }
+    )
+    await runtime.process_update(
+        {
+            "update_id": 106,
+            "business_message": customer_message(
+                message_id=5, sender_business_bot={"id": 999}
+            ),
+        }
+    )
+    assert len(feishu.sent) == 3
+    conversation = await dbassistant.get_conversation_by_chat("conn-a", 501)
+    messages = await dbassistant.recent_conversation_messages(
+        conversation["conversation_id"], limit=10
+    )
+    incoming = [item for item in messages if item["telegram_message_id"] == 2][0]
+    assert incoming["metadata"]["feishu_notification_status"] == "sent"
+
+
 async def test_handoff_sensitive_non_text_provider_error_and_expired_window(app_modules):
     dbassistant, service, runtime, _knowledge = await prepare_runtime(app_modules)
     await dbassistant.update_account_settings(
@@ -726,6 +862,7 @@ class FakeResponse:
     def __init__(self, status: int, payload: dict) -> None:
         self.status = status
         self.payload = payload
+        self.headers: dict[str, str] = {}
 
     async def __aenter__(self):
         return self
@@ -780,6 +917,32 @@ async def test_telegram_api_retries_429_and_5xx(app_modules, monkeypatch):
     assert sleeps == [3, 2]
 
 
+async def test_feishu_webhook_signature_and_success_contract(app_modules):
+    service = app_modules.load("wbb.services.business_assistant")
+    response = FakeResponse(200, {"code": 0, "msg": "success"})
+
+    class CaptureSession:
+        def __init__(self) -> None:
+            self.payload = None
+
+        def post(self, _url, *, json, timeout):
+            self.payload = json
+            assert timeout.total == 8
+            return response
+
+    session = CaptureSession()
+    client = service.FeishuWebhookClient(session)
+    await client.send(
+        "https://open.feishu.cn/open-apis/bot/v2/hook/12345678-abcd-4321-abcd-123456789abc",
+        "Telegram 新消息提醒",
+        signing_secret="secret",
+    )
+    assert session.payload["msg_type"] == "text"
+    assert session.payload["content"]["text"] == "Telegram 新消息提醒"
+    assert session.payload["timestamp"]
+    assert session.payload["sign"]
+
+
 async def _login_and_change_password(client: TestClient) -> str:
     login = await client.post(
         "/api/admin/v1/auth/login",
@@ -838,10 +1001,51 @@ async def test_business_assistant_api_permission_csrf_audit_and_clear(app_module
         saved = await client.put(
             "/api/admin/v1/business-assistant/settings",
             headers={"X-CSRF-Token": csrf},
-            json={"connection_id": "conn-a", "assistant_enabled": True, "confirm": True},
+            json={
+                "connection_id": "conn-a",
+                "assistant_enabled": True,
+                "feishu_webhook_enabled": True,
+                "feishu_webhook_url": "https://open.feishu.cn/open-apis/bot/v2/hook/12345678-abcd-4321-abcd-123456789abc",
+                "feishu_webhook_signing_secret": "api-secret",
+                "confirm": True,
+            },
         )
         assert saved.status == 200
-        assert (await saved.json())["data"]["assistant_enabled"] is True
+        saved_data = (await saved.json())["data"]
+        assert saved_data["assistant_enabled"] is True
+        assert saved_data["feishu_webhook_url"] == ""
+        assert saved_data["feishu_webhook_signing_secret"] == ""
+        assert saved_data["feishu_webhook_configured"] is True
+        stored_settings = await dbassistant.get_account_settings("conn-a")
+        assert stored_settings["feishu_webhook_signing_secret"] == "api-secret"
+        listed_connections = await client.get(
+            "/api/admin/v1/business-assistant/connections?page=1&page_size=100"
+        )
+        listed_settings = (await listed_connections.json())["data"]["items"][0][
+            "settings"
+        ]
+        assert listed_settings["feishu_webhook_url"] == ""
+        assert listed_settings["feishu_webhook_signing_secret"] == ""
+
+        class FakeRuntime:
+            def __init__(self) -> None:
+                self.calls = []
+
+            async def test_feishu_webhook(self, connection_id, overrides):
+                self.calls.append((connection_id, overrides))
+                return {"ok": True, "account": "店主"}
+
+        module = app_modules.load("wbb.modules.business_assistant")
+        fake_runtime = FakeRuntime()
+        module._runtime = fake_runtime
+        tested = await client.post(
+            "/api/admin/v1/business-assistant/feishu-webhook/test",
+            headers={"X-CSRF-Token": csrf},
+            json={"connection_id": "conn-a", "confirm": True},
+        )
+        assert tested.status == 200
+        assert (await tested.json())["data"]["ok"] is True
+        assert fake_runtime.calls == [("conn-a", {})]
 
         created = await client.post(
             "/api/admin/v1/business-assistant/knowledge",

+ 39 - 2
wbb/admin/api.py

@@ -44,6 +44,7 @@ from wbb.admin.telegram_webapp import (
 from wbb.services.bot_permissions import api_permission, has_permission
 from wbb.services.business_assistant import (
     AssistantProviderError,
+    FeishuWebhookError,
     runtime_overview,
 )
 from wbb.services.chat_management import (
@@ -112,6 +113,7 @@ from wbb.utils.dbassistant import (
     list_knowledge_entries,
     list_knowledge_sources,
     pause_conversation,
+    public_account_settings,
     publish_knowledge_candidate,
     reject_knowledge_candidate,
     resume_conversation,
@@ -735,6 +737,10 @@ class AdminApi:
             f"{API_PREFIX}/business-assistant/settings",
             self.business_assistant_settings_update,
         )
+        router.add_post(
+            f"{API_PREFIX}/business-assistant/feishu-webhook/test",
+            self.business_assistant_feishu_webhook_test,
+        )
         router.add_get(
             f"{API_PREFIX}/business-assistant/knowledge",
             self.business_assistant_knowledge,
@@ -2430,7 +2436,7 @@ class AdminApi:
             raise ApiProblem("connection_required", "请先选择一个 Business 连接。")
         if not await get_business_connection(connection_id):
             raise AssistantDataError("connection_not_found", "未找到 Business 连接。")
-        return success(await get_account_settings(connection_id))
+        return success(public_account_settings(await get_account_settings(connection_id)))
 
     async def business_assistant_settings_update(
         self, request: web.Request
@@ -2447,7 +2453,38 @@ class AdminApi:
             target_id=connection_id,
             summary="更新智能接待设置",
         )
-        return success(await update_account_settings(connection_id, body))
+        saved = await update_account_settings(connection_id, body)
+        return success(public_account_settings(saved))
+
+    async def business_assistant_feishu_webhook_test(
+        self, request: web.Request
+    ) -> web.Response:
+        body = await json_body(request)
+        require_confirmation(body)
+        connection_id = str(body.pop("connection_id", "")).strip()
+        body.pop("confirm", None)
+        if not connection_id:
+            raise ApiProblem("connection_required", "请先选择一个 Business 连接。")
+        allowed = {"feishu_webhook_url", "feishu_webhook_signing_secret"}
+        overrides = {key: value for key, value in body.items() if key in allowed}
+        from wbb.modules.business_assistant import get_business_assistant_runtime
+
+        runtime = get_business_assistant_runtime()
+        if runtime is None:
+            raise ApiProblem(
+                "assistant_runtime_unavailable", "智能接待运行时尚未启动。", status=503
+            )
+        try:
+            result = await runtime.test_feishu_webhook(connection_id, overrides)
+        except FeishuWebhookError as exc:
+            raise ApiProblem("feishu_webhook_error", str(exc), status=502) from exc
+        set_audit(
+            request,
+            "business_assistant.feishu_webhook.test",
+            target_id=connection_id,
+            summary="测试飞书群机器人 Webhook",
+        )
+        return success(result)
 
     async def business_assistant_knowledge(self, request: web.Request) -> web.Response:
         connection_id = str(request.query.get("connection_id") or "").strip()

+ 197 - 1
wbb/services/business_assistant.py

@@ -1,17 +1,21 @@
 from __future__ import annotations
 
 import asyncio
+import base64
+import hashlib
+import hmac
 import json
 import re
 from contextlib import suppress
 from datetime import UTC, datetime
 from typing import Any
 
-from aiohttp import ClientError, ClientSession
+from aiohttp import ClientError, ClientSession, ClientTimeout
 
 import wbb
 from wbb.utils.dbassistant import (
     append_conversation_message,
+    claim_feishu_message_notification,
     claim_update,
     classify_handoff,
     dead_letter_update,
@@ -31,10 +35,12 @@ from wbb.utils.dbassistant import (
     list_enabled_business_sources,
     load_update_offset,
     mark_digest_sent,
+    mark_feishu_message_notification,
     mark_source_event_deleted,
     mark_update_done,
     mark_update_failed,
     match_knowledge,
+    normalize_account_settings,
     pause_conversation_for_human,
     recent_conversation_messages,
     record_source_event,
@@ -79,6 +85,74 @@ class AssistantProviderError(RuntimeError):
     pass
 
 
+class FeishuWebhookError(RuntimeError):
+    pass
+
+
+class FeishuWebhookClient:
+    def __init__(self, session: ClientSession) -> None:
+        self._session = session
+
+    @staticmethod
+    def _signature(secret: str, timestamp: int) -> str:
+        sign_key = f"{timestamp}\n{secret}".encode()
+        digest = hmac.new(sign_key, digestmod=hashlib.sha256).digest()
+        return base64.b64encode(digest).decode()
+
+    async def send(self, webhook_url: str, text: str, *, signing_secret: str = "") -> None:
+        payload: dict[str, Any] = {
+            "msg_type": "text",
+            "content": {"text": str(text)[:4000]},
+        }
+        if signing_secret:
+            timestamp = int(utc_now().timestamp())
+            payload.update(
+                {
+                    "timestamp": str(timestamp),
+                    "sign": self._signature(signing_secret, timestamp),
+                }
+            )
+        last_error: FeishuWebhookError | None = None
+        for attempt in range(3):
+            try:
+                async with self._session.post(
+                    webhook_url,
+                    json=payload,
+                    timeout=ClientTimeout(total=8),
+                ) as response:
+                    data = await response.json(content_type=None)
+            except (ClientError, TimeoutError, ValueError) as exc:
+                last_error = FeishuWebhookError("无法连接飞书群机器人 Webhook。")
+                if attempt < 2:
+                    await asyncio.sleep(2**attempt)
+                    continue
+                raise last_error from exc
+            raw_code = (
+                data.get("code", data.get("StatusCode", -1))
+                if isinstance(data, dict)
+                else -1
+            )
+            try:
+                code = int(raw_code)
+            except (TypeError, ValueError):
+                code = -1
+            if response.status == 200 and code == 0:
+                return
+            last_error = FeishuWebhookError(
+                f"飞书群机器人 Webhook 请求失败(HTTP {response.status},code={code})。"
+            )
+            if (response.status == 429 or response.status >= 500) and attempt < 2:
+                retry_after = response.headers.get("Retry-After", "")
+                try:
+                    wait_seconds = max(1, min(int(retry_after), 30))
+                except (TypeError, ValueError):
+                    wait_seconds = 2**attempt
+                await asyncio.sleep(wait_seconds)
+                continue
+            raise last_error
+        raise last_error or FeishuWebhookError("飞书群机器人 Webhook 请求失败。")
+
+
 class TelegramBusinessApi:
     def __init__(self, token: str, session: ClientSession) -> None:
         self._base_url = f"https://api.telegram.org/bot{token}"
@@ -588,6 +662,7 @@ class BusinessAssistantRuntime:
         provider: OpenAICompatibleAssistant | None = None,
     ) -> None:
         self.api = TelegramBusinessApi(token, session)
+        self.feishu = FeishuWebhookClient(session)
         self.provider = provider or provider_from_wbb(session)
         self.ingestion = KnowledgeIngestionService(self.provider)
         self._poll_task: asyncio.Task[None] | None = None
@@ -888,6 +963,14 @@ class BusinessAssistantRuntime:
             sender_id=int(sender.get("id") or 0),
             metadata={"content_type": "text" if text else "unsupported"},
         )
+        await self._notify_feishu_customer_message(
+            connection=connection,
+            conversation=conversation,
+            message=message,
+            sender=sender,
+            settings=settings,
+            text=text,
+        )
         if not connection.get("is_enabled") or not settings.get("assistant_enabled"):
             return
         current = await get_conversation(conversation["conversation_id"]) or conversation
@@ -992,6 +1075,119 @@ class BusinessAssistantRuntime:
             },
         )
 
+    async def _notify_feishu_customer_message(
+        self,
+        *,
+        connection: dict[str, Any],
+        conversation: dict[str, Any],
+        message: dict[str, Any],
+        sender: dict[str, Any],
+        settings: dict[str, Any],
+        text: str,
+    ) -> None:
+        webhook_url = str(settings.get("feishu_webhook_url") or "")
+        if not settings.get("feishu_webhook_enabled") or not webhook_url:
+            return
+        conversation_id = str(conversation["conversation_id"])
+        message_id = int(message.get("message_id") or 0)
+        if not await claim_feishu_message_notification(conversation_id, message_id):
+            return
+        customer_name = " ".join(
+            str(item).strip()
+            for item in (sender.get("first_name"), sender.get("last_name"))
+            if str(item or "").strip()
+        )
+        if not customer_name:
+            customer_name = (
+                f"@{sender.get('username')}"
+                if sender.get("username")
+                else str((message.get("chat") or {}).get("id") or "未知")
+            )
+        account = connection.get("user") or {}
+        account_name = " ".join(
+            str(item).strip()
+            for item in (account.get("first_name"), account.get("last_name"))
+            if str(item or "").strip()
+        ) or (f"@{account.get('username')}" if account.get("username") else "Business 账号")
+        message_type = "文本" if text else "非文本"
+        preview_text = str(
+            message.get("text") or message.get("caption") or ""
+        ).strip()
+        lines = [
+            "Telegram 新消息提醒",
+            f"账号:{account_name}",
+            f"客户:{customer_name}",
+            f"类型:{message_type}",
+        ]
+        if settings.get("feishu_message_preview_enabled"):
+            lines.append(f"内容:{preview_text[:500] if preview_text else '[非文本消息]'}")
+        try:
+            message_time = datetime.fromtimestamp(
+                int(message.get("date") or 0), UTC
+            ).strftime("%Y-%m-%d %H:%M:%S UTC")
+        except (OSError, OverflowError, TypeError, ValueError):
+            message_time = "未知"
+        lines.extend(
+            [
+                f"时间:{message_time}",
+                f"消息标识:chat_id={(message.get('chat') or {}).get('id')} / message_id={message_id}",
+            ]
+        )
+        try:
+            await self.feishu.send(
+                webhook_url,
+                "\n".join(lines),
+                signing_secret=str(settings.get("feishu_webhook_signing_secret") or ""),
+            )
+        except FeishuWebhookError as exc:
+            await mark_feishu_message_notification(
+                conversation_id, message_id, status="failed", error=str(exc)
+            )
+            await update_runtime_status(
+                {
+                    "feishu_notification_last_error": str(exc)[:1000],
+                    "feishu_notification_last_error_at": utc_now(),
+                }
+            )
+            return
+        await mark_feishu_message_notification(
+            conversation_id, message_id, status="sent"
+        )
+        await update_runtime_status(
+            {
+                "feishu_notification_last_error": "",
+                "feishu_notification_last_sent_at": utc_now(),
+            }
+        )
+
+    async def test_feishu_webhook(
+        self, connection_id: str, overrides: dict[str, Any] | None = None
+    ) -> dict[str, Any]:
+        connection = await self._resolve_connection(connection_id)
+        current = await get_account_settings(connection_id)
+        settings = normalize_account_settings(overrides or {}, previous=current)
+        webhook_url = str(settings.get("feishu_webhook_url") or "")
+        if not webhook_url:
+            raise FeishuWebhookError("请先填写飞书群机器人 Webhook 地址。")
+        account = connection.get("user") or {}
+        account_name = " ".join(
+            str(item).strip()
+            for item in (account.get("first_name"), account.get("last_name"))
+            if str(item or "").strip()
+        ) or (f"@{account.get('username')}" if account.get("username") else "Business 账号")
+        await self.feishu.send(
+            webhook_url,
+            "\n".join(
+                [
+                    "Telegram 新消息提醒",
+                    "飞书群机器人连接测试成功。",
+                    f"账号:{account_name}",
+                ]
+            ),
+            signing_secret=str(settings.get("feishu_webhook_signing_secret") or ""),
+        )
+        return {"ok": True, "account": account_name}
+
     async def _handoff(
         self,
         connection: dict[str, Any],

+ 136 - 1
wbb/utils/dbassistant.py

@@ -6,6 +6,7 @@ import json
 import re
 from datetime import UTC, datetime, timedelta
 from typing import Any
+from urllib.parse import urlsplit
 from uuid import uuid4
 from zoneinfo import ZoneInfo, ZoneInfoNotFoundError
 
@@ -44,6 +45,15 @@ DEFAULT_ACCOUNT_SETTINGS: dict[str, Any] = {
     "digest_time": "09:00",
     "handoff_message": "这个问题需要人工确认,我已经通知负责人,请稍候。",
     "unsupported_message": "已收到你的消息,这类内容需要人工处理,我已经通知负责人。",
+    "feishu_webhook_enabled": False,
+    "feishu_webhook_url": "",
+    "feishu_webhook_signing_secret": "",
+    "feishu_message_preview_enabled": False,
+}
+
+SECRET_ACCOUNT_SETTINGS = {
+    "feishu_webhook_url",
+    "feishu_webhook_signing_secret",
 }
 
 VALID_NOTIFICATION_DESTINATIONS = {"owner", "ops", "both"}
@@ -135,6 +145,33 @@ def _bounded_int(
     return parsed
 
 
+def _normalize_feishu_webhook_url(value: Any) -> str:
+    url = clean_text(value, max_length=1000)
+    if not url:
+        return ""
+    try:
+        parsed = urlsplit(url)
+        port = parsed.port
+    except ValueError as exc:
+        raise AssistantDataError("invalid_setting", "飞书 Webhook 地址无效。") from exc
+    allowed_hosts = {"open.feishu.cn", "open.larksuite.com"}
+    if (
+        parsed.scheme != "https"
+        or (parsed.hostname or "").casefold() not in allowed_hosts
+        or port not in {None, 443}
+        or parsed.username
+        or parsed.password
+        or parsed.query
+        or parsed.fragment
+        or not re.fullmatch(r"/open-apis/bot/v2/hook/[A-Za-z0-9-]{8,128}", parsed.path)
+    ):
+        raise AssistantDataError(
+            "invalid_setting",
+            "仅支持飞书或 Lark 官方 HTTPS 群机器人 Webhook 地址。",
+        )
+    return url
+
+
 def _normalize_string_list(value: Any, *, max_items: int, max_length: int) -> list[str]:
     values = value if isinstance(value, (list, tuple, set)) else str(value or "").split(",")
     normalized: list[str] = []
@@ -383,7 +420,9 @@ async def list_business_connections(
     )
     items = []
     async for item in cursor:
-        item["settings"] = await get_account_settings(item["connection_id"])
+        item["settings"] = public_account_settings(
+            await get_account_settings(item["connection_id"])
+        )
         items.append(item)
     return items, total
 
@@ -398,6 +437,25 @@ async def get_account_settings(connection_id: str) -> dict[str, Any]:
     }
 
 
+def public_account_settings(settings: dict[str, Any]) -> dict[str, Any]:
+    public = {
+        key: value
+        for key, value in settings.items()
+        if key not in SECRET_ACCOUNT_SETTINGS
+    }
+    public.update(
+        {
+            "feishu_webhook_url": "",
+            "feishu_webhook_signing_secret": "",
+            "feishu_webhook_configured": bool(settings.get("feishu_webhook_url")),
+            "feishu_signing_secret_configured": bool(
+                settings.get("feishu_webhook_signing_secret")
+            ),
+        }
+    )
+    return public
+
+
 def normalize_account_settings(
     values: dict[str, Any], *, previous: dict[str, Any] | None = None
 ) -> dict[str, Any]:
@@ -454,6 +512,24 @@ def normalize_account_settings(
     for key, max_length in (("handoff_message", 1000), ("unsupported_message", 1000)):
         if key in values:
             current[key] = clean_text(values.get(key), max_length=max_length, required=True)
+    if "feishu_webhook_url" in values:
+        current["feishu_webhook_url"] = _normalize_feishu_webhook_url(
+            values.get("feishu_webhook_url")
+        )
+    if "feishu_webhook_signing_secret" in values:
+        current["feishu_webhook_signing_secret"] = clean_text(
+            values.get("feishu_webhook_signing_secret"), max_length=512
+        )
+    if "feishu_webhook_enabled" in values:
+        current["feishu_webhook_enabled"] = bool(values.get("feishu_webhook_enabled"))
+    if "feishu_message_preview_enabled" in values:
+        current["feishu_message_preview_enabled"] = bool(
+            values.get("feishu_message_preview_enabled")
+        )
+    if current["feishu_webhook_enabled"] and not current["feishu_webhook_url"]:
+        raise AssistantDataError(
+            "invalid_setting", "开启飞书消息提醒前,请先填写群机器人 Webhook 地址。"
+        )
     return {key: current[key] for key in DEFAULT_ACCOUNT_SETTINGS}
 
 
@@ -1306,6 +1382,65 @@ async def append_conversation_message(
     return document
 
 
+async def claim_feishu_message_notification(
+    conversation_id: str, telegram_message_id: int
+) -> bool:
+    message_key = f"{conversation_id}:incoming:{int(telegram_message_id)}"
+    stale_before = utc_now() - timedelta(minutes=2)
+    result = await messagesdb.find_one_and_update(
+        _scope(
+            {
+                "message_key": message_key,
+                "$or": [
+                    {
+                        "metadata.feishu_notification_status": {
+                            "$nin": ["sending", "sent"]
+                        }
+                    },
+                    {
+                        "metadata.feishu_notification_status": "sending",
+                        "metadata.feishu_notification_started_at": {
+                            "$lte": stale_before
+                        },
+                    },
+                ],
+            }
+        ),
+        {
+            "$set": {
+                "metadata.feishu_notification_status": "sending",
+                "metadata.feishu_notification_started_at": utc_now(),
+            }
+        },
+        return_document=True,
+    )
+    return bool(result)
+
+
+async def mark_feishu_message_notification(
+    conversation_id: str,
+    telegram_message_id: int,
+    *,
+    status: str,
+    error: str = "",
+) -> None:
+    if status not in {"sent", "failed"}:
+        raise AssistantDataError("invalid_status", "飞书通知状态无效。")
+    message_key = f"{conversation_id}:incoming:{int(telegram_message_id)}"
+    await messagesdb.update_one(
+        _scope(message_key=message_key),
+        {
+            "$set": {
+                "metadata.feishu_notification_status": status,
+                "metadata.feishu_notification_error": clean_text(
+                    error, max_length=1000
+                ),
+                "metadata.feishu_notification_finished_at": utc_now(),
+            }
+        },
+    )
+
+
 async def recent_conversation_messages(
     conversation_id: str, *, limit: int = 12
 ) -> list[dict[str, Any]]: