business_assistant.py 35 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893
  1. from __future__ import annotations
  2. import asyncio
  3. import json
  4. import re
  5. from contextlib import suppress
  6. from datetime import UTC, datetime
  7. from typing import Any
  8. from aiohttp import ClientError, ClientSession
  9. import wbb
  10. from wbb.utils.dbassistant import (
  11. append_conversation_message,
  12. claim_update,
  13. classify_handoff,
  14. dead_letter_update,
  15. due_digest_connections,
  16. get_account_settings,
  17. get_business_connection,
  18. get_conversation,
  19. get_conversation_by_chat,
  20. get_or_create_conversation,
  21. list_business_connections,
  22. load_update_offset,
  23. mark_digest_sent,
  24. mark_update_done,
  25. mark_update_failed,
  26. match_knowledge,
  27. pause_conversation_for_human,
  28. recent_conversation_messages,
  29. reserve_ai_usage,
  30. resume_conversation,
  31. runtime_status,
  32. save_update_offset,
  33. set_conversation_handoff,
  34. touch_business_connection,
  35. update_conversation_summary,
  36. update_runtime_status,
  37. upsert_business_connection,
  38. usage_metrics,
  39. utc_now,
  40. )
  41. BUSINESS_ALLOWED_UPDATES = [
  42. "business_connection",
  43. "business_message",
  44. "edited_business_message",
  45. "deleted_business_messages",
  46. ]
  47. BUSINESS_REPLY_WINDOW_SECONDS = 24 * 60 * 60
  48. class TelegramBotApiError(RuntimeError):
  49. def __init__(
  50. self,
  51. message: str,
  52. *,
  53. error_code: int = 0,
  54. retry_after: int = 0,
  55. ) -> None:
  56. super().__init__(message)
  57. self.error_code = int(error_code or 0)
  58. self.retry_after = int(retry_after or 0)
  59. class AssistantProviderError(RuntimeError):
  60. pass
  61. class TelegramBusinessApi:
  62. def __init__(self, token: str, session: ClientSession) -> None:
  63. self._base_url = f"https://api.telegram.org/bot{token}"
  64. self._session = session
  65. async def call(
  66. self,
  67. method: str,
  68. payload: dict[str, Any] | None = None,
  69. ) -> Any:
  70. last_error: TelegramBotApiError | None = None
  71. for attempt in range(3):
  72. try:
  73. async with self._session.post(
  74. f"{self._base_url}/{method}", json=payload or {}
  75. ) as response:
  76. data = await response.json(content_type=None)
  77. except (ClientError, TimeoutError, ValueError) as exc:
  78. last_error = TelegramBotApiError("无法连接 Telegram Bot API。")
  79. if attempt < 2:
  80. await asyncio.sleep(2**attempt)
  81. continue
  82. raise last_error from exc
  83. if response.status == 200 and isinstance(data, dict) and data.get("ok"):
  84. return data.get("result")
  85. parameters = data.get("parameters") if isinstance(data, dict) else {}
  86. last_error = TelegramBotApiError(
  87. str(data.get("description") or "Telegram Bot API 请求失败。")
  88. if isinstance(data, dict)
  89. else "Telegram Bot API 请求失败。",
  90. error_code=int(data.get("error_code") or response.status)
  91. if isinstance(data, dict)
  92. else response.status,
  93. retry_after=int((parameters or {}).get("retry_after") or 0),
  94. )
  95. retryable = last_error.error_code == 429 or response.status >= 500
  96. if retryable and attempt < 2:
  97. await asyncio.sleep(last_error.retry_after or 2**attempt)
  98. continue
  99. raise last_error
  100. raise last_error or TelegramBotApiError("Telegram Bot API 请求失败。")
  101. async def get_me(self) -> dict[str, Any]:
  102. result = await self.call("getMe")
  103. return result if isinstance(result, dict) else {}
  104. async def get_webhook_info(self) -> dict[str, Any]:
  105. result = await self.call("getWebhookInfo")
  106. return result if isinstance(result, dict) else {}
  107. async def get_business_connection(self, connection_id: str) -> dict[str, Any]:
  108. result = await self.call(
  109. "getBusinessConnection", {"business_connection_id": str(connection_id)}
  110. )
  111. return result if isinstance(result, dict) else {}
  112. async def get_updates(
  113. self, *, offset: int, poll_timeout: int = 30
  114. ) -> list[dict[str, Any]]:
  115. result = await self.call(
  116. "getUpdates",
  117. {
  118. "offset": int(offset),
  119. "timeout": int(poll_timeout),
  120. "limit": 100,
  121. "allowed_updates": BUSINESS_ALLOWED_UPDATES,
  122. },
  123. )
  124. return [item for item in (result or []) if isinstance(item, dict)]
  125. async def send_message(
  126. self,
  127. chat_id: int,
  128. text: str,
  129. *,
  130. business_connection_id: str = "",
  131. reply_markup: dict[str, Any] | None = None,
  132. ) -> dict[str, Any]:
  133. payload: dict[str, Any] = {
  134. "chat_id": int(chat_id),
  135. "text": str(text)[:4096],
  136. }
  137. if business_connection_id:
  138. payload["business_connection_id"] = str(business_connection_id)
  139. if reply_markup:
  140. payload["reply_markup"] = reply_markup
  141. result = await self.call("sendMessage", payload)
  142. return result if isinstance(result, dict) else {}
  143. async def read_business_message(
  144. self,
  145. connection_id: str,
  146. chat_id: int,
  147. message_id: int,
  148. ) -> None:
  149. await self.call(
  150. "readBusinessMessage",
  151. {
  152. "business_connection_id": str(connection_id),
  153. "chat_id": int(chat_id),
  154. "message_id": int(message_id),
  155. },
  156. )
  157. def _strip_json_fence(value: str) -> str:
  158. text = value.strip()
  159. if text.startswith("```"):
  160. text = re.sub(r"^```(?:json)?\s*", "", text, flags=re.IGNORECASE)
  161. text = re.sub(r"\s*```$", "", text)
  162. return text.strip()
  163. def business_reply_window_open(
  164. message: dict[str, Any], *, now: datetime | None = None
  165. ) -> bool:
  166. try:
  167. sent_at = datetime.fromtimestamp(int(message.get("date") or 0), UTC)
  168. except (OSError, OverflowError, TypeError, ValueError):
  169. return False
  170. if sent_at.timestamp() <= 0:
  171. return False
  172. age = (now or utc_now()) - sent_at
  173. return age.total_seconds() <= BUSINESS_REPLY_WINDOW_SECONDS
  174. def parse_ai_decision(value: str, *, allowed_entry_ids: set[str]) -> dict[str, Any]:
  175. try:
  176. payload = json.loads(_strip_json_fence(value))
  177. except (json.JSONDecodeError, TypeError) as exc:
  178. raise AssistantProviderError("模型没有返回有效的 JSON 结果。") from exc
  179. if not isinstance(payload, dict):
  180. raise AssistantProviderError("模型结果必须是 JSON 对象。")
  181. action = str(payload.get("action") or "")
  182. if action not in {"answer", "clarify", "handoff"}:
  183. raise AssistantProviderError("模型返回了不支持的接待动作。")
  184. reply = str(payload.get("reply") or "").strip()
  185. if action in {"answer", "clarify"} and not reply:
  186. raise AssistantProviderError("模型没有提供回复内容。")
  187. if len(reply) > 4000:
  188. raise AssistantProviderError("模型回复内容过长。")
  189. matched_ids = [
  190. str(item)
  191. for item in payload.get("matched_entry_ids") or []
  192. if str(item) in allowed_entry_ids
  193. ]
  194. if action == "answer" and not matched_ids:
  195. raise AssistantProviderError("业务回答没有引用知识条目。")
  196. return {
  197. "action": action,
  198. "reply": reply,
  199. "handoff_reason": str(payload.get("handoff_reason") or "ai_handoff")[:200],
  200. "matched_entry_ids": matched_ids,
  201. "summary": str(payload.get("summary") or "")[:4000],
  202. }
  203. class OpenAICompatibleAssistant:
  204. def __init__(
  205. self,
  206. session: ClientSession,
  207. *,
  208. base_url: str,
  209. api_key: str,
  210. model: str,
  211. timeout_seconds: int = 30,
  212. max_output_tokens: int = 600,
  213. ) -> None:
  214. self._session = session
  215. self.base_url = str(base_url).rstrip("/")
  216. self.api_key = str(api_key)
  217. self.model = str(model)
  218. self.timeout_seconds = max(5, min(int(timeout_seconds), 120))
  219. self.max_output_tokens = max(100, min(int(max_output_tokens), 2000))
  220. @property
  221. def configured(self) -> bool:
  222. return bool(self.base_url and self.api_key and self.model)
  223. async def decide(
  224. self,
  225. *,
  226. settings: dict[str, Any],
  227. conversation: dict[str, Any],
  228. messages: list[dict[str, Any]],
  229. knowledge: list[dict[str, Any]],
  230. customer_text: str,
  231. ) -> dict[str, Any]:
  232. if not self.configured:
  233. raise AssistantProviderError("OpenAI 兼容模型尚未配置。")
  234. knowledge_payload = [
  235. {
  236. "entry_id": item["entry_id"],
  237. "question": item["question"],
  238. "answer": item["answer"],
  239. "tags": item.get("tags") or [],
  240. }
  241. for item in knowledge
  242. ]
  243. recent_payload = [
  244. {"direction": item.get("direction"), "text": item.get("text", "")}
  245. for item in messages[-12:]
  246. ]
  247. system_prompt = (
  248. f"{settings.get('system_prompt')}\n"
  249. "必须遵守:业务事实只能来自 knowledge;不得自行补全价格、承诺、退款或政策。"
  250. "仅输出一个 JSON 对象,不使用 Markdown。字段为 action、reply、handoff_reason、"
  251. "matched_entry_ids、summary。action 只能是 answer、clarify、handoff。"
  252. "answer 必须填写实际引用的 matched_entry_ids;无法可靠回答时使用 handoff。"
  253. )
  254. request_body = {
  255. "model": self.model,
  256. "temperature": 0.2,
  257. "max_tokens": self.max_output_tokens,
  258. "messages": [
  259. {"role": "system", "content": system_prompt},
  260. {
  261. "role": "user",
  262. "content": json.dumps(
  263. {
  264. "language": settings.get("language"),
  265. "tone": settings.get("tone"),
  266. "existing_summary": conversation.get("summary", ""),
  267. "recent_messages": recent_payload,
  268. "knowledge": knowledge_payload,
  269. "customer_message": customer_text,
  270. },
  271. ensure_ascii=False,
  272. ),
  273. },
  274. ],
  275. }
  276. try:
  277. async with asyncio.timeout(self.timeout_seconds):
  278. async with self._session.post(
  279. f"{self.base_url}/chat/completions",
  280. headers={
  281. "Authorization": f"Bearer {self.api_key}",
  282. "Content-Type": "application/json",
  283. },
  284. json=request_body,
  285. ) as response:
  286. payload = await response.json(content_type=None)
  287. except (ClientError, TimeoutError, ValueError) as exc:
  288. raise AssistantProviderError("OpenAI 兼容服务暂时不可用。") from exc
  289. if response.status != 200:
  290. raise AssistantProviderError(f"OpenAI 兼容服务返回 HTTP {response.status}。")
  291. try:
  292. content = payload["choices"][0]["message"]["content"]
  293. except (KeyError, IndexError, TypeError) as exc:
  294. raise AssistantProviderError("OpenAI 兼容服务响应格式无效。") from exc
  295. return parse_ai_decision(
  296. str(content),
  297. allowed_entry_ids={str(item["entry_id"]) for item in knowledge},
  298. )
  299. async def test_connection(self) -> dict[str, Any]:
  300. if not self.configured:
  301. raise AssistantProviderError("OpenAI 兼容模型尚未配置。")
  302. request_body = {
  303. "model": self.model,
  304. "temperature": 0,
  305. "max_tokens": 12,
  306. "messages": [{"role": "user", "content": "只回复 OK"}],
  307. }
  308. try:
  309. async with asyncio.timeout(self.timeout_seconds):
  310. async with self._session.post(
  311. f"{self.base_url}/chat/completions",
  312. headers={
  313. "Authorization": f"Bearer {self.api_key}",
  314. "Content-Type": "application/json",
  315. },
  316. json=request_body,
  317. ) as response:
  318. payload = await response.json(content_type=None)
  319. except (ClientError, TimeoutError, ValueError) as exc:
  320. raise AssistantProviderError("无法连接 OpenAI 兼容服务。") from exc
  321. if response.status != 200:
  322. raise AssistantProviderError(f"模型测试失败:HTTP {response.status}。")
  323. return {"ok": True, "model": self.model, "response_id": payload.get("id", "")}
  324. def provider_from_wbb(session: ClientSession) -> OpenAICompatibleAssistant:
  325. return OpenAICompatibleAssistant(
  326. session,
  327. base_url=str(getattr(wbb, "BUSINESS_ASSISTANT_OPENAI_BASE_URL", "")),
  328. api_key=str(getattr(wbb, "BUSINESS_ASSISTANT_OPENAI_API_KEY", "")),
  329. model=str(getattr(wbb, "BUSINESS_ASSISTANT_OPENAI_MODEL", "")),
  330. timeout_seconds=int(
  331. getattr(wbb, "BUSINESS_ASSISTANT_OPENAI_TIMEOUT_SECONDS", 30)
  332. ),
  333. max_output_tokens=int(
  334. getattr(wbb, "BUSINESS_ASSISTANT_OPENAI_MAX_OUTPUT_TOKENS", 600)
  335. ),
  336. )
  337. class BusinessAssistantRuntime:
  338. def __init__(
  339. self,
  340. *,
  341. token: str,
  342. session: ClientSession,
  343. provider: OpenAICompatibleAssistant | None = None,
  344. ) -> None:
  345. self.api = TelegramBusinessApi(token, session)
  346. self.provider = provider or provider_from_wbb(session)
  347. self._poll_task: asyncio.Task[None] | None = None
  348. self._digest_task: asyncio.Task[None] | None = None
  349. self._stopping = asyncio.Event()
  350. async def start(self) -> dict[str, Any]:
  351. if self._poll_task and not self._poll_task.done():
  352. return await runtime_status()
  353. self._stopping.clear()
  354. await update_runtime_status(
  355. {
  356. "polling_state": "starting",
  357. "model_configured": self.provider.configured,
  358. "last_error": "",
  359. }
  360. )
  361. try:
  362. me = await self.api.get_me()
  363. webhook = await self.api.get_webhook_info()
  364. except TelegramBotApiError as exc:
  365. return await update_runtime_status(
  366. {"polling_state": "error", "last_error": str(exc)}
  367. )
  368. webhook_url = str(webhook.get("url") or "")
  369. supported = bool(me.get("can_connect_to_business"))
  370. await update_runtime_status(
  371. {
  372. "business_mode_supported": supported,
  373. "webhook_conflict": bool(webhook_url),
  374. "webhook_url": webhook_url,
  375. "bot_username": str(me.get("username") or ""),
  376. }
  377. )
  378. if webhook_url:
  379. return await update_runtime_status(
  380. {
  381. "polling_state": "blocked",
  382. "last_error": "检测到 Telegram webhook,Business 长轮询未启动。",
  383. }
  384. )
  385. if not supported:
  386. return await update_runtime_status(
  387. {
  388. "polling_state": "blocked",
  389. "last_error": "请先在 BotFather 为机器人开启 Business Mode。",
  390. }
  391. )
  392. await self._refresh_connections()
  393. self._poll_task = asyncio.create_task(
  394. self._poll_loop(), name="business-assistant-poller"
  395. )
  396. self._digest_task = asyncio.create_task(
  397. self._digest_loop(), name="business-assistant-digest"
  398. )
  399. return await update_runtime_status(
  400. {"polling_state": "running", "last_error": "", "started_at": utc_now()}
  401. )
  402. async def _refresh_connections(self) -> None:
  403. page = 1
  404. while True:
  405. connections, total = await list_business_connections(
  406. page=page, page_size=100
  407. )
  408. for connection in connections:
  409. connection_id = str(connection["connection_id"])
  410. try:
  411. payload = await self.api.get_business_connection(connection_id)
  412. await upsert_business_connection(payload)
  413. except TelegramBotApiError as exc:
  414. await touch_business_connection(
  415. connection_id,
  416. error=str(exc),
  417. is_enabled=False if exc.error_code in {400, 403} else None,
  418. )
  419. if not connections or page * 100 >= total:
  420. return
  421. page += 1
  422. async def stop(self) -> None:
  423. self._stopping.set()
  424. tasks = [task for task in (self._poll_task, self._digest_task) if task]
  425. for task in tasks:
  426. task.cancel()
  427. if tasks:
  428. await asyncio.gather(*tasks, return_exceptions=True)
  429. self._poll_task = None
  430. self._digest_task = None
  431. await update_runtime_status(
  432. {"polling_state": "stopped", "stopped_at": utc_now()}
  433. )
  434. async def _poll_loop(self) -> None:
  435. offset = await load_update_offset()
  436. backoff = 1
  437. while not self._stopping.is_set():
  438. try:
  439. updates = await self.api.get_updates(offset=offset, poll_timeout=30)
  440. backoff = 1
  441. await update_runtime_status(
  442. {"polling_state": "running", "last_poll_at": utc_now(), "last_error": ""}
  443. )
  444. for update in updates:
  445. update_id = int(update.get("update_id") or 0)
  446. if update_id <= 0:
  447. continue
  448. await self.process_update(update)
  449. offset = max(offset, update_id + 1)
  450. await save_update_offset(offset)
  451. except asyncio.CancelledError:
  452. raise
  453. except TelegramBotApiError as exc:
  454. wait_for = exc.retry_after or backoff
  455. await update_runtime_status(
  456. {"polling_state": "retrying", "last_error": str(exc)}
  457. )
  458. await asyncio.sleep(min(max(wait_for, 1), 60))
  459. backoff = min(backoff * 2, 60)
  460. except Exception as exc:
  461. await update_runtime_status(
  462. {"polling_state": "retrying", "last_error": str(exc)[:1000]}
  463. )
  464. await asyncio.sleep(backoff)
  465. backoff = min(backoff * 2, 60)
  466. async def process_update(self, update: dict[str, Any]) -> None:
  467. update_id = int(update.get("update_id") or 0)
  468. if not await claim_update(update_id):
  469. return
  470. last_error = ""
  471. for attempt in range(1, 4):
  472. try:
  473. if isinstance(update.get("business_connection"), dict):
  474. await upsert_business_connection(update["business_connection"])
  475. elif isinstance(update.get("business_message"), dict):
  476. await self._handle_business_message(update["business_message"])
  477. elif isinstance(update.get("edited_business_message"), dict):
  478. await self._handle_edited_message(update["edited_business_message"])
  479. elif isinstance(update.get("deleted_business_messages"), dict):
  480. await self._handle_deleted_messages(update["deleted_business_messages"])
  481. await mark_update_done(update_id)
  482. return
  483. except TelegramBotApiError as exc:
  484. last_error = str(exc)
  485. if exc.retry_after:
  486. await asyncio.sleep(min(exc.retry_after, 60))
  487. except Exception as exc:
  488. last_error = str(exc)
  489. if attempt < 3:
  490. await asyncio.sleep(attempt)
  491. await mark_update_failed(update_id, last_error)
  492. await dead_letter_update(update, last_error or "unknown update error")
  493. async def _resolve_connection(self, connection_id: str) -> dict[str, Any]:
  494. connection = await get_business_connection(connection_id)
  495. if connection:
  496. return connection
  497. payload = await self.api.get_business_connection(connection_id)
  498. return await upsert_business_connection(payload)
  499. async def _handle_business_message(self, message: dict[str, Any]) -> None:
  500. connection_id = str(message.get("business_connection_id") or "")
  501. if not connection_id or message.get("is_from_offline"):
  502. return
  503. if message.get("sender_business_bot"):
  504. return
  505. connection = await self._resolve_connection(connection_id)
  506. await touch_business_connection(connection_id)
  507. chat = message.get("chat") if isinstance(message.get("chat"), dict) else {}
  508. sender = message.get("from") if isinstance(message.get("from"), dict) else {}
  509. chat_id = int(chat.get("id") or 0)
  510. message_id = int(message.get("message_id") or 0)
  511. if not chat_id or not message_id:
  512. return
  513. owner_id = int((connection.get("user") or {}).get("id") or 0)
  514. owner_reply = int(sender.get("id") or 0) == owner_id
  515. conversation = await get_or_create_conversation(
  516. connection_id,
  517. chat_id,
  518. customer=None if owner_reply else sender,
  519. )
  520. if owner_reply:
  521. settings = await get_account_settings(connection_id)
  522. await append_conversation_message(
  523. conversation["conversation_id"],
  524. direction="human",
  525. telegram_message_id=message_id,
  526. text=str(message.get("text") or message.get("caption") or ""),
  527. sender_id=owner_id,
  528. )
  529. await pause_conversation_for_human(
  530. conversation["conversation_id"],
  531. hours=int(settings.get("human_pause_hours") or 24),
  532. )
  533. return
  534. await self._handle_customer_message(
  535. connection=connection,
  536. conversation=conversation,
  537. message=message,
  538. sender=sender,
  539. )
  540. async def _handle_customer_message(
  541. self,
  542. *,
  543. connection: dict[str, Any],
  544. conversation: dict[str, Any],
  545. message: dict[str, Any],
  546. sender: dict[str, Any],
  547. ) -> None:
  548. connection_id = str(connection["connection_id"])
  549. chat_id = int((message.get("chat") or {}).get("id") or 0)
  550. message_id = int(message.get("message_id") or 0)
  551. text = str(message.get("text") or "").strip()
  552. settings = await get_account_settings(connection_id)
  553. await append_conversation_message(
  554. conversation["conversation_id"],
  555. direction="incoming",
  556. telegram_message_id=message_id,
  557. text=text or "[非文本消息]",
  558. sender_id=int(sender.get("id") or 0),
  559. metadata={"content_type": "text" if text else "unsupported"},
  560. )
  561. if not connection.get("is_enabled") or not settings.get("assistant_enabled"):
  562. return
  563. current = await get_conversation(conversation["conversation_id"]) or conversation
  564. if current.get("status") == "human_paused":
  565. paused_until = current.get("paused_until")
  566. if isinstance(paused_until, datetime) and (
  567. paused_until.replace(tzinfo=UTC) if paused_until.tzinfo is None else paused_until
  568. ) <= utc_now():
  569. current = await resume_conversation(current["conversation_id"])
  570. else:
  571. return
  572. if current.get("status") in {"handoff", "closed"}:
  573. return
  574. if not business_reply_window_open(message):
  575. await self._handoff(
  576. connection,
  577. current,
  578. settings,
  579. "reply_window_expired",
  580. send_customer_notice=False,
  581. )
  582. return
  583. rights = connection.get("rights") or {}
  584. if not rights.get("can_reply"):
  585. await self._handoff(
  586. connection, current, settings, "can_reply_missing", send_customer_notice=False
  587. )
  588. return
  589. if rights.get("can_read_messages"):
  590. with suppress(TelegramBotApiError):
  591. await self.api.read_business_message(connection_id, chat_id, message_id)
  592. if not text:
  593. await self._handoff(
  594. connection,
  595. current,
  596. settings,
  597. "unsupported_message",
  598. customer_notice=str(settings.get("unsupported_message")),
  599. )
  600. return
  601. classification = classify_handoff(text)
  602. if classification:
  603. await self._handoff(connection, current, settings, classification)
  604. return
  605. knowledge = await match_knowledge(connection_id, text)
  606. if not knowledge:
  607. await self._handoff(connection, current, settings, "knowledge_not_found")
  608. return
  609. allowed, quota_reason = await reserve_ai_usage(
  610. connection_id, chat_id, settings
  611. )
  612. if not allowed:
  613. await self._handoff(connection, current, settings, quota_reason)
  614. return
  615. messages = await recent_conversation_messages(current["conversation_id"], limit=12)
  616. try:
  617. decision = await self.provider.decide(
  618. settings=settings,
  619. conversation=current,
  620. messages=messages,
  621. knowledge=knowledge,
  622. customer_text=text,
  623. )
  624. except AssistantProviderError as exc:
  625. await update_runtime_status(
  626. {"provider_last_error": str(exc)[:1000], "provider_last_error_at": utc_now()}
  627. )
  628. await self._handoff(
  629. connection,
  630. current,
  631. settings,
  632. "provider_error",
  633. )
  634. return
  635. if decision["summary"]:
  636. await update_conversation_summary(
  637. current["conversation_id"], decision["summary"]
  638. )
  639. if decision["action"] == "handoff":
  640. await self._handoff(
  641. connection,
  642. current,
  643. settings,
  644. decision["handoff_reason"],
  645. summary=decision["summary"],
  646. )
  647. return
  648. response = await self.api.send_message(
  649. chat_id,
  650. decision["reply"],
  651. business_connection_id=connection_id,
  652. )
  653. await append_conversation_message(
  654. current["conversation_id"],
  655. direction="assistant",
  656. telegram_message_id=int(response.get("message_id") or 0),
  657. text=decision["reply"],
  658. sender_id=int((response.get("sender_business_bot") or {}).get("id") or 0),
  659. metadata={
  660. "action": decision["action"],
  661. "matched_entry_ids": decision["matched_entry_ids"],
  662. },
  663. )
  664. async def _handoff(
  665. self,
  666. connection: dict[str, Any],
  667. conversation: dict[str, Any],
  668. settings: dict[str, Any],
  669. reason: str,
  670. *,
  671. summary: str = "",
  672. customer_notice: str = "",
  673. send_customer_notice: bool = True,
  674. ) -> None:
  675. updated = await set_conversation_handoff(
  676. conversation["conversation_id"], reason, summary=summary
  677. )
  678. notice = customer_notice or str(settings.get("handoff_message") or "")
  679. if send_customer_notice and notice and connection.get("rights", {}).get("can_reply"):
  680. try:
  681. response = await self.api.send_message(
  682. int(updated["chat_id"]),
  683. notice,
  684. business_connection_id=str(connection["connection_id"]),
  685. )
  686. await append_conversation_message(
  687. updated["conversation_id"],
  688. direction="assistant",
  689. telegram_message_id=int(response.get("message_id") or 0),
  690. text=notice,
  691. metadata={"action": "handoff", "reason": reason},
  692. )
  693. except TelegramBotApiError as exc:
  694. await update_runtime_status(
  695. {
  696. "customer_notice_last_error": str(exc)[:1000],
  697. "customer_notice_last_error_at": utc_now(),
  698. }
  699. )
  700. await self.notify_handoff(connection, updated, settings, reason)
  701. async def notify_handoff(
  702. self,
  703. connection: dict[str, Any],
  704. conversation: dict[str, Any],
  705. settings: dict[str, Any],
  706. reason: str,
  707. ) -> None:
  708. customer = conversation.get("customer") or {}
  709. display = " ".join(
  710. item for item in (customer.get("first_name"), customer.get("last_name")) if item
  711. ) or (f"@{customer.get('username')}" if customer.get("username") else str(conversation["chat_id"]))
  712. latest = await recent_conversation_messages(conversation["conversation_id"], limit=1)
  713. latest_text = str(latest[-1].get("text") or "")[:1000] if latest else ""
  714. text = (
  715. "需要人工接待\n"
  716. f"账号:{(connection.get('user') or {}).get('first_name') or connection['connection_id']}\n"
  717. f"客户:{display}\n"
  718. f"原因:{reason}\n"
  719. f"原消息:chat_id={conversation['chat_id']} / "
  720. f"message_id={latest[-1].get('telegram_message_id') if latest else '未知'}\n"
  721. f"摘要:{conversation.get('summary') or latest_text or '暂无'}"
  722. )
  723. markup = {
  724. "inline_keyboard": [
  725. [
  726. {
  727. "text": "恢复自动回复",
  728. "callback_data": f"ba:resume:{conversation['conversation_id']}",
  729. }
  730. ]
  731. ]
  732. }
  733. destinations: list[int] = []
  734. destination = str(settings.get("notification_destination") or "owner")
  735. if destination in {"owner", "both"} and int(connection.get("user_chat_id") or 0):
  736. destinations.append(int(connection["user_chat_id"]))
  737. if destination in {"ops", "both"} and int(settings.get("ops_group_id") or 0):
  738. destinations.append(int(settings["ops_group_id"]))
  739. for chat_id in dict.fromkeys(destinations):
  740. try:
  741. await self.api.send_message(chat_id, text, reply_markup=markup)
  742. except TelegramBotApiError as exc:
  743. await update_runtime_status(
  744. {
  745. "notification_last_error": str(exc)[:1000],
  746. "notification_last_error_at": utc_now(),
  747. }
  748. )
  749. async def _handle_edited_message(self, message: dict[str, Any]) -> None:
  750. connection_id = str(message.get("business_connection_id") or "")
  751. chat_id = int((message.get("chat") or {}).get("id") or 0)
  752. if not connection_id or not chat_id:
  753. return
  754. await touch_business_connection(connection_id)
  755. conversation = await get_conversation_by_chat(connection_id, chat_id)
  756. if not conversation:
  757. return
  758. await append_conversation_message(
  759. conversation["conversation_id"],
  760. direction="incoming",
  761. telegram_message_id=-int(message.get("message_id") or 0),
  762. text=str(message.get("text") or "[消息已编辑]"),
  763. sender_id=int((message.get("from") or {}).get("id") or 0),
  764. metadata={"edited_message_id": int(message.get("message_id") or 0)},
  765. )
  766. async def _handle_deleted_messages(self, payload: dict[str, Any]) -> None:
  767. connection_id = str(payload.get("business_connection_id") or "")
  768. chat_id = int((payload.get("chat") or {}).get("id") or 0)
  769. if connection_id:
  770. await touch_business_connection(connection_id)
  771. conversation = await get_conversation_by_chat(connection_id, chat_id)
  772. if not conversation:
  773. return
  774. await append_conversation_message(
  775. conversation["conversation_id"],
  776. direction="incoming",
  777. telegram_message_id=-abs(int((payload.get("message_ids") or [0])[0] or 0)) - 1,
  778. text="[消息已删除]",
  779. metadata={"deleted_message_ids": payload.get("message_ids") or []},
  780. )
  781. async def _digest_loop(self) -> None:
  782. while not self._stopping.is_set():
  783. try:
  784. for item in await due_digest_connections():
  785. connection = item["connection"]
  786. settings = item["settings"]
  787. metrics = await usage_metrics(
  788. connection_id=str(connection["connection_id"]),
  789. day=str(item["day"]),
  790. )
  791. text = (
  792. f"智能接待运营简报 · {item['day']}\n"
  793. f"客户数:{metrics['customers']}\n"
  794. f"AI 调用:{metrics['ai_calls']}\n"
  795. f"待人工:{metrics['handoffs']}\n"
  796. f"人工暂停:{metrics['human_paused']}"
  797. )
  798. destinations: list[int] = []
  799. destination = str(settings.get("notification_destination") or "owner")
  800. if destination in {"owner", "both"} and connection.get("user_chat_id"):
  801. destinations.append(int(connection["user_chat_id"]))
  802. if destination in {"ops", "both"} and settings.get("ops_group_id"):
  803. destinations.append(int(settings["ops_group_id"]))
  804. sent = False
  805. for chat_id in dict.fromkeys(destinations):
  806. try:
  807. await self.api.send_message(chat_id, text)
  808. sent = True
  809. except TelegramBotApiError as exc:
  810. await update_runtime_status(
  811. {
  812. "digest_last_error": str(exc)[:1000],
  813. "digest_last_error_at": utc_now(),
  814. }
  815. )
  816. if sent:
  817. await mark_digest_sent(connection["connection_id"], item["day"])
  818. except asyncio.CancelledError:
  819. raise
  820. except Exception as exc:
  821. await update_runtime_status({"digest_last_error": str(exc)[:1000]})
  822. await asyncio.sleep(60)
  823. async def preview_answer(self, connection_id: str, text: str) -> dict[str, Any]:
  824. settings = await get_account_settings(connection_id)
  825. knowledge = await match_knowledge(connection_id, text)
  826. classification = classify_handoff(text)
  827. if classification or not knowledge:
  828. return {
  829. "action": "handoff",
  830. "reason": classification or "knowledge_not_found",
  831. "matched_entries": knowledge,
  832. }
  833. conversation = {"summary": ""}
  834. decision = await self.provider.decide(
  835. settings=settings,
  836. conversation=conversation,
  837. messages=[],
  838. knowledge=knowledge,
  839. customer_text=text,
  840. )
  841. return {**decision, "matched_entries": knowledge}
  842. async def runtime_overview() -> dict[str, Any]:
  843. status = await runtime_status()
  844. status["usage"] = await usage_metrics()
  845. return status