494 lines
20 KiB
Python
494 lines
20 KiB
Python
"""LLM-based knowledge extraction (wiki-capture style).
|
||
|
||
Turns a conversation into a *declarative* knowledge note instead of dumping the
|
||
raw chat. Follows the obsidian-wiki ``wiki-capture`` skill: classify the note,
|
||
keep only high-value findings/decisions, drop greetings / intermediate
|
||
reasoning / verbatim transcripts, mark inferred content, and emit structured
|
||
JSON that we render into the Context / Finding / Reasoning / Implications /
|
||
Related body.
|
||
|
||
Falls back to the rule-based :func:`deerflow.knowledge.extractor.extract_thread`
|
||
when the LLM is unavailable or returns something unusable, so capture never
|
||
hard-fails.
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import json
|
||
import logging
|
||
import re
|
||
from typing import Any
|
||
|
||
from deerflow.knowledge import extractor
|
||
from deerflow.knowledge.schemas import NoteDraft, SourceDraft
|
||
from deerflow.knowledge.skill_templates import WIKI_CAPTURE_GUIDANCE
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
_VALID_CATEGORIES = {"synthesis", "concepts", "references", "decision", "journal"}
|
||
|
||
|
||
class LLMRejectedIngest(Exception):
|
||
"""The LLM judged the conversation as having no knowledge value.
|
||
|
||
Raised (instead of returning ``[]``) when the model explicitly sets
|
||
``should_ingest: false`` with no notes — so callers can *skip* ingestion
|
||
entirely rather than falling back to the rule-based verbatim dump.
|
||
"""
|
||
|
||
|
||
def _is_falsey(value: Any) -> bool:
|
||
if value is False:
|
||
return True
|
||
if isinstance(value, str):
|
||
return value.strip().lower() in ("false", "0", "no")
|
||
return False
|
||
|
||
_SYSTEM_PROMPT = """\
|
||
你是知识库整理专家,遵循 obsidian-wiki 的 wiki-capture 规范,把一次对话提炼成\
|
||
**陈述式知识(declarative knowledge)**,而不是聊天流水账。
|
||
|
||
要求:
|
||
- 只保留高价值内容:明确结论、方案、决策、关键事实、被引用的来源。
|
||
- 丢弃:寒暄、客套、中间推理过程、重复内容、无意义的原文堆砌。
|
||
- 不要写“用户问了什么、AI 回答了什么”,要写“结论是什么、为什么成立”。
|
||
- 推断内容用 ^[inferred] 标记,不确定/冲突用 ^[ambiguous] 标记。
|
||
- 分类 category 取值之一:synthesis(分析/方案/结论)、concepts(概念/框架)、\
|
||
references(外部来源汇总)、decision(架构/设计决策)、journal(会话纪要)。
|
||
- 用中文输出。
|
||
|
||
一次对话可能只讲一个主题,也可能覆盖多个相互独立的主题。请把它拆成 1~4 页知识:
|
||
- 默认只产出 **1 页**;只有当对话明确包含**多个彼此独立、可各自成文**的主题时才拆成多页。
|
||
- 每一页聚焦单一主题,标题与 summary 必须彼此区分,不要把同一结论拆到多页。
|
||
- 不值得沉淀时返回空的 notes 数组。
|
||
|
||
只输出一个 JSON 对象,不要任何额外文字或 Markdown 代码块,结构如下:
|
||
{
|
||
"should_ingest": true,
|
||
"reason": "为什么值得/不值得沉淀",
|
||
"notes": [
|
||
{
|
||
"title": "简洁的知识标题(不超过 40 字)",
|
||
"summary": "1-2 句话概括这页知识",
|
||
"category": "synthesis",
|
||
"context": "问题背景或触发场景,1-3 句",
|
||
"findings": ["结论/要点(陈述式,可含 ^[inferred])", "..."],
|
||
"reasoning": "为什么成立、取舍、不确定性",
|
||
"implications": ["后续建议/风险/下一步", "..."],
|
||
"tags": ["标签1", "标签2"],
|
||
"entities": [{"name": "实体名", "type": "technology|concept|person|org|product|other"}],
|
||
"relations": [{"from": "实体A", "to": "实体B", "type": "uses|relates_to|part_of|depends_on"}],
|
||
"confidence": 0.0
|
||
}
|
||
]
|
||
}"""
|
||
|
||
|
||
_DOC_SYSTEM_PROMPT = """\
|
||
你是知识库整理专家,遵循 obsidian-wiki 的 wiki-capture 规范,把一篇**文档**提炼成\
|
||
**陈述式知识(declarative knowledge)**,而不是照抄原文。
|
||
|
||
要求:
|
||
- 只保留高价值内容:核心结论、关键事实、定义、流程、数据、决策与其依据。
|
||
- 丢弃:目录、页眉页脚、版权声明、重复的样板文字、无意义的排版噪声。
|
||
- 用自己的话重组为结构化知识,不要逐段复制原文。
|
||
- 推断内容用 ^[inferred] 标记,不确定/冲突用 ^[ambiguous] 标记。
|
||
- 分类 category 取值之一:synthesis(分析/方案/结论)、concepts(概念/框架)、\
|
||
references(外部来源汇总)、decision(架构/设计决策)、journal(纪要)。
|
||
- 用中文输出。
|
||
|
||
一篇文档可能只讲一个主题,也可能覆盖多个相互独立的主题。请把它拆成 1~4 页知识:
|
||
- 默认只产出 **1 页**;只有当文档明确包含**多个彼此独立、可各自成文**的主题时才拆成多页。
|
||
- 每一页聚焦单一主题,标题与 summary 必须彼此区分。
|
||
- 文档没有可沉淀的知识时返回空的 notes 数组。
|
||
|
||
只输出一个 JSON 对象,不要任何额外文字或 Markdown 代码块,结构与对话提炼一致:
|
||
{
|
||
"should_ingest": true,
|
||
"reason": "为什么值得/不值得沉淀",
|
||
"notes": [
|
||
{
|
||
"title": "简洁的知识标题(不超过 40 字)",
|
||
"summary": "1-2 句话概括这页知识",
|
||
"category": "synthesis",
|
||
"context": "文档背景或适用场景,1-3 句",
|
||
"findings": ["结论/要点(陈述式,可含 ^[inferred])", "..."],
|
||
"reasoning": "为什么成立、取舍、不确定性",
|
||
"implications": ["后续建议/风险/下一步", "..."],
|
||
"tags": ["标签1", "标签2"],
|
||
"entities": [{"name": "实体名", "type": "technology|concept|person|org|product|other"}],
|
||
"relations": [{"from": "实体A", "to": "实体B", "type": "uses|relates_to|part_of|depends_on"}],
|
||
"confidence": 0.0
|
||
}
|
||
]
|
||
}"""
|
||
|
||
|
||
def _build_transcript(turns: list[dict[str, str]], sources: list[SourceDraft], *, max_chars: int) -> str:
|
||
parts: list[str] = []
|
||
for i, t in enumerate(turns, 1):
|
||
if t.get("user"):
|
||
parts.append(f"[用户 {i}] {t['user']}")
|
||
if t.get("ai"):
|
||
parts.append(f"[助手 {i}] {t['ai']}")
|
||
if sources:
|
||
parts.append("\n[来源]")
|
||
for s in sources[:20]:
|
||
label = s.title or s.url or s.tool_name or "来源"
|
||
parts.append(f"- {label}" + (f" ({s.url})" if s.url else ""))
|
||
text = "\n\n".join(parts)
|
||
if len(text) > max_chars:
|
||
text = text[:max_chars] + "\n…(已截断)"
|
||
return text
|
||
|
||
|
||
def _parse_json(raw: str) -> dict[str, Any] | None:
|
||
if not raw:
|
||
return None
|
||
text = raw.strip()
|
||
# Strip ```json fences if present.
|
||
fence = re.search(r"```(?:json)?\s*(.*?)```", text, re.DOTALL)
|
||
if fence:
|
||
text = fence.group(1).strip()
|
||
# Otherwise grab the outermost { ... }.
|
||
if not text.startswith("{"):
|
||
start = text.find("{")
|
||
end = text.rfind("}")
|
||
if start >= 0 and end > start:
|
||
text = text[start : end + 1]
|
||
try:
|
||
data = json.loads(text)
|
||
return data if isinstance(data, dict) else None
|
||
except json.JSONDecodeError:
|
||
logger.warning("Knowledge LLM returned non-JSON output; falling back to rule-based")
|
||
return None
|
||
|
||
|
||
def _render_body(data: dict[str, Any], title: str, related: list[str], sources: list[SourceDraft]) -> str:
|
||
def _bullets(items: Any) -> str:
|
||
if isinstance(items, str):
|
||
items = [items]
|
||
if not isinstance(items, list):
|
||
return ""
|
||
return "\n".join(f"- {str(x).strip()}" for x in items if str(x).strip())
|
||
|
||
context = str(data.get("context") or "").strip() or "- 来自一次对话沉淀。^[inferred]"
|
||
findings = _bullets(data.get("findings")) or "_暂无可提炼的结论。_"
|
||
reasoning = str(data.get("reasoning") or "").strip() or "(未提供推理说明)"
|
||
implications = _bullets(data.get("implications")) or "- 如内容过期或冲突,请编辑或归档。"
|
||
|
||
related_lines = [f"- {link}" for link in related]
|
||
source_lines = [f"- {s.title or s.url or '来源'}" + (f" — {s.url}" if s.url else "") for s in sources[:20]]
|
||
related_section = extractor._related_section_md(related_lines, source_lines)
|
||
|
||
return (
|
||
f"# {title}\n\n"
|
||
f"## Context\n\n{context}\n\n"
|
||
f"## Finding / Decision\n\n{findings}\n\n"
|
||
f"## Reasoning\n\n{reasoning}\n\n"
|
||
f"## Implications\n\n{implications}\n\n"
|
||
f"{related_section}"
|
||
)
|
||
|
||
|
||
def _notes_from_data(data: dict[str, Any]) -> list[dict[str, Any]]:
|
||
"""Normalize the LLM JSON into a list of per-note dicts.
|
||
|
||
Accepts the multi-note shape ``{"notes": [...]}`` as well as the legacy
|
||
single-note shape (``{"title": ..., "findings": ...}``) for backward
|
||
compatibility. Notes without any title/findings are dropped.
|
||
"""
|
||
raw_notes = data.get("notes")
|
||
if isinstance(raw_notes, list):
|
||
candidates = [n for n in raw_notes if isinstance(n, dict)]
|
||
elif data.get("title") or data.get("findings") or data.get("context"):
|
||
candidates = [data] # legacy single-object output
|
||
else:
|
||
candidates = []
|
||
return [n for n in candidates if str(n.get("title") or "").strip() or n.get("findings")]
|
||
|
||
|
||
def _build_note_draft(
|
||
note_data: dict[str, Any],
|
||
*,
|
||
source_type: str,
|
||
source_id: str | None,
|
||
source_key: str,
|
||
fallback_title: str,
|
||
title_override: str | None,
|
||
tags: list[str] | None,
|
||
project_name: str | None,
|
||
sources: list[SourceDraft],
|
||
) -> NoteDraft:
|
||
"""Render one note dict into a :class:`NoteDraft`.
|
||
|
||
Source-agnostic: the caller supplies ``source_type`` / ``source_id`` /
|
||
``source_key`` and a ``fallback_title`` (used when neither the user override
|
||
nor the LLM produced a usable title), so both thread and document capture
|
||
share one renderer.
|
||
"""
|
||
category = str(note_data.get("category") or "synthesis").strip().lower()
|
||
if category not in _VALID_CATEGORIES:
|
||
category = "synthesis"
|
||
|
||
final_title = (title_override or "").strip() or str(note_data.get("title") or "").strip() or fallback_title
|
||
summary = str(note_data.get("summary") or "").strip()
|
||
related = extractor.related_links(project_name)
|
||
body = _render_body(note_data, final_title, related, sources)
|
||
|
||
base_tags = list(tags or [])
|
||
llm_tags = note_data.get("tags") if isinstance(note_data.get("tags"), list) else []
|
||
for tag in [*llm_tags, project_name or "zncm", "knowledge-base"]:
|
||
tag = str(tag).strip()
|
||
if tag and tag not in base_tags:
|
||
base_tags.append(tag)
|
||
|
||
try:
|
||
confidence = float(note_data.get("confidence") or 0.7)
|
||
except (TypeError, ValueError):
|
||
confidence = 0.7
|
||
|
||
entities = note_data.get("entities") if isinstance(note_data.get("entities"), list) else []
|
||
relations = note_data.get("relations") if isinstance(note_data.get("relations"), list) else []
|
||
|
||
return NoteDraft(
|
||
title=final_title,
|
||
summary=summary or final_title,
|
||
content_md=body,
|
||
category=category,
|
||
source_type=source_type,
|
||
source_id=source_id,
|
||
tags=base_tags,
|
||
confidence=max(0.0, min(confidence, 1.0)),
|
||
related=related,
|
||
sources=sources,
|
||
entities=[e for e in entities if isinstance(e, dict) and e.get("name")],
|
||
relations=[r for r in relations if isinstance(r, dict) and r.get("from") and r.get("to")],
|
||
source_key=source_key,
|
||
)
|
||
|
||
|
||
async def extract_thread_llm_multi(
|
||
messages: list[dict[str, Any]],
|
||
*,
|
||
thread_id: str,
|
||
thread_title: str | None = None,
|
||
mode: str = "summary",
|
||
title: str | None = None,
|
||
tags: list[str] | None = None,
|
||
project_name: str | None = "zncm",
|
||
sources: list[SourceDraft] | None = None,
|
||
max_messages: int = 80,
|
||
model_name: str | None = None,
|
||
max_input_chars: int = 12000,
|
||
allow_multi: bool = True,
|
||
max_notes: int = 4,
|
||
system_prompt: str | None = None,
|
||
) -> list[NoteDraft]:
|
||
"""Distill a thread into one *or several* cross-linkable knowledge notes.
|
||
|
||
The LLM may split a conversation that covers multiple independent topics into
|
||
several declarative notes. Returns ``[]`` when there is nothing to sediment
|
||
or the LLM path fails — callers should fall back to the rule-based extractor.
|
||
When ``allow_multi`` is false only the first note is kept. ``system_prompt``
|
||
overrides the built-in distillation prompt with an operator-edited template
|
||
(used as-is, since the template is expected to be self-contained).
|
||
"""
|
||
turns = extractor.split_turns(messages, max_messages=max_messages)
|
||
if not turns:
|
||
return []
|
||
sources = sources or []
|
||
|
||
try:
|
||
from langchain_core.messages import HumanMessage, SystemMessage
|
||
|
||
from deerflow.models import create_chat_model
|
||
|
||
model = create_chat_model(name=model_name or None, thinking_enabled=False)
|
||
transcript = _build_transcript(turns, sources, max_chars=max_input_chars)
|
||
system_text = (system_prompt or "").strip() or (_SYSTEM_PROMPT + "\n\n" + WIKI_CAPTURE_GUIDANCE)
|
||
prompt = [
|
||
SystemMessage(content=system_text),
|
||
HumanMessage(content=f"以下是一次对话,请按规范提炼为知识 JSON:\n\n{transcript}"),
|
||
]
|
||
response = await model.ainvoke(prompt)
|
||
raw = response.content if isinstance(response.content, str) else str(response.content)
|
||
data = _parse_json(raw)
|
||
except Exception:
|
||
logger.exception("[knowledge] thread %s — LLM extraction call failed; using rule-based fallback", thread_id)
|
||
return []
|
||
|
||
if not data:
|
||
logger.warning("[knowledge] thread %s — LLM returned no parsable JSON; using rule-based fallback", thread_id)
|
||
return []
|
||
|
||
note_dicts = _notes_from_data(data)
|
||
if not note_dicts:
|
||
reason = str(data.get("reason") or "")[:200]
|
||
if _is_falsey(data.get("should_ingest")):
|
||
# The model deliberately judged this conversation worthless (greeting,
|
||
# meta chatter, no findings). Skip ingestion entirely — do NOT fall
|
||
# back to the rule-based verbatim dump.
|
||
logger.warning("[knowledge] thread %s SKIPPED — model judged not worth ingesting: %s", thread_id, reason)
|
||
raise LLMRejectedIngest(reason or "no knowledge value")
|
||
logger.warning(
|
||
"[knowledge] thread %s — LLM produced no notes (should_ingest=%s, reason=%r); using rule-based fallback",
|
||
thread_id,
|
||
data.get("should_ingest"),
|
||
reason,
|
||
)
|
||
return []
|
||
if not allow_multi:
|
||
note_dicts = note_dicts[:1]
|
||
else:
|
||
note_dicts = note_dicts[: max(1, max_notes)]
|
||
|
||
fallback_title = extractor._auto_title(thread_title, turns)
|
||
drafts: list[NoteDraft] = []
|
||
for i, nd in enumerate(note_dicts):
|
||
# The user-supplied title override only applies to a single-note result;
|
||
# multi-note splits keep each note's own LLM title.
|
||
override = title if (len(note_dicts) == 1) else None
|
||
suffix = "" if len(note_dicts) == 1 else f"#{i}"
|
||
drafts.append(
|
||
_build_note_draft(
|
||
nd,
|
||
source_type="thread",
|
||
source_id=thread_id,
|
||
source_key=f"thread:{thread_id}{suffix}",
|
||
fallback_title=fallback_title,
|
||
title_override=override,
|
||
tags=tags,
|
||
project_name=project_name,
|
||
sources=sources,
|
||
)
|
||
)
|
||
return drafts
|
||
|
||
|
||
async def extract_document_llm_multi(
|
||
text: str,
|
||
*,
|
||
doc_id: str,
|
||
doc_title: str | None = None,
|
||
title: str | None = None,
|
||
tags: list[str] | None = None,
|
||
project_name: str | None = "zncm",
|
||
sources: list[SourceDraft] | None = None,
|
||
model_name: str | None = None,
|
||
max_input_chars: int = 12000,
|
||
allow_multi: bool = True,
|
||
max_notes: int = 4,
|
||
system_prompt: str | None = None,
|
||
) -> list[NoteDraft]:
|
||
"""Distill an imported **document** into one *or several* knowledge notes.
|
||
|
||
Mirrors :func:`extract_thread_llm_multi` but takes raw document text (already
|
||
converted to Markdown/plain text) instead of a conversation. Returns ``[]``
|
||
when the LLM path fails (callers fall back to a verbatim note); raises
|
||
:class:`LLMRejectedIngest` when the model deliberately judges the document
|
||
worthless. Resulting notes carry ``source_type="file"``.
|
||
"""
|
||
text = (text or "").strip()
|
||
if not text:
|
||
return []
|
||
sources = sources or []
|
||
|
||
try:
|
||
from langchain_core.messages import HumanMessage, SystemMessage
|
||
|
||
from deerflow.models import create_chat_model
|
||
|
||
model = create_chat_model(name=model_name or None, thinking_enabled=False)
|
||
body = text if len(text) <= max_input_chars else text[:max_input_chars] + "\n…(已截断)"
|
||
header = f"文档标题:{doc_title}\n\n" if doc_title else ""
|
||
system_text = (system_prompt or "").strip() or (_DOC_SYSTEM_PROMPT + "\n\n" + WIKI_CAPTURE_GUIDANCE)
|
||
prompt = [
|
||
SystemMessage(content=system_text),
|
||
HumanMessage(content=f"以下是一篇文档,请按规范提炼为知识 JSON:\n\n{header}{body}"),
|
||
]
|
||
response = await model.ainvoke(prompt)
|
||
raw = response.content if isinstance(response.content, str) else str(response.content)
|
||
data = _parse_json(raw)
|
||
except Exception:
|
||
logger.exception("[knowledge] document %s — LLM extraction call failed; using rule-based fallback", doc_id)
|
||
return []
|
||
|
||
if not data:
|
||
logger.warning("[knowledge] document %s — LLM returned no parsable JSON; using rule-based fallback", doc_id)
|
||
return []
|
||
|
||
note_dicts = _notes_from_data(data)
|
||
if not note_dicts:
|
||
reason = str(data.get("reason") or "")[:200]
|
||
if _is_falsey(data.get("should_ingest")):
|
||
logger.warning("[knowledge] document %s SKIPPED — model judged not worth ingesting: %s", doc_id, reason)
|
||
raise LLMRejectedIngest(reason or "no knowledge value")
|
||
logger.warning(
|
||
"[knowledge] document %s — LLM produced no notes (should_ingest=%s, reason=%r); using rule-based fallback",
|
||
doc_id,
|
||
data.get("should_ingest"),
|
||
reason,
|
||
)
|
||
return []
|
||
note_dicts = note_dicts[:1] if not allow_multi else note_dicts[: max(1, max_notes)]
|
||
|
||
fallback_title = (doc_title or title or "文档知识").strip() or "文档知识"
|
||
drafts: list[NoteDraft] = []
|
||
for i, nd in enumerate(note_dicts):
|
||
override = title if (len(note_dicts) == 1) else None
|
||
suffix = "" if len(note_dicts) == 1 else f"#{i}"
|
||
drafts.append(
|
||
_build_note_draft(
|
||
nd,
|
||
source_type="file",
|
||
source_id=doc_id,
|
||
source_key=f"file:{doc_id}{suffix}",
|
||
fallback_title=fallback_title,
|
||
title_override=override,
|
||
tags=tags,
|
||
project_name=project_name,
|
||
sources=sources,
|
||
)
|
||
)
|
||
return drafts
|
||
|
||
|
||
async def extract_thread_llm(
|
||
messages: list[dict[str, Any]],
|
||
*,
|
||
thread_id: str,
|
||
thread_title: str | None = None,
|
||
mode: str = "summary",
|
||
title: str | None = None,
|
||
tags: list[str] | None = None,
|
||
project_name: str | None = "zncm",
|
||
sources: list[SourceDraft] | None = None,
|
||
max_messages: int = 80,
|
||
model_name: str | None = None,
|
||
max_input_chars: int = 12000,
|
||
) -> NoteDraft | None:
|
||
"""Distill a thread into a single :class:`NoteDraft` (backward-compatible).
|
||
|
||
Thin wrapper over :func:`extract_thread_llm_multi` that returns the first
|
||
note, or ``None`` when nothing was produced.
|
||
"""
|
||
try:
|
||
drafts = await extract_thread_llm_multi(
|
||
messages,
|
||
thread_id=thread_id,
|
||
thread_title=thread_title,
|
||
mode=mode,
|
||
title=title,
|
||
tags=tags,
|
||
project_name=project_name,
|
||
sources=sources,
|
||
max_messages=max_messages,
|
||
model_name=model_name,
|
||
max_input_chars=max_input_chars,
|
||
allow_multi=False,
|
||
)
|
||
except LLMRejectedIngest:
|
||
return None
|
||
return drafts[0] if drafts else None
|