memu.py 4.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113
  1. from __future__ import annotations
  2. import json
  3. from typing import Any
  4. from uuid import uuid4
  5. from ..config import settings
  6. from ..schemas import ChatResponse, SystemDescriptor
  7. from .base import AdapterUnavailable, MemoryAgent
  8. class MemUAgent(MemoryAgent):
  9. id = "memu"
  10. def __init__(self) -> None:
  11. super().__init__()
  12. self._service: Any | None = None
  13. self._import_error: str | None = None
  14. try:
  15. from memu import MemoryService, MemUService # type: ignore
  16. self._service_classes = (MemoryService, MemUService)
  17. except Exception as exc:
  18. try:
  19. from memu import MemoryService # type: ignore
  20. self._service_classes = (MemoryService,)
  21. except Exception:
  22. self._service_classes = ()
  23. self._import_error = str(exc)
  24. embedding_ready = bool(
  25. settings.embedding_api_key
  26. and settings.embedding_base_url
  27. and settings.embedding_model
  28. )
  29. ready = bool(self._service_classes and settings.memu_enabled and embedding_ready)
  30. self.descriptor = SystemDescriptor(
  31. id="memu",
  32. name="memU",
  33. paradigm="主动式异步记忆管线",
  34. description="把对话摄入、意图提取、候选更新和主动检索拆成后台记忆服务。",
  35. available=ready,
  36. mode="real-sdk" if ready else "unavailable",
  37. status="ready" if ready else "not-configured",
  38. package="memU 官方源代码",
  39. setup_hint=None if ready else (
  40. "当前仅保留 memU 适配骨架,请保持 MEMU_ENABLED=false。"
  41. "完成队列、预算、审批和停止开关后再启用。"
  42. ),
  43. )
  44. async def _ensure(self) -> Any:
  45. if not self._service_classes or not settings.memu_enabled:
  46. raise AdapterUnavailable(self.descriptor.setup_hint or "memU 不可用")
  47. if self._service is None:
  48. service_class = self._service_classes[0]
  49. self._service = service_class(
  50. llm_profiles={
  51. "default": {
  52. "provider": "openai",
  53. "base_url": settings.llm_base_url,
  54. "api_key": settings.llm_api_key,
  55. "chat_model": settings.llm_model,
  56. "client_backend": "sdk",
  57. },
  58. "embedding": {
  59. "provider": "openai",
  60. "base_url": settings.embedding_base_url,
  61. "api_key": settings.embedding_api_key,
  62. "embed_model": settings.embedding_model,
  63. "client_backend": "sdk",
  64. },
  65. },
  66. database_config={
  67. "metadata_store": {
  68. "provider": "sqlite",
  69. "dsn": f"sqlite:///{settings.data_dir / 'memu.sqlite3'}",
  70. },
  71. },
  72. )
  73. return self._service
  74. async def chat(self, message: str) -> ChatResponse:
  75. service = await self._ensure()
  76. conversation_dir = settings.data_dir / "memu" / "conversations"
  77. conversation_dir.mkdir(parents=True, exist_ok=True)
  78. resource_path = conversation_dir / f"{uuid4().hex}.json"
  79. resource_path.write_text(
  80. json.dumps([{"role": "user", "content": message}], ensure_ascii=False),
  81. encoding="utf-8",
  82. )
  83. memorize = service.memorize(
  84. resource_url=str(resource_path),
  85. modality="conversation",
  86. user={"user_id": settings.workspace_id, "content": message},
  87. )
  88. if hasattr(memorize, "__await__"):
  89. await memorize
  90. context = service.retrieve(
  91. queries=[{"role": "user", "content": {"text": message}}],
  92. where={"user_id": settings.workspace_id},
  93. )
  94. if hasattr(context, "__await__"):
  95. context = await context
  96. await self._audit("MEMU/memorize+retrieve", details={"context": str(context)[:1200]})
  97. return ChatResponse(
  98. system="memu",
  99. answer=f"memU 已完成异步记忆摄入与主动检索。\n\n当前上下文:\n{context}",
  100. mode="real-sdk",
  101. memory_context=await self.memories(),
  102. memory_events=[{"event": "memorize", "candidate": True, "context": str(context)[:2000]}],
  103. audit_events=await self.audit(),
  104. )