| 12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460 |
- from __future__ import annotations
- import asyncio
- import base64
- import hashlib
- import hmac
- import json
- import re
- from contextlib import suppress
- from datetime import UTC, datetime
- from typing import Any
- from aiohttp import ClientError, ClientSession, ClientTimeout
- import wbb
- from wbb.utils.dbassistant import (
- append_conversation_message,
- claim_feishu_message_notification,
- claim_update,
- classify_handoff,
- dead_letter_update,
- due_digest_connections,
- get_account_settings,
- get_business_connection,
- get_conversation,
- get_conversation_by_chat,
- get_knowledge_candidate,
- get_knowledge_source,
- get_or_create_conversation,
- get_source_event,
- get_source_event_by_key,
- invalidate_source_event_candidates,
- list_business_connections,
- list_business_sources,
- list_enabled_business_sources,
- load_update_offset,
- mark_digest_sent,
- mark_feishu_message_notification,
- mark_source_event_deleted,
- mark_update_done,
- mark_update_failed,
- match_knowledge,
- normalize_account_settings,
- pause_conversation_for_human,
- recent_conversation_messages,
- record_source_event,
- replace_event_candidates,
- reserve_ai_usage,
- resume_conversation,
- runtime_status,
- save_update_offset,
- set_conversation_handoff,
- touch_business_connection,
- update_conversation_summary,
- update_runtime_status,
- update_source_event_extraction,
- upsert_business_connection,
- usage_metrics,
- utc_now,
- )
- BUSINESS_ALLOWED_UPDATES = [
- "business_connection",
- "business_message",
- "edited_business_message",
- "deleted_business_messages",
- ]
- BUSINESS_REPLY_WINDOW_SECONDS = 24 * 60 * 60
- class TelegramBotApiError(RuntimeError):
- def __init__(
- self,
- message: str,
- *,
- error_code: int = 0,
- retry_after: int = 0,
- ) -> None:
- super().__init__(message)
- self.error_code = int(error_code or 0)
- self.retry_after = int(retry_after or 0)
- class AssistantProviderError(RuntimeError):
- pass
- class FeishuWebhookError(RuntimeError):
- pass
- class FeishuWebhookClient:
- def __init__(self, session: ClientSession) -> None:
- self._session = session
- @staticmethod
- def _signature(secret: str, timestamp: int) -> str:
- sign_key = f"{timestamp}\n{secret}".encode()
- digest = hmac.new(sign_key, digestmod=hashlib.sha256).digest()
- return base64.b64encode(digest).decode()
- async def send(self, webhook_url: str, text: str, *, signing_secret: str = "") -> None:
- payload: dict[str, Any] = {
- "msg_type": "text",
- "content": {"text": str(text)[:4000]},
- }
- await self._send_payload(
- webhook_url, payload, signing_secret=signing_secret
- )
- async def send_post(
- self,
- webhook_url: str,
- *,
- title: str,
- content: list[list[dict[str, str]]],
- signing_secret: str = "",
- ) -> None:
- payload: dict[str, Any] = {
- "msg_type": "post",
- "content": {
- "post": {
- "zh_cn": {
- "title": str(title)[:100],
- "content": content,
- }
- }
- },
- }
- await self._send_payload(
- webhook_url, payload, signing_secret=signing_secret
- )
- async def _send_payload(
- self,
- webhook_url: str,
- payload: dict[str, Any],
- *,
- signing_secret: str = "",
- ) -> None:
- if signing_secret:
- timestamp = int(utc_now().timestamp())
- payload.update(
- {
- "timestamp": str(timestamp),
- "sign": self._signature(signing_secret, timestamp),
- }
- )
- last_error: FeishuWebhookError | None = None
- for attempt in range(3):
- try:
- async with self._session.post(
- webhook_url,
- json=payload,
- timeout=ClientTimeout(total=8),
- ) as response:
- data = await response.json(content_type=None)
- except (ClientError, TimeoutError, ValueError) as exc:
- last_error = FeishuWebhookError("无法连接飞书群机器人 Webhook。")
- if attempt < 2:
- await asyncio.sleep(2**attempt)
- continue
- raise last_error from exc
- raw_code = (
- data.get("code", data.get("StatusCode", -1))
- if isinstance(data, dict)
- else -1
- )
- try:
- code = int(raw_code)
- except (TypeError, ValueError):
- code = -1
- if response.status == 200 and code == 0:
- return
- last_error = FeishuWebhookError(
- f"飞书群机器人 Webhook 请求失败(HTTP {response.status},code={code})。"
- )
- if (response.status == 429 or response.status >= 500) and attempt < 2:
- retry_after = response.headers.get("Retry-After", "")
- try:
- wait_seconds = max(1, min(int(retry_after), 30))
- except (TypeError, ValueError):
- wait_seconds = 2**attempt
- await asyncio.sleep(wait_seconds)
- continue
- raise last_error
- raise last_error or FeishuWebhookError("飞书群机器人 Webhook 请求失败。")
- class TelegramBusinessApi:
- def __init__(self, token: str, session: ClientSession) -> None:
- self._base_url = f"https://api.telegram.org/bot{token}"
- self._session = session
- async def call(
- self,
- method: str,
- payload: dict[str, Any] | None = None,
- ) -> Any:
- last_error: TelegramBotApiError | None = None
- for attempt in range(3):
- try:
- async with self._session.post(
- f"{self._base_url}/{method}", json=payload or {}
- ) as response:
- data = await response.json(content_type=None)
- except (ClientError, TimeoutError, ValueError) as exc:
- last_error = TelegramBotApiError("无法连接 Telegram Bot API。")
- if attempt < 2:
- await asyncio.sleep(2**attempt)
- continue
- raise last_error from exc
- if response.status == 200 and isinstance(data, dict) and data.get("ok"):
- return data.get("result")
- parameters = data.get("parameters") if isinstance(data, dict) else {}
- last_error = TelegramBotApiError(
- str(data.get("description") or "Telegram Bot API 请求失败。")
- if isinstance(data, dict)
- else "Telegram Bot API 请求失败。",
- error_code=int(data.get("error_code") or response.status)
- if isinstance(data, dict)
- else response.status,
- retry_after=int((parameters or {}).get("retry_after") or 0),
- )
- retryable = last_error.error_code == 429 or response.status >= 500
- if retryable and attempt < 2:
- await asyncio.sleep(last_error.retry_after or 2**attempt)
- continue
- raise last_error
- raise last_error or TelegramBotApiError("Telegram Bot API 请求失败。")
- async def get_me(self) -> dict[str, Any]:
- result = await self.call("getMe")
- return result if isinstance(result, dict) else {}
- async def get_webhook_info(self) -> dict[str, Any]:
- result = await self.call("getWebhookInfo")
- return result if isinstance(result, dict) else {}
- async def get_business_connection(self, connection_id: str) -> dict[str, Any]:
- result = await self.call(
- "getBusinessConnection", {"business_connection_id": str(connection_id)}
- )
- return result if isinstance(result, dict) else {}
- async def get_updates(
- self, *, offset: int, poll_timeout: int = 30
- ) -> list[dict[str, Any]]:
- result = await self.call(
- "getUpdates",
- {
- "offset": int(offset),
- "timeout": int(poll_timeout),
- "limit": 100,
- "allowed_updates": BUSINESS_ALLOWED_UPDATES,
- },
- )
- return [item for item in (result or []) if isinstance(item, dict)]
- async def send_message(
- self,
- chat_id: int,
- text: str,
- *,
- business_connection_id: str = "",
- reply_markup: dict[str, Any] | None = None,
- ) -> dict[str, Any]:
- payload: dict[str, Any] = {
- "chat_id": int(chat_id),
- "text": str(text)[:4096],
- }
- if business_connection_id:
- payload["business_connection_id"] = str(business_connection_id)
- if reply_markup:
- payload["reply_markup"] = reply_markup
- result = await self.call("sendMessage", payload)
- return result if isinstance(result, dict) else {}
- async def read_business_message(
- self,
- connection_id: str,
- chat_id: int,
- message_id: int,
- ) -> None:
- await self.call(
- "readBusinessMessage",
- {
- "business_connection_id": str(connection_id),
- "chat_id": int(chat_id),
- "message_id": int(message_id),
- },
- )
- def _strip_json_fence(value: str) -> str:
- text = value.strip()
- if text.startswith("```"):
- text = re.sub(r"^```(?:json)?\s*", "", text, flags=re.IGNORECASE)
- text = re.sub(r"\s*```$", "", text)
- return text.strip()
- def business_reply_window_open(
- message: dict[str, Any], *, now: datetime | None = None
- ) -> bool:
- try:
- sent_at = datetime.fromtimestamp(int(message.get("date") or 0), UTC)
- except (OSError, OverflowError, TypeError, ValueError):
- return False
- if sent_at.timestamp() <= 0:
- return False
- age = (now or utc_now()) - sent_at
- return age.total_seconds() <= BUSINESS_REPLY_WINDOW_SECONDS
- def parse_ai_decision(value: str, *, allowed_entry_ids: set[str]) -> dict[str, Any]:
- try:
- payload = json.loads(_strip_json_fence(value))
- except (json.JSONDecodeError, TypeError) as exc:
- raise AssistantProviderError("模型没有返回有效的 JSON 结果。") from exc
- if not isinstance(payload, dict):
- raise AssistantProviderError("模型结果必须是 JSON 对象。")
- action = str(payload.get("action") or "")
- if action not in {"answer", "clarify", "handoff"}:
- raise AssistantProviderError("模型返回了不支持的接待动作。")
- reply = str(payload.get("reply") or "").strip()
- if action in {"answer", "clarify"} and not reply:
- raise AssistantProviderError("模型没有提供回复内容。")
- if len(reply) > 4000:
- raise AssistantProviderError("模型回复内容过长。")
- matched_ids = [
- str(item)
- for item in payload.get("matched_entry_ids") or []
- if str(item) in allowed_entry_ids
- ]
- if action == "answer" and not matched_ids:
- raise AssistantProviderError("业务回答没有引用知识条目。")
- return {
- "action": action,
- "reply": reply,
- "handoff_reason": str(payload.get("handoff_reason") or "ai_handoff")[:200],
- "matched_entry_ids": matched_ids,
- "summary": str(payload.get("summary") or "")[:4000],
- }
- def _short_string_list(value: Any, *, items: int, length: int) -> list[str]:
- raw_items = value if isinstance(value, list) else []
- result: list[str] = []
- for item in raw_items:
- text = " ".join(str(item or "").strip().split())[:length]
- if text and text.casefold() not in {existing.casefold() for existing in result}:
- result.append(text)
- return result[:items]
- def parse_knowledge_extraction(value: str) -> list[dict[str, Any]]:
- try:
- payload = json.loads(_strip_json_fence(value))
- except (json.JSONDecodeError, TypeError) as exc:
- raise AssistantProviderError("模型没有返回有效的知识提取 JSON。") from exc
- if not isinstance(payload, dict) or not isinstance(payload.get("items"), list):
- raise AssistantProviderError("知识提取结果必须包含 items 数组。")
- result: list[dict[str, Any]] = []
- for raw_item in payload["items"][:5]:
- if not isinstance(raw_item, dict):
- raise AssistantProviderError("知识候选必须是 JSON 对象。")
- question = " ".join(str(raw_item.get("question") or "").strip().split())[:300]
- answer = str(raw_item.get("answer") or "").strip()[:4000]
- if not question or not answer:
- raise AssistantProviderError("知识候选缺少标准问题或答案。")
- try:
- confidence = float(raw_item.get("confidence") or 0)
- except (TypeError, ValueError):
- confidence = 0
- result.append(
- {
- "question": question,
- "aliases": _short_string_list(
- raw_item.get("aliases"), items=20, length=200
- ),
- "keywords": _short_string_list(
- raw_item.get("keywords"), items=30, length=50
- ),
- "answer": answer,
- "tags": _short_string_list(raw_item.get("tags"), items=20, length=30),
- "confidence": max(0.0, min(confidence, 1.0)),
- }
- )
- return result
- def redact_business_knowledge_text(value: str) -> str:
- text = re.sub(
- r"(?<![\w.])[A-Z0-9._%+-]+@[A-Z0-9.-]+\.[A-Z]{2,}(?!\w)",
- "[邮箱已脱敏]",
- str(value),
- flags=re.IGNORECASE,
- )
- text = re.sub(r"(?<!\d)(?:\+?\d[\d\s-]{7,}\d)(?!\d)", "[电话已脱敏]", text)
- return re.sub(r"(?<!\w)@[A-Za-z][A-Za-z0-9_]{4,31}", "[用户名已脱敏]", text)
- class OpenAICompatibleAssistant:
- def __init__(
- self,
- session: ClientSession,
- *,
- base_url: str,
- api_key: str,
- model: str,
- timeout_seconds: int = 30,
- max_output_tokens: int = 600,
- ) -> None:
- self._session = session
- self.base_url = str(base_url).rstrip("/")
- self.api_key = str(api_key)
- self.model = str(model)
- self.timeout_seconds = max(5, min(int(timeout_seconds), 120))
- self.max_output_tokens = max(100, min(int(max_output_tokens), 2000))
- @property
- def configured(self) -> bool:
- return bool(self.base_url and self.api_key and self.model)
- async def decide(
- self,
- *,
- settings: dict[str, Any],
- conversation: dict[str, Any],
- messages: list[dict[str, Any]],
- knowledge: list[dict[str, Any]],
- customer_text: str,
- ) -> dict[str, Any]:
- if not self.configured:
- raise AssistantProviderError("OpenAI 兼容模型尚未配置。")
- knowledge_payload = [
- {
- "entry_id": item["entry_id"],
- "question": item["question"],
- "answer": item["answer"],
- "tags": item.get("tags") or [],
- }
- for item in knowledge
- ]
- recent_payload = [
- {"direction": item.get("direction"), "text": item.get("text", "")}
- for item in messages[-12:]
- ]
- system_prompt = (
- f"{settings.get('system_prompt')}\n"
- "必须遵守:业务事实只能来自 knowledge;不得自行补全价格、承诺、退款或政策。"
- "仅输出一个 JSON 对象,不使用 Markdown。字段为 action、reply、handoff_reason、"
- "matched_entry_ids、summary。action 只能是 answer、clarify、handoff。"
- "answer 必须填写实际引用的 matched_entry_ids;无法可靠回答时使用 handoff。"
- )
- request_body = {
- "model": self.model,
- "temperature": 0.2,
- "max_tokens": self.max_output_tokens,
- "messages": [
- {"role": "system", "content": system_prompt},
- {
- "role": "user",
- "content": json.dumps(
- {
- "language": settings.get("language"),
- "tone": settings.get("tone"),
- "existing_summary": conversation.get("summary", ""),
- "recent_messages": recent_payload,
- "knowledge": knowledge_payload,
- "customer_message": customer_text,
- },
- ensure_ascii=False,
- ),
- },
- ],
- }
- try:
- async with asyncio.timeout(self.timeout_seconds):
- async with self._session.post(
- f"{self.base_url}/chat/completions",
- headers={
- "Authorization": f"Bearer {self.api_key}",
- "Content-Type": "application/json",
- },
- json=request_body,
- ) as response:
- payload = await response.json(content_type=None)
- except (ClientError, TimeoutError, ValueError) as exc:
- raise AssistantProviderError("OpenAI 兼容服务暂时不可用。") from exc
- if response.status != 200:
- raise AssistantProviderError(f"OpenAI 兼容服务返回 HTTP {response.status}。")
- try:
- content = payload["choices"][0]["message"]["content"]
- except (KeyError, IndexError, TypeError) as exc:
- raise AssistantProviderError("OpenAI 兼容服务响应格式无效。") from exc
- return parse_ai_decision(
- str(content),
- allowed_entry_ids={str(item["entry_id"]) for item in knowledge},
- )
- async def extract_knowledge(
- self,
- *,
- source_type: str,
- source_title: str,
- content: str,
- question_context: str = "",
- ) -> list[dict[str, Any]]:
- if not self.configured:
- raise AssistantProviderError("OpenAI 兼容模型尚未配置。")
- if source_type == "business":
- safe_content = redact_business_knowledge_text(content)
- safe_question = redact_business_knowledge_text(question_context)
- source_rule = (
- "这是 Business 人工接待对话。customer_question 只能作为标准问题参考,"
- "答案事实只能来自 human_reply。不得把客户自述、联系方式、用户名、"
- "投诉、退款要求或未经确认的承诺写入知识。"
- )
- source_payload = {
- "customer_question": safe_question,
- "human_reply": safe_content,
- }
- else:
- source_rule = (
- "这是频道或群组内容。只提取明确、稳定、可复用的业务事实;"
- "忽略寒暄、广告口号、个人观点、临时活动、投诉、退款争议和未经确认的承诺。"
- "没有稳定事实时返回空 items。"
- )
- source_payload = {"content": content}
- system_prompt = (
- "你是知识库整理器,只提取输入中明确存在的事实,不得推测或补全。"
- f"{source_rule}"
- "只输出一个 JSON 对象,不使用 Markdown,格式为 "
- '{"items":[{"question":"","aliases":[],"keywords":[],"answer":"",'
- '"tags":[],"confidence":0.0}]}。items 最多 5 条。'
- )
- request_body = {
- "model": self.model,
- "temperature": 0.1,
- "max_tokens": self.max_output_tokens,
- "messages": [
- {"role": "system", "content": system_prompt},
- {
- "role": "user",
- "content": json.dumps(
- {
- "source_type": source_type,
- "source_title": source_title,
- **source_payload,
- },
- ensure_ascii=False,
- ),
- },
- ],
- }
- try:
- async with asyncio.timeout(self.timeout_seconds):
- async with self._session.post(
- f"{self.base_url}/chat/completions",
- headers={
- "Authorization": f"Bearer {self.api_key}",
- "Content-Type": "application/json",
- },
- json=request_body,
- ) as response:
- payload = await response.json(content_type=None)
- except (ClientError, TimeoutError, ValueError) as exc:
- raise AssistantProviderError("知识提取模型暂时不可用。") from exc
- if response.status != 200:
- raise AssistantProviderError(f"知识提取模型返回 HTTP {response.status}。")
- try:
- result = payload["choices"][0]["message"]["content"]
- except (KeyError, IndexError, TypeError) as exc:
- raise AssistantProviderError("知识提取模型响应格式无效。") from exc
- return parse_knowledge_extraction(str(result))
- async def test_connection(self) -> dict[str, Any]:
- if not self.configured:
- raise AssistantProviderError("OpenAI 兼容模型尚未配置。")
- request_body = {
- "model": self.model,
- "temperature": 0,
- "max_tokens": 12,
- "messages": [{"role": "user", "content": "只回复 OK"}],
- }
- try:
- async with asyncio.timeout(self.timeout_seconds):
- async with self._session.post(
- f"{self.base_url}/chat/completions",
- headers={
- "Authorization": f"Bearer {self.api_key}",
- "Content-Type": "application/json",
- },
- json=request_body,
- ) as response:
- payload = await response.json(content_type=None)
- except (ClientError, TimeoutError, ValueError) as exc:
- raise AssistantProviderError("无法连接 OpenAI 兼容服务。") from exc
- if response.status != 200:
- raise AssistantProviderError(f"模型测试失败:HTTP {response.status}。")
- return {"ok": True, "model": self.model, "response_id": payload.get("id", "")}
- def provider_from_wbb(session: ClientSession) -> OpenAICompatibleAssistant:
- return OpenAICompatibleAssistant(
- session,
- base_url=str(getattr(wbb, "BUSINESS_ASSISTANT_OPENAI_BASE_URL", "")),
- api_key=str(getattr(wbb, "BUSINESS_ASSISTANT_OPENAI_API_KEY", "")),
- model=str(getattr(wbb, "BUSINESS_ASSISTANT_OPENAI_MODEL", "")),
- timeout_seconds=int(
- getattr(wbb, "BUSINESS_ASSISTANT_OPENAI_TIMEOUT_SECONDS", 30)
- ),
- max_output_tokens=int(
- getattr(wbb, "BUSINESS_ASSISTANT_OPENAI_MAX_OUTPUT_TOKENS", 600)
- ),
- )
- class KnowledgeIngestionService:
- def __init__(self, provider: OpenAICompatibleAssistant) -> None:
- self.provider = provider
- async def ingest(
- self,
- source: dict[str, Any],
- *,
- event_key: str,
- content: str,
- metadata: dict[str, Any],
- question_context: str = "",
- auto_publish: bool = False,
- force: bool = False,
- ) -> list[dict[str, Any]]:
- event, changed = await record_source_event(
- source,
- event_key=event_key,
- content=content,
- metadata={
- **metadata,
- "question_context": question_context,
- "trusted_author": bool(metadata.get("trusted_author")),
- },
- )
- if not changed and not force:
- return []
- try:
- items = await self.provider.extract_knowledge(
- source_type=str(metadata.get("effective_source_type") or source["source_type"]),
- source_title=str(source.get("title") or ""),
- content=content,
- question_context=question_context,
- )
- return await replace_event_candidates(
- source, event, items, auto_publish=auto_publish
- )
- except Exception as exc:
- await invalidate_source_event_candidates(
- str(event["event_id"]), reason="extraction_failed"
- )
- await update_source_event_extraction(
- str(event["event_id"]), status="failed", error=str(exc)
- )
- raise
- async def regenerate_candidate(self, candidate_id: str) -> list[dict[str, Any]]:
- candidate = await get_knowledge_candidate(candidate_id)
- if not candidate:
- raise AssistantProviderError("未找到知识候选。")
- source = await get_knowledge_source(str(candidate["source_id"]))
- event = await get_source_event(str(candidate["event_id"]))
- if not source or not event or event.get("deleted"):
- raise AssistantProviderError("知识来源或原始事件已失效。")
- metadata = event.get("metadata") if isinstance(event.get("metadata"), dict) else {}
- auto_publish = bool(
- source.get("publication_mode") == "auto" and metadata.get("trusted_author")
- )
- return await self.ingest(
- source,
- event_key=str(event["event_key"]),
- content=str(event.get("content") or ""),
- metadata=metadata,
- question_context=str(metadata.get("question_context") or ""),
- auto_publish=auto_publish,
- force=True,
- )
- class BusinessAssistantRuntime:
- def __init__(
- self,
- *,
- token: str,
- session: ClientSession,
- provider: OpenAICompatibleAssistant | None = None,
- ) -> None:
- self.api = TelegramBusinessApi(token, session)
- self.feishu = FeishuWebhookClient(session)
- self.provider = provider or provider_from_wbb(session)
- self.ingestion = KnowledgeIngestionService(self.provider)
- self._poll_task: asyncio.Task[None] | None = None
- self._digest_task: asyncio.Task[None] | None = None
- self._stopping = asyncio.Event()
- async def start(self) -> dict[str, Any]:
- if self._poll_task and not self._poll_task.done():
- return await runtime_status()
- self._stopping.clear()
- await update_runtime_status(
- {
- "polling_state": "starting",
- "model_configured": self.provider.configured,
- "last_error": "",
- }
- )
- try:
- me = await self.api.get_me()
- webhook = await self.api.get_webhook_info()
- except TelegramBotApiError as exc:
- return await update_runtime_status(
- {"polling_state": "error", "last_error": str(exc)}
- )
- webhook_url = str(webhook.get("url") or "")
- supported = bool(me.get("can_connect_to_business"))
- await update_runtime_status(
- {
- "business_mode_supported": supported,
- "webhook_conflict": bool(webhook_url),
- "webhook_url": webhook_url,
- "bot_username": str(me.get("username") or ""),
- }
- )
- if webhook_url:
- return await update_runtime_status(
- {
- "polling_state": "blocked",
- "last_error": "检测到 Telegram webhook,Business 长轮询未启动。",
- }
- )
- if not supported:
- return await update_runtime_status(
- {
- "polling_state": "blocked",
- "last_error": "请先在 BotFather 为机器人开启 Business Mode。",
- }
- )
- await self._refresh_connections()
- self._poll_task = asyncio.create_task(
- self._poll_loop(), name="business-assistant-poller"
- )
- self._digest_task = asyncio.create_task(
- self._digest_loop(), name="business-assistant-digest"
- )
- return await update_runtime_status(
- {"polling_state": "running", "last_error": "", "started_at": utc_now()}
- )
- async def _refresh_connections(self) -> None:
- page = 1
- while True:
- connections, total = await list_business_connections(
- page=page, page_size=100
- )
- for connection in connections:
- connection_id = str(connection["connection_id"])
- try:
- payload = await self.api.get_business_connection(connection_id)
- await upsert_business_connection(payload)
- except TelegramBotApiError as exc:
- await touch_business_connection(
- connection_id,
- error=str(exc),
- is_enabled=False if exc.error_code in {400, 403} else None,
- )
- if not connections or page * 100 >= total:
- return
- page += 1
- async def stop(self) -> None:
- self._stopping.set()
- tasks = [task for task in (self._poll_task, self._digest_task) if task]
- for task in tasks:
- task.cancel()
- if tasks:
- await asyncio.gather(*tasks, return_exceptions=True)
- self._poll_task = None
- self._digest_task = None
- await update_runtime_status(
- {"polling_state": "stopped", "stopped_at": utc_now()}
- )
- async def _poll_loop(self) -> None:
- offset = await load_update_offset()
- backoff = 1
- while not self._stopping.is_set():
- try:
- updates = await self.api.get_updates(offset=offset, poll_timeout=30)
- backoff = 1
- await update_runtime_status(
- {"polling_state": "running", "last_poll_at": utc_now(), "last_error": ""}
- )
- for update in updates:
- update_id = int(update.get("update_id") or 0)
- if update_id <= 0:
- continue
- await self.process_update(update)
- offset = max(offset, update_id + 1)
- await save_update_offset(offset)
- except asyncio.CancelledError:
- raise
- except TelegramBotApiError as exc:
- wait_for = exc.retry_after or backoff
- await update_runtime_status(
- {"polling_state": "retrying", "last_error": str(exc)}
- )
- await asyncio.sleep(min(max(wait_for, 1), 60))
- backoff = min(backoff * 2, 60)
- except Exception as exc:
- await update_runtime_status(
- {"polling_state": "retrying", "last_error": str(exc)[:1000]}
- )
- await asyncio.sleep(backoff)
- backoff = min(backoff * 2, 60)
- async def process_update(self, update: dict[str, Any]) -> None:
- update_id = int(update.get("update_id") or 0)
- if not await claim_update(update_id):
- return
- last_error = ""
- for attempt in range(1, 4):
- try:
- if isinstance(update.get("business_connection"), dict):
- await upsert_business_connection(update["business_connection"])
- elif isinstance(update.get("business_message"), dict):
- await self._handle_business_message(update["business_message"])
- elif isinstance(update.get("edited_business_message"), dict):
- await self._handle_edited_message(update["edited_business_message"])
- elif isinstance(update.get("deleted_business_messages"), dict):
- await self._handle_deleted_messages(update["deleted_business_messages"])
- await mark_update_done(update_id)
- return
- except TelegramBotApiError as exc:
- last_error = str(exc)
- if exc.retry_after:
- await asyncio.sleep(min(exc.retry_after, 60))
- except Exception as exc:
- last_error = str(exc)
- if attempt < 3:
- await asyncio.sleep(attempt)
- await mark_update_failed(update_id, last_error)
- await dead_letter_update(update, last_error or "unknown update error")
- async def _resolve_connection(self, connection_id: str) -> dict[str, Any]:
- connection = await get_business_connection(connection_id)
- if connection:
- return connection
- payload = await self.api.get_business_connection(connection_id)
- return await upsert_business_connection(payload)
- async def _handle_business_message(self, message: dict[str, Any]) -> None:
- connection_id = str(message.get("business_connection_id") or "")
- if not connection_id or message.get("is_from_offline"):
- return
- if message.get("sender_business_bot"):
- return
- connection = await self._resolve_connection(connection_id)
- await touch_business_connection(connection_id)
- chat = message.get("chat") if isinstance(message.get("chat"), dict) else {}
- sender = message.get("from") if isinstance(message.get("from"), dict) else {}
- chat_id = int(chat.get("id") or 0)
- message_id = int(message.get("message_id") or 0)
- if not chat_id or not message_id:
- return
- owner_id = int((connection.get("user") or {}).get("id") or 0)
- owner_reply = int(sender.get("id") or 0) == owner_id
- conversation = await get_or_create_conversation(
- connection_id,
- chat_id,
- customer=None if owner_reply else sender,
- )
- if owner_reply:
- settings = await get_account_settings(connection_id)
- human_text = str(message.get("text") or message.get("caption") or "").strip()
- await append_conversation_message(
- conversation["conversation_id"],
- direction="human",
- telegram_message_id=message_id,
- text=human_text,
- sender_id=owner_id,
- )
- await pause_conversation_for_human(
- conversation["conversation_id"],
- hours=int(settings.get("human_pause_hours") or 24),
- )
- if human_text:
- await self._learn_from_human_reply(
- connection_id=connection_id,
- conversation_id=str(conversation["conversation_id"]),
- chat_id=chat_id,
- message_id=message_id,
- owner_id=owner_id,
- human_text=human_text,
- )
- return
- await self._handle_customer_message(
- connection=connection,
- conversation=conversation,
- message=message,
- sender=sender,
- )
- async def _learn_from_human_reply(
- self,
- *,
- connection_id: str,
- conversation_id: str,
- chat_id: int,
- message_id: int,
- owner_id: int,
- human_text: str,
- ) -> None:
- sources = await list_enabled_business_sources(connection_id)
- if not sources:
- return
- messages = await recent_conversation_messages(conversation_id, limit=20)
- customer_question = next(
- (
- str(item.get("text") or "")
- for item in reversed(messages)
- if item.get("direction") == "incoming"
- and str(item.get("text") or "") != "[非文本消息]"
- ),
- "",
- )
- if not customer_question:
- return
- for source in sources:
- event_key = f"business:{chat_id}:{message_id}"
- existing_event = await get_source_event_by_key(
- str(source["source_id"]), event_key
- )
- existing_metadata = (
- existing_event.get("metadata")
- if existing_event and isinstance(existing_event.get("metadata"), dict)
- else {}
- )
- source_question = str(
- existing_metadata.get("question_context") or customer_question
- )
- if any(
- term in f"{source_question}\n{human_text}".casefold()
- for term in ("承诺", "保证", "投诉", "退款", "退钱", "赔偿", "律师", "起诉")
- ):
- continue
- try:
- await self.ingestion.ingest(
- source,
- event_key=event_key,
- content=human_text,
- question_context=source_question,
- metadata={
- "effective_source_type": "business",
- "chat_id": int(chat_id),
- "message_id": int(message_id),
- "author_id": int(owner_id),
- "trusted_author": True,
- },
- auto_publish=source.get("publication_mode") == "auto",
- )
- except Exception as exc:
- await update_runtime_status(
- {
- "knowledge_ingestion_last_error": str(exc)[:1000],
- "knowledge_ingestion_last_error_at": utc_now(),
- }
- )
- async def _handle_customer_message(
- self,
- *,
- connection: dict[str, Any],
- conversation: dict[str, Any],
- message: dict[str, Any],
- sender: dict[str, Any],
- ) -> None:
- connection_id = str(connection["connection_id"])
- chat_id = int((message.get("chat") or {}).get("id") or 0)
- message_id = int(message.get("message_id") or 0)
- text = str(message.get("text") or "").strip()
- settings = await get_account_settings(connection_id)
- await append_conversation_message(
- conversation["conversation_id"],
- direction="incoming",
- telegram_message_id=message_id,
- text=text or "[非文本消息]",
- sender_id=int(sender.get("id") or 0),
- metadata={"content_type": "text" if text else "unsupported"},
- )
- await self._notify_feishu_customer_message(
- connection=connection,
- conversation=conversation,
- message=message,
- sender=sender,
- settings=settings,
- text=text,
- )
- if not connection.get("is_enabled") or not settings.get("assistant_enabled"):
- return
- current = await get_conversation(conversation["conversation_id"]) or conversation
- if current.get("status") == "human_paused":
- paused_until = current.get("paused_until")
- if isinstance(paused_until, datetime) and (
- paused_until.replace(tzinfo=UTC) if paused_until.tzinfo is None else paused_until
- ) <= utc_now():
- current = await resume_conversation(current["conversation_id"])
- else:
- return
- if current.get("status") in {"handoff", "closed"}:
- return
- if not business_reply_window_open(message):
- await self._handoff(
- connection,
- current,
- settings,
- "reply_window_expired",
- send_customer_notice=False,
- )
- return
- rights = connection.get("rights") or {}
- if not rights.get("can_reply"):
- await self._handoff(
- connection, current, settings, "can_reply_missing", send_customer_notice=False
- )
- return
- if rights.get("can_read_messages"):
- with suppress(TelegramBotApiError):
- await self.api.read_business_message(connection_id, chat_id, message_id)
- if not text:
- await self._handoff(
- connection,
- current,
- settings,
- "unsupported_message",
- customer_notice=str(settings.get("unsupported_message")),
- )
- return
- classification = classify_handoff(text)
- if classification:
- await self._handoff(connection, current, settings, classification)
- return
- knowledge = await match_knowledge(connection_id, text)
- if not knowledge:
- await self._handoff(connection, current, settings, "knowledge_not_found")
- return
- allowed, quota_reason = await reserve_ai_usage(
- connection_id, chat_id, settings
- )
- if not allowed:
- await self._handoff(connection, current, settings, quota_reason)
- return
- messages = await recent_conversation_messages(current["conversation_id"], limit=12)
- try:
- decision = await self.provider.decide(
- settings=settings,
- conversation=current,
- messages=messages,
- knowledge=knowledge,
- customer_text=text,
- )
- except AssistantProviderError as exc:
- await update_runtime_status(
- {"provider_last_error": str(exc)[:1000], "provider_last_error_at": utc_now()}
- )
- await self._handoff(
- connection,
- current,
- settings,
- "provider_error",
- )
- return
- if decision["summary"]:
- await update_conversation_summary(
- current["conversation_id"], decision["summary"]
- )
- if decision["action"] == "handoff":
- await self._handoff(
- connection,
- current,
- settings,
- decision["handoff_reason"],
- summary=decision["summary"],
- )
- return
- response = await self.api.send_message(
- chat_id,
- decision["reply"],
- business_connection_id=connection_id,
- )
- await append_conversation_message(
- current["conversation_id"],
- direction="assistant",
- telegram_message_id=int(response.get("message_id") or 0),
- text=decision["reply"],
- sender_id=int((response.get("sender_business_bot") or {}).get("id") or 0),
- metadata={
- "action": decision["action"],
- "matched_entry_ids": decision["matched_entry_ids"],
- },
- )
- async def _notify_feishu_customer_message(
- self,
- *,
- connection: dict[str, Any],
- conversation: dict[str, Any],
- message: dict[str, Any],
- sender: dict[str, Any],
- settings: dict[str, Any],
- text: str,
- ) -> None:
- webhook_url = str(settings.get("feishu_webhook_url") or "")
- if not settings.get("feishu_webhook_enabled") or not webhook_url:
- return
- conversation_id = str(conversation["conversation_id"])
- message_id = int(message.get("message_id") or 0)
- if not await claim_feishu_message_notification(conversation_id, message_id):
- return
- user_name = " ".join(
- str(item).strip()
- for item in (sender.get("first_name"), sender.get("last_name"))
- if str(item or "").strip()
- )
- if not user_name:
- user_name = (
- f"@{sender.get('username')}"
- if sender.get("username")
- else str((message.get("chat") or {}).get("id") or "未知")
- )
- account = connection.get("user") or {}
- account_name = " ".join(
- str(item).strip()
- for item in (account.get("first_name"), account.get("last_name"))
- if str(item or "").strip()
- ) or (f"@{account.get('username')}" if account.get("username") else "Business 账号")
- message_type = "文本" if text else "非文本"
- preview_text = str(
- message.get("text") or message.get("caption") or ""
- ).strip()
- content: list[list[dict[str, str]]] = [
- [{"tag": "text", "text": f"账号:{account_name}"}],
- [{"tag": "text", "text": "用户:"}],
- [{"tag": "text", "text": f"类型:{message_type}"}],
- ]
- telegram_username = str(sender.get("username") or "").strip().lstrip("@")
- if re.fullmatch(r"[A-Za-z0-9_]{1,64}", telegram_username):
- content[1].append(
- {
- "tag": "a",
- "text": user_name,
- "href": f"https://t.me/{telegram_username}",
- }
- )
- else:
- content[1].append({"tag": "text", "text": user_name})
- if settings.get("feishu_message_preview_enabled"):
- content.append(
- [
- {
- "tag": "text",
- "text": f"内容:{preview_text[:500] if preview_text else '[非文本消息]'}",
- }
- ]
- )
- try:
- message_time = datetime.fromtimestamp(
- int(message.get("date") or 0), UTC
- ).strftime("%Y-%m-%d %H:%M:%S UTC")
- except (OSError, OverflowError, TypeError, ValueError):
- message_time = "未知"
- content.append([{"tag": "text", "text": f"时间:{message_time}"}])
- try:
- await self.feishu.send_post(
- webhook_url,
- title="消息提醒",
- content=content,
- signing_secret=str(settings.get("feishu_webhook_signing_secret") or ""),
- )
- except FeishuWebhookError as exc:
- await mark_feishu_message_notification(
- conversation_id, message_id, status="failed", error=str(exc)
- )
- await update_runtime_status(
- {
- "feishu_notification_last_error": str(exc)[:1000],
- "feishu_notification_last_error_at": utc_now(),
- }
- )
- return
- await mark_feishu_message_notification(
- conversation_id, message_id, status="sent"
- )
- await update_runtime_status(
- {
- "feishu_notification_last_error": "",
- "feishu_notification_last_sent_at": utc_now(),
- }
- )
- async def test_feishu_webhook(
- self, connection_id: str, overrides: dict[str, Any] | None = None
- ) -> dict[str, Any]:
- connection = await self._resolve_connection(connection_id)
- current = await get_account_settings(connection_id)
- settings = normalize_account_settings(overrides or {}, previous=current)
- webhook_url = str(settings.get("feishu_webhook_url") or "")
- if not webhook_url:
- raise FeishuWebhookError("请先填写飞书群机器人 Webhook 地址。")
- account = connection.get("user") or {}
- account_name = " ".join(
- str(item).strip()
- for item in (account.get("first_name"), account.get("last_name"))
- if str(item or "").strip()
- ) or (f"@{account.get('username')}" if account.get("username") else "Business 账号")
- await self.feishu.send(
- webhook_url,
- "\n".join(
- [
- "消息提醒",
- "飞书群机器人连接测试成功。",
- f"账号:{account_name}",
- ]
- ),
- signing_secret=str(settings.get("feishu_webhook_signing_secret") or ""),
- )
- return {"ok": True, "account": account_name}
- async def _handoff(
- self,
- connection: dict[str, Any],
- conversation: dict[str, Any],
- settings: dict[str, Any],
- reason: str,
- *,
- summary: str = "",
- customer_notice: str = "",
- send_customer_notice: bool = True,
- ) -> None:
- updated = await set_conversation_handoff(
- conversation["conversation_id"], reason, summary=summary
- )
- notice = customer_notice or str(settings.get("handoff_message") or "")
- if send_customer_notice and notice and connection.get("rights", {}).get("can_reply"):
- try:
- response = await self.api.send_message(
- int(updated["chat_id"]),
- notice,
- business_connection_id=str(connection["connection_id"]),
- )
- await append_conversation_message(
- updated["conversation_id"],
- direction="assistant",
- telegram_message_id=int(response.get("message_id") or 0),
- text=notice,
- metadata={"action": "handoff", "reason": reason},
- )
- except TelegramBotApiError as exc:
- await update_runtime_status(
- {
- "customer_notice_last_error": str(exc)[:1000],
- "customer_notice_last_error_at": utc_now(),
- }
- )
- await self.notify_handoff(connection, updated, settings, reason)
- async def notify_handoff(
- self,
- connection: dict[str, Any],
- conversation: dict[str, Any],
- settings: dict[str, Any],
- reason: str,
- ) -> None:
- customer = conversation.get("customer") or {}
- display = " ".join(
- item for item in (customer.get("first_name"), customer.get("last_name")) if item
- ) or (f"@{customer.get('username')}" if customer.get("username") else str(conversation["chat_id"]))
- latest = await recent_conversation_messages(conversation["conversation_id"], limit=1)
- latest_text = str(latest[-1].get("text") or "")[:1000] if latest else ""
- text = (
- "需要人工接待\n"
- f"账号:{(connection.get('user') or {}).get('first_name') or connection['connection_id']}\n"
- f"客户:{display}\n"
- f"原因:{reason}\n"
- f"原消息:chat_id={conversation['chat_id']} / "
- f"message_id={latest[-1].get('telegram_message_id') if latest else '未知'}\n"
- f"摘要:{conversation.get('summary') or latest_text or '暂无'}"
- )
- markup = {
- "inline_keyboard": [
- [
- {
- "text": "恢复自动回复",
- "callback_data": f"ba:resume:{conversation['conversation_id']}",
- }
- ]
- ]
- }
- destinations: list[int] = []
- destination = str(settings.get("notification_destination") or "owner")
- if destination in {"owner", "both"} and int(connection.get("user_chat_id") or 0):
- destinations.append(int(connection["user_chat_id"]))
- if destination in {"ops", "both"} and int(settings.get("ops_group_id") or 0):
- destinations.append(int(settings["ops_group_id"]))
- for chat_id in dict.fromkeys(destinations):
- try:
- await self.api.send_message(chat_id, text, reply_markup=markup)
- except TelegramBotApiError as exc:
- await update_runtime_status(
- {
- "notification_last_error": str(exc)[:1000],
- "notification_last_error_at": utc_now(),
- }
- )
- async def _handle_edited_message(self, message: dict[str, Any]) -> None:
- connection_id = str(message.get("business_connection_id") or "")
- chat_id = int((message.get("chat") or {}).get("id") or 0)
- if not connection_id or not chat_id or message.get("sender_business_bot"):
- return
- await touch_business_connection(connection_id)
- conversation = await get_conversation_by_chat(connection_id, chat_id)
- if not conversation:
- return
- connection = await self._resolve_connection(connection_id)
- sender_id = int((message.get("from") or {}).get("id") or 0)
- owner_id = int((connection.get("user") or {}).get("id") or 0)
- message_id = int(message.get("message_id") or 0)
- text = str(message.get("text") or message.get("caption") or "").strip()
- if sender_id == owner_id:
- settings = await get_account_settings(connection_id)
- await append_conversation_message(
- conversation["conversation_id"],
- direction="human",
- telegram_message_id=-message_id,
- text=text or "[人工消息已编辑]",
- sender_id=owner_id,
- metadata={"edited_message_id": message_id},
- )
- await pause_conversation_for_human(
- conversation["conversation_id"],
- hours=int(settings.get("human_pause_hours") or 24),
- )
- if text:
- await self._learn_from_human_reply(
- connection_id=connection_id,
- conversation_id=str(conversation["conversation_id"]),
- chat_id=chat_id,
- message_id=message_id,
- owner_id=owner_id,
- human_text=text,
- )
- return
- await append_conversation_message(
- conversation["conversation_id"],
- direction="incoming",
- telegram_message_id=-message_id,
- text=text or "[消息已编辑]",
- sender_id=sender_id,
- metadata={"edited_message_id": message_id},
- )
- async def _handle_deleted_messages(self, payload: dict[str, Any]) -> None:
- connection_id = str(payload.get("business_connection_id") or "")
- chat_id = int((payload.get("chat") or {}).get("id") or 0)
- if connection_id:
- await touch_business_connection(connection_id)
- conversation = await get_conversation_by_chat(connection_id, chat_id)
- if not conversation:
- return
- for source in await list_business_sources(connection_id):
- for message_id in payload.get("message_ids") or []:
- await mark_source_event_deleted(
- str(source["source_id"]), f"business:{chat_id}:{int(message_id)}"
- )
- await append_conversation_message(
- conversation["conversation_id"],
- direction="incoming",
- telegram_message_id=-abs(int((payload.get("message_ids") or [0])[0] or 0)) - 1,
- text="[消息已删除]",
- metadata={"deleted_message_ids": payload.get("message_ids") or []},
- )
- async def _digest_loop(self) -> None:
- while not self._stopping.is_set():
- try:
- for item in await due_digest_connections():
- connection = item["connection"]
- settings = item["settings"]
- metrics = await usage_metrics(
- connection_id=str(connection["connection_id"]),
- day=str(item["day"]),
- )
- text = (
- f"智能接待运营简报 · {item['day']}\n"
- f"客户数:{metrics['customers']}\n"
- f"AI 调用:{metrics['ai_calls']}\n"
- f"待人工:{metrics['handoffs']}\n"
- f"人工暂停:{metrics['human_paused']}"
- )
- destinations: list[int] = []
- destination = str(settings.get("notification_destination") or "owner")
- if destination in {"owner", "both"} and connection.get("user_chat_id"):
- destinations.append(int(connection["user_chat_id"]))
- if destination in {"ops", "both"} and settings.get("ops_group_id"):
- destinations.append(int(settings["ops_group_id"]))
- sent = False
- for chat_id in dict.fromkeys(destinations):
- try:
- await self.api.send_message(chat_id, text)
- sent = True
- except TelegramBotApiError as exc:
- await update_runtime_status(
- {
- "digest_last_error": str(exc)[:1000],
- "digest_last_error_at": utc_now(),
- }
- )
- if sent:
- await mark_digest_sent(connection["connection_id"], item["day"])
- except asyncio.CancelledError:
- raise
- except Exception as exc:
- await update_runtime_status({"digest_last_error": str(exc)[:1000]})
- await asyncio.sleep(60)
- async def preview_answer(self, connection_id: str, text: str) -> dict[str, Any]:
- settings = await get_account_settings(connection_id)
- knowledge = await match_knowledge(connection_id, text)
- classification = classify_handoff(text)
- if classification or not knowledge:
- return {
- "action": "handoff",
- "reason": classification or "knowledge_not_found",
- "matched_entries": knowledge,
- }
- conversation = {"summary": ""}
- decision = await self.provider.decide(
- settings=settings,
- conversation=conversation,
- messages=[],
- knowledge=knowledge,
- customer_text=text,
- )
- return {**decision, "matched_entries": knowledge}
- async def runtime_overview() -> dict[str, Any]:
- status = await runtime_status()
- status["usage"] = await usage_metrics()
- return status
|