from __future__ import annotations import json from typing import Any from uuid import uuid4 from ..config import settings from ..schemas import ChatResponse, SystemDescriptor from .base import AdapterUnavailable, MemoryAgent class MemUAgent(MemoryAgent): id = "memu" def __init__(self) -> None: super().__init__() self._service: Any | None = None self._import_error: str | None = None try: from memu import MemoryService, MemUService # type: ignore self._service_classes = (MemoryService, MemUService) except Exception as exc: try: from memu import MemoryService # type: ignore self._service_classes = (MemoryService,) except Exception: self._service_classes = () self._import_error = str(exc) embedding_ready = bool( settings.embedding_api_key and settings.embedding_base_url and settings.embedding_model ) ready = bool(self._service_classes and settings.memu_enabled and embedding_ready) self.descriptor = SystemDescriptor( id="memu", name="memU", paradigm="主动式异步记忆管线", description="把对话摄入、意图提取、候选更新和主动检索拆成后台记忆服务。", available=ready, mode="real-sdk" if ready else "unavailable", status="ready" if ready else "not-configured", package="memU 官方源代码", setup_hint=None if ready else ( "当前仅保留 memU 适配骨架,请保持 MEMU_ENABLED=false。" "完成队列、预算、审批和停止开关后再启用。" ), ) async def _ensure(self) -> Any: if not self._service_classes or not settings.memu_enabled: raise AdapterUnavailable(self.descriptor.setup_hint or "memU 不可用") if self._service is None: service_class = self._service_classes[0] self._service = service_class( llm_profiles={ "default": { "provider": "openai", "base_url": settings.llm_base_url, "api_key": settings.llm_api_key, "chat_model": settings.llm_model, "client_backend": "sdk", }, "embedding": { "provider": "openai", "base_url": settings.embedding_base_url, "api_key": settings.embedding_api_key, "embed_model": settings.embedding_model, "client_backend": "sdk", }, }, database_config={ "metadata_store": { "provider": "sqlite", "dsn": f"sqlite:///{settings.data_dir / 'memu.sqlite3'}", }, }, ) return self._service async def chat(self, message: str) -> ChatResponse: service = await self._ensure() conversation_dir = settings.data_dir / "memu" / "conversations" conversation_dir.mkdir(parents=True, exist_ok=True) resource_path = conversation_dir / f"{uuid4().hex}.json" resource_path.write_text( json.dumps([{"role": "user", "content": message}], ensure_ascii=False), encoding="utf-8", ) memorize = service.memorize( resource_url=str(resource_path), modality="conversation", user={"user_id": settings.workspace_id, "content": message}, ) if hasattr(memorize, "__await__"): await memorize context = service.retrieve( queries=[{"role": "user", "content": {"text": message}}], where={"user_id": settings.workspace_id}, ) if hasattr(context, "__await__"): context = await context await self._audit("MEMU/memorize+retrieve", details={"context": str(context)[:1200]}) return ChatResponse( system="memu", answer=f"memU 已完成异步记忆摄入与主动检索。\n\n当前上下文:\n{context}", mode="real-sdk", memory_context=await self.memories(), memory_events=[{"event": "memorize", "candidate": True, "context": str(context)[:2000]}], audit_events=await self.audit(), )