reme.py 18 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427
  1. from __future__ import annotations
  2. import hashlib
  3. import re
  4. from datetime import datetime, timezone
  5. from typing import Any
  6. from ..embeddings import embedding_client, embedding_fingerprint
  7. from ..config import settings
  8. from ..schemas import ChatResponse, MemoryItem, SystemDescriptor
  9. from .base import AdapterUnavailable, MemoryAgent
  10. class ReMeAgent(MemoryAgent):
  11. """用 Markdown 文件保存记忆,并提供文件与向量搜索。"""
  12. id = "reme"
  13. _PASSIVE_JOB_OVERRIDES = {
  14. # 关闭后台任务,改为每次对话显式更新索引,避免隐式改文件和模型调用。
  15. "index_update_loop": {"backend": "base", "enable_serve": False},
  16. "resource_watch_loop": {"backend": "base", "enable_serve": False},
  17. "digest_watch_loop": {"backend": "base", "enable_serve": False},
  18. "dream_cron": {"backend": "base", "enable_serve": False},
  19. }
  20. # 扩展常见工程词,补强 ReMe 的 BM25 召回;语义召回单独使用 pgvector。
  21. _QUERY_EXPANSIONS = {
  22. "部署": ("deploy", "发布", "上线", "测试环境", "索引构建", "批次"),
  23. "上线": ("deploy", "发布", "测试环境", "索引构建"),
  24. "问题": ("超时", "错误", "失败", "异常", "报错", "故障", "冲突"),
  25. "故障": ("超时", "错误", "失败", "异常", "报错", "冲突"),
  26. "错误": ("失败", "异常", "报错", "故障"),
  27. "超时": ("timeout", "索引", "构建", "分批", "批次"),
  28. "连接": ("connection", "connect", "API", "数据库", "服务"),
  29. "索引": ("index", "向量", "构建", "分批", "批次"),
  30. "检索": ("retrieval", "RAG", "向量", "关键词", "重排"),
  31. "知识库": ("RAG", "文档", "向量", "检索", "重排"),
  32. "上次": ("之前", "最近", "历史"),
  33. }
  34. _QUERY_STOPWORDS = frozenset(
  35. {
  36. "什么",
  37. "哪些",
  38. "哪个",
  39. "怎么",
  40. "如何",
  41. "是否",
  42. "有没有",
  43. "请问",
  44. "告诉我",
  45. "根据",
  46. "当前",
  47. "相关",
  48. "一下",
  49. "遇到",
  50. "遇到了",
  51. "发生",
  52. "发生了",
  53. "的",
  54. "吗",
  55. "呢",
  56. "?",
  57. "?",
  58. }
  59. )
  60. _EXPLICIT_WRITE_RE = re.compile(r"记住|记下|请记录|请保存")
  61. _DECLARATIVE_RE = re.compile(
  62. r"我(?:偏好|喜欢|习惯|要求|正在|在开发|决定|选择|使用)|"
  63. r"项目(?:目前|现在|已经|使用|采用|数据库|前端|后端)|"
  64. r"(?:前端|后端|数据库|对话模型|向量模型)(?:使用|采用|选择|是|改成)|"
  65. r"更新一下|以后|上次|之前|曾经|长期规则|代码要求"
  66. )
  67. _QUESTION_RE = re.compile(
  68. r"[??]$|(?:什么|哪些|哪个|怎么|如何|是否|有没有|多少|为什么).*(?:[??]|$)"
  69. )
  70. def __init__(self) -> None:
  71. """加载 ReMe 依赖并记录当前是否可用。"""
  72. super().__init__()
  73. self._service: Any | None = None
  74. try:
  75. from reme.reme import ReMe # type: ignore
  76. from reme.config import resolve_app_config # type: ignore
  77. self._class = ReMe
  78. self._resolve_config = resolve_app_config
  79. except Exception:
  80. self._class = None
  81. self._resolve_config = None
  82. ready = bool(self._class and settings.reme_enabled)
  83. self.descriptor = SystemDescriptor(
  84. id="reme",
  85. name="ReMe",
  86. paradigm="文件即记忆",
  87. description="把长期记忆组织为可读 Markdown 文件,并融合 ReMe 检索与 pgvector 召回。",
  88. available=ready,
  89. mode="real-sdk" if ready else "unavailable",
  90. status="ready" if ready else "not-configured",
  91. package="reme-ai[core]>=0.4,<0.5",
  92. setup_hint=None if ready else "安装 reme-ai[core],并设置 REME_ENABLED=true。",
  93. )
  94. async def memories(self, query: str | None = None) -> list[MemoryItem]:
  95. """把 ReMe 的 Markdown 文件转换为统一记忆模型。"""
  96. root = settings.data_dir / "reme"
  97. if not root.exists():
  98. return []
  99. items: list[MemoryItem] = []
  100. # 对外接口只认识 MemoryItem,因此需要把文件逐个转换。
  101. for path in sorted(root.glob("daily/**/*.md")):
  102. content = path.read_text(encoding="utf-8", errors="replace").strip()
  103. if not content:
  104. continue
  105. if query and query.lower() not in f"{path} {content}".lower():
  106. continue
  107. updated = datetime.fromtimestamp(path.stat().st_mtime, tz=timezone.utc)
  108. digest = hashlib.sha1(str(path).encode("utf-8")).hexdigest()[:16]
  109. items.append(MemoryItem(
  110. id=f"reme_{digest}",
  111. system=self.id,
  112. workspace_id=settings.workspace_id,
  113. content=content,
  114. memory_type="file-memory",
  115. source=str(path.relative_to(root)),
  116. confidence=1.0,
  117. created_at=updated,
  118. updated_at=updated,
  119. metadata={"path": str(path.relative_to(root)), "retrieval": "hybrid"},
  120. ))
  121. return items
  122. async def delete_memory(self, memory_id: str) -> bool:
  123. """按记忆 ID 删除文件,并同步清理索引和操作记录。"""
  124. root = settings.data_dir / "reme"
  125. deleted = False
  126. for path in root.glob("daily/**/*.md"):
  127. digest = hashlib.sha1(str(path).encode("utf-8")).hexdigest()[:16]
  128. if f"reme_{digest}" != memory_id:
  129. continue
  130. path.unlink(missing_ok=True)
  131. deleted = True
  132. break
  133. # 文件删掉后,对应向量和文件索引也要一起更新。
  134. if deleted:
  135. await self.repo.delete_embedding(self.id, memory_id)
  136. if deleted and self._class and settings.reme_enabled:
  137. service = await self._ensure()
  138. await service.run_job("reindex")
  139. if deleted:
  140. await self._audit("DELETE/FileMemory", target=memory_id)
  141. return deleted
  142. async def _ensure(self) -> Any:
  143. """按需启动 ReMe 服务,同一进程内重复使用。"""
  144. if not self._class or not settings.reme_enabled:
  145. raise AdapterUnavailable(self.descriptor.setup_hint or "ReMe 不可用")
  146. if self._service is None:
  147. working_dir = str(settings.data_dir / "reme")
  148. config = self._resolve_config(
  149. log_config=False,
  150. workspace_dir=working_dir,
  151. enable_logo=False,
  152. log_to_console=False,
  153. service={"backend": "http"},
  154. jobs=self._PASSIVE_JOB_OVERRIDES,
  155. )
  156. self._service = self._class(**config)
  157. await self._service.start()
  158. return self._service
  159. @classmethod
  160. def _should_memorize(cls, message: str) -> bool:
  161. """判断当前输入是否值得写入长期文件。"""
  162. normalized = " ".join(message.strip().split())
  163. if not normalized:
  164. return False
  165. if cls._EXPLICIT_WRITE_RE.search(normalized):
  166. return True
  167. if cls._QUESTION_RE.search(normalized):
  168. return False
  169. return bool(cls._DECLARATIVE_RE.search(normalized))
  170. @classmethod
  171. def _heuristic_queries(cls, message: str) -> list[str]:
  172. """为文件检索生成短查询和工程词扩展。"""
  173. normalized = message.strip()
  174. queries: list[str] = [normalized] if normalized else []
  175. # 保留 RAG、FastAPI、PostgreSQL 等可直接匹配的领域词。
  176. queries.extend(re.findall(r"[A-Za-z][A-Za-z0-9_.:/-]{1,}", normalized))
  177. for trigger, expansions in cls._QUERY_EXPANSIONS.items():
  178. # 用户用词较少时,补上文件里可能出现的近义说法。
  179. if trigger in normalized:
  180. queries.extend((trigger, *expansions))
  181. # 中文短词作为兜底,过滤会干扰词面检索的疑问词。
  182. for chunk in re.findall(r"[\u4e00-\u9fff]{2,}", normalized):
  183. if chunk not in cls._QUERY_STOPWORDS:
  184. queries.append(chunk)
  185. if len(chunk) > 2:
  186. queries.extend(
  187. chunk[index : index + 2]
  188. for index in range(len(chunk) - 1)
  189. if chunk[index : index + 2] not in cls._QUERY_STOPWORDS
  190. )
  191. return cls._dedupe_queries(queries)
  192. @staticmethod
  193. def _dedupe_queries(queries: list[str]) -> list[str]:
  194. """清理重复或空查询,并限制单次搜索数量。"""
  195. unique: list[str] = []
  196. seen: set[str] = set()
  197. for query in queries:
  198. cleaned = " ".join(str(query).strip().split())
  199. if not cleaned:
  200. continue
  201. key = cleaned.casefold()
  202. if key in seen:
  203. continue
  204. seen.add(key)
  205. unique.append(cleaned)
  206. return unique[:12]
  207. async def _rewrite_query(self, message: str) -> list[str]:
  208. """合并规则扩展和可选的模型查询改写。"""
  209. queries = self._heuristic_queries(message)
  210. try:
  211. rewritten = await self.llm.json(
  212. (
  213. "你是记忆检索查询改写器。只输出 JSON,不要回答用户问题。"
  214. "把用户问题改写成 2 到 5 个适合混合文件检索的短查询,"
  215. "保留专有名词,并补充可能出现在记忆里的中英文同义词。"
  216. '格式必须是 {"queries": ["..."]}。'
  217. ),
  218. f"用户问题:{message}",
  219. )
  220. generated = rewritten.get("queries", [])
  221. if isinstance(generated, list):
  222. queries.extend(str(item) for item in generated if str(item).strip())
  223. except Exception:
  224. # 改写失败时继续使用规则查询,不影响基础检索。
  225. pass
  226. return self._dedupe_queries(queries)
  227. @staticmethod
  228. def _merge_search_results(results: list[str]) -> str:
  229. """合并多个查询结果并按正文去重。"""
  230. merged: list[str] = []
  231. seen: set[str] = set()
  232. for result in results:
  233. for block in re.split(r"(?=^========== )", result.strip(), flags=re.MULTILINE):
  234. block = block.strip()
  235. if not block:
  236. continue
  237. fingerprint_source = block
  238. if block.startswith("==========") and "\n" in block:
  239. fingerprint_source = block.split("\n", 1)[1]
  240. fingerprint = re.sub(r"\[(?:vector_)?score=[^\]]+\]", "", fingerprint_source)
  241. fingerprint = " ".join(fingerprint.split()).casefold()
  242. if fingerprint in seen:
  243. continue
  244. seen.add(fingerprint)
  245. merged.append(block)
  246. return "\n\n".join(merged)
  247. async def _search_memory(self, service: Any, queries: list[str]) -> tuple[str, list[str]]:
  248. """逐个查询 ReMe 文件,并返回有结果的查询词。"""
  249. results: list[str] = []
  250. used_queries: list[str] = []
  251. for query in queries:
  252. search_result = await service.run_job("search", query=query, limit=5)
  253. answer = str(getattr(search_result, "answer", "") or "").strip()
  254. if answer:
  255. results.append(answer)
  256. used_queries.append(query)
  257. return self._merge_search_results(results), used_queries
  258. async def _sync_embeddings(self) -> dict[str, Any]:
  259. """更新已变化文件的向量,并清理已删除文件。"""
  260. if not embedding_client.configured:
  261. return {"enabled": False, "embedded": 0, "removed": 0}
  262. items = await self.memories()
  263. current_ids = {item.id for item in items}
  264. previous_hashes = await self.repo.embedding_hashes(self.id)
  265. # 文件不存在后,旧向量也要删除。
  266. stale_ids = set(previous_hashes) - current_ids
  267. for memory_id in stale_ids:
  268. await self.repo.delete_embedding(self.id, memory_id)
  269. # 只处理新增或内容有变化的文件。
  270. pending = [
  271. item
  272. for item in items
  273. if previous_hashes.get(item.id) != embedding_fingerprint(item.content)
  274. ]
  275. if not pending:
  276. return {"enabled": True, "embedded": 0, "removed": len(stale_ids)}
  277. vectors = await embedding_client.embed([item.content for item in pending])
  278. for item, vector in zip(pending, vectors):
  279. await self.repo.upsert_embedding(
  280. system=self.id,
  281. memory_id=item.id,
  282. content=item.content,
  283. source=item.source,
  284. content_hash=embedding_fingerprint(item.content),
  285. embedding=vector,
  286. )
  287. return {"enabled": True, "embedded": len(pending), "removed": len(stale_ids)}
  288. async def _semantic_search(self, message: str) -> tuple[str, int]:
  289. """查找内容相近的文件,并过滤分数过低的结果。"""
  290. if not embedding_client.configured:
  291. return "", 0
  292. try:
  293. query_vector = (await embedding_client.embed([message]))[0]
  294. matches = await self.repo.search_embeddings(self.id, query_vector, limit=5)
  295. except Exception:
  296. # 这一路失败时,主流程仍可使用 ReMe 文件搜索。
  297. return "", 0
  298. matches = [
  299. match
  300. for match in matches
  301. if float(match.get("score", 0.0)) >= settings.embedding_min_score
  302. ]
  303. blocks = [
  304. (
  305. f"========== {match['source']} [vector_score={float(match['score']):.4f}] =========="
  306. f"\n{match['content']}"
  307. )
  308. for match in matches
  309. ]
  310. return self._merge_search_results(blocks), len(matches)
  311. async def _answer_from_memory(self, message: str, evidence: str) -> str:
  312. """根据搜索结果生成回答,模型不可用时直接返回原文。"""
  313. if not evidence:
  314. return "没有检索到与这个问题相关的历史记忆。"
  315. try:
  316. return await self.llm.chat(
  317. (
  318. "你是一个使用文件记忆的 AI 助手。"
  319. "只能依据提供的记忆证据回答,不要编造证据中没有的事实。"
  320. "如果证据不足,要明确说明不确定。直接回答用户问题,简洁自然。"
  321. ),
  322. f"用户问题:{message}\n\n记忆证据:\n{evidence}",
  323. )
  324. except Exception:
  325. # 保留原始内容比吞掉已经找到的结果更方便排查。
  326. return f"根据检索到的记忆:\n\n{evidence}"
  327. async def chat(self, message: str) -> ChatResponse:
  328. """完成文件写入、索引更新、搜索、回答和操作记录。"""
  329. service = await self._ensure()
  330. write_candidate = self._should_memorize(message)
  331. memory_result: Any | None = None
  332. # 问题和临时请求不写文件,只处理值得长期保存的陈述。
  333. if write_candidate:
  334. memory_result = await service.run_job(
  335. "auto_memory",
  336. messages=[{"name": "user", "role": "user", "content": message}],
  337. session_id=settings.workspace_id,
  338. memory_hint=(
  339. "只记录稳定偏好、项目事实、进度变化、可复用的故障经验和长期流程规则。"
  340. "问题、寒暄和临时请求不应写入长期记忆。"
  341. ),
  342. )
  343. # 写入后立即更新索引,保证本轮搜索能看到刚保存的内容。
  344. reindex_result = await service.run_job("reindex")
  345. try:
  346. embedding_sync = await self._sync_embeddings()
  347. except Exception as exc:
  348. # 向量服务异常时仍保留 ReMe 自带的文件检索。
  349. embedding_sync = {"enabled": True, "embedded": 0, "removed": 0, "error": str(exc)}
  350. queries = await self._rewrite_query(message)
  351. # 文件搜索和相近内容搜索各跑一次,最后合并并去重。
  352. lexical_search, used_queries = await self._search_memory(service, queries)
  353. semantic_search, semantic_count = await self._semantic_search(message)
  354. search = self._merge_search_results([lexical_search, semantic_search])
  355. answer = await self._answer_from_memory(message, search)
  356. # 记录查询词和搜索摘要,出现误匹配时可回看整个过程。
  357. await self._audit(
  358. "REME/auto_memory+search",
  359. details={
  360. "write_candidate": write_candidate,
  361. "memory": str(getattr(memory_result, "answer", "skipped"))[:800],
  362. "reindex": str(reindex_result.metadata)[:800],
  363. "embedding_sync": embedding_sync,
  364. "queries": queries,
  365. "used_queries": used_queries,
  366. "semantic_results": semantic_count,
  367. "search": str(search)[:1000],
  368. },
  369. )
  370. return ChatResponse(
  371. system="reme",
  372. answer=answer,
  373. mode="real-sdk",
  374. memory_context=await self.memories(),
  375. memory_events=[
  376. {
  377. "event": "file-memory",
  378. "write_candidate": write_candidate,
  379. "queries": queries,
  380. "retrieval": "hybrid",
  381. "semantic_results": semantic_count,
  382. "result": str(search)[:2000],
  383. }
  384. ],
  385. audit_events=await self.audit(),
  386. )
  387. async def reset(self) -> None:
  388. """清空公共数据,关闭服务并删除 ReMe 工作目录。"""
  389. await super().reset()
  390. # 先关闭服务,避免删除目录后仍有后台句柄占用文件。
  391. if self._service is not None:
  392. await self._service.close()
  393. self._service = None
  394. working_dir = settings.data_dir / "reme"
  395. if working_dir.exists():
  396. import shutil
  397. shutil.rmtree(working_dir)