| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113 |
- 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(),
- )
|