"""In-memory LLMWiki metadata store for tests and DB-less deployments.""" from __future__ import annotations from datetime import UTC, datetime from typing import Any from uuid import uuid4 from deerflow.persistence.llmwiki.base import LlmWikiScope, LlmWikiStore def _now() -> str: return datetime.now(UTC).isoformat() class MemoryLlmWikiStore(LlmWikiStore): def __init__(self) -> None: self._rows: dict[str, dict[str, Any]] = {} self._conversation_deposits: dict[str, dict[str, Any]] = {} async def create_mapping(self, *, weknora_id: str, owner_user_id: str, name: str, description: str, kb_type: str) -> dict[str, Any]: if any(row["weknora_id"] == weknora_id for row in self._rows.values()): raise ValueError("WeKnora knowledge base is already mapped") now = _now() row = { "id": str(uuid4()), "weknora_id": weknora_id, "owner_user_id": owner_user_id, "name": name, "description": description, "kb_type": kb_type, "publication_status": "private", "published_by": None, "reviewed_by": None, "published_at": None, "wiki_index_enabled": True, "external_search_enabled": False, "external_search_updated_at": None, "created_at": now, "updated_at": now, } self._rows[row["id"]] = row return dict(row) async def list_visible(self, user_id: str, *, scope: LlmWikiScope = "all", is_admin: bool = False) -> list[dict[str, Any]]: rows = list(self._rows.values()) if scope == "personal": rows = rows if is_admin else [row for row in rows if row["owner_user_id"] == user_id] elif scope == "public": rows = [row for row in rows if row["publication_status"] == "published"] else: rows = rows if is_admin else [row for row in rows if row["owner_user_id"] == user_id or row["publication_status"] == "published"] return [dict(row) for row in sorted(rows, key=lambda row: row["created_at"], reverse=True)] async def get_authorized(self, mapping_id: str, user_id: str, *, write: bool, is_admin: bool) -> dict[str, Any] | None: row = self._rows.get(mapping_id) if row is None: return None if is_admin or row["owner_user_id"] == user_id: return dict(row) if not write and row["publication_status"] == "published": return dict(row) return None async def get_by_weknora_id(self, weknora_id: str) -> dict[str, Any] | None: for row in self._rows.values(): if row["weknora_id"] == weknora_id: return dict(row) return None async def update_mapping( self, mapping_id: str, *, name: str | None = None, description: str | None = None, wiki_index_enabled: bool | None = None, external_search_enabled: bool | None = None, ) -> dict[str, Any] | None: row = self._rows.get(mapping_id) if row is None: return None if name is not None: row["name"] = name if description is not None: row["description"] = description if wiki_index_enabled is not None: row["wiki_index_enabled"] = wiki_index_enabled if external_search_enabled is not None: row["external_search_enabled"] = external_search_enabled row["external_search_updated_at"] = _now() row["updated_at"] = _now() return dict(row) async def adopt_conversation_deposit_mapping( self, mapping_id: str, *, name: str, description: str, ) -> dict[str, Any] | None: row = self._rows.get(mapping_id) if row is None: return None now = _now() row.update( { "owner_user_id": "system", "name": name, "description": description, "publication_status": "published", "published_by": "system", "reviewed_by": None, "published_at": row.get("published_at") or now, "updated_at": now, } ) return dict(row) async def delete_mapping(self, mapping_id: str) -> bool: return self._rows.pop(mapping_id, None) is not None async def request_publish(self, mapping_id: str, actor_user_id: str, *, is_admin: bool) -> dict[str, Any] | None: row = self._rows.get(mapping_id) if row is None or (not is_admin and row["owner_user_id"] != actor_user_id): return None row["publication_status"] = "published" row["published_by"] = actor_user_id row["reviewed_by"] = None row["published_at"] = _now() row["updated_at"] = _now() return dict(row) async def review_publish(self, mapping_id: str, reviewer_user_id: str, *, approved: bool) -> dict[str, Any] | None: row = self._rows.get(mapping_id) if row is None or row["publication_status"] != "pending": return None row["publication_status"] = "published" if approved else "rejected" row["reviewed_by"] = reviewer_user_id row["published_at"] = _now() if approved else None row["updated_at"] = _now() return dict(row) async def unpublish(self, mapping_id: str, actor_user_id: str, *, is_admin: bool) -> dict[str, Any] | None: row = self._rows.get(mapping_id) if row is None or (not is_admin and row["owner_user_id"] != actor_user_id): return None row["publication_status"] = "private" row["reviewed_by"] = actor_user_id if is_admin else row.get("reviewed_by") row["published_at"] = None row["updated_at"] = _now() return dict(row) async def list_publications(self, *, status: str | None = None) -> list[dict[str, Any]]: rows = [row for row in self._rows.values() if row["publication_status"] != "private"] if status: rows = [row for row in rows if row["publication_status"] == status] return [dict(row) for row in sorted(rows, key=lambda row: row["updated_at"], reverse=True)] async def get_conversation_deposit_mapping(self) -> dict[str, Any] | None: for row in self._rows.values(): if row.get("owner_user_id") == "system" and row.get("name") == "对话沉淀" and row.get("publication_status") == "published": return dict(row) return None async def get_conversation_deposit( self, *, thread_id: str, assistant_message_id: str, ) -> dict[str, Any] | None: key = f"{thread_id}:{assistant_message_id}" row = self._conversation_deposits.get(key) return dict(row) if row else None async def upsert_conversation_deposit( self, *, thread_id: str, human_message_id: str | None, assistant_message_id: str, knowledge_base_mapping_id: str, weknora_knowledge_id: str | None, title: str, created_by_user_id: str | None, ) -> dict[str, Any]: key = f"{thread_id}:{assistant_message_id}" now = _now() row = self._conversation_deposits.get(key) if row is None: row = { "id": str(uuid4()), "thread_id": thread_id, "human_message_id": human_message_id, "assistant_message_id": assistant_message_id, "knowledge_base_mapping_id": knowledge_base_mapping_id, "weknora_knowledge_id": weknora_knowledge_id, "title": title, "created_by_user_id": created_by_user_id, "created_at": now, "updated_at": now, } self._conversation_deposits[key] = row else: row.update( { "human_message_id": human_message_id, "knowledge_base_mapping_id": knowledge_base_mapping_id, "weknora_knowledge_id": weknora_knowledge_id, "title": title, "updated_at": now, } ) return dict(row) async def delete_conversation_deposits( self, *, thread_id: str, message_ids: list[str] | None = None, ) -> list[dict[str, Any]]: ids = set(message_ids or []) deleted: list[dict[str, Any]] = [] for key, row in list(self._conversation_deposits.items()): if row.get("thread_id") != thread_id: continue if ids and row.get("human_message_id") not in ids and row.get("assistant_message_id") not in ids: continue deleted.append(dict(row)) self._conversation_deposits.pop(key, None) return deleted