deerflow-code/offline-backend-20260512/backend/packages/harness/deerflow/knowledge/obsidian_adapter.py
2026-09-07 18:24:55 +08:00

257 lines
10 KiB
Python

"""ObsidianWikiAdapter — the obsidian-wiki vault facade.
Owns the Markdown vault: writes notes with wiki-capture frontmatter/body, keeps
``index.md`` / ``log.md`` / ``hot.md`` / ``.manifest.json`` in sync after every
write, and produces wiki-export graph artifacts. Output stays compatible with
``ar9av/obsidian-wiki`` without ever invoking its CLI.
This class is storage-only (no DB, no LLM). The :class:`KnowledgeService`
combines it with the repository + extractor.
"""
from __future__ import annotations
from pathlib import Path
from typing import Any
from deerflow.knowledge import exporter
from deerflow.knowledge.manifest import Manifest, content_hash
from deerflow.knowledge.markdown_vault import MarkdownVault
from deerflow.knowledge.schemas import NoteDraft
_INDEX_PAGES_HEADER = "## Pages"
_HOT_HEADER = "## Recent"
_HOT_MAX = 20
class ObsidianWikiAdapter:
"""Read/write an obsidian-wiki compatible vault."""
def __init__(self, vault_root: Path | str, *, project_name: str = "zncm", export_html: bool = True, export_json: bool = True) -> None:
self.vault = MarkdownVault(vault_root)
self.root = self.vault.root
self.manifest = Manifest(self.root)
self.project_name = project_name
self.export_html = export_html
self.export_json = export_json
def ensure_ready(self) -> None:
self.vault.ensure_scaffold(project_name=self.project_name)
# -- writes -----------------------------------------------------------
def write_note(
self,
note_id: str,
draft: NoteDraft,
*,
created: str,
updated: str,
status: str,
) -> str:
"""Write the note file, sync vault metadata, return the vault path."""
self.ensure_ready()
rel = self.vault.note_relpath(category=draft.category, project_name=self.project_name, slug=draft.title, note_id=note_id)
sources_fm = self._frontmatter_sources(draft)
relationships = [{"target": link, "type": "related"} for link in (draft.related or [])[:8]]
frontmatter = self.vault.build_frontmatter(
note_id=note_id,
title=draft.title,
category=draft.category,
tags=draft.tags,
sources=sources_fm,
summary=draft.summary,
created=created,
updated=updated,
base_confidence=draft.confidence,
lifecycle=status,
relationships=relationships,
)
self.vault.write_note(rel, frontmatter, draft.content_md)
self.manifest.record_page(
note_id=note_id,
vault_path=rel,
source_key=draft.source_key,
hash_value=content_hash(draft.content_md),
updated=updated,
)
self._update_index(note_id, rel, draft.title, draft.summary)
self._append_log("CAPTURE" if draft.source_type in ("thread", "search", "tool") else "INGEST", note_id, draft.title, updated)
self._update_hot(draft.title, rel, updated)
return rel
def update_note_file(
self,
note_id: str,
vault_path: str | None,
*,
title: str,
content_md: str,
tags: list[str],
status: str,
summary: str,
category: str,
created: str,
updated: str,
) -> str:
"""Rewrite a note file in place (frontmatter refreshed), sync metadata."""
self.ensure_ready()
rel = vault_path
if not rel:
rel = self.vault.note_relpath(category=category, project_name=self.project_name, slug=title, note_id=note_id)
# Preserve original `sources`/`provenance`/`category` frontmatter when present.
existing = self.vault.read_note(rel)
sources_fm = existing[0].get("sources", []) if existing else []
provenance = existing[0].get("provenance") if existing else None
relationships = existing[0].get("relationships") if existing else []
effective_category = (existing[0].get("category") if existing else None) or category
frontmatter = self.vault.build_frontmatter(
note_id=note_id,
title=title,
category=effective_category,
tags=tags,
sources=sources_fm,
summary=summary,
created=created,
updated=updated,
base_confidence=float((existing[0].get("base_confidence") if existing else 0.0) or 0.0),
lifecycle=status,
provenance=provenance,
relationships=relationships,
)
self.vault.write_note(rel, frontmatter, content_md)
self.manifest.record_page(
note_id=note_id,
vault_path=rel,
source_key=self.manifest.load().get("pages", {}).get(note_id, {}).get("source"),
hash_value=content_hash(content_md),
updated=updated,
)
self._update_index(note_id, rel, title, summary)
self._append_log("UPDATE", note_id, title, updated)
self._update_hot(title, rel, updated)
return rel
def archive_note(self, note_id: str, vault_path: str | None, *, updated: str) -> None:
"""Mark a note archived in its frontmatter + log (soft delete)."""
if vault_path:
parsed = self.vault.read_note(vault_path)
if parsed:
fm, body = parsed
fm["lifecycle"] = "archived"
fm["updated"] = updated
self.vault.write_note(vault_path, fm, body)
self._remove_from_index(note_id)
self._append_log("ARCHIVE", note_id, "", updated)
def export_graph(
self,
notes: list[dict[str, Any]],
*,
entities: list[dict[str, Any]] | None = None,
relations: list[dict[str, Any]] | None = None,
updated: str,
blacklist: set[str] | None = None,
) -> dict[str, Any]:
self.ensure_ready()
result = exporter.export_all(
self.root,
notes,
entities=entities,
relations=relations,
write_html=self.export_html,
write_json=self.export_json,
blacklist=blacklist,
)
self._append_log("EXPORT", "graph", f"nodes={result['stats']['nodes']} edges={result['stats']['edges']}", updated)
return result
def export_file_path(self, filename: str) -> Path | None:
allowed = {"graph.json", "graph.graphml", "cypher.txt", "graph.html"}
if filename not in allowed:
return None
path = (self.root / "wiki-export" / filename).resolve()
try:
path.relative_to(self.root)
except ValueError:
return None
return path if path.exists() else None
# -- frontmatter source serialization ---------------------------------
@staticmethod
def _frontmatter_sources(draft: NoteDraft) -> list[str]:
out: list[str] = []
if draft.source_type == "thread" and draft.source_id:
out.append(f"thread:{draft.source_id}")
for s in draft.sources[:20]:
if s.url:
out.append(f"url:{s.url}")
elif s.thread_id:
out.append(f"thread:{s.thread_id}")
elif s.tool_name:
out.append(f"tool:{s.tool_name}")
# De-dup preserving order.
seen: set[str] = set()
deduped = []
for item in out:
if item not in seen:
seen.add(item)
deduped.append(item)
return deduped
# -- index.md ---------------------------------------------------------
def _index_path(self) -> Path:
return self.root / "index.md"
def _update_index(self, note_id: str, rel: str, title: str, summary: str | None) -> None:
path = self._index_path()
text = path.read_text(encoding="utf-8") if path.exists() else "# Knowledge Index\n\n"
link = rel[:-3] if rel.endswith(".md") else rel
summary_txt = (summary or "").strip().replace("\n", " ")
if len(summary_txt) > 120:
summary_txt = summary_txt[:120] + "…"
line = f"- [[{link}|{title}]]" + (f" — {summary_txt}" if summary_txt else "") + f" <!-- id:{note_id} -->"
lines = text.splitlines()
# Drop any existing line for this note id.
lines = [ln for ln in lines if f"<!-- id:{note_id} -->" not in ln]
if _INDEX_PAGES_HEADER not in text:
lines.append("")
lines.append(_INDEX_PAGES_HEADER)
lines.append("")
# Insert after the Pages header.
try:
insert_at = lines.index(_INDEX_PAGES_HEADER) + 1
except ValueError:
insert_at = len(lines)
lines.insert(insert_at, line)
path.write_text("\n".join(lines).rstrip() + "\n", encoding="utf-8")
def _remove_from_index(self, note_id: str) -> None:
path = self._index_path()
if not path.exists():
return
lines = [ln for ln in path.read_text(encoding="utf-8").splitlines() if f"<!-- id:{note_id} -->" not in ln]
path.write_text("\n".join(lines).rstrip() + "\n", encoding="utf-8")
# -- log.md -----------------------------------------------------------
def _append_log(self, action: str, note_id: str, label: str, updated: str) -> None:
path = self.root / "log.md"
text = path.read_text(encoding="utf-8") if path.exists() else "# Activity Log\n\n"
entry = f"- {updated} {action} {note_id}" + (f" — {label}" if label else "")
path.write_text(text.rstrip() + "\n" + entry + "\n", encoding="utf-8")
# -- hot.md -----------------------------------------------------------
def _update_hot(self, title: str, rel: str, updated: str) -> None:
path = self.root / "hot.md"
link = rel[:-3] if rel.endswith(".md") else rel
new_line = f"- {updated} [[{link}|{title}]]"
existing_lines: list[str] = []
if path.exists():
for ln in path.read_text(encoding="utf-8").splitlines():
if ln.startswith("- ") and "[[" in ln:
existing_lines.append(ln)
existing_lines = [ln for ln in existing_lines if f"[[{link}" not in ln]
existing_lines.insert(0, new_line)
existing_lines = existing_lines[:_HOT_MAX]
body = "# Hot\n\n_近期活动摘要。_\n\n" + _HOT_HEADER + "\n\n" + "\n".join(existing_lines) + "\n"
path.write_text(body, encoding="utf-8")