deerflow-code/offline-backend-20260512/backend/app/gateway/llmwiki_deposit.py
2026-09-07 18:24:55 +08:00

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)