294 lines
10 KiB
Python
294 lines
10 KiB
Python
"""Asynchronous conversation sedimentation into WeKnora-backed LLMWiki."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import logging
|
|
from dataclasses import dataclass
|
|
from typing import Any
|
|
|
|
from deerflow.integrations.weknora.client import WeKnoraClient, WeKnoraError
|
|
from deerflow.integrations.weknora.runtime import build_weknora_client, get_resolved_llmwiki_runtime
|
|
from deerflow.persistence.llmwiki import LlmWikiStore
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
CONVERSATION_DEPOSIT_KB_NAME = "对话沉淀"
|
|
CONVERSATION_DEPOSIT_OWNER_USER_ID = "system"
|
|
CONVERSATION_DEPOSIT_KB_DESCRIPTION = "cmzs 全局问答沉淀知识库。所有用户可见,不默认参与问答检索。"
|
|
_DEPOSIT_KB_LOCK = asyncio.Lock()
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class ConversationTurnDeposit:
|
|
thread_id: str
|
|
question: str
|
|
answer: str
|
|
human_message_id: str | None
|
|
assistant_message_id: str
|
|
assistant_id: str | None
|
|
knowledge_base_ids: list[str]
|
|
created_by_user_id: str | None
|
|
|
|
|
|
def _store(app: Any) -> LlmWikiStore | None:
|
|
return getattr(app.state, "llmwiki_store", None)
|
|
|
|
|
|
def _client(app: Any) -> WeKnoraClient | None:
|
|
runtime = get_resolved_llmwiki_runtime(getattr(app.state, "config", None))
|
|
if not runtime.weknora_enabled:
|
|
return None
|
|
try:
|
|
return build_weknora_client(runtime)
|
|
except RuntimeError:
|
|
logger.debug("LLMWiki conversation deposit skipped: WeKnora client is unavailable", exc_info=True)
|
|
return None
|
|
|
|
|
|
def _shorten(text: str, max_chars: int) -> str:
|
|
value = " ".join(str(text or "").split())
|
|
if len(value) <= max_chars:
|
|
return value
|
|
return f"{value[: max_chars - 1].rstrip()}…"
|
|
|
|
|
|
def _title_for_turn(turn: ConversationTurnDeposit) -> str:
|
|
question = _shorten(turn.question, 48)
|
|
return f"对话沉淀 - {question or turn.assistant_message_id[:12]}"
|
|
|
|
|
|
def _format_markdown(turn: ConversationTurnDeposit, title: str) -> str:
|
|
kb_line = "、".join(turn.knowledge_base_ids) if turn.knowledge_base_ids else "未显式选择知识库"
|
|
return f"""---
|
|
source: cmzs
|
|
type: conversation-deposit
|
|
thread_id: {turn.thread_id}
|
|
human_message_id: {turn.human_message_id or ""}
|
|
assistant_message_id: {turn.assistant_message_id}
|
|
assistant_id: {turn.assistant_id or ""}
|
|
knowledge_base_ids: {kb_line}
|
|
created_by_user_id: {turn.created_by_user_id or ""}
|
|
---
|
|
|
|
# {title}
|
|
|
|
> 本页由 cmzs 在一轮问答结束后自动沉淀,用于后续 LLMWiki 检索、Wiki 索引与图谱构建。
|
|
|
|
## 摘要
|
|
|
|
{_shorten(turn.answer, 360) or "本轮回答为空。"}
|
|
|
|
## 用户问题
|
|
|
|
{turn.question.strip() or "(空)"}
|
|
|
|
## 助手回答
|
|
|
|
{turn.answer.strip() or "(空)"}
|
|
|
|
## 结构化线索
|
|
|
|
- 来源线程:`{turn.thread_id}`
|
|
- 用户消息:`{turn.human_message_id or ""}`
|
|
- 回答消息:`{turn.assistant_message_id}`
|
|
- 智能体:`{turn.assistant_id or "通用问答"}`
|
|
- 本轮显式知识库:{kb_line}
|
|
|
|
## LLMWiki 标签
|
|
|
|
#对话沉淀 #cmzs #问答知识
|
|
"""
|
|
|
|
|
|
async def _sync_conversation_deposit_description(
|
|
store: LlmWikiStore,
|
|
client: WeKnoraClient,
|
|
row: dict[str, Any],
|
|
) -> dict[str, Any]:
|
|
if str(row.get("description") or "") == CONVERSATION_DEPOSIT_KB_DESCRIPTION:
|
|
return row
|
|
try:
|
|
await client.update_knowledge_base(
|
|
str(row["weknora_id"]),
|
|
{"description": CONVERSATION_DEPOSIT_KB_DESCRIPTION},
|
|
)
|
|
except WeKnoraError:
|
|
logger.warning("Failed to update the conversation deposit knowledge-base description", exc_info=True)
|
|
return row
|
|
updated = await store.update_mapping(
|
|
str(row["id"]),
|
|
description=CONVERSATION_DEPOSIT_KB_DESCRIPTION,
|
|
)
|
|
return updated or row
|
|
|
|
|
|
async def ensure_conversation_deposit_mapping(
|
|
store: LlmWikiStore,
|
|
client: WeKnoraClient,
|
|
*,
|
|
remote_knowledge_bases: list[dict[str, Any]] | None = None,
|
|
) -> dict[str, Any] | None:
|
|
existing = await store.get_conversation_deposit_mapping()
|
|
if existing is not None:
|
|
return await _sync_conversation_deposit_description(store, client, existing)
|
|
|
|
async with _DEPOSIT_KB_LOCK:
|
|
existing = await store.get_conversation_deposit_mapping()
|
|
if existing is not None:
|
|
return await _sync_conversation_deposit_description(store, client, existing)
|
|
|
|
remote_id = ""
|
|
try:
|
|
remote_items = remote_knowledge_bases
|
|
if remote_items is None:
|
|
remote_items = await client.list_knowledge_bases()
|
|
for remote in remote_items:
|
|
if str(remote.get("name") or "").strip() == CONVERSATION_DEPOSIT_KB_NAME and remote.get("id"):
|
|
remote_id = str(remote["id"])
|
|
break
|
|
except WeKnoraError:
|
|
logger.debug("Could not scan WeKnora knowledge bases before creating deposit base", exc_info=True)
|
|
|
|
if remote_id:
|
|
mapped = await store.get_by_weknora_id(remote_id)
|
|
if mapped is not None:
|
|
try:
|
|
await client.update_knowledge_base(
|
|
remote_id,
|
|
{"description": CONVERSATION_DEPOSIT_KB_DESCRIPTION},
|
|
)
|
|
except WeKnoraError:
|
|
logger.warning("Failed to update the existing conversation deposit knowledge-base description", exc_info=True)
|
|
return await store.adopt_conversation_deposit_mapping(
|
|
str(mapped["id"]),
|
|
name=CONVERSATION_DEPOSIT_KB_NAME,
|
|
description=CONVERSATION_DEPOSIT_KB_DESCRIPTION,
|
|
)
|
|
try:
|
|
remote = await client.update_knowledge_base(
|
|
remote_id,
|
|
{"description": CONVERSATION_DEPOSIT_KB_DESCRIPTION},
|
|
)
|
|
except WeKnoraError:
|
|
logger.warning("Failed to update the existing conversation deposit knowledge-base description", exc_info=True)
|
|
remote = {
|
|
"id": remote_id,
|
|
"name": CONVERSATION_DEPOSIT_KB_NAME,
|
|
"description": CONVERSATION_DEPOSIT_KB_DESCRIPTION,
|
|
"type": "document",
|
|
}
|
|
else:
|
|
try:
|
|
remote = await client.create_knowledge_base(
|
|
name=CONVERSATION_DEPOSIT_KB_NAME,
|
|
description=CONVERSATION_DEPOSIT_KB_DESCRIPTION,
|
|
kb_type="document",
|
|
wiki_enabled=True,
|
|
)
|
|
except WeKnoraError:
|
|
logger.warning("Failed to create WeKnora conversation deposit knowledge base", exc_info=True)
|
|
return None
|
|
|
|
try:
|
|
row = await store.create_mapping(
|
|
weknora_id=str(remote["id"]),
|
|
owner_user_id=CONVERSATION_DEPOSIT_OWNER_USER_ID,
|
|
name=CONVERSATION_DEPOSIT_KB_NAME,
|
|
description=CONVERSATION_DEPOSIT_KB_DESCRIPTION,
|
|
kb_type="document",
|
|
)
|
|
published = await store.request_publish(str(row["id"]), CONVERSATION_DEPOSIT_OWNER_USER_ID, is_admin=True)
|
|
return published or row
|
|
except Exception:
|
|
logger.exception("Failed to persist conversation deposit knowledge-base mapping")
|
|
return None
|
|
|
|
|
|
async def ensure_conversation_deposit_knowledge_base(
|
|
app: Any,
|
|
*,
|
|
owner_user_id: str | None,
|
|
) -> dict[str, Any] | None:
|
|
store = _store(app)
|
|
client = _client(app)
|
|
if store is None or client is None:
|
|
return None
|
|
return await ensure_conversation_deposit_mapping(store, client)
|
|
|
|
|
|
async def deposit_conversation_turn(app: Any, turn: ConversationTurnDeposit) -> None:
|
|
if not turn.thread_id or not turn.assistant_message_id:
|
|
return
|
|
if not turn.question.strip() and not turn.answer.strip():
|
|
return
|
|
|
|
store = _store(app)
|
|
client = _client(app)
|
|
if store is None or client is None:
|
|
return
|
|
|
|
kb = await ensure_conversation_deposit_knowledge_base(app, owner_user_id=turn.created_by_user_id)
|
|
if kb is None:
|
|
return
|
|
|
|
previous = await store.get_conversation_deposit(
|
|
thread_id=turn.thread_id,
|
|
assistant_message_id=turn.assistant_message_id,
|
|
)
|
|
previous_doc_id = str((previous or {}).get("weknora_knowledge_id") or "")
|
|
if previous_doc_id:
|
|
try:
|
|
await client.delete_document(previous_doc_id)
|
|
except WeKnoraError as exc:
|
|
if exc.status_code != 404:
|
|
logger.warning("Failed to delete old conversation deposit document %s", previous_doc_id, exc_info=True)
|
|
|
|
title = _title_for_turn(turn)
|
|
content = _format_markdown(turn, title)
|
|
try:
|
|
remote_doc = await client.create_manual_document(
|
|
str(kb["weknora_id"]),
|
|
title=title,
|
|
content=content,
|
|
)
|
|
except WeKnoraError:
|
|
logger.warning("Failed to create conversation deposit document", exc_info=True)
|
|
return
|
|
|
|
await store.upsert_conversation_deposit(
|
|
thread_id=turn.thread_id,
|
|
human_message_id=turn.human_message_id,
|
|
assistant_message_id=turn.assistant_message_id,
|
|
knowledge_base_mapping_id=str(kb["id"]),
|
|
weknora_knowledge_id=str(remote_doc.get("id") or remote_doc.get("knowledge_id") or ""),
|
|
title=title,
|
|
created_by_user_id=turn.created_by_user_id,
|
|
)
|
|
|
|
|
|
async def delete_conversation_deposit_documents(
|
|
app: Any,
|
|
*,
|
|
thread_id: str,
|
|
message_ids: list[str] | None = None,
|
|
) -> None:
|
|
store = _store(app)
|
|
client = _client(app)
|
|
if store is None or client is None:
|
|
return
|
|
try:
|
|
rows = await store.delete_conversation_deposits(thread_id=thread_id, message_ids=message_ids)
|
|
except Exception:
|
|
logger.debug("Failed to delete local conversation deposit rows", exc_info=True)
|
|
return
|
|
for row in rows:
|
|
doc_id = str(row.get("weknora_knowledge_id") or "")
|
|
if not doc_id:
|
|
continue
|
|
try:
|
|
await client.delete_document(doc_id)
|
|
except WeKnoraError as exc:
|
|
if exc.status_code != 404:
|
|
logger.warning("Failed to delete WeKnora conversation deposit document %s", doc_id, exc_info=True)
|