test_business_assistant.py 46 KB

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