business_assistant.py 57 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460
  1. from __future__ import annotations
  2. import asyncio
  3. import base64
  4. import hashlib
  5. import hmac
  6. import json
  7. import re
  8. from contextlib import suppress
  9. from datetime import UTC, datetime
  10. from typing import Any
  11. from aiohttp import ClientError, ClientSession, ClientTimeout
  12. import wbb
  13. from wbb.utils.dbassistant import (
  14. append_conversation_message,
  15. claim_feishu_message_notification,
  16. claim_update,
  17. classify_handoff,
  18. dead_letter_update,
  19. due_digest_connections,
  20. get_account_settings,
  21. get_business_connection,
  22. get_conversation,
  23. get_conversation_by_chat,
  24. get_knowledge_candidate,
  25. get_knowledge_source,
  26. get_or_create_conversation,
  27. get_source_event,
  28. get_source_event_by_key,
  29. invalidate_source_event_candidates,
  30. list_business_connections,
  31. list_business_sources,
  32. list_enabled_business_sources,
  33. load_update_offset,
  34. mark_digest_sent,
  35. mark_feishu_message_notification,
  36. mark_source_event_deleted,
  37. mark_update_done,
  38. mark_update_failed,
  39. match_knowledge,
  40. normalize_account_settings,
  41. pause_conversation_for_human,
  42. recent_conversation_messages,
  43. record_source_event,
  44. replace_event_candidates,
  45. reserve_ai_usage,
  46. resume_conversation,
  47. runtime_status,
  48. save_update_offset,
  49. set_conversation_handoff,
  50. touch_business_connection,
  51. update_conversation_summary,
  52. update_runtime_status,
  53. update_source_event_extraction,
  54. upsert_business_connection,
  55. usage_metrics,
  56. utc_now,
  57. )
  58. BUSINESS_ALLOWED_UPDATES = [
  59. "business_connection",
  60. "business_message",
  61. "edited_business_message",
  62. "deleted_business_messages",
  63. ]
  64. BUSINESS_REPLY_WINDOW_SECONDS = 24 * 60 * 60
  65. class TelegramBotApiError(RuntimeError):
  66. def __init__(
  67. self,
  68. message: str,
  69. *,
  70. error_code: int = 0,
  71. retry_after: int = 0,
  72. ) -> None:
  73. super().__init__(message)
  74. self.error_code = int(error_code or 0)
  75. self.retry_after = int(retry_after or 0)
  76. class AssistantProviderError(RuntimeError):
  77. pass
  78. class FeishuWebhookError(RuntimeError):
  79. pass
  80. class FeishuWebhookClient:
  81. def __init__(self, session: ClientSession) -> None:
  82. self._session = session
  83. @staticmethod
  84. def _signature(secret: str, timestamp: int) -> str:
  85. sign_key = f"{timestamp}\n{secret}".encode()
  86. digest = hmac.new(sign_key, digestmod=hashlib.sha256).digest()
  87. return base64.b64encode(digest).decode()
  88. async def send(self, webhook_url: str, text: str, *, signing_secret: str = "") -> None:
  89. payload: dict[str, Any] = {
  90. "msg_type": "text",
  91. "content": {"text": str(text)[:4000]},
  92. }
  93. await self._send_payload(
  94. webhook_url, payload, signing_secret=signing_secret
  95. )
  96. async def send_post(
  97. self,
  98. webhook_url: str,
  99. *,
  100. title: str,
  101. content: list[list[dict[str, str]]],
  102. signing_secret: str = "",
  103. ) -> None:
  104. payload: dict[str, Any] = {
  105. "msg_type": "post",
  106. "content": {
  107. "post": {
  108. "zh_cn": {
  109. "title": str(title)[:100],
  110. "content": content,
  111. }
  112. }
  113. },
  114. }
  115. await self._send_payload(
  116. webhook_url, payload, signing_secret=signing_secret
  117. )
  118. async def _send_payload(
  119. self,
  120. webhook_url: str,
  121. payload: dict[str, Any],
  122. *,
  123. signing_secret: str = "",
  124. ) -> None:
  125. if signing_secret:
  126. timestamp = int(utc_now().timestamp())
  127. payload.update(
  128. {
  129. "timestamp": str(timestamp),
  130. "sign": self._signature(signing_secret, timestamp),
  131. }
  132. )
  133. last_error: FeishuWebhookError | None = None
  134. for attempt in range(3):
  135. try:
  136. async with self._session.post(
  137. webhook_url,
  138. json=payload,
  139. timeout=ClientTimeout(total=8),
  140. ) as response:
  141. data = await response.json(content_type=None)
  142. except (ClientError, TimeoutError, ValueError) as exc:
  143. last_error = FeishuWebhookError("无法连接飞书群机器人 Webhook。")
  144. if attempt < 2:
  145. await asyncio.sleep(2**attempt)
  146. continue
  147. raise last_error from exc
  148. raw_code = (
  149. data.get("code", data.get("StatusCode", -1))
  150. if isinstance(data, dict)
  151. else -1
  152. )
  153. try:
  154. code = int(raw_code)
  155. except (TypeError, ValueError):
  156. code = -1
  157. if response.status == 200 and code == 0:
  158. return
  159. last_error = FeishuWebhookError(
  160. f"飞书群机器人 Webhook 请求失败(HTTP {response.status},code={code})。"
  161. )
  162. if (response.status == 429 or response.status >= 500) and attempt < 2:
  163. retry_after = response.headers.get("Retry-After", "")
  164. try:
  165. wait_seconds = max(1, min(int(retry_after), 30))
  166. except (TypeError, ValueError):
  167. wait_seconds = 2**attempt
  168. await asyncio.sleep(wait_seconds)
  169. continue
  170. raise last_error
  171. raise last_error or FeishuWebhookError("飞书群机器人 Webhook 请求失败。")
  172. class TelegramBusinessApi:
  173. def __init__(self, token: str, session: ClientSession) -> None:
  174. self._base_url = f"https://api.telegram.org/bot{token}"
  175. self._session = session
  176. async def call(
  177. self,
  178. method: str,
  179. payload: dict[str, Any] | None = None,
  180. ) -> Any:
  181. last_error: TelegramBotApiError | None = None
  182. for attempt in range(3):
  183. try:
  184. async with self._session.post(
  185. f"{self._base_url}/{method}", json=payload or {}
  186. ) as response:
  187. data = await response.json(content_type=None)
  188. except (ClientError, TimeoutError, ValueError) as exc:
  189. last_error = TelegramBotApiError("无法连接 Telegram Bot API。")
  190. if attempt < 2:
  191. await asyncio.sleep(2**attempt)
  192. continue
  193. raise last_error from exc
  194. if response.status == 200 and isinstance(data, dict) and data.get("ok"):
  195. return data.get("result")
  196. parameters = data.get("parameters") if isinstance(data, dict) else {}
  197. last_error = TelegramBotApiError(
  198. str(data.get("description") or "Telegram Bot API 请求失败。")
  199. if isinstance(data, dict)
  200. else "Telegram Bot API 请求失败。",
  201. error_code=int(data.get("error_code") or response.status)
  202. if isinstance(data, dict)
  203. else response.status,
  204. retry_after=int((parameters or {}).get("retry_after") or 0),
  205. )
  206. retryable = last_error.error_code == 429 or response.status >= 500
  207. if retryable and attempt < 2:
  208. await asyncio.sleep(last_error.retry_after or 2**attempt)
  209. continue
  210. raise last_error
  211. raise last_error or TelegramBotApiError("Telegram Bot API 请求失败。")
  212. async def get_me(self) -> dict[str, Any]:
  213. result = await self.call("getMe")
  214. return result if isinstance(result, dict) else {}
  215. async def get_webhook_info(self) -> dict[str, Any]:
  216. result = await self.call("getWebhookInfo")
  217. return result if isinstance(result, dict) else {}
  218. async def get_business_connection(self, connection_id: str) -> dict[str, Any]:
  219. result = await self.call(
  220. "getBusinessConnection", {"business_connection_id": str(connection_id)}
  221. )
  222. return result if isinstance(result, dict) else {}
  223. async def get_updates(
  224. self, *, offset: int, poll_timeout: int = 30
  225. ) -> list[dict[str, Any]]:
  226. result = await self.call(
  227. "getUpdates",
  228. {
  229. "offset": int(offset),
  230. "timeout": int(poll_timeout),
  231. "limit": 100,
  232. "allowed_updates": BUSINESS_ALLOWED_UPDATES,
  233. },
  234. )
  235. return [item for item in (result or []) if isinstance(item, dict)]
  236. async def send_message(
  237. self,
  238. chat_id: int,
  239. text: str,
  240. *,
  241. business_connection_id: str = "",
  242. reply_markup: dict[str, Any] | None = None,
  243. ) -> dict[str, Any]:
  244. payload: dict[str, Any] = {
  245. "chat_id": int(chat_id),
  246. "text": str(text)[:4096],
  247. }
  248. if business_connection_id:
  249. payload["business_connection_id"] = str(business_connection_id)
  250. if reply_markup:
  251. payload["reply_markup"] = reply_markup
  252. result = await self.call("sendMessage", payload)
  253. return result if isinstance(result, dict) else {}
  254. async def read_business_message(
  255. self,
  256. connection_id: str,
  257. chat_id: int,
  258. message_id: int,
  259. ) -> None:
  260. await self.call(
  261. "readBusinessMessage",
  262. {
  263. "business_connection_id": str(connection_id),
  264. "chat_id": int(chat_id),
  265. "message_id": int(message_id),
  266. },
  267. )
  268. def _strip_json_fence(value: str) -> str:
  269. text = value.strip()
  270. if text.startswith("```"):
  271. text = re.sub(r"^```(?:json)?\s*", "", text, flags=re.IGNORECASE)
  272. text = re.sub(r"\s*```$", "", text)
  273. return text.strip()
  274. def business_reply_window_open(
  275. message: dict[str, Any], *, now: datetime | None = None
  276. ) -> bool:
  277. try:
  278. sent_at = datetime.fromtimestamp(int(message.get("date") or 0), UTC)
  279. except (OSError, OverflowError, TypeError, ValueError):
  280. return False
  281. if sent_at.timestamp() <= 0:
  282. return False
  283. age = (now or utc_now()) - sent_at
  284. return age.total_seconds() <= BUSINESS_REPLY_WINDOW_SECONDS
  285. def parse_ai_decision(value: str, *, allowed_entry_ids: set[str]) -> dict[str, Any]:
  286. try:
  287. payload = json.loads(_strip_json_fence(value))
  288. except (json.JSONDecodeError, TypeError) as exc:
  289. raise AssistantProviderError("模型没有返回有效的 JSON 结果。") from exc
  290. if not isinstance(payload, dict):
  291. raise AssistantProviderError("模型结果必须是 JSON 对象。")
  292. action = str(payload.get("action") or "")
  293. if action not in {"answer", "clarify", "handoff"}:
  294. raise AssistantProviderError("模型返回了不支持的接待动作。")
  295. reply = str(payload.get("reply") or "").strip()
  296. if action in {"answer", "clarify"} and not reply:
  297. raise AssistantProviderError("模型没有提供回复内容。")
  298. if len(reply) > 4000:
  299. raise AssistantProviderError("模型回复内容过长。")
  300. matched_ids = [
  301. str(item)
  302. for item in payload.get("matched_entry_ids") or []
  303. if str(item) in allowed_entry_ids
  304. ]
  305. if action == "answer" and not matched_ids:
  306. raise AssistantProviderError("业务回答没有引用知识条目。")
  307. return {
  308. "action": action,
  309. "reply": reply,
  310. "handoff_reason": str(payload.get("handoff_reason") or "ai_handoff")[:200],
  311. "matched_entry_ids": matched_ids,
  312. "summary": str(payload.get("summary") or "")[:4000],
  313. }
  314. def _short_string_list(value: Any, *, items: int, length: int) -> list[str]:
  315. raw_items = value if isinstance(value, list) else []
  316. result: list[str] = []
  317. for item in raw_items:
  318. text = " ".join(str(item or "").strip().split())[:length]
  319. if text and text.casefold() not in {existing.casefold() for existing in result}:
  320. result.append(text)
  321. return result[:items]
  322. def parse_knowledge_extraction(value: str) -> list[dict[str, Any]]:
  323. try:
  324. payload = json.loads(_strip_json_fence(value))
  325. except (json.JSONDecodeError, TypeError) as exc:
  326. raise AssistantProviderError("模型没有返回有效的知识提取 JSON。") from exc
  327. if not isinstance(payload, dict) or not isinstance(payload.get("items"), list):
  328. raise AssistantProviderError("知识提取结果必须包含 items 数组。")
  329. result: list[dict[str, Any]] = []
  330. for raw_item in payload["items"][:5]:
  331. if not isinstance(raw_item, dict):
  332. raise AssistantProviderError("知识候选必须是 JSON 对象。")
  333. question = " ".join(str(raw_item.get("question") or "").strip().split())[:300]
  334. answer = str(raw_item.get("answer") or "").strip()[:4000]
  335. if not question or not answer:
  336. raise AssistantProviderError("知识候选缺少标准问题或答案。")
  337. try:
  338. confidence = float(raw_item.get("confidence") or 0)
  339. except (TypeError, ValueError):
  340. confidence = 0
  341. result.append(
  342. {
  343. "question": question,
  344. "aliases": _short_string_list(
  345. raw_item.get("aliases"), items=20, length=200
  346. ),
  347. "keywords": _short_string_list(
  348. raw_item.get("keywords"), items=30, length=50
  349. ),
  350. "answer": answer,
  351. "tags": _short_string_list(raw_item.get("tags"), items=20, length=30),
  352. "confidence": max(0.0, min(confidence, 1.0)),
  353. }
  354. )
  355. return result
  356. def redact_business_knowledge_text(value: str) -> str:
  357. text = re.sub(
  358. r"(?<![\w.])[A-Z0-9._%+-]+@[A-Z0-9.-]+\.[A-Z]{2,}(?!\w)",
  359. "[邮箱已脱敏]",
  360. str(value),
  361. flags=re.IGNORECASE,
  362. )
  363. text = re.sub(r"(?<!\d)(?:\+?\d[\d\s-]{7,}\d)(?!\d)", "[电话已脱敏]", text)
  364. return re.sub(r"(?<!\w)@[A-Za-z][A-Za-z0-9_]{4,31}", "[用户名已脱敏]", text)
  365. class OpenAICompatibleAssistant:
  366. def __init__(
  367. self,
  368. session: ClientSession,
  369. *,
  370. base_url: str,
  371. api_key: str,
  372. model: str,
  373. timeout_seconds: int = 30,
  374. max_output_tokens: int = 600,
  375. ) -> None:
  376. self._session = session
  377. self.base_url = str(base_url).rstrip("/")
  378. self.api_key = str(api_key)
  379. self.model = str(model)
  380. self.timeout_seconds = max(5, min(int(timeout_seconds), 120))
  381. self.max_output_tokens = max(100, min(int(max_output_tokens), 2000))
  382. @property
  383. def configured(self) -> bool:
  384. return bool(self.base_url and self.api_key and self.model)
  385. async def decide(
  386. self,
  387. *,
  388. settings: dict[str, Any],
  389. conversation: dict[str, Any],
  390. messages: list[dict[str, Any]],
  391. knowledge: list[dict[str, Any]],
  392. customer_text: str,
  393. ) -> dict[str, Any]:
  394. if not self.configured:
  395. raise AssistantProviderError("OpenAI 兼容模型尚未配置。")
  396. knowledge_payload = [
  397. {
  398. "entry_id": item["entry_id"],
  399. "question": item["question"],
  400. "answer": item["answer"],
  401. "tags": item.get("tags") or [],
  402. }
  403. for item in knowledge
  404. ]
  405. recent_payload = [
  406. {"direction": item.get("direction"), "text": item.get("text", "")}
  407. for item in messages[-12:]
  408. ]
  409. system_prompt = (
  410. f"{settings.get('system_prompt')}\n"
  411. "必须遵守:业务事实只能来自 knowledge;不得自行补全价格、承诺、退款或政策。"
  412. "仅输出一个 JSON 对象,不使用 Markdown。字段为 action、reply、handoff_reason、"
  413. "matched_entry_ids、summary。action 只能是 answer、clarify、handoff。"
  414. "answer 必须填写实际引用的 matched_entry_ids;无法可靠回答时使用 handoff。"
  415. )
  416. request_body = {
  417. "model": self.model,
  418. "temperature": 0.2,
  419. "max_tokens": self.max_output_tokens,
  420. "messages": [
  421. {"role": "system", "content": system_prompt},
  422. {
  423. "role": "user",
  424. "content": json.dumps(
  425. {
  426. "language": settings.get("language"),
  427. "tone": settings.get("tone"),
  428. "existing_summary": conversation.get("summary", ""),
  429. "recent_messages": recent_payload,
  430. "knowledge": knowledge_payload,
  431. "customer_message": customer_text,
  432. },
  433. ensure_ascii=False,
  434. ),
  435. },
  436. ],
  437. }
  438. try:
  439. async with asyncio.timeout(self.timeout_seconds):
  440. async with self._session.post(
  441. f"{self.base_url}/chat/completions",
  442. headers={
  443. "Authorization": f"Bearer {self.api_key}",
  444. "Content-Type": "application/json",
  445. },
  446. json=request_body,
  447. ) as response:
  448. payload = await response.json(content_type=None)
  449. except (ClientError, TimeoutError, ValueError) as exc:
  450. raise AssistantProviderError("OpenAI 兼容服务暂时不可用。") from exc
  451. if response.status != 200:
  452. raise AssistantProviderError(f"OpenAI 兼容服务返回 HTTP {response.status}。")
  453. try:
  454. content = payload["choices"][0]["message"]["content"]
  455. except (KeyError, IndexError, TypeError) as exc:
  456. raise AssistantProviderError("OpenAI 兼容服务响应格式无效。") from exc
  457. return parse_ai_decision(
  458. str(content),
  459. allowed_entry_ids={str(item["entry_id"]) for item in knowledge},
  460. )
  461. async def extract_knowledge(
  462. self,
  463. *,
  464. source_type: str,
  465. source_title: str,
  466. content: str,
  467. question_context: str = "",
  468. ) -> list[dict[str, Any]]:
  469. if not self.configured:
  470. raise AssistantProviderError("OpenAI 兼容模型尚未配置。")
  471. if source_type == "business":
  472. safe_content = redact_business_knowledge_text(content)
  473. safe_question = redact_business_knowledge_text(question_context)
  474. source_rule = (
  475. "这是 Business 人工接待对话。customer_question 只能作为标准问题参考,"
  476. "答案事实只能来自 human_reply。不得把客户自述、联系方式、用户名、"
  477. "投诉、退款要求或未经确认的承诺写入知识。"
  478. )
  479. source_payload = {
  480. "customer_question": safe_question,
  481. "human_reply": safe_content,
  482. }
  483. else:
  484. source_rule = (
  485. "这是频道或群组内容。只提取明确、稳定、可复用的业务事实;"
  486. "忽略寒暄、广告口号、个人观点、临时活动、投诉、退款争议和未经确认的承诺。"
  487. "没有稳定事实时返回空 items。"
  488. )
  489. source_payload = {"content": content}
  490. system_prompt = (
  491. "你是知识库整理器,只提取输入中明确存在的事实,不得推测或补全。"
  492. f"{source_rule}"
  493. "只输出一个 JSON 对象,不使用 Markdown,格式为 "
  494. '{"items":[{"question":"","aliases":[],"keywords":[],"answer":"",'
  495. '"tags":[],"confidence":0.0}]}。items 最多 5 条。'
  496. )
  497. request_body = {
  498. "model": self.model,
  499. "temperature": 0.1,
  500. "max_tokens": self.max_output_tokens,
  501. "messages": [
  502. {"role": "system", "content": system_prompt},
  503. {
  504. "role": "user",
  505. "content": json.dumps(
  506. {
  507. "source_type": source_type,
  508. "source_title": source_title,
  509. **source_payload,
  510. },
  511. ensure_ascii=False,
  512. ),
  513. },
  514. ],
  515. }
  516. try:
  517. async with asyncio.timeout(self.timeout_seconds):
  518. async with self._session.post(
  519. f"{self.base_url}/chat/completions",
  520. headers={
  521. "Authorization": f"Bearer {self.api_key}",
  522. "Content-Type": "application/json",
  523. },
  524. json=request_body,
  525. ) as response:
  526. payload = await response.json(content_type=None)
  527. except (ClientError, TimeoutError, ValueError) as exc:
  528. raise AssistantProviderError("知识提取模型暂时不可用。") from exc
  529. if response.status != 200:
  530. raise AssistantProviderError(f"知识提取模型返回 HTTP {response.status}。")
  531. try:
  532. result = payload["choices"][0]["message"]["content"]
  533. except (KeyError, IndexError, TypeError) as exc:
  534. raise AssistantProviderError("知识提取模型响应格式无效。") from exc
  535. return parse_knowledge_extraction(str(result))
  536. async def test_connection(self) -> dict[str, Any]:
  537. if not self.configured:
  538. raise AssistantProviderError("OpenAI 兼容模型尚未配置。")
  539. request_body = {
  540. "model": self.model,
  541. "temperature": 0,
  542. "max_tokens": 12,
  543. "messages": [{"role": "user", "content": "只回复 OK"}],
  544. }
  545. try:
  546. async with asyncio.timeout(self.timeout_seconds):
  547. async with self._session.post(
  548. f"{self.base_url}/chat/completions",
  549. headers={
  550. "Authorization": f"Bearer {self.api_key}",
  551. "Content-Type": "application/json",
  552. },
  553. json=request_body,
  554. ) as response:
  555. payload = await response.json(content_type=None)
  556. except (ClientError, TimeoutError, ValueError) as exc:
  557. raise AssistantProviderError("无法连接 OpenAI 兼容服务。") from exc
  558. if response.status != 200:
  559. raise AssistantProviderError(f"模型测试失败:HTTP {response.status}。")
  560. return {"ok": True, "model": self.model, "response_id": payload.get("id", "")}
  561. def provider_from_wbb(session: ClientSession) -> OpenAICompatibleAssistant:
  562. return OpenAICompatibleAssistant(
  563. session,
  564. base_url=str(getattr(wbb, "BUSINESS_ASSISTANT_OPENAI_BASE_URL", "")),
  565. api_key=str(getattr(wbb, "BUSINESS_ASSISTANT_OPENAI_API_KEY", "")),
  566. model=str(getattr(wbb, "BUSINESS_ASSISTANT_OPENAI_MODEL", "")),
  567. timeout_seconds=int(
  568. getattr(wbb, "BUSINESS_ASSISTANT_OPENAI_TIMEOUT_SECONDS", 30)
  569. ),
  570. max_output_tokens=int(
  571. getattr(wbb, "BUSINESS_ASSISTANT_OPENAI_MAX_OUTPUT_TOKENS", 600)
  572. ),
  573. )
  574. class KnowledgeIngestionService:
  575. def __init__(self, provider: OpenAICompatibleAssistant) -> None:
  576. self.provider = provider
  577. async def ingest(
  578. self,
  579. source: dict[str, Any],
  580. *,
  581. event_key: str,
  582. content: str,
  583. metadata: dict[str, Any],
  584. question_context: str = "",
  585. auto_publish: bool = False,
  586. force: bool = False,
  587. ) -> list[dict[str, Any]]:
  588. event, changed = await record_source_event(
  589. source,
  590. event_key=event_key,
  591. content=content,
  592. metadata={
  593. **metadata,
  594. "question_context": question_context,
  595. "trusted_author": bool(metadata.get("trusted_author")),
  596. },
  597. )
  598. if not changed and not force:
  599. return []
  600. try:
  601. items = await self.provider.extract_knowledge(
  602. source_type=str(metadata.get("effective_source_type") or source["source_type"]),
  603. source_title=str(source.get("title") or ""),
  604. content=content,
  605. question_context=question_context,
  606. )
  607. return await replace_event_candidates(
  608. source, event, items, auto_publish=auto_publish
  609. )
  610. except Exception as exc:
  611. await invalidate_source_event_candidates(
  612. str(event["event_id"]), reason="extraction_failed"
  613. )
  614. await update_source_event_extraction(
  615. str(event["event_id"]), status="failed", error=str(exc)
  616. )
  617. raise
  618. async def regenerate_candidate(self, candidate_id: str) -> list[dict[str, Any]]:
  619. candidate = await get_knowledge_candidate(candidate_id)
  620. if not candidate:
  621. raise AssistantProviderError("未找到知识候选。")
  622. source = await get_knowledge_source(str(candidate["source_id"]))
  623. event = await get_source_event(str(candidate["event_id"]))
  624. if not source or not event or event.get("deleted"):
  625. raise AssistantProviderError("知识来源或原始事件已失效。")
  626. metadata = event.get("metadata") if isinstance(event.get("metadata"), dict) else {}
  627. auto_publish = bool(
  628. source.get("publication_mode") == "auto" and metadata.get("trusted_author")
  629. )
  630. return await self.ingest(
  631. source,
  632. event_key=str(event["event_key"]),
  633. content=str(event.get("content") or ""),
  634. metadata=metadata,
  635. question_context=str(metadata.get("question_context") or ""),
  636. auto_publish=auto_publish,
  637. force=True,
  638. )
  639. class BusinessAssistantRuntime:
  640. def __init__(
  641. self,
  642. *,
  643. token: str,
  644. session: ClientSession,
  645. provider: OpenAICompatibleAssistant | None = None,
  646. ) -> None:
  647. self.api = TelegramBusinessApi(token, session)
  648. self.feishu = FeishuWebhookClient(session)
  649. self.provider = provider or provider_from_wbb(session)
  650. self.ingestion = KnowledgeIngestionService(self.provider)
  651. self._poll_task: asyncio.Task[None] | None = None
  652. self._digest_task: asyncio.Task[None] | None = None
  653. self._stopping = asyncio.Event()
  654. async def start(self) -> dict[str, Any]:
  655. if self._poll_task and not self._poll_task.done():
  656. return await runtime_status()
  657. self._stopping.clear()
  658. await update_runtime_status(
  659. {
  660. "polling_state": "starting",
  661. "model_configured": self.provider.configured,
  662. "last_error": "",
  663. }
  664. )
  665. try:
  666. me = await self.api.get_me()
  667. webhook = await self.api.get_webhook_info()
  668. except TelegramBotApiError as exc:
  669. return await update_runtime_status(
  670. {"polling_state": "error", "last_error": str(exc)}
  671. )
  672. webhook_url = str(webhook.get("url") or "")
  673. supported = bool(me.get("can_connect_to_business"))
  674. await update_runtime_status(
  675. {
  676. "business_mode_supported": supported,
  677. "webhook_conflict": bool(webhook_url),
  678. "webhook_url": webhook_url,
  679. "bot_username": str(me.get("username") or ""),
  680. }
  681. )
  682. if webhook_url:
  683. return await update_runtime_status(
  684. {
  685. "polling_state": "blocked",
  686. "last_error": "检测到 Telegram webhook,Business 长轮询未启动。",
  687. }
  688. )
  689. if not supported:
  690. return await update_runtime_status(
  691. {
  692. "polling_state": "blocked",
  693. "last_error": "请先在 BotFather 为机器人开启 Business Mode。",
  694. }
  695. )
  696. await self._refresh_connections()
  697. self._poll_task = asyncio.create_task(
  698. self._poll_loop(), name="business-assistant-poller"
  699. )
  700. self._digest_task = asyncio.create_task(
  701. self._digest_loop(), name="business-assistant-digest"
  702. )
  703. return await update_runtime_status(
  704. {"polling_state": "running", "last_error": "", "started_at": utc_now()}
  705. )
  706. async def _refresh_connections(self) -> None:
  707. page = 1
  708. while True:
  709. connections, total = await list_business_connections(
  710. page=page, page_size=100
  711. )
  712. for connection in connections:
  713. connection_id = str(connection["connection_id"])
  714. try:
  715. payload = await self.api.get_business_connection(connection_id)
  716. await upsert_business_connection(payload)
  717. except TelegramBotApiError as exc:
  718. await touch_business_connection(
  719. connection_id,
  720. error=str(exc),
  721. is_enabled=False if exc.error_code in {400, 403} else None,
  722. )
  723. if not connections or page * 100 >= total:
  724. return
  725. page += 1
  726. async def stop(self) -> None:
  727. self._stopping.set()
  728. tasks = [task for task in (self._poll_task, self._digest_task) if task]
  729. for task in tasks:
  730. task.cancel()
  731. if tasks:
  732. await asyncio.gather(*tasks, return_exceptions=True)
  733. self._poll_task = None
  734. self._digest_task = None
  735. await update_runtime_status(
  736. {"polling_state": "stopped", "stopped_at": utc_now()}
  737. )
  738. async def _poll_loop(self) -> None:
  739. offset = await load_update_offset()
  740. backoff = 1
  741. while not self._stopping.is_set():
  742. try:
  743. updates = await self.api.get_updates(offset=offset, poll_timeout=30)
  744. backoff = 1
  745. await update_runtime_status(
  746. {"polling_state": "running", "last_poll_at": utc_now(), "last_error": ""}
  747. )
  748. for update in updates:
  749. update_id = int(update.get("update_id") or 0)
  750. if update_id <= 0:
  751. continue
  752. await self.process_update(update)
  753. offset = max(offset, update_id + 1)
  754. await save_update_offset(offset)
  755. except asyncio.CancelledError:
  756. raise
  757. except TelegramBotApiError as exc:
  758. wait_for = exc.retry_after or backoff
  759. await update_runtime_status(
  760. {"polling_state": "retrying", "last_error": str(exc)}
  761. )
  762. await asyncio.sleep(min(max(wait_for, 1), 60))
  763. backoff = min(backoff * 2, 60)
  764. except Exception as exc:
  765. await update_runtime_status(
  766. {"polling_state": "retrying", "last_error": str(exc)[:1000]}
  767. )
  768. await asyncio.sleep(backoff)
  769. backoff = min(backoff * 2, 60)
  770. async def process_update(self, update: dict[str, Any]) -> None:
  771. update_id = int(update.get("update_id") or 0)
  772. if not await claim_update(update_id):
  773. return
  774. last_error = ""
  775. for attempt in range(1, 4):
  776. try:
  777. if isinstance(update.get("business_connection"), dict):
  778. await upsert_business_connection(update["business_connection"])
  779. elif isinstance(update.get("business_message"), dict):
  780. await self._handle_business_message(update["business_message"])
  781. elif isinstance(update.get("edited_business_message"), dict):
  782. await self._handle_edited_message(update["edited_business_message"])
  783. elif isinstance(update.get("deleted_business_messages"), dict):
  784. await self._handle_deleted_messages(update["deleted_business_messages"])
  785. await mark_update_done(update_id)
  786. return
  787. except TelegramBotApiError as exc:
  788. last_error = str(exc)
  789. if exc.retry_after:
  790. await asyncio.sleep(min(exc.retry_after, 60))
  791. except Exception as exc:
  792. last_error = str(exc)
  793. if attempt < 3:
  794. await asyncio.sleep(attempt)
  795. await mark_update_failed(update_id, last_error)
  796. await dead_letter_update(update, last_error or "unknown update error")
  797. async def _resolve_connection(self, connection_id: str) -> dict[str, Any]:
  798. connection = await get_business_connection(connection_id)
  799. if connection:
  800. return connection
  801. payload = await self.api.get_business_connection(connection_id)
  802. return await upsert_business_connection(payload)
  803. async def _handle_business_message(self, message: dict[str, Any]) -> None:
  804. connection_id = str(message.get("business_connection_id") or "")
  805. if not connection_id or message.get("is_from_offline"):
  806. return
  807. if message.get("sender_business_bot"):
  808. return
  809. connection = await self._resolve_connection(connection_id)
  810. await touch_business_connection(connection_id)
  811. chat = message.get("chat") if isinstance(message.get("chat"), dict) else {}
  812. sender = message.get("from") if isinstance(message.get("from"), dict) else {}
  813. chat_id = int(chat.get("id") or 0)
  814. message_id = int(message.get("message_id") or 0)
  815. if not chat_id or not message_id:
  816. return
  817. owner_id = int((connection.get("user") or {}).get("id") or 0)
  818. owner_reply = int(sender.get("id") or 0) == owner_id
  819. conversation = await get_or_create_conversation(
  820. connection_id,
  821. chat_id,
  822. customer=None if owner_reply else sender,
  823. )
  824. if owner_reply:
  825. settings = await get_account_settings(connection_id)
  826. human_text = str(message.get("text") or message.get("caption") or "").strip()
  827. await append_conversation_message(
  828. conversation["conversation_id"],
  829. direction="human",
  830. telegram_message_id=message_id,
  831. text=human_text,
  832. sender_id=owner_id,
  833. )
  834. await pause_conversation_for_human(
  835. conversation["conversation_id"],
  836. hours=int(settings.get("human_pause_hours") or 24),
  837. )
  838. if human_text:
  839. await self._learn_from_human_reply(
  840. connection_id=connection_id,
  841. conversation_id=str(conversation["conversation_id"]),
  842. chat_id=chat_id,
  843. message_id=message_id,
  844. owner_id=owner_id,
  845. human_text=human_text,
  846. )
  847. return
  848. await self._handle_customer_message(
  849. connection=connection,
  850. conversation=conversation,
  851. message=message,
  852. sender=sender,
  853. )
  854. async def _learn_from_human_reply(
  855. self,
  856. *,
  857. connection_id: str,
  858. conversation_id: str,
  859. chat_id: int,
  860. message_id: int,
  861. owner_id: int,
  862. human_text: str,
  863. ) -> None:
  864. sources = await list_enabled_business_sources(connection_id)
  865. if not sources:
  866. return
  867. messages = await recent_conversation_messages(conversation_id, limit=20)
  868. customer_question = next(
  869. (
  870. str(item.get("text") or "")
  871. for item in reversed(messages)
  872. if item.get("direction") == "incoming"
  873. and str(item.get("text") or "") != "[非文本消息]"
  874. ),
  875. "",
  876. )
  877. if not customer_question:
  878. return
  879. for source in sources:
  880. event_key = f"business:{chat_id}:{message_id}"
  881. existing_event = await get_source_event_by_key(
  882. str(source["source_id"]), event_key
  883. )
  884. existing_metadata = (
  885. existing_event.get("metadata")
  886. if existing_event and isinstance(existing_event.get("metadata"), dict)
  887. else {}
  888. )
  889. source_question = str(
  890. existing_metadata.get("question_context") or customer_question
  891. )
  892. if any(
  893. term in f"{source_question}\n{human_text}".casefold()
  894. for term in ("承诺", "保证", "投诉", "退款", "退钱", "赔偿", "律师", "起诉")
  895. ):
  896. continue
  897. try:
  898. await self.ingestion.ingest(
  899. source,
  900. event_key=event_key,
  901. content=human_text,
  902. question_context=source_question,
  903. metadata={
  904. "effective_source_type": "business",
  905. "chat_id": int(chat_id),
  906. "message_id": int(message_id),
  907. "author_id": int(owner_id),
  908. "trusted_author": True,
  909. },
  910. auto_publish=source.get("publication_mode") == "auto",
  911. )
  912. except Exception as exc:
  913. await update_runtime_status(
  914. {
  915. "knowledge_ingestion_last_error": str(exc)[:1000],
  916. "knowledge_ingestion_last_error_at": utc_now(),
  917. }
  918. )
  919. async def _handle_customer_message(
  920. self,
  921. *,
  922. connection: dict[str, Any],
  923. conversation: dict[str, Any],
  924. message: dict[str, Any],
  925. sender: dict[str, Any],
  926. ) -> None:
  927. connection_id = str(connection["connection_id"])
  928. chat_id = int((message.get("chat") or {}).get("id") or 0)
  929. message_id = int(message.get("message_id") or 0)
  930. text = str(message.get("text") or "").strip()
  931. settings = await get_account_settings(connection_id)
  932. await append_conversation_message(
  933. conversation["conversation_id"],
  934. direction="incoming",
  935. telegram_message_id=message_id,
  936. text=text or "[非文本消息]",
  937. sender_id=int(sender.get("id") or 0),
  938. metadata={"content_type": "text" if text else "unsupported"},
  939. )
  940. await self._notify_feishu_customer_message(
  941. connection=connection,
  942. conversation=conversation,
  943. message=message,
  944. sender=sender,
  945. settings=settings,
  946. text=text,
  947. )
  948. if not connection.get("is_enabled") or not settings.get("assistant_enabled"):
  949. return
  950. current = await get_conversation(conversation["conversation_id"]) or conversation
  951. if current.get("status") == "human_paused":
  952. paused_until = current.get("paused_until")
  953. if isinstance(paused_until, datetime) and (
  954. paused_until.replace(tzinfo=UTC) if paused_until.tzinfo is None else paused_until
  955. ) <= utc_now():
  956. current = await resume_conversation(current["conversation_id"])
  957. else:
  958. return
  959. if current.get("status") in {"handoff", "closed"}:
  960. return
  961. if not business_reply_window_open(message):
  962. await self._handoff(
  963. connection,
  964. current,
  965. settings,
  966. "reply_window_expired",
  967. send_customer_notice=False,
  968. )
  969. return
  970. rights = connection.get("rights") or {}
  971. if not rights.get("can_reply"):
  972. await self._handoff(
  973. connection, current, settings, "can_reply_missing", send_customer_notice=False
  974. )
  975. return
  976. if rights.get("can_read_messages"):
  977. with suppress(TelegramBotApiError):
  978. await self.api.read_business_message(connection_id, chat_id, message_id)
  979. if not text:
  980. await self._handoff(
  981. connection,
  982. current,
  983. settings,
  984. "unsupported_message",
  985. customer_notice=str(settings.get("unsupported_message")),
  986. )
  987. return
  988. classification = classify_handoff(text)
  989. if classification:
  990. await self._handoff(connection, current, settings, classification)
  991. return
  992. knowledge = await match_knowledge(connection_id, text)
  993. if not knowledge:
  994. await self._handoff(connection, current, settings, "knowledge_not_found")
  995. return
  996. allowed, quota_reason = await reserve_ai_usage(
  997. connection_id, chat_id, settings
  998. )
  999. if not allowed:
  1000. await self._handoff(connection, current, settings, quota_reason)
  1001. return
  1002. messages = await recent_conversation_messages(current["conversation_id"], limit=12)
  1003. try:
  1004. decision = await self.provider.decide(
  1005. settings=settings,
  1006. conversation=current,
  1007. messages=messages,
  1008. knowledge=knowledge,
  1009. customer_text=text,
  1010. )
  1011. except AssistantProviderError as exc:
  1012. await update_runtime_status(
  1013. {"provider_last_error": str(exc)[:1000], "provider_last_error_at": utc_now()}
  1014. )
  1015. await self._handoff(
  1016. connection,
  1017. current,
  1018. settings,
  1019. "provider_error",
  1020. )
  1021. return
  1022. if decision["summary"]:
  1023. await update_conversation_summary(
  1024. current["conversation_id"], decision["summary"]
  1025. )
  1026. if decision["action"] == "handoff":
  1027. await self._handoff(
  1028. connection,
  1029. current,
  1030. settings,
  1031. decision["handoff_reason"],
  1032. summary=decision["summary"],
  1033. )
  1034. return
  1035. response = await self.api.send_message(
  1036. chat_id,
  1037. decision["reply"],
  1038. business_connection_id=connection_id,
  1039. )
  1040. await append_conversation_message(
  1041. current["conversation_id"],
  1042. direction="assistant",
  1043. telegram_message_id=int(response.get("message_id") or 0),
  1044. text=decision["reply"],
  1045. sender_id=int((response.get("sender_business_bot") or {}).get("id") or 0),
  1046. metadata={
  1047. "action": decision["action"],
  1048. "matched_entry_ids": decision["matched_entry_ids"],
  1049. },
  1050. )
  1051. async def _notify_feishu_customer_message(
  1052. self,
  1053. *,
  1054. connection: dict[str, Any],
  1055. conversation: dict[str, Any],
  1056. message: dict[str, Any],
  1057. sender: dict[str, Any],
  1058. settings: dict[str, Any],
  1059. text: str,
  1060. ) -> None:
  1061. webhook_url = str(settings.get("feishu_webhook_url") or "")
  1062. if not settings.get("feishu_webhook_enabled") or not webhook_url:
  1063. return
  1064. conversation_id = str(conversation["conversation_id"])
  1065. message_id = int(message.get("message_id") or 0)
  1066. if not await claim_feishu_message_notification(conversation_id, message_id):
  1067. return
  1068. user_name = " ".join(
  1069. str(item).strip()
  1070. for item in (sender.get("first_name"), sender.get("last_name"))
  1071. if str(item or "").strip()
  1072. )
  1073. if not user_name:
  1074. user_name = (
  1075. f"@{sender.get('username')}"
  1076. if sender.get("username")
  1077. else str((message.get("chat") or {}).get("id") or "未知")
  1078. )
  1079. account = connection.get("user") or {}
  1080. account_name = " ".join(
  1081. str(item).strip()
  1082. for item in (account.get("first_name"), account.get("last_name"))
  1083. if str(item or "").strip()
  1084. ) or (f"@{account.get('username')}" if account.get("username") else "Business 账号")
  1085. message_type = "文本" if text else "非文本"
  1086. preview_text = str(
  1087. message.get("text") or message.get("caption") or ""
  1088. ).strip()
  1089. content: list[list[dict[str, str]]] = [
  1090. [{"tag": "text", "text": f"账号:{account_name}"}],
  1091. [{"tag": "text", "text": "用户:"}],
  1092. [{"tag": "text", "text": f"类型:{message_type}"}],
  1093. ]
  1094. telegram_username = str(sender.get("username") or "").strip().lstrip("@")
  1095. if re.fullmatch(r"[A-Za-z0-9_]{1,64}", telegram_username):
  1096. content[1].append(
  1097. {
  1098. "tag": "a",
  1099. "text": user_name,
  1100. "href": f"https://t.me/{telegram_username}",
  1101. }
  1102. )
  1103. else:
  1104. content[1].append({"tag": "text", "text": user_name})
  1105. if settings.get("feishu_message_preview_enabled"):
  1106. content.append(
  1107. [
  1108. {
  1109. "tag": "text",
  1110. "text": f"内容:{preview_text[:500] if preview_text else '[非文本消息]'}",
  1111. }
  1112. ]
  1113. )
  1114. try:
  1115. message_time = datetime.fromtimestamp(
  1116. int(message.get("date") or 0), UTC
  1117. ).strftime("%Y-%m-%d %H:%M:%S UTC")
  1118. except (OSError, OverflowError, TypeError, ValueError):
  1119. message_time = "未知"
  1120. content.append([{"tag": "text", "text": f"时间:{message_time}"}])
  1121. try:
  1122. await self.feishu.send_post(
  1123. webhook_url,
  1124. title="消息提醒",
  1125. content=content,
  1126. signing_secret=str(settings.get("feishu_webhook_signing_secret") or ""),
  1127. )
  1128. except FeishuWebhookError as exc:
  1129. await mark_feishu_message_notification(
  1130. conversation_id, message_id, status="failed", error=str(exc)
  1131. )
  1132. await update_runtime_status(
  1133. {
  1134. "feishu_notification_last_error": str(exc)[:1000],
  1135. "feishu_notification_last_error_at": utc_now(),
  1136. }
  1137. )
  1138. return
  1139. await mark_feishu_message_notification(
  1140. conversation_id, message_id, status="sent"
  1141. )
  1142. await update_runtime_status(
  1143. {
  1144. "feishu_notification_last_error": "",
  1145. "feishu_notification_last_sent_at": utc_now(),
  1146. }
  1147. )
  1148. async def test_feishu_webhook(
  1149. self, connection_id: str, overrides: dict[str, Any] | None = None
  1150. ) -> dict[str, Any]:
  1151. connection = await self._resolve_connection(connection_id)
  1152. current = await get_account_settings(connection_id)
  1153. settings = normalize_account_settings(overrides or {}, previous=current)
  1154. webhook_url = str(settings.get("feishu_webhook_url") or "")
  1155. if not webhook_url:
  1156. raise FeishuWebhookError("请先填写飞书群机器人 Webhook 地址。")
  1157. account = connection.get("user") or {}
  1158. account_name = " ".join(
  1159. str(item).strip()
  1160. for item in (account.get("first_name"), account.get("last_name"))
  1161. if str(item or "").strip()
  1162. ) or (f"@{account.get('username')}" if account.get("username") else "Business 账号")
  1163. await self.feishu.send(
  1164. webhook_url,
  1165. "\n".join(
  1166. [
  1167. "消息提醒",
  1168. "飞书群机器人连接测试成功。",
  1169. f"账号:{account_name}",
  1170. ]
  1171. ),
  1172. signing_secret=str(settings.get("feishu_webhook_signing_secret") or ""),
  1173. )
  1174. return {"ok": True, "account": account_name}
  1175. async def _handoff(
  1176. self,
  1177. connection: dict[str, Any],
  1178. conversation: dict[str, Any],
  1179. settings: dict[str, Any],
  1180. reason: str,
  1181. *,
  1182. summary: str = "",
  1183. customer_notice: str = "",
  1184. send_customer_notice: bool = True,
  1185. ) -> None:
  1186. updated = await set_conversation_handoff(
  1187. conversation["conversation_id"], reason, summary=summary
  1188. )
  1189. notice = customer_notice or str(settings.get("handoff_message") or "")
  1190. if send_customer_notice and notice and connection.get("rights", {}).get("can_reply"):
  1191. try:
  1192. response = await self.api.send_message(
  1193. int(updated["chat_id"]),
  1194. notice,
  1195. business_connection_id=str(connection["connection_id"]),
  1196. )
  1197. await append_conversation_message(
  1198. updated["conversation_id"],
  1199. direction="assistant",
  1200. telegram_message_id=int(response.get("message_id") or 0),
  1201. text=notice,
  1202. metadata={"action": "handoff", "reason": reason},
  1203. )
  1204. except TelegramBotApiError as exc:
  1205. await update_runtime_status(
  1206. {
  1207. "customer_notice_last_error": str(exc)[:1000],
  1208. "customer_notice_last_error_at": utc_now(),
  1209. }
  1210. )
  1211. await self.notify_handoff(connection, updated, settings, reason)
  1212. async def notify_handoff(
  1213. self,
  1214. connection: dict[str, Any],
  1215. conversation: dict[str, Any],
  1216. settings: dict[str, Any],
  1217. reason: str,
  1218. ) -> None:
  1219. customer = conversation.get("customer") or {}
  1220. display = " ".join(
  1221. item for item in (customer.get("first_name"), customer.get("last_name")) if item
  1222. ) or (f"@{customer.get('username')}" if customer.get("username") else str(conversation["chat_id"]))
  1223. latest = await recent_conversation_messages(conversation["conversation_id"], limit=1)
  1224. latest_text = str(latest[-1].get("text") or "")[:1000] if latest else ""
  1225. text = (
  1226. "需要人工接待\n"
  1227. f"账号:{(connection.get('user') or {}).get('first_name') or connection['connection_id']}\n"
  1228. f"客户:{display}\n"
  1229. f"原因:{reason}\n"
  1230. f"原消息:chat_id={conversation['chat_id']} / "
  1231. f"message_id={latest[-1].get('telegram_message_id') if latest else '未知'}\n"
  1232. f"摘要:{conversation.get('summary') or latest_text or '暂无'}"
  1233. )
  1234. markup = {
  1235. "inline_keyboard": [
  1236. [
  1237. {
  1238. "text": "恢复自动回复",
  1239. "callback_data": f"ba:resume:{conversation['conversation_id']}",
  1240. }
  1241. ]
  1242. ]
  1243. }
  1244. destinations: list[int] = []
  1245. destination = str(settings.get("notification_destination") or "owner")
  1246. if destination in {"owner", "both"} and int(connection.get("user_chat_id") or 0):
  1247. destinations.append(int(connection["user_chat_id"]))
  1248. if destination in {"ops", "both"} and int(settings.get("ops_group_id") or 0):
  1249. destinations.append(int(settings["ops_group_id"]))
  1250. for chat_id in dict.fromkeys(destinations):
  1251. try:
  1252. await self.api.send_message(chat_id, text, reply_markup=markup)
  1253. except TelegramBotApiError as exc:
  1254. await update_runtime_status(
  1255. {
  1256. "notification_last_error": str(exc)[:1000],
  1257. "notification_last_error_at": utc_now(),
  1258. }
  1259. )
  1260. async def _handle_edited_message(self, message: dict[str, Any]) -> None:
  1261. connection_id = str(message.get("business_connection_id") or "")
  1262. chat_id = int((message.get("chat") or {}).get("id") or 0)
  1263. if not connection_id or not chat_id or message.get("sender_business_bot"):
  1264. return
  1265. await touch_business_connection(connection_id)
  1266. conversation = await get_conversation_by_chat(connection_id, chat_id)
  1267. if not conversation:
  1268. return
  1269. connection = await self._resolve_connection(connection_id)
  1270. sender_id = int((message.get("from") or {}).get("id") or 0)
  1271. owner_id = int((connection.get("user") or {}).get("id") or 0)
  1272. message_id = int(message.get("message_id") or 0)
  1273. text = str(message.get("text") or message.get("caption") or "").strip()
  1274. if sender_id == owner_id:
  1275. settings = await get_account_settings(connection_id)
  1276. await append_conversation_message(
  1277. conversation["conversation_id"],
  1278. direction="human",
  1279. telegram_message_id=-message_id,
  1280. text=text or "[人工消息已编辑]",
  1281. sender_id=owner_id,
  1282. metadata={"edited_message_id": message_id},
  1283. )
  1284. await pause_conversation_for_human(
  1285. conversation["conversation_id"],
  1286. hours=int(settings.get("human_pause_hours") or 24),
  1287. )
  1288. if text:
  1289. await self._learn_from_human_reply(
  1290. connection_id=connection_id,
  1291. conversation_id=str(conversation["conversation_id"]),
  1292. chat_id=chat_id,
  1293. message_id=message_id,
  1294. owner_id=owner_id,
  1295. human_text=text,
  1296. )
  1297. return
  1298. await append_conversation_message(
  1299. conversation["conversation_id"],
  1300. direction="incoming",
  1301. telegram_message_id=-message_id,
  1302. text=text or "[消息已编辑]",
  1303. sender_id=sender_id,
  1304. metadata={"edited_message_id": message_id},
  1305. )
  1306. async def _handle_deleted_messages(self, payload: dict[str, Any]) -> None:
  1307. connection_id = str(payload.get("business_connection_id") or "")
  1308. chat_id = int((payload.get("chat") or {}).get("id") or 0)
  1309. if connection_id:
  1310. await touch_business_connection(connection_id)
  1311. conversation = await get_conversation_by_chat(connection_id, chat_id)
  1312. if not conversation:
  1313. return
  1314. for source in await list_business_sources(connection_id):
  1315. for message_id in payload.get("message_ids") or []:
  1316. await mark_source_event_deleted(
  1317. str(source["source_id"]), f"business:{chat_id}:{int(message_id)}"
  1318. )
  1319. await append_conversation_message(
  1320. conversation["conversation_id"],
  1321. direction="incoming",
  1322. telegram_message_id=-abs(int((payload.get("message_ids") or [0])[0] or 0)) - 1,
  1323. text="[消息已删除]",
  1324. metadata={"deleted_message_ids": payload.get("message_ids") or []},
  1325. )
  1326. async def _digest_loop(self) -> None:
  1327. while not self._stopping.is_set():
  1328. try:
  1329. for item in await due_digest_connections():
  1330. connection = item["connection"]
  1331. settings = item["settings"]
  1332. metrics = await usage_metrics(
  1333. connection_id=str(connection["connection_id"]),
  1334. day=str(item["day"]),
  1335. )
  1336. text = (
  1337. f"智能接待运营简报 · {item['day']}\n"
  1338. f"客户数:{metrics['customers']}\n"
  1339. f"AI 调用:{metrics['ai_calls']}\n"
  1340. f"待人工:{metrics['handoffs']}\n"
  1341. f"人工暂停:{metrics['human_paused']}"
  1342. )
  1343. destinations: list[int] = []
  1344. destination = str(settings.get("notification_destination") or "owner")
  1345. if destination in {"owner", "both"} and connection.get("user_chat_id"):
  1346. destinations.append(int(connection["user_chat_id"]))
  1347. if destination in {"ops", "both"} and settings.get("ops_group_id"):
  1348. destinations.append(int(settings["ops_group_id"]))
  1349. sent = False
  1350. for chat_id in dict.fromkeys(destinations):
  1351. try:
  1352. await self.api.send_message(chat_id, text)
  1353. sent = True
  1354. except TelegramBotApiError as exc:
  1355. await update_runtime_status(
  1356. {
  1357. "digest_last_error": str(exc)[:1000],
  1358. "digest_last_error_at": utc_now(),
  1359. }
  1360. )
  1361. if sent:
  1362. await mark_digest_sent(connection["connection_id"], item["day"])
  1363. except asyncio.CancelledError:
  1364. raise
  1365. except Exception as exc:
  1366. await update_runtime_status({"digest_last_error": str(exc)[:1000]})
  1367. await asyncio.sleep(60)
  1368. async def preview_answer(self, connection_id: str, text: str) -> dict[str, Any]:
  1369. settings = await get_account_settings(connection_id)
  1370. knowledge = await match_knowledge(connection_id, text)
  1371. classification = classify_handoff(text)
  1372. if classification or not knowledge:
  1373. return {
  1374. "action": "handoff",
  1375. "reason": classification or "knowledge_not_found",
  1376. "matched_entries": knowledge,
  1377. }
  1378. conversation = {"summary": ""}
  1379. decision = await self.provider.decide(
  1380. settings=settings,
  1381. conversation=conversation,
  1382. messages=[],
  1383. knowledge=knowledge,
  1384. customer_text=text,
  1385. )
  1386. return {**decision, "matched_entries": knowledge}
  1387. async def runtime_overview() -> dict[str, Any]:
  1388. status = await runtime_status()
  1389. status["usage"] = await usage_metrics()
  1390. return status