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()