test_business_assistant.py 44 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239
  1. from __future__ import annotations
  2. import json
  3. from datetime import UTC, datetime, timedelta
  4. from types import SimpleNamespace
  5. import pytest
  6. from aiohttp import CookieJar
  7. from aiohttp.test_utils import TestClient, TestServer
  8. from pyrogram.enums import ChatMemberStatus
  9. def connection_payload(connection_id: str = "conn-a", *, can_reply: bool = True) -> dict:
  10. return {
  11. "id": connection_id,
  12. "user": {"id": 900, "username": "owner", "first_name": "店主"},
  13. "user_chat_id": 900,
  14. "date": int(datetime.now(UTC).timestamp()),
  15. "rights": {"can_reply": can_reply, "can_read_messages": True},
  16. "is_enabled": True,
  17. }
  18. def customer_message(
  19. *,
  20. message_id: int,
  21. text: str | None = "营业时间是什么?",
  22. chat_id: int = 501,
  23. sender_id: int = 501,
  24. sent_at: datetime | None = None,
  25. **values,
  26. ) -> dict:
  27. message = {
  28. "business_connection_id": "conn-a",
  29. "message_id": message_id,
  30. "date": int((sent_at or datetime.now(UTC)).timestamp()),
  31. "chat": {"id": chat_id, "type": "private"},
  32. "from": {"id": sender_id, "username": f"user{sender_id}", "first_name": "客户"},
  33. **values,
  34. }
  35. if text is not None:
  36. message["text"] = text
  37. return message
  38. class FakeProvider:
  39. configured = True
  40. def __init__(self, *, error: Exception | None = None) -> None:
  41. self.error = error
  42. self.calls: list[dict] = []
  43. self.extraction_calls: list[dict] = []
  44. async def decide(self, **values):
  45. self.calls.append(values)
  46. if self.error:
  47. raise self.error
  48. entry_id = str(values["knowledge"][0]["entry_id"])
  49. return {
  50. "action": "answer",
  51. "reply": "每天 09:00 至 18:00 营业。",
  52. "handoff_reason": "",
  53. "matched_entry_ids": [entry_id],
  54. "summary": "客户询问营业时间。",
  55. }
  56. async def test_connection(self):
  57. return {"ok": True, "model": "test-model", "response_id": "response-1"}
  58. async def extract_knowledge(self, **values):
  59. self.extraction_calls.append(values)
  60. if self.error:
  61. raise self.error
  62. return [
  63. {
  64. "question": values.get("question_context") or "如何办理?",
  65. "aliases": [],
  66. "keywords": ["办理"],
  67. "answer": values["content"],
  68. "tags": ["自动采集"],
  69. "confidence": 0.92,
  70. }
  71. ]
  72. class FakeBusinessApi:
  73. def __init__(self) -> None:
  74. self.sent: list[dict] = []
  75. self.read: list[tuple[str, int, int]] = []
  76. async def send_message(
  77. self,
  78. chat_id: int,
  79. text: str,
  80. *,
  81. business_connection_id: str = "",
  82. reply_markup: dict | None = None,
  83. ) -> dict:
  84. self.sent.append(
  85. {
  86. "chat_id": chat_id,
  87. "text": text,
  88. "business_connection_id": business_connection_id,
  89. "reply_markup": reply_markup,
  90. }
  91. )
  92. return {
  93. "message_id": 1000 + len(self.sent),
  94. "sender_business_bot": {"id": 999},
  95. }
  96. async def read_business_message(
  97. self, connection_id: str, chat_id: int, message_id: int
  98. ) -> None:
  99. self.read.append((connection_id, chat_id, message_id))
  100. class FakeFeishuWebhook:
  101. def __init__(self) -> None:
  102. self.sent: list[dict] = []
  103. async def send(self, webhook_url: str, text: str, *, signing_secret: str = "") -> None:
  104. self.sent.append(
  105. {
  106. "webhook_url": webhook_url,
  107. "text": text,
  108. "signing_secret": signing_secret,
  109. }
  110. )
  111. async def send_card(
  112. self,
  113. webhook_url: str,
  114. *,
  115. title: str,
  116. lines: list[str],
  117. link_url: str = "",
  118. signing_secret: str = "",
  119. ) -> None:
  120. self.sent.append(
  121. {
  122. "webhook_url": webhook_url,
  123. "title": title,
  124. "lines": lines,
  125. "text": "\n".join(lines),
  126. "link_url": link_url,
  127. "signing_secret": signing_secret,
  128. }
  129. )
  130. async def prepare_runtime(app_modules, *, provider: FakeProvider | None = None):
  131. dbassistant = app_modules.load("wbb.utils.dbassistant")
  132. service = app_modules.load("wbb.services.business_assistant")
  133. await dbassistant.upsert_business_connection(connection_payload())
  134. await dbassistant.update_account_settings("conn-a", {"assistant_enabled": True})
  135. knowledge = await dbassistant.create_knowledge_entry(
  136. "conn-a",
  137. {
  138. "question": "营业时间是什么?",
  139. "aliases": ["几点营业"],
  140. "keywords": ["营业时间", "营业"],
  141. "answer": "每天 09:00 至 18:00 营业。",
  142. "priority": 10,
  143. },
  144. )
  145. selected_provider = provider or FakeProvider()
  146. runtime = service.BusinessAssistantRuntime(
  147. token="123456:" + "A" * 30,
  148. session=SimpleNamespace(),
  149. provider=selected_provider,
  150. )
  151. runtime.api = FakeBusinessApi()
  152. return dbassistant, service, runtime, knowledge
  153. async def test_connection_settings_knowledge_quota_and_bot_scope(app_modules):
  154. dbassistant = app_modules.load("wbb.utils.dbassistant")
  155. await dbassistant.upsert_business_connection(connection_payload("conn-a"))
  156. await dbassistant.upsert_business_connection(connection_payload("conn-b"))
  157. settings = await dbassistant.update_account_settings(
  158. "conn-a",
  159. {
  160. "assistant_enabled": True,
  161. "account_daily_limit": 2,
  162. "customer_daily_limit": 1,
  163. "timezone": "Asia/Shanghai",
  164. },
  165. )
  166. first = await dbassistant.create_knowledge_entry(
  167. "conn-a",
  168. {
  169. "question": "营业时间",
  170. "aliases": ["几点开门"],
  171. "keywords": ["开门"],
  172. "answer": "九点开门。",
  173. "priority": 20,
  174. },
  175. )
  176. await dbassistant.create_knowledge_entry(
  177. "conn-b",
  178. {"question": "营业时间", "keywords": ["开门"], "answer": "十点开门。"},
  179. )
  180. await app_modules.wbb.db.business_assistant_knowledge.insert_one(
  181. {
  182. "bot_id": "foreign-bot",
  183. "entry_id": "foreign-entry",
  184. "connection_id": "conn-a",
  185. "question": "营业时间",
  186. "aliases": [],
  187. "keywords": ["开门"],
  188. "answer": "不应跨 Bot 返回。",
  189. "priority": 1000,
  190. "enabled": True,
  191. "updated_at": datetime.now(UTC),
  192. }
  193. )
  194. matched = await dbassistant.match_knowledge("conn-a", "你们几点开门?")
  195. assert [item["entry_id"] for item in matched] == [first["entry_id"]]
  196. items, total = await dbassistant.list_knowledge_entries("conn-a")
  197. assert total == 1
  198. assert items[0]["answer"] == "九点开门。"
  199. await app_modules.wbb.db.business_assistant_connections.insert_one(
  200. {
  201. "bot_id": "foreign-bot",
  202. "connection_id": "foreign-connection",
  203. "updated_at": datetime.now(UTC),
  204. }
  205. )
  206. connections, total = await dbassistant.list_business_connections(page_size=100)
  207. assert total == 2
  208. assert {item["connection_id"] for item in connections} == {"conn-a", "conn-b"}
  209. assert await dbassistant.reserve_ai_usage("conn-a", 101, settings) == (True, "")
  210. assert await dbassistant.reserve_ai_usage("conn-a", 101, settings) == (
  211. False,
  212. "customer_daily_limit",
  213. )
  214. assert await dbassistant.reserve_ai_usage("conn-a", 102, settings) == (True, "")
  215. assert await dbassistant.reserve_ai_usage("conn-a", 103, settings) == (
  216. False,
  217. "account_daily_limit",
  218. )
  219. usage = await dbassistant.usage_metrics(connection_id="conn-a")
  220. assert usage["ai_calls"] == 2
  221. assert usage["customers"] == 2
  222. with pytest.raises(dbassistant.AssistantDataError) as error:
  223. dbassistant.normalize_account_settings({"timezone": "invalid/timezone"})
  224. assert error.value.code == "invalid_setting"
  225. async def test_feishu_settings_are_validated_and_never_exposed(app_modules):
  226. dbassistant = app_modules.load("wbb.utils.dbassistant")
  227. await dbassistant.upsert_business_connection(connection_payload())
  228. webhook_url = "https://open.feishu.cn/open-apis/bot/v2/hook/12345678-abcd-4321-abcd-123456789abc"
  229. stored = await dbassistant.update_account_settings(
  230. "conn-a",
  231. {
  232. "feishu_webhook_enabled": True,
  233. "feishu_webhook_url": webhook_url,
  234. "feishu_webhook_signing_secret": "signing-secret",
  235. "feishu_message_preview_enabled": True,
  236. },
  237. )
  238. assert stored["feishu_webhook_url"] == webhook_url
  239. assert stored["feishu_webhook_signing_secret"] == "signing-secret"
  240. public = dbassistant.public_account_settings(stored)
  241. assert public["feishu_webhook_url"] == ""
  242. assert public["feishu_webhook_signing_secret"] == ""
  243. assert public["feishu_webhook_configured"] is True
  244. assert public["feishu_signing_secret_configured"] is True
  245. connections, _total = await dbassistant.list_business_connections()
  246. assert connections[0]["settings"]["feishu_webhook_url"] == ""
  247. assert connections[0]["settings"]["feishu_webhook_configured"] is True
  248. with pytest.raises(dbassistant.AssistantDataError):
  249. dbassistant.normalize_account_settings(
  250. {
  251. "feishu_webhook_enabled": True,
  252. "feishu_webhook_url": "https://example.com/internal-hook",
  253. }
  254. )
  255. def test_ai_contract_and_reply_window_validation(app_modules):
  256. service = app_modules.load("wbb.services.business_assistant")
  257. parsed = service.parse_ai_decision(
  258. json.dumps(
  259. {
  260. "action": "answer",
  261. "reply": "标准答案",
  262. "matched_entry_ids": ["entry-1", "foreign"],
  263. "summary": "摘要",
  264. }
  265. ),
  266. allowed_entry_ids={"entry-1"},
  267. )
  268. assert parsed["matched_entry_ids"] == ["entry-1"]
  269. with pytest.raises(service.AssistantProviderError):
  270. service.parse_ai_decision(
  271. '{"action":"answer","reply":"没有引用"}', allowed_entry_ids={"entry-1"}
  272. )
  273. now = datetime.now(UTC)
  274. assert service.business_reply_window_open(
  275. {"date": int((now - timedelta(hours=23)).timestamp())}, now=now
  276. )
  277. assert not service.business_reply_window_open(
  278. {"date": int((now - timedelta(hours=25)).timestamp())}, now=now
  279. )
  280. assert not service.business_reply_window_open({}, now=now)
  281. extracted = service.parse_knowledge_extraction(
  282. json.dumps(
  283. {
  284. "items": [
  285. {
  286. "question": "如何配送?",
  287. "aliases": ["配送范围"],
  288. "keywords": ["配送"],
  289. "answer": "仅支持市区配送。",
  290. "tags": ["配送"],
  291. "confidence": 0.9,
  292. }
  293. ]
  294. }
  295. )
  296. )
  297. assert extracted[0]["confidence"] == 0.9
  298. assert (
  299. service.redact_business_knowledge_text(
  300. "电话 138 0013 8000,邮箱 user@example.com,Telegram @private_user"
  301. )
  302. == "电话 [电话已脱敏],邮箱 [邮箱已脱敏],Telegram [用户名已脱敏]"
  303. )
  304. async def test_source_event_candidates_publish_edit_delete_and_isolate(app_modules):
  305. dbassistant = app_modules.load("wbb.utils.dbassistant")
  306. await dbassistant.upsert_business_connection(connection_payload("conn-a"))
  307. await dbassistant.upsert_business_connection(connection_payload("conn-b"))
  308. source = await dbassistant.create_knowledge_source(
  309. "conn-a",
  310. {
  311. "source_type": "channel",
  312. "chat_id": -100100,
  313. "title": "官方频道",
  314. "publication_mode": "auto",
  315. "author_policy": "admins_or_allowlist",
  316. },
  317. )
  318. event, changed = await dbassistant.record_source_event(
  319. source,
  320. event_key="telegram:-100100:10",
  321. content="每天九点营业。",
  322. metadata={"chat_id": -100100, "message_id": 10, "trusted_author": True},
  323. )
  324. assert changed is True
  325. candidates = await dbassistant.replace_event_candidates(
  326. source,
  327. event,
  328. [
  329. {
  330. "question": "几点营业?",
  331. "aliases": ["营业时间"],
  332. "keywords": ["营业"],
  333. "answer": "每天九点营业。",
  334. "tags": ["营业"],
  335. "confidence": 0.95,
  336. }
  337. ],
  338. auto_publish=True,
  339. )
  340. candidate = candidates[0]
  341. first_entry_id = candidate["knowledge_entry_id"]
  342. first_entry = await app_modules.wbb.db.business_assistant_knowledge.find_one(
  343. {"bot_id": "primary", "entry_id": first_entry_id}
  344. )
  345. assert candidate["status"] == "published"
  346. assert first_entry["enabled"] is True
  347. edited_event, changed = await dbassistant.record_source_event(
  348. source,
  349. event_key="telegram:-100100:10",
  350. content="每天十点营业。",
  351. metadata={"chat_id": -100100, "message_id": 10, "trusted_author": True},
  352. )
  353. assert changed is True
  354. edited = await dbassistant.replace_event_candidates(
  355. source,
  356. edited_event,
  357. [
  358. {
  359. "question": "几点营业?",
  360. "aliases": [],
  361. "keywords": ["营业"],
  362. "answer": "每天十点营业。",
  363. "tags": ["营业"],
  364. "confidence": 0.96,
  365. }
  366. ],
  367. auto_publish=True,
  368. )
  369. assert edited[0]["knowledge_entry_id"] == first_entry_id
  370. updated_entry = await app_modules.wbb.db.business_assistant_knowledge.find_one(
  371. {"bot_id": "primary", "entry_id": first_entry_id}
  372. )
  373. assert updated_entry["answer"] == "每天十点营业。"
  374. assert await dbassistant.mark_source_event_deleted(
  375. source["source_id"], "telegram:-100100:10"
  376. )
  377. stale = await dbassistant.get_knowledge_candidate(candidate["candidate_id"])
  378. disabled_entry = await app_modules.wbb.db.business_assistant_knowledge.find_one(
  379. {"bot_id": "primary", "entry_id": first_entry_id}
  380. )
  381. assert stale["status"] == "stale"
  382. assert disabled_entry["enabled"] is False
  383. sources_a, total_a = await dbassistant.list_knowledge_sources("conn-a")
  384. sources_b, total_b = await dbassistant.list_knowledge_sources("conn-b")
  385. assert total_a == 1 and sources_a[0]["source_id"] == source["source_id"]
  386. assert total_b == 0 and sources_b == []
  387. event_indexes = await app_modules.wbb.db.business_assistant_source_events.index_information()
  388. assert event_indexes["expires_at_1"]["expireAfterSeconds"] == 0
  389. async def test_business_human_reply_learns_only_when_source_enabled(app_modules):
  390. dbassistant, _service, runtime, _knowledge = await prepare_runtime(app_modules)
  391. await runtime.process_update(
  392. {
  393. "update_id": 9,
  394. "business_message": customer_message(
  395. message_id=9, text="配送到哪里?", chat_id=509, sender_id=509
  396. ),
  397. }
  398. )
  399. await runtime.process_update(
  400. {
  401. "update_id": 10,
  402. "business_message": customer_message(
  403. message_id=10, text="仅支持市区配送。", chat_id=509, sender_id=900
  404. ),
  405. }
  406. )
  407. assert runtime.provider.extraction_calls == []
  408. await dbassistant.create_knowledge_source(
  409. "conn-a",
  410. {
  411. "source_type": "business",
  412. "title": "人工接待对话",
  413. "publication_mode": "auto",
  414. },
  415. )
  416. await runtime.process_update(
  417. {
  418. "update_id": 15,
  419. "business_message": customer_message(
  420. message_id=15, text="支持哪些区域?", chat_id=515, sender_id=515
  421. ),
  422. }
  423. )
  424. await runtime.process_update(
  425. {
  426. "update_id": 16,
  427. "business_message": customer_message(
  428. message_id=16, text="仅支持市区配送。", chat_id=515, sender_id=900
  429. ),
  430. }
  431. )
  432. assert len(runtime.provider.extraction_calls) == 1
  433. candidates, total = await dbassistant.list_knowledge_candidates(
  434. "conn-a", status="published"
  435. )
  436. assert total == 1
  437. assert candidates[0]["answer"] == "仅支持市区配送。"
  438. entry_id = candidates[0]["knowledge_entry_id"]
  439. await runtime.process_update(
  440. {
  441. "update_id": 17,
  442. "edited_business_message": customer_message(
  443. message_id=16, text="仅二环内支持配送。", chat_id=515, sender_id=900
  444. ),
  445. }
  446. )
  447. edited = await dbassistant.get_knowledge_candidate(candidates[0]["candidate_id"])
  448. assert edited["answer"] == "仅二环内支持配送。"
  449. assert edited["knowledge_entry_id"] == entry_id
  450. await runtime.process_update(
  451. {
  452. "update_id": 18,
  453. "deleted_business_messages": {
  454. "business_connection_id": "conn-a",
  455. "chat": {"id": 515},
  456. "message_ids": [16],
  457. },
  458. }
  459. )
  460. stale = await dbassistant.get_knowledge_candidate(candidates[0]["candidate_id"])
  461. entry = await app_modules.wbb.db.business_assistant_knowledge.find_one(
  462. {"bot_id": "primary", "entry_id": entry_id}
  463. )
  464. assert stale["status"] == "stale"
  465. assert entry["enabled"] is False
  466. async def test_group_collection_requires_trusted_author_for_auto_publish(app_modules):
  467. dbassistant, _service, runtime, _knowledge = await prepare_runtime(app_modules)
  468. module = app_modules.load("wbb.modules.business_assistant")
  469. module._runtime = runtime
  470. await dbassistant.create_knowledge_source(
  471. "conn-a",
  472. {
  473. "source_type": "group",
  474. "chat_id": -100300,
  475. "title": "运营讨论群",
  476. "publication_mode": "auto",
  477. "author_policy": "admins_or_allowlist",
  478. },
  479. )
  480. def group_message(message_id: int, author_id: int, text: str):
  481. return SimpleNamespace(
  482. id=message_id,
  483. text=text,
  484. caption=None,
  485. outgoing=False,
  486. is_automatic_forward=False,
  487. date=datetime.now(UTC),
  488. chat=SimpleNamespace(id=-100300),
  489. from_user=SimpleNamespace(id=author_id, is_bot=False),
  490. )
  491. ordinary = group_message(1, 301, "市区支持当日配送。")
  492. await module.process_knowledge_source_message(ordinary)
  493. pending, pending_total = await dbassistant.list_knowledge_candidates(
  494. "conn-a", status="pending"
  495. )
  496. assert pending_total == 1
  497. assert pending[0]["status"] == "pending"
  498. app_modules.app.members[(-100300, 302)] = SimpleNamespace(
  499. status=ChatMemberStatus.ADMINISTRATOR
  500. )
  501. admin_message = group_message(2, 302, "每周一至周五提供配送。")
  502. await module.process_knowledge_source_message(admin_message)
  503. published, published_total = await dbassistant.list_knowledge_candidates(
  504. "conn-a", status="published"
  505. )
  506. assert published_total == 1
  507. assert published[0]["source_snapshot"]["author_id"] == 302
  508. extraction_count = len(runtime.provider.extraction_calls)
  509. await module.process_knowledge_source_message(admin_message)
  510. assert len(runtime.provider.extraction_calls) == extraction_count
  511. async def test_business_updates_are_idempotent_and_manual_reply_pauses(app_modules):
  512. dbassistant, _service, runtime, knowledge = await prepare_runtime(app_modules)
  513. update = {
  514. "update_id": 11,
  515. "business_message": customer_message(message_id=1),
  516. }
  517. await runtime.process_update(update)
  518. await runtime.process_update(update)
  519. assert len(runtime.api.sent) == 1
  520. assert runtime.api.sent[0]["business_connection_id"] == "conn-a"
  521. assert runtime.api.read == [("conn-a", 501, 1)]
  522. assert len(runtime.provider.calls) == 1
  523. assert runtime.provider.calls[0]["knowledge"][0]["entry_id"] == knowledge["entry_id"]
  524. conversation = await dbassistant.get_conversation_by_chat("conn-a", 501)
  525. assert conversation["summary"] == "客户询问营业时间。"
  526. messages = await dbassistant.recent_conversation_messages(
  527. conversation["conversation_id"], limit=10
  528. )
  529. assert [item["direction"] for item in messages] == ["incoming", "assistant"]
  530. assert messages[0]["expires_at"] - messages[0]["created_at"] == timedelta(days=30)
  531. message_indexes = await app_modules.wbb.db.business_assistant_messages.index_information()
  532. assert message_indexes["expires_at_1"]["expireAfterSeconds"] == 0
  533. await runtime.process_update(
  534. {
  535. "update_id": 12,
  536. "business_message": customer_message(
  537. message_id=2, text="账号本人回复", sender_id=900
  538. ),
  539. }
  540. )
  541. paused = await dbassistant.get_conversation(conversation["conversation_id"])
  542. assert paused["status"] == "human_paused"
  543. assert paused["handoff_reason"] == "human_reply_detected"
  544. assert paused["customer"]["id"] == 501
  545. await runtime.process_update(
  546. {
  547. "update_id": 13,
  548. "business_message": customer_message(
  549. message_id=3,
  550. sender_business_bot={"id": 999},
  551. ),
  552. }
  553. )
  554. await runtime.process_update(
  555. {
  556. "update_id": 14,
  557. "business_message": customer_message(
  558. message_id=4,
  559. is_from_offline=True,
  560. ),
  561. }
  562. )
  563. assert len(runtime.api.sent) == 1
  564. async def test_every_user_message_notifies_feishu_once(app_modules):
  565. dbassistant, _service, runtime, _knowledge = await prepare_runtime(app_modules)
  566. webhook_url = "https://open.feishu.cn/open-apis/bot/v2/hook/12345678-abcd-4321-abcd-123456789abc"
  567. await dbassistant.update_account_settings(
  568. "conn-a",
  569. {
  570. "assistant_enabled": False,
  571. "feishu_webhook_enabled": True,
  572. "feishu_webhook_url": webhook_url,
  573. "feishu_webhook_signing_secret": "secret",
  574. "feishu_message_preview_enabled": False,
  575. },
  576. )
  577. feishu = FakeFeishuWebhook()
  578. runtime.feishu = feishu
  579. await runtime.process_update(
  580. {
  581. "update_id": 101,
  582. "business_message": customer_message(
  583. message_id=1,
  584. sent_at=datetime(2026, 8, 24, 2, 3, 4, tzinfo=UTC),
  585. ),
  586. }
  587. )
  588. await runtime.process_update(
  589. {"update_id": 102, "business_message": customer_message(message_id=1)}
  590. )
  591. assert len(feishu.sent) == 1
  592. assert feishu.sent[0]["webhook_url"] == webhook_url
  593. assert feishu.sent[0]["signing_secret"] == "secret"
  594. assert "营业时间是什么" not in feishu.sent[0]["text"]
  595. assert feishu.sent[0]["title"] == "客户"
  596. assert feishu.sent[0]["lines"][0] == "消息已隐藏"
  597. assert "账号:" not in feishu.sent[0]["text"]
  598. assert "用户:" not in feishu.sent[0]["text"]
  599. assert feishu.sent[0]["link_url"] == "https://t.me/user501"
  600. assert "2026-08-24 10:03:04" in feishu.sent[0]["lines"]
  601. await dbassistant.update_account_settings(
  602. "conn-a", {"feishu_message_preview_enabled": True}
  603. )
  604. await runtime.process_update(
  605. {
  606. "update_id": 103,
  607. "business_message": customer_message(
  608. message_id=2, text="第二条用户消息"
  609. ),
  610. }
  611. )
  612. assert len(feishu.sent) == 2
  613. assert "第二条用户消息" in feishu.sent[1]["text"]
  614. await runtime.process_update(
  615. {
  616. "update_id": 1030,
  617. "business_message": customer_message(
  618. message_id=20,
  619. text=None,
  620. caption="图片说明",
  621. photo=[{"file_id": "photo"}],
  622. ),
  623. }
  624. )
  625. assert len(feishu.sent) == 3
  626. assert "图片说明" in feishu.sent[2]["text"]
  627. await runtime.process_update(
  628. {
  629. "update_id": 1031,
  630. "business_message": customer_message(
  631. message_id=21,
  632. chat_id=502,
  633. sender_id=502,
  634. **{"from": {"id": 502, "first_name": "无用户名用户"}},
  635. ),
  636. }
  637. )
  638. assert len(feishu.sent) == 4
  639. assert feishu.sent[3]["title"] == "无用户名用户"
  640. assert feishu.sent[3]["link_url"] == ""
  641. await runtime.process_update(
  642. {
  643. "update_id": 104,
  644. "business_message": customer_message(
  645. message_id=3, text="账号本人回复", sender_id=900
  646. ),
  647. }
  648. )
  649. await runtime.process_update(
  650. {
  651. "update_id": 105,
  652. "business_message": customer_message(message_id=4, is_from_offline=True),
  653. }
  654. )
  655. await runtime.process_update(
  656. {
  657. "update_id": 106,
  658. "business_message": customer_message(
  659. message_id=5, sender_business_bot={"id": 999}
  660. ),
  661. }
  662. )
  663. assert len(feishu.sent) == 4
  664. conversation = await dbassistant.get_conversation_by_chat("conn-a", 501)
  665. messages = await dbassistant.recent_conversation_messages(
  666. conversation["conversation_id"], limit=10
  667. )
  668. incoming = [item for item in messages if item["telegram_message_id"] == 2][0]
  669. assert incoming["metadata"]["feishu_notification_status"] == "sent"
  670. async def test_handoff_sensitive_non_text_provider_error_and_expired_window(app_modules):
  671. dbassistant, service, runtime, _knowledge = await prepare_runtime(app_modules)
  672. await dbassistant.update_account_settings(
  673. "conn-a",
  674. {
  675. "notification_destination": "both",
  676. "ops_group_id": -100123,
  677. "account_daily_limit": 1,
  678. "customer_daily_limit": 1,
  679. },
  680. )
  681. await runtime.process_update(
  682. {
  683. "update_id": 21,
  684. "business_message": customer_message(
  685. message_id=1, text="我要退款", chat_id=601, sender_id=601
  686. ),
  687. }
  688. )
  689. sensitive = await dbassistant.get_conversation_by_chat("conn-a", 601)
  690. assert sensitive["status"] == "handoff"
  691. assert sensitive["handoff_reason"] == "sensitive_request"
  692. assert runtime.api.sent[-1]["business_connection_id"] == ""
  693. assert "message_id=1" in runtime.api.sent[-1]["text"]
  694. assert {item["chat_id"] for item in runtime.api.sent[-2:]} == {900, -100123}
  695. await runtime.process_update(
  696. {
  697. "update_id": 22,
  698. "business_message": customer_message(
  699. message_id=2, text=None, chat_id=602, sender_id=602, photo=[{"file_id": "x"}]
  700. ),
  701. }
  702. )
  703. unsupported = await dbassistant.get_conversation_by_chat("conn-a", 602)
  704. assert unsupported["handoff_reason"] == "unsupported_message"
  705. runtime.provider = FakeProvider(error=service.AssistantProviderError("模型超时"))
  706. await runtime.process_update(
  707. {
  708. "update_id": 23,
  709. "business_message": customer_message(
  710. message_id=3, chat_id=603, sender_id=603
  711. ),
  712. }
  713. )
  714. provider_failed = await dbassistant.get_conversation_by_chat("conn-a", 603)
  715. assert provider_failed["handoff_reason"] == "provider_error"
  716. await runtime.process_update(
  717. {
  718. "update_id": 230,
  719. "business_message": customer_message(
  720. message_id=30, chat_id=605, sender_id=605
  721. ),
  722. }
  723. )
  724. quota_handoff = await dbassistant.get_conversation_by_chat("conn-a", 605)
  725. assert quota_handoff["handoff_reason"] == "account_daily_limit"
  726. sent_before = len(runtime.api.sent)
  727. await runtime.process_update(
  728. {
  729. "update_id": 24,
  730. "business_message": customer_message(
  731. message_id=4,
  732. chat_id=604,
  733. sender_id=604,
  734. sent_at=datetime.now(UTC) - timedelta(hours=25),
  735. ),
  736. }
  737. )
  738. expired = await dbassistant.get_conversation_by_chat("conn-a", 604)
  739. assert expired["handoff_reason"] == "reply_window_expired"
  740. new_sends = runtime.api.sent[sent_before:]
  741. assert len(new_sends) == 2
  742. assert all(item["business_connection_id"] == "" for item in new_sends)
  743. async def test_failed_update_goes_to_dead_letter_without_blocking(app_modules, monkeypatch):
  744. dbassistant, _service, runtime, _knowledge = await prepare_runtime(app_modules)
  745. async def fail(_message):
  746. raise RuntimeError("persistent failure")
  747. async def no_sleep(_seconds):
  748. return None
  749. runtime._handle_business_message = fail
  750. monkeypatch.setattr("wbb.services.business_assistant.asyncio.sleep", no_sleep)
  751. await runtime.process_update(
  752. {"update_id": 31, "business_message": customer_message(message_id=1)}
  753. )
  754. dead_letter = await app_modules.wbb.db.business_assistant_dead_letters.find_one(
  755. {"bot_id": "primary", "update_id": 31}
  756. )
  757. update = await app_modules.wbb.db.business_assistant_updates.find_one(
  758. {"bot_id": "primary", "update_id": 31}
  759. )
  760. assert dead_letter["error"] == "persistent failure"
  761. assert update["status"] == "done"
  762. assert await dbassistant.claim_update(31) is False
  763. async def test_connection_updates_revoke_permissions_and_persist_offset(app_modules):
  764. service = app_modules.load("wbb.services.business_assistant")
  765. dbassistant = app_modules.load("wbb.utils.dbassistant")
  766. runtime = service.BusinessAssistantRuntime(
  767. token="123456:" + "A" * 30,
  768. session=SimpleNamespace(),
  769. provider=FakeProvider(),
  770. )
  771. runtime.api = FakeBusinessApi()
  772. await runtime.process_update(
  773. {"update_id": 41, "business_connection": connection_payload(can_reply=True)}
  774. )
  775. connected = await dbassistant.get_business_connection("conn-a")
  776. assert connected["is_enabled"] is True
  777. assert connected["rights"]["can_reply"] is True
  778. await dbassistant.update_account_settings("conn-a", {"assistant_enabled": True})
  779. revoked = connection_payload(can_reply=False)
  780. revoked["is_enabled"] = False
  781. await runtime.process_update({"update_id": 42, "business_connection": revoked})
  782. disconnected = await dbassistant.get_business_connection("conn-a")
  783. assert disconnected["is_enabled"] is False
  784. assert "can_reply" not in disconnected["rights"]
  785. await runtime.process_update(
  786. {"update_id": 420, "business_message": customer_message(message_id=20)}
  787. )
  788. assert runtime.api.sent == []
  789. await dbassistant.save_update_offset(43)
  790. assert await dbassistant.load_update_offset() == 43
  791. await dbassistant.save_update_offset(44)
  792. assert await dbassistant.load_update_offset() == 44
  793. class StartupApi:
  794. def __init__(self, *, supported: bool, webhook_url: str = "") -> None:
  795. self.supported = supported
  796. self.webhook_url = webhook_url
  797. async def get_me(self):
  798. return {"username": "business_bot", "can_connect_to_business": self.supported}
  799. async def get_webhook_info(self):
  800. return {"url": self.webhook_url}
  801. async def test_runtime_blocks_webhook_conflict_and_disabled_business_mode(app_modules):
  802. service = app_modules.load("wbb.services.business_assistant")
  803. runtime = service.BusinessAssistantRuntime(
  804. token="123456:" + "A" * 30,
  805. session=SimpleNamespace(),
  806. provider=FakeProvider(),
  807. )
  808. runtime.api = StartupApi(supported=True, webhook_url="https://hooks.example.com/tg")
  809. webhook_status = await runtime.start()
  810. assert webhook_status["polling_state"] == "blocked"
  811. assert webhook_status["webhook_conflict"] is True
  812. assert runtime._poll_task is None
  813. runtime.api = StartupApi(supported=False)
  814. business_status = await runtime.start()
  815. assert business_status["polling_state"] == "blocked"
  816. assert business_status["business_mode_supported"] is False
  817. assert "BotFather" in business_status["last_error"]
  818. class FakeResponse:
  819. def __init__(self, status: int, payload: dict) -> None:
  820. self.status = status
  821. self.payload = payload
  822. self.headers: dict[str, str] = {}
  823. async def __aenter__(self):
  824. return self
  825. async def __aexit__(self, *_args):
  826. return None
  827. async def json(self, **_kwargs):
  828. return self.payload
  829. class RetrySession:
  830. def __init__(self, responses: list[FakeResponse]) -> None:
  831. self.responses = responses
  832. self.calls = 0
  833. def post(self, *_args, **_kwargs):
  834. response = self.responses[self.calls]
  835. self.calls += 1
  836. return response
  837. async def test_telegram_api_retries_429_and_5xx(app_modules, monkeypatch):
  838. service = app_modules.load("wbb.services.business_assistant")
  839. session = RetrySession(
  840. [
  841. FakeResponse(
  842. 429,
  843. {
  844. "ok": False,
  845. "error_code": 429,
  846. "description": "Too Many Requests",
  847. "parameters": {"retry_after": 3},
  848. },
  849. ),
  850. FakeResponse(
  851. 502,
  852. {"ok": False, "error_code": 502, "description": "Bad Gateway"},
  853. ),
  854. FakeResponse(200, {"ok": True, "result": {"id": 999}}),
  855. ]
  856. )
  857. sleeps: list[int] = []
  858. async def record_sleep(seconds):
  859. sleeps.append(seconds)
  860. monkeypatch.setattr("wbb.services.business_assistant.asyncio.sleep", record_sleep)
  861. api = service.TelegramBusinessApi("123456:" + "A" * 30, session)
  862. assert await api.get_me() == {"id": 999}
  863. assert session.calls == 3
  864. assert sleeps == [3, 2]
  865. async def test_feishu_webhook_signature_and_success_contract(app_modules):
  866. service = app_modules.load("wbb.services.business_assistant")
  867. response = FakeResponse(200, {"code": 0, "msg": "success"})
  868. class CaptureSession:
  869. def __init__(self) -> None:
  870. self.payload = None
  871. def post(self, _url, *, json, timeout):
  872. self.payload = json
  873. assert timeout.total == 8
  874. return response
  875. session = CaptureSession()
  876. client = service.FeishuWebhookClient(session)
  877. await client.send(
  878. "https://open.feishu.cn/open-apis/bot/v2/hook/12345678-abcd-4321-abcd-123456789abc",
  879. "消息提醒",
  880. signing_secret="secret",
  881. )
  882. assert session.payload["msg_type"] == "text"
  883. assert session.payload["content"]["text"] == "消息提醒"
  884. assert session.payload["timestamp"]
  885. assert session.payload["sign"]
  886. card_lines = ["账号:店主", "用户:测试用户"]
  887. await client.send_card(
  888. "https://open.feishu.cn/open-apis/bot/v2/hook/12345678-abcd-4321-abcd-123456789abc",
  889. title="消息提醒",
  890. lines=card_lines,
  891. link_url="https://t.me/test_user",
  892. )
  893. assert session.payload == {
  894. "msg_type": "interactive",
  895. "card": {
  896. "config": {"wide_screen_mode": True},
  897. "header": {
  898. "template": "blue",
  899. "title": {"tag": "plain_text", "content": "消息提醒"},
  900. },
  901. "elements": [
  902. {
  903. "tag": "div",
  904. "text": {"tag": "plain_text", "content": line},
  905. }
  906. for line in card_lines
  907. ]
  908. + [
  909. {"tag": "hr"},
  910. {
  911. "tag": "action",
  912. "actions": [
  913. {
  914. "tag": "button",
  915. "text": {
  916. "tag": "plain_text",
  917. "content": "打开 Telegram 对话",
  918. },
  919. "type": "primary",
  920. "url": "https://t.me/test_user",
  921. }
  922. ],
  923. },
  924. ],
  925. },
  926. }
  927. await client.send_card(
  928. "https://open.feishu.cn/open-apis/bot/v2/hook/12345678-abcd-4321-abcd-123456789abc",
  929. title="消息提醒",
  930. lines=["无链接"],
  931. )
  932. assert session.payload == {
  933. "msg_type": "interactive",
  934. "card": {
  935. "config": {"wide_screen_mode": True},
  936. "header": {
  937. "template": "blue",
  938. "title": {"tag": "plain_text", "content": "消息提醒"},
  939. },
  940. "elements": [
  941. {
  942. "tag": "div",
  943. "text": {"tag": "plain_text", "content": "无链接"},
  944. }
  945. ],
  946. },
  947. }
  948. async def _login_and_change_password(client: TestClient) -> str:
  949. login = await client.post(
  950. "/api/admin/v1/auth/login",
  951. json={"username": "admin", "password": "qwe0.123456"},
  952. )
  953. login_data = (await login.json())["data"]
  954. changed = await client.put(
  955. "/api/admin/v1/auth/password",
  956. headers={"X-CSRF-Token": login_data["csrf_token"]},
  957. json={
  958. "current_password": "qwe0.123456",
  959. "new_password": "changed-pass-123",
  960. },
  961. )
  962. return (await changed.json())["data"]["csrf_token"]
  963. async def test_business_assistant_api_permission_csrf_audit_and_clear(app_modules):
  964. dbassistant = app_modules.load("wbb.utils.dbassistant")
  965. await dbassistant.upsert_business_connection(connection_payload())
  966. conversation = await dbassistant.get_or_create_conversation(
  967. "conn-a", 701, customer={"id": 701, "first_name": "待清理客户"}
  968. )
  969. await dbassistant.append_conversation_message(
  970. conversation["conversation_id"],
  971. direction="incoming",
  972. telegram_message_id=1,
  973. text="请清理",
  974. )
  975. await dbassistant.update_conversation_summary(conversation["conversation_id"], "待清理摘要")
  976. admin_api = app_modules.load("wbb.admin.api")
  977. application = admin_api.build_admin_application()
  978. await application["admin_api"].initialize()
  979. client = TestClient(TestServer(application), cookie_jar=CookieJar(unsafe=True))
  980. await client.start_server()
  981. try:
  982. csrf = await _login_and_change_password(client)
  983. app_modules.wbb.BOT_PERMISSIONS = set()
  984. denied = await client.get("/api/admin/v1/business-assistant/status")
  985. assert denied.status == 403
  986. app_modules.wbb.BOT_PERMISSIONS = {"business_assistant.manage"}
  987. no_csrf = await client.put(
  988. "/api/admin/v1/business-assistant/settings",
  989. json={"connection_id": "conn-a", "assistant_enabled": True, "confirm": True},
  990. )
  991. assert no_csrf.status == 403
  992. unconfirmed = await client.put(
  993. "/api/admin/v1/business-assistant/settings",
  994. headers={"X-CSRF-Token": csrf},
  995. json={"connection_id": "conn-a", "assistant_enabled": True},
  996. )
  997. assert unconfirmed.status == 409
  998. saved = await client.put(
  999. "/api/admin/v1/business-assistant/settings",
  1000. headers={"X-CSRF-Token": csrf},
  1001. json={
  1002. "connection_id": "conn-a",
  1003. "assistant_enabled": True,
  1004. "feishu_webhook_enabled": True,
  1005. "feishu_webhook_url": "https://open.feishu.cn/open-apis/bot/v2/hook/12345678-abcd-4321-abcd-123456789abc",
  1006. "feishu_webhook_signing_secret": "api-secret",
  1007. "confirm": True,
  1008. },
  1009. )
  1010. assert saved.status == 200
  1011. saved_data = (await saved.json())["data"]
  1012. assert saved_data["assistant_enabled"] is True
  1013. assert saved_data["feishu_webhook_url"] == ""
  1014. assert saved_data["feishu_webhook_signing_secret"] == ""
  1015. assert saved_data["feishu_webhook_configured"] is True
  1016. stored_settings = await dbassistant.get_account_settings("conn-a")
  1017. assert stored_settings["feishu_webhook_signing_secret"] == "api-secret"
  1018. listed_connections = await client.get(
  1019. "/api/admin/v1/business-assistant/connections?page=1&page_size=100"
  1020. )
  1021. listed_settings = (await listed_connections.json())["data"]["items"][0][
  1022. "settings"
  1023. ]
  1024. assert listed_settings["feishu_webhook_url"] == ""
  1025. assert listed_settings["feishu_webhook_signing_secret"] == ""
  1026. class FakeRuntime:
  1027. def __init__(self) -> None:
  1028. self.calls = []
  1029. async def test_feishu_webhook(self, connection_id, overrides):
  1030. self.calls.append((connection_id, overrides))
  1031. return {"ok": True, "account": "店主"}
  1032. module = app_modules.load("wbb.modules.business_assistant")
  1033. fake_runtime = FakeRuntime()
  1034. module._runtime = fake_runtime
  1035. tested = await client.post(
  1036. "/api/admin/v1/business-assistant/feishu-webhook/test",
  1037. headers={"X-CSRF-Token": csrf},
  1038. json={"connection_id": "conn-a", "confirm": True},
  1039. )
  1040. assert tested.status == 200
  1041. assert (await tested.json())["data"]["ok"] is True
  1042. assert fake_runtime.calls == [("conn-a", {})]
  1043. created = await client.post(
  1044. "/api/admin/v1/business-assistant/knowledge",
  1045. headers={"X-CSRF-Token": csrf},
  1046. json={
  1047. "connection_id": "conn-a",
  1048. "question": "配送范围",
  1049. "keywords": ["配送"],
  1050. "answer": "仅限市区。",
  1051. "confirm": True,
  1052. },
  1053. )
  1054. assert created.status == 201
  1055. listed = await client.get(
  1056. "/api/admin/v1/business-assistant/knowledge?connection_id=conn-a"
  1057. )
  1058. assert (await listed.json())["data"]["total"] == 1
  1059. created_source = await client.post(
  1060. "/api/admin/v1/business-assistant/knowledge-sources",
  1061. headers={"X-CSRF-Token": csrf},
  1062. json={
  1063. "connection_id": "conn-a",
  1064. "source_type": "business",
  1065. "title": "人工对话",
  1066. "publication_mode": "review",
  1067. "confirm": True,
  1068. },
  1069. )
  1070. assert created_source.status == 201
  1071. source = (await created_source.json())["data"]
  1072. event, _ = await dbassistant.record_source_event(
  1073. source,
  1074. event_key="business:701:8",
  1075. content="市区可配送。",
  1076. metadata={"chat_id": 701, "message_id": 8, "trusted_author": True},
  1077. )
  1078. candidates = await dbassistant.replace_event_candidates(
  1079. source,
  1080. event,
  1081. [
  1082. {
  1083. "question": "配送范围?",
  1084. "keywords": ["配送"],
  1085. "answer": "市区可配送。",
  1086. "confidence": 0.9,
  1087. }
  1088. ],
  1089. auto_publish=False,
  1090. )
  1091. candidate_page = await client.get(
  1092. "/api/admin/v1/business-assistant/knowledge-candidates"
  1093. "?connection_id=conn-a&status=pending"
  1094. )
  1095. assert candidate_page.status == 200
  1096. assert (await candidate_page.json())["data"]["total"] == 1
  1097. approved = await client.post(
  1098. "/api/admin/v1/business-assistant/knowledge-candidates/"
  1099. f"{candidates[0]['candidate_id']}/actions",
  1100. headers={"X-CSRF-Token": csrf},
  1101. json={"action": "approve", "confirm": True},
  1102. )
  1103. assert approved.status == 200
  1104. assert (await approved.json())["data"]["status"] == "published"
  1105. cleared = await client.post(
  1106. f"/api/admin/v1/business-assistant/conversations/{conversation['conversation_id']}/actions",
  1107. headers={"X-CSRF-Token": csrf},
  1108. json={"action": "clear", "confirm": True},
  1109. )
  1110. assert cleared.status == 200
  1111. detail = await dbassistant.conversation_detail(conversation["conversation_id"])
  1112. assert detail["summary"] == ""
  1113. assert detail["messages"] == []
  1114. audit = await app_modules.wbb.db.admin_audit_logs.find_one(
  1115. {"action": "business_assistant.conversation.clear"}
  1116. )
  1117. assert audit["success"] is True
  1118. finally:
  1119. await client.close()