| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239 |
- from __future__ import annotations
- import json
- from datetime import UTC, datetime, timedelta
- from types import SimpleNamespace
- import pytest
- from aiohttp import CookieJar
- from aiohttp.test_utils import TestClient, TestServer
- from pyrogram.enums import ChatMemberStatus
- def connection_payload(connection_id: str = "conn-a", *, can_reply: bool = True) -> dict:
- return {
- "id": connection_id,
- "user": {"id": 900, "username": "owner", "first_name": "店主"},
- "user_chat_id": 900,
- "date": int(datetime.now(UTC).timestamp()),
- "rights": {"can_reply": can_reply, "can_read_messages": True},
- "is_enabled": True,
- }
- def customer_message(
- *,
- message_id: int,
- text: str | None = "营业时间是什么?",
- chat_id: int = 501,
- sender_id: int = 501,
- sent_at: datetime | None = None,
- **values,
- ) -> dict:
- message = {
- "business_connection_id": "conn-a",
- "message_id": message_id,
- "date": int((sent_at or datetime.now(UTC)).timestamp()),
- "chat": {"id": chat_id, "type": "private"},
- "from": {"id": sender_id, "username": f"user{sender_id}", "first_name": "客户"},
- **values,
- }
- if text is not None:
- message["text"] = text
- return message
- class FakeProvider:
- configured = True
- def __init__(self, *, error: Exception | None = None) -> None:
- self.error = error
- self.calls: list[dict] = []
- self.extraction_calls: list[dict] = []
- async def decide(self, **values):
- self.calls.append(values)
- if self.error:
- raise self.error
- entry_id = str(values["knowledge"][0]["entry_id"])
- return {
- "action": "answer",
- "reply": "每天 09:00 至 18:00 营业。",
- "handoff_reason": "",
- "matched_entry_ids": [entry_id],
- "summary": "客户询问营业时间。",
- }
- async def test_connection(self):
- return {"ok": True, "model": "test-model", "response_id": "response-1"}
- async def extract_knowledge(self, **values):
- self.extraction_calls.append(values)
- if self.error:
- raise self.error
- return [
- {
- "question": values.get("question_context") or "如何办理?",
- "aliases": [],
- "keywords": ["办理"],
- "answer": values["content"],
- "tags": ["自动采集"],
- "confidence": 0.92,
- }
- ]
- class FakeBusinessApi:
- def __init__(self) -> None:
- self.sent: list[dict] = []
- self.read: list[tuple[str, int, int]] = []
- async def send_message(
- self,
- chat_id: int,
- text: str,
- *,
- business_connection_id: str = "",
- reply_markup: dict | None = None,
- ) -> dict:
- self.sent.append(
- {
- "chat_id": chat_id,
- "text": text,
- "business_connection_id": business_connection_id,
- "reply_markup": reply_markup,
- }
- )
- return {
- "message_id": 1000 + len(self.sent),
- "sender_business_bot": {"id": 999},
- }
- async def read_business_message(
- self, connection_id: str, chat_id: int, message_id: int
- ) -> None:
- 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 send_card(
- self,
- webhook_url: str,
- *,
- title: str,
- lines: list[str],
- link_url: str = "",
- signing_secret: str = "",
- ) -> None:
- self.sent.append(
- {
- "webhook_url": webhook_url,
- "title": title,
- "lines": lines,
- "text": "\n".join(lines),
- "link_url": link_url,
- "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")
- await dbassistant.upsert_business_connection(connection_payload())
- await dbassistant.update_account_settings("conn-a", {"assistant_enabled": True})
- knowledge = await dbassistant.create_knowledge_entry(
- "conn-a",
- {
- "question": "营业时间是什么?",
- "aliases": ["几点营业"],
- "keywords": ["营业时间", "营业"],
- "answer": "每天 09:00 至 18:00 营业。",
- "priority": 10,
- },
- )
- selected_provider = provider or FakeProvider()
- runtime = service.BusinessAssistantRuntime(
- token="123456:" + "A" * 30,
- session=SimpleNamespace(),
- provider=selected_provider,
- )
- runtime.api = FakeBusinessApi()
- return dbassistant, service, runtime, knowledge
- async def test_connection_settings_knowledge_quota_and_bot_scope(app_modules):
- dbassistant = app_modules.load("wbb.utils.dbassistant")
- await dbassistant.upsert_business_connection(connection_payload("conn-a"))
- await dbassistant.upsert_business_connection(connection_payload("conn-b"))
- settings = await dbassistant.update_account_settings(
- "conn-a",
- {
- "assistant_enabled": True,
- "account_daily_limit": 2,
- "customer_daily_limit": 1,
- "timezone": "Asia/Shanghai",
- },
- )
- first = await dbassistant.create_knowledge_entry(
- "conn-a",
- {
- "question": "营业时间",
- "aliases": ["几点开门"],
- "keywords": ["开门"],
- "answer": "九点开门。",
- "priority": 20,
- },
- )
- await dbassistant.create_knowledge_entry(
- "conn-b",
- {"question": "营业时间", "keywords": ["开门"], "answer": "十点开门。"},
- )
- await app_modules.wbb.db.business_assistant_knowledge.insert_one(
- {
- "bot_id": "foreign-bot",
- "entry_id": "foreign-entry",
- "connection_id": "conn-a",
- "question": "营业时间",
- "aliases": [],
- "keywords": ["开门"],
- "answer": "不应跨 Bot 返回。",
- "priority": 1000,
- "enabled": True,
- "updated_at": datetime.now(UTC),
- }
- )
- matched = await dbassistant.match_knowledge("conn-a", "你们几点开门?")
- assert [item["entry_id"] for item in matched] == [first["entry_id"]]
- items, total = await dbassistant.list_knowledge_entries("conn-a")
- assert total == 1
- assert items[0]["answer"] == "九点开门。"
- await app_modules.wbb.db.business_assistant_connections.insert_one(
- {
- "bot_id": "foreign-bot",
- "connection_id": "foreign-connection",
- "updated_at": datetime.now(UTC),
- }
- )
- connections, total = await dbassistant.list_business_connections(page_size=100)
- assert total == 2
- assert {item["connection_id"] for item in connections} == {"conn-a", "conn-b"}
- assert await dbassistant.reserve_ai_usage("conn-a", 101, settings) == (True, "")
- assert await dbassistant.reserve_ai_usage("conn-a", 101, settings) == (
- False,
- "customer_daily_limit",
- )
- assert await dbassistant.reserve_ai_usage("conn-a", 102, settings) == (True, "")
- assert await dbassistant.reserve_ai_usage("conn-a", 103, settings) == (
- False,
- "account_daily_limit",
- )
- usage = await dbassistant.usage_metrics(connection_id="conn-a")
- assert usage["ai_calls"] == 2
- assert usage["customers"] == 2
- with pytest.raises(dbassistant.AssistantDataError) as error:
- dbassistant.normalize_account_settings({"timezone": "invalid/timezone"})
- 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(
- json.dumps(
- {
- "action": "answer",
- "reply": "标准答案",
- "matched_entry_ids": ["entry-1", "foreign"],
- "summary": "摘要",
- }
- ),
- allowed_entry_ids={"entry-1"},
- )
- assert parsed["matched_entry_ids"] == ["entry-1"]
- with pytest.raises(service.AssistantProviderError):
- service.parse_ai_decision(
- '{"action":"answer","reply":"没有引用"}', allowed_entry_ids={"entry-1"}
- )
- now = datetime.now(UTC)
- assert service.business_reply_window_open(
- {"date": int((now - timedelta(hours=23)).timestamp())}, now=now
- )
- assert not service.business_reply_window_open(
- {"date": int((now - timedelta(hours=25)).timestamp())}, now=now
- )
- assert not service.business_reply_window_open({}, now=now)
- extracted = service.parse_knowledge_extraction(
- json.dumps(
- {
- "items": [
- {
- "question": "如何配送?",
- "aliases": ["配送范围"],
- "keywords": ["配送"],
- "answer": "仅支持市区配送。",
- "tags": ["配送"],
- "confidence": 0.9,
- }
- ]
- }
- )
- )
- assert extracted[0]["confidence"] == 0.9
- assert (
- service.redact_business_knowledge_text(
- "电话 138 0013 8000,邮箱 user@example.com,Telegram @private_user"
- )
- == "电话 [电话已脱敏],邮箱 [邮箱已脱敏],Telegram [用户名已脱敏]"
- )
- async def test_source_event_candidates_publish_edit_delete_and_isolate(app_modules):
- dbassistant = app_modules.load("wbb.utils.dbassistant")
- await dbassistant.upsert_business_connection(connection_payload("conn-a"))
- await dbassistant.upsert_business_connection(connection_payload("conn-b"))
- source = await dbassistant.create_knowledge_source(
- "conn-a",
- {
- "source_type": "channel",
- "chat_id": -100100,
- "title": "官方频道",
- "publication_mode": "auto",
- "author_policy": "admins_or_allowlist",
- },
- )
- event, changed = await dbassistant.record_source_event(
- source,
- event_key="telegram:-100100:10",
- content="每天九点营业。",
- metadata={"chat_id": -100100, "message_id": 10, "trusted_author": True},
- )
- assert changed is True
- candidates = await dbassistant.replace_event_candidates(
- source,
- event,
- [
- {
- "question": "几点营业?",
- "aliases": ["营业时间"],
- "keywords": ["营业"],
- "answer": "每天九点营业。",
- "tags": ["营业"],
- "confidence": 0.95,
- }
- ],
- auto_publish=True,
- )
- candidate = candidates[0]
- first_entry_id = candidate["knowledge_entry_id"]
- first_entry = await app_modules.wbb.db.business_assistant_knowledge.find_one(
- {"bot_id": "primary", "entry_id": first_entry_id}
- )
- assert candidate["status"] == "published"
- assert first_entry["enabled"] is True
- edited_event, changed = await dbassistant.record_source_event(
- source,
- event_key="telegram:-100100:10",
- content="每天十点营业。",
- metadata={"chat_id": -100100, "message_id": 10, "trusted_author": True},
- )
- assert changed is True
- edited = await dbassistant.replace_event_candidates(
- source,
- edited_event,
- [
- {
- "question": "几点营业?",
- "aliases": [],
- "keywords": ["营业"],
- "answer": "每天十点营业。",
- "tags": ["营业"],
- "confidence": 0.96,
- }
- ],
- auto_publish=True,
- )
- assert edited[0]["knowledge_entry_id"] == first_entry_id
- updated_entry = await app_modules.wbb.db.business_assistant_knowledge.find_one(
- {"bot_id": "primary", "entry_id": first_entry_id}
- )
- assert updated_entry["answer"] == "每天十点营业。"
- assert await dbassistant.mark_source_event_deleted(
- source["source_id"], "telegram:-100100:10"
- )
- stale = await dbassistant.get_knowledge_candidate(candidate["candidate_id"])
- disabled_entry = await app_modules.wbb.db.business_assistant_knowledge.find_one(
- {"bot_id": "primary", "entry_id": first_entry_id}
- )
- assert stale["status"] == "stale"
- assert disabled_entry["enabled"] is False
- sources_a, total_a = await dbassistant.list_knowledge_sources("conn-a")
- sources_b, total_b = await dbassistant.list_knowledge_sources("conn-b")
- assert total_a == 1 and sources_a[0]["source_id"] == source["source_id"]
- assert total_b == 0 and sources_b == []
- event_indexes = await app_modules.wbb.db.business_assistant_source_events.index_information()
- assert event_indexes["expires_at_1"]["expireAfterSeconds"] == 0
- async def test_business_human_reply_learns_only_when_source_enabled(app_modules):
- dbassistant, _service, runtime, _knowledge = await prepare_runtime(app_modules)
- await runtime.process_update(
- {
- "update_id": 9,
- "business_message": customer_message(
- message_id=9, text="配送到哪里?", chat_id=509, sender_id=509
- ),
- }
- )
- await runtime.process_update(
- {
- "update_id": 10,
- "business_message": customer_message(
- message_id=10, text="仅支持市区配送。", chat_id=509, sender_id=900
- ),
- }
- )
- assert runtime.provider.extraction_calls == []
- await dbassistant.create_knowledge_source(
- "conn-a",
- {
- "source_type": "business",
- "title": "人工接待对话",
- "publication_mode": "auto",
- },
- )
- await runtime.process_update(
- {
- "update_id": 15,
- "business_message": customer_message(
- message_id=15, text="支持哪些区域?", chat_id=515, sender_id=515
- ),
- }
- )
- await runtime.process_update(
- {
- "update_id": 16,
- "business_message": customer_message(
- message_id=16, text="仅支持市区配送。", chat_id=515, sender_id=900
- ),
- }
- )
- assert len(runtime.provider.extraction_calls) == 1
- candidates, total = await dbassistant.list_knowledge_candidates(
- "conn-a", status="published"
- )
- assert total == 1
- assert candidates[0]["answer"] == "仅支持市区配送。"
- entry_id = candidates[0]["knowledge_entry_id"]
- await runtime.process_update(
- {
- "update_id": 17,
- "edited_business_message": customer_message(
- message_id=16, text="仅二环内支持配送。", chat_id=515, sender_id=900
- ),
- }
- )
- edited = await dbassistant.get_knowledge_candidate(candidates[0]["candidate_id"])
- assert edited["answer"] == "仅二环内支持配送。"
- assert edited["knowledge_entry_id"] == entry_id
- await runtime.process_update(
- {
- "update_id": 18,
- "deleted_business_messages": {
- "business_connection_id": "conn-a",
- "chat": {"id": 515},
- "message_ids": [16],
- },
- }
- )
- stale = await dbassistant.get_knowledge_candidate(candidates[0]["candidate_id"])
- entry = await app_modules.wbb.db.business_assistant_knowledge.find_one(
- {"bot_id": "primary", "entry_id": entry_id}
- )
- assert stale["status"] == "stale"
- assert entry["enabled"] is False
- async def test_group_collection_requires_trusted_author_for_auto_publish(app_modules):
- dbassistant, _service, runtime, _knowledge = await prepare_runtime(app_modules)
- module = app_modules.load("wbb.modules.business_assistant")
- module._runtime = runtime
- await dbassistant.create_knowledge_source(
- "conn-a",
- {
- "source_type": "group",
- "chat_id": -100300,
- "title": "运营讨论群",
- "publication_mode": "auto",
- "author_policy": "admins_or_allowlist",
- },
- )
- def group_message(message_id: int, author_id: int, text: str):
- return SimpleNamespace(
- id=message_id,
- text=text,
- caption=None,
- outgoing=False,
- is_automatic_forward=False,
- date=datetime.now(UTC),
- chat=SimpleNamespace(id=-100300),
- from_user=SimpleNamespace(id=author_id, is_bot=False),
- )
- ordinary = group_message(1, 301, "市区支持当日配送。")
- await module.process_knowledge_source_message(ordinary)
- pending, pending_total = await dbassistant.list_knowledge_candidates(
- "conn-a", status="pending"
- )
- assert pending_total == 1
- assert pending[0]["status"] == "pending"
- app_modules.app.members[(-100300, 302)] = SimpleNamespace(
- status=ChatMemberStatus.ADMINISTRATOR
- )
- admin_message = group_message(2, 302, "每周一至周五提供配送。")
- await module.process_knowledge_source_message(admin_message)
- published, published_total = await dbassistant.list_knowledge_candidates(
- "conn-a", status="published"
- )
- assert published_total == 1
- assert published[0]["source_snapshot"]["author_id"] == 302
- extraction_count = len(runtime.provider.extraction_calls)
- await module.process_knowledge_source_message(admin_message)
- assert len(runtime.provider.extraction_calls) == extraction_count
- async def test_business_updates_are_idempotent_and_manual_reply_pauses(app_modules):
- dbassistant, _service, runtime, knowledge = await prepare_runtime(app_modules)
- update = {
- "update_id": 11,
- "business_message": customer_message(message_id=1),
- }
- await runtime.process_update(update)
- await runtime.process_update(update)
- assert len(runtime.api.sent) == 1
- assert runtime.api.sent[0]["business_connection_id"] == "conn-a"
- assert runtime.api.read == [("conn-a", 501, 1)]
- assert len(runtime.provider.calls) == 1
- assert runtime.provider.calls[0]["knowledge"][0]["entry_id"] == knowledge["entry_id"]
- conversation = await dbassistant.get_conversation_by_chat("conn-a", 501)
- assert conversation["summary"] == "客户询问营业时间。"
- messages = await dbassistant.recent_conversation_messages(
- conversation["conversation_id"], limit=10
- )
- assert [item["direction"] for item in messages] == ["incoming", "assistant"]
- assert messages[0]["expires_at"] - messages[0]["created_at"] == timedelta(days=30)
- message_indexes = await app_modules.wbb.db.business_assistant_messages.index_information()
- assert message_indexes["expires_at_1"]["expireAfterSeconds"] == 0
- await runtime.process_update(
- {
- "update_id": 12,
- "business_message": customer_message(
- message_id=2, text="账号本人回复", sender_id=900
- ),
- }
- )
- paused = await dbassistant.get_conversation(conversation["conversation_id"])
- assert paused["status"] == "human_paused"
- assert paused["handoff_reason"] == "human_reply_detected"
- assert paused["customer"]["id"] == 501
- await runtime.process_update(
- {
- "update_id": 13,
- "business_message": customer_message(
- message_id=3,
- sender_business_bot={"id": 999},
- ),
- }
- )
- await runtime.process_update(
- {
- "update_id": 14,
- "business_message": customer_message(
- message_id=4,
- is_from_offline=True,
- ),
- }
- )
- assert len(runtime.api.sent) == 1
- async def test_every_user_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,
- sent_at=datetime(2026, 8, 24, 2, 3, 4, tzinfo=UTC),
- ),
- }
- )
- 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"]
- assert feishu.sent[0]["title"] == "客户"
- assert feishu.sent[0]["lines"][0] == "消息已隐藏"
- assert "账号:" not in feishu.sent[0]["text"]
- assert "用户:" not in feishu.sent[0]["text"]
- assert feishu.sent[0]["link_url"] == "https://t.me/user501"
- assert "2026-08-24 10:03:04" in feishu.sent[0]["lines"]
- 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"]
- await runtime.process_update(
- {
- "update_id": 1031,
- "business_message": customer_message(
- message_id=21,
- chat_id=502,
- sender_id=502,
- **{"from": {"id": 502, "first_name": "无用户名用户"}},
- ),
- }
- )
- assert len(feishu.sent) == 4
- assert feishu.sent[3]["title"] == "无用户名用户"
- assert feishu.sent[3]["link_url"] == ""
- 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) == 4
- 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(
- "conn-a",
- {
- "notification_destination": "both",
- "ops_group_id": -100123,
- "account_daily_limit": 1,
- "customer_daily_limit": 1,
- },
- )
- await runtime.process_update(
- {
- "update_id": 21,
- "business_message": customer_message(
- message_id=1, text="我要退款", chat_id=601, sender_id=601
- ),
- }
- )
- sensitive = await dbassistant.get_conversation_by_chat("conn-a", 601)
- assert sensitive["status"] == "handoff"
- assert sensitive["handoff_reason"] == "sensitive_request"
- assert runtime.api.sent[-1]["business_connection_id"] == ""
- assert "message_id=1" in runtime.api.sent[-1]["text"]
- assert {item["chat_id"] for item in runtime.api.sent[-2:]} == {900, -100123}
- await runtime.process_update(
- {
- "update_id": 22,
- "business_message": customer_message(
- message_id=2, text=None, chat_id=602, sender_id=602, photo=[{"file_id": "x"}]
- ),
- }
- )
- unsupported = await dbassistant.get_conversation_by_chat("conn-a", 602)
- assert unsupported["handoff_reason"] == "unsupported_message"
- runtime.provider = FakeProvider(error=service.AssistantProviderError("模型超时"))
- await runtime.process_update(
- {
- "update_id": 23,
- "business_message": customer_message(
- message_id=3, chat_id=603, sender_id=603
- ),
- }
- )
- provider_failed = await dbassistant.get_conversation_by_chat("conn-a", 603)
- assert provider_failed["handoff_reason"] == "provider_error"
- await runtime.process_update(
- {
- "update_id": 230,
- "business_message": customer_message(
- message_id=30, chat_id=605, sender_id=605
- ),
- }
- )
- quota_handoff = await dbassistant.get_conversation_by_chat("conn-a", 605)
- assert quota_handoff["handoff_reason"] == "account_daily_limit"
- sent_before = len(runtime.api.sent)
- await runtime.process_update(
- {
- "update_id": 24,
- "business_message": customer_message(
- message_id=4,
- chat_id=604,
- sender_id=604,
- sent_at=datetime.now(UTC) - timedelta(hours=25),
- ),
- }
- )
- expired = await dbassistant.get_conversation_by_chat("conn-a", 604)
- assert expired["handoff_reason"] == "reply_window_expired"
- new_sends = runtime.api.sent[sent_before:]
- assert len(new_sends) == 2
- assert all(item["business_connection_id"] == "" for item in new_sends)
- async def test_failed_update_goes_to_dead_letter_without_blocking(app_modules, monkeypatch):
- dbassistant, _service, runtime, _knowledge = await prepare_runtime(app_modules)
- async def fail(_message):
- raise RuntimeError("persistent failure")
- async def no_sleep(_seconds):
- return None
- runtime._handle_business_message = fail
- monkeypatch.setattr("wbb.services.business_assistant.asyncio.sleep", no_sleep)
- await runtime.process_update(
- {"update_id": 31, "business_message": customer_message(message_id=1)}
- )
- dead_letter = await app_modules.wbb.db.business_assistant_dead_letters.find_one(
- {"bot_id": "primary", "update_id": 31}
- )
- update = await app_modules.wbb.db.business_assistant_updates.find_one(
- {"bot_id": "primary", "update_id": 31}
- )
- assert dead_letter["error"] == "persistent failure"
- assert update["status"] == "done"
- assert await dbassistant.claim_update(31) is False
- async def test_connection_updates_revoke_permissions_and_persist_offset(app_modules):
- service = app_modules.load("wbb.services.business_assistant")
- dbassistant = app_modules.load("wbb.utils.dbassistant")
- runtime = service.BusinessAssistantRuntime(
- token="123456:" + "A" * 30,
- session=SimpleNamespace(),
- provider=FakeProvider(),
- )
- runtime.api = FakeBusinessApi()
- await runtime.process_update(
- {"update_id": 41, "business_connection": connection_payload(can_reply=True)}
- )
- connected = await dbassistant.get_business_connection("conn-a")
- assert connected["is_enabled"] is True
- assert connected["rights"]["can_reply"] is True
- await dbassistant.update_account_settings("conn-a", {"assistant_enabled": True})
- revoked = connection_payload(can_reply=False)
- revoked["is_enabled"] = False
- await runtime.process_update({"update_id": 42, "business_connection": revoked})
- disconnected = await dbassistant.get_business_connection("conn-a")
- assert disconnected["is_enabled"] is False
- assert "can_reply" not in disconnected["rights"]
- await runtime.process_update(
- {"update_id": 420, "business_message": customer_message(message_id=20)}
- )
- assert runtime.api.sent == []
- await dbassistant.save_update_offset(43)
- assert await dbassistant.load_update_offset() == 43
- await dbassistant.save_update_offset(44)
- assert await dbassistant.load_update_offset() == 44
- class StartupApi:
- def __init__(self, *, supported: bool, webhook_url: str = "") -> None:
- self.supported = supported
- self.webhook_url = webhook_url
- async def get_me(self):
- return {"username": "business_bot", "can_connect_to_business": self.supported}
- async def get_webhook_info(self):
- return {"url": self.webhook_url}
- async def test_runtime_blocks_webhook_conflict_and_disabled_business_mode(app_modules):
- service = app_modules.load("wbb.services.business_assistant")
- runtime = service.BusinessAssistantRuntime(
- token="123456:" + "A" * 30,
- session=SimpleNamespace(),
- provider=FakeProvider(),
- )
- runtime.api = StartupApi(supported=True, webhook_url="https://hooks.example.com/tg")
- webhook_status = await runtime.start()
- assert webhook_status["polling_state"] == "blocked"
- assert webhook_status["webhook_conflict"] is True
- assert runtime._poll_task is None
- runtime.api = StartupApi(supported=False)
- business_status = await runtime.start()
- assert business_status["polling_state"] == "blocked"
- assert business_status["business_mode_supported"] is False
- assert "BotFather" in business_status["last_error"]
- 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
- async def __aexit__(self, *_args):
- return None
- async def json(self, **_kwargs):
- return self.payload
- class RetrySession:
- def __init__(self, responses: list[FakeResponse]) -> None:
- self.responses = responses
- self.calls = 0
- def post(self, *_args, **_kwargs):
- response = self.responses[self.calls]
- self.calls += 1
- return response
- async def test_telegram_api_retries_429_and_5xx(app_modules, monkeypatch):
- service = app_modules.load("wbb.services.business_assistant")
- session = RetrySession(
- [
- FakeResponse(
- 429,
- {
- "ok": False,
- "error_code": 429,
- "description": "Too Many Requests",
- "parameters": {"retry_after": 3},
- },
- ),
- FakeResponse(
- 502,
- {"ok": False, "error_code": 502, "description": "Bad Gateway"},
- ),
- FakeResponse(200, {"ok": True, "result": {"id": 999}}),
- ]
- )
- sleeps: list[int] = []
- async def record_sleep(seconds):
- sleeps.append(seconds)
- monkeypatch.setattr("wbb.services.business_assistant.asyncio.sleep", record_sleep)
- api = service.TelegramBusinessApi("123456:" + "A" * 30, session)
- assert await api.get_me() == {"id": 999}
- assert session.calls == 3
- 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",
- "消息提醒",
- signing_secret="secret",
- )
- assert session.payload["msg_type"] == "text"
- assert session.payload["content"]["text"] == "消息提醒"
- assert session.payload["timestamp"]
- assert session.payload["sign"]
- card_lines = ["账号:店主", "用户:测试用户"]
- await client.send_card(
- "https://open.feishu.cn/open-apis/bot/v2/hook/12345678-abcd-4321-abcd-123456789abc",
- title="消息提醒",
- lines=card_lines,
- link_url="https://t.me/test_user",
- )
- assert session.payload == {
- "msg_type": "interactive",
- "card": {
- "config": {"wide_screen_mode": True},
- "header": {
- "template": "blue",
- "title": {"tag": "plain_text", "content": "消息提醒"},
- },
- "elements": [
- {
- "tag": "div",
- "text": {"tag": "plain_text", "content": line},
- }
- for line in card_lines
- ]
- + [
- {"tag": "hr"},
- {
- "tag": "action",
- "actions": [
- {
- "tag": "button",
- "text": {
- "tag": "plain_text",
- "content": "打开 Telegram 对话",
- },
- "type": "primary",
- "url": "https://t.me/test_user",
- }
- ],
- },
- ],
- },
- }
- await client.send_card(
- "https://open.feishu.cn/open-apis/bot/v2/hook/12345678-abcd-4321-abcd-123456789abc",
- title="消息提醒",
- lines=["无链接"],
- )
- assert session.payload == {
- "msg_type": "interactive",
- "card": {
- "config": {"wide_screen_mode": True},
- "header": {
- "template": "blue",
- "title": {"tag": "plain_text", "content": "消息提醒"},
- },
- "elements": [
- {
- "tag": "div",
- "text": {"tag": "plain_text", "content": "无链接"},
- }
- ],
- },
- }
- async def _login_and_change_password(client: TestClient) -> str:
- login = await client.post(
- "/api/admin/v1/auth/login",
- json={"username": "admin", "password": "qwe0.123456"},
- )
- login_data = (await login.json())["data"]
- changed = await client.put(
- "/api/admin/v1/auth/password",
- headers={"X-CSRF-Token": login_data["csrf_token"]},
- json={
- "current_password": "qwe0.123456",
- "new_password": "changed-pass-123",
- },
- )
- return (await changed.json())["data"]["csrf_token"]
- async def test_business_assistant_api_permission_csrf_audit_and_clear(app_modules):
- dbassistant = app_modules.load("wbb.utils.dbassistant")
- await dbassistant.upsert_business_connection(connection_payload())
- conversation = await dbassistant.get_or_create_conversation(
- "conn-a", 701, customer={"id": 701, "first_name": "待清理客户"}
- )
- await dbassistant.append_conversation_message(
- conversation["conversation_id"],
- direction="incoming",
- telegram_message_id=1,
- text="请清理",
- )
- await dbassistant.update_conversation_summary(conversation["conversation_id"], "待清理摘要")
- admin_api = app_modules.load("wbb.admin.api")
- application = admin_api.build_admin_application()
- await application["admin_api"].initialize()
- client = TestClient(TestServer(application), cookie_jar=CookieJar(unsafe=True))
- await client.start_server()
- try:
- csrf = await _login_and_change_password(client)
- app_modules.wbb.BOT_PERMISSIONS = set()
- denied = await client.get("/api/admin/v1/business-assistant/status")
- assert denied.status == 403
- app_modules.wbb.BOT_PERMISSIONS = {"business_assistant.manage"}
- no_csrf = await client.put(
- "/api/admin/v1/business-assistant/settings",
- json={"connection_id": "conn-a", "assistant_enabled": True, "confirm": True},
- )
- assert no_csrf.status == 403
- unconfirmed = await client.put(
- "/api/admin/v1/business-assistant/settings",
- headers={"X-CSRF-Token": csrf},
- json={"connection_id": "conn-a", "assistant_enabled": True},
- )
- assert unconfirmed.status == 409
- saved = await client.put(
- "/api/admin/v1/business-assistant/settings",
- headers={"X-CSRF-Token": csrf},
- 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
- 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",
- headers={"X-CSRF-Token": csrf},
- json={
- "connection_id": "conn-a",
- "question": "配送范围",
- "keywords": ["配送"],
- "answer": "仅限市区。",
- "confirm": True,
- },
- )
- assert created.status == 201
- listed = await client.get(
- "/api/admin/v1/business-assistant/knowledge?connection_id=conn-a"
- )
- assert (await listed.json())["data"]["total"] == 1
- created_source = await client.post(
- "/api/admin/v1/business-assistant/knowledge-sources",
- headers={"X-CSRF-Token": csrf},
- json={
- "connection_id": "conn-a",
- "source_type": "business",
- "title": "人工对话",
- "publication_mode": "review",
- "confirm": True,
- },
- )
- assert created_source.status == 201
- source = (await created_source.json())["data"]
- event, _ = await dbassistant.record_source_event(
- source,
- event_key="business:701:8",
- content="市区可配送。",
- metadata={"chat_id": 701, "message_id": 8, "trusted_author": True},
- )
- candidates = await dbassistant.replace_event_candidates(
- source,
- event,
- [
- {
- "question": "配送范围?",
- "keywords": ["配送"],
- "answer": "市区可配送。",
- "confidence": 0.9,
- }
- ],
- auto_publish=False,
- )
- candidate_page = await client.get(
- "/api/admin/v1/business-assistant/knowledge-candidates"
- "?connection_id=conn-a&status=pending"
- )
- assert candidate_page.status == 200
- assert (await candidate_page.json())["data"]["total"] == 1
- approved = await client.post(
- "/api/admin/v1/business-assistant/knowledge-candidates/"
- f"{candidates[0]['candidate_id']}/actions",
- headers={"X-CSRF-Token": csrf},
- json={"action": "approve", "confirm": True},
- )
- assert approved.status == 200
- assert (await approved.json())["data"]["status"] == "published"
- cleared = await client.post(
- f"/api/admin/v1/business-assistant/conversations/{conversation['conversation_id']}/actions",
- headers={"X-CSRF-Token": csrf},
- json={"action": "clear", "confirm": True},
- )
- assert cleared.status == 200
- detail = await dbassistant.conversation_detail(conversation["conversation_id"])
- assert detail["summary"] == ""
- assert detail["messages"] == []
- audit = await app_modules.wbb.db.admin_audit_logs.find_one(
- {"action": "business_assistant.conversation.clear"}
- )
- assert audit["success"] is True
- finally:
- await client.close()
|