from __future__ import annotations import hashlib import re from datetime import datetime, timezone from typing import Any from ..embeddings import embedding_client, embedding_fingerprint from ..config import settings from ..schemas import ChatResponse, MemoryItem, SystemDescriptor from .base import AdapterUnavailable, MemoryAgent class ReMeAgent(MemoryAgent): """用 Markdown 文件保存记忆,并提供文件与向量搜索。""" id = "reme" _PASSIVE_JOB_OVERRIDES = { # 关闭后台任务,改为每次对话显式更新索引,避免隐式改文件和模型调用。 "index_update_loop": {"backend": "base", "enable_serve": False}, "resource_watch_loop": {"backend": "base", "enable_serve": False}, "digest_watch_loop": {"backend": "base", "enable_serve": False}, "dream_cron": {"backend": "base", "enable_serve": False}, } # 扩展常见工程词,补强 ReMe 的 BM25 召回;语义召回单独使用 pgvector。 _QUERY_EXPANSIONS = { "部署": ("deploy", "发布", "上线", "测试环境", "索引构建", "批次"), "上线": ("deploy", "发布", "测试环境", "索引构建"), "问题": ("超时", "错误", "失败", "异常", "报错", "故障", "冲突"), "故障": ("超时", "错误", "失败", "异常", "报错", "冲突"), "错误": ("失败", "异常", "报错", "故障"), "超时": ("timeout", "索引", "构建", "分批", "批次"), "连接": ("connection", "connect", "API", "数据库", "服务"), "索引": ("index", "向量", "构建", "分批", "批次"), "检索": ("retrieval", "RAG", "向量", "关键词", "重排"), "知识库": ("RAG", "文档", "向量", "检索", "重排"), "上次": ("之前", "最近", "历史"), } _QUERY_STOPWORDS = frozenset( { "什么", "哪些", "哪个", "怎么", "如何", "是否", "有没有", "请问", "告诉我", "根据", "当前", "相关", "一下", "遇到", "遇到了", "发生", "发生了", "的", "吗", "呢", "?", "?", } ) _EXPLICIT_WRITE_RE = re.compile(r"记住|记下|请记录|请保存") _DECLARATIVE_RE = re.compile( r"我(?:偏好|喜欢|习惯|要求|正在|在开发|决定|选择|使用)|" r"项目(?:目前|现在|已经|使用|采用|数据库|前端|后端)|" r"(?:前端|后端|数据库|对话模型|向量模型)(?:使用|采用|选择|是|改成)|" r"更新一下|以后|上次|之前|曾经|长期规则|代码要求" ) _QUESTION_RE = re.compile( r"[??]$|(?:什么|哪些|哪个|怎么|如何|是否|有没有|多少|为什么).*(?:[??]|$)" ) def __init__(self) -> None: """加载 ReMe 依赖并记录当前是否可用。""" super().__init__() self._service: Any | None = None try: from reme.reme import ReMe # type: ignore from reme.config import resolve_app_config # type: ignore self._class = ReMe self._resolve_config = resolve_app_config except Exception: self._class = None self._resolve_config = None ready = bool(self._class and settings.reme_enabled) self.descriptor = SystemDescriptor( id="reme", name="ReMe", paradigm="文件即记忆", description="把长期记忆组织为可读 Markdown 文件,并融合 ReMe 检索与 pgvector 召回。", available=ready, mode="real-sdk" if ready else "unavailable", status="ready" if ready else "not-configured", package="reme-ai[core]>=0.4,<0.5", setup_hint=None if ready else "安装 reme-ai[core],并设置 REME_ENABLED=true。", ) async def memories(self, query: str | None = None) -> list[MemoryItem]: """把 ReMe 的 Markdown 文件转换为统一记忆模型。""" root = settings.data_dir / "reme" if not root.exists(): return [] items: list[MemoryItem] = [] # 对外接口只认识 MemoryItem,因此需要把文件逐个转换。 for path in sorted(root.glob("daily/**/*.md")): content = path.read_text(encoding="utf-8", errors="replace").strip() if not content: continue if query and query.lower() not in f"{path} {content}".lower(): continue updated = datetime.fromtimestamp(path.stat().st_mtime, tz=timezone.utc) digest = hashlib.sha1(str(path).encode("utf-8")).hexdigest()[:16] items.append(MemoryItem( id=f"reme_{digest}", system=self.id, workspace_id=settings.workspace_id, content=content, memory_type="file-memory", source=str(path.relative_to(root)), confidence=1.0, created_at=updated, updated_at=updated, metadata={"path": str(path.relative_to(root)), "retrieval": "hybrid"}, )) return items async def delete_memory(self, memory_id: str) -> bool: """按记忆 ID 删除文件,并同步清理索引和操作记录。""" root = settings.data_dir / "reme" deleted = False for path in root.glob("daily/**/*.md"): digest = hashlib.sha1(str(path).encode("utf-8")).hexdigest()[:16] if f"reme_{digest}" != memory_id: continue path.unlink(missing_ok=True) deleted = True break # 文件删掉后,对应向量和文件索引也要一起更新。 if deleted: await self.repo.delete_embedding(self.id, memory_id) if deleted and self._class and settings.reme_enabled: service = await self._ensure() await service.run_job("reindex") if deleted: await self._audit("DELETE/FileMemory", target=memory_id) return deleted async def _ensure(self) -> Any: """按需启动 ReMe 服务,同一进程内重复使用。""" if not self._class or not settings.reme_enabled: raise AdapterUnavailable(self.descriptor.setup_hint or "ReMe 不可用") if self._service is None: working_dir = str(settings.data_dir / "reme") config = self._resolve_config( log_config=False, workspace_dir=working_dir, enable_logo=False, log_to_console=False, service={"backend": "http"}, jobs=self._PASSIVE_JOB_OVERRIDES, ) self._service = self._class(**config) await self._service.start() return self._service @classmethod def _should_memorize(cls, message: str) -> bool: """判断当前输入是否值得写入长期文件。""" normalized = " ".join(message.strip().split()) if not normalized: return False if cls._EXPLICIT_WRITE_RE.search(normalized): return True if cls._QUESTION_RE.search(normalized): return False return bool(cls._DECLARATIVE_RE.search(normalized)) @classmethod def _heuristic_queries(cls, message: str) -> list[str]: """为文件检索生成短查询和工程词扩展。""" normalized = message.strip() queries: list[str] = [normalized] if normalized else [] # 保留 RAG、FastAPI、PostgreSQL 等可直接匹配的领域词。 queries.extend(re.findall(r"[A-Za-z][A-Za-z0-9_.:/-]{1,}", normalized)) for trigger, expansions in cls._QUERY_EXPANSIONS.items(): # 用户用词较少时,补上文件里可能出现的近义说法。 if trigger in normalized: queries.extend((trigger, *expansions)) # 中文短词作为兜底,过滤会干扰词面检索的疑问词。 for chunk in re.findall(r"[\u4e00-\u9fff]{2,}", normalized): if chunk not in cls._QUERY_STOPWORDS: queries.append(chunk) if len(chunk) > 2: queries.extend( chunk[index : index + 2] for index in range(len(chunk) - 1) if chunk[index : index + 2] not in cls._QUERY_STOPWORDS ) return cls._dedupe_queries(queries) @staticmethod def _dedupe_queries(queries: list[str]) -> list[str]: """清理重复或空查询,并限制单次搜索数量。""" unique: list[str] = [] seen: set[str] = set() for query in queries: cleaned = " ".join(str(query).strip().split()) if not cleaned: continue key = cleaned.casefold() if key in seen: continue seen.add(key) unique.append(cleaned) return unique[:12] async def _rewrite_query(self, message: str) -> list[str]: """合并规则扩展和可选的模型查询改写。""" queries = self._heuristic_queries(message) try: rewritten = await self.llm.json( ( "你是记忆检索查询改写器。只输出 JSON,不要回答用户问题。" "把用户问题改写成 2 到 5 个适合混合文件检索的短查询," "保留专有名词,并补充可能出现在记忆里的中英文同义词。" '格式必须是 {"queries": ["..."]}。' ), f"用户问题:{message}", ) generated = rewritten.get("queries", []) if isinstance(generated, list): queries.extend(str(item) for item in generated if str(item).strip()) except Exception: # 改写失败时继续使用规则查询,不影响基础检索。 pass return self._dedupe_queries(queries) @staticmethod def _merge_search_results(results: list[str]) -> str: """合并多个查询结果并按正文去重。""" merged: list[str] = [] seen: set[str] = set() for result in results: for block in re.split(r"(?=^========== )", result.strip(), flags=re.MULTILINE): block = block.strip() if not block: continue fingerprint_source = block if block.startswith("==========") and "\n" in block: fingerprint_source = block.split("\n", 1)[1] fingerprint = re.sub(r"\[(?:vector_)?score=[^\]]+\]", "", fingerprint_source) fingerprint = " ".join(fingerprint.split()).casefold() if fingerprint in seen: continue seen.add(fingerprint) merged.append(block) return "\n\n".join(merged) async def _search_memory(self, service: Any, queries: list[str]) -> tuple[str, list[str]]: """逐个查询 ReMe 文件,并返回有结果的查询词。""" results: list[str] = [] used_queries: list[str] = [] for query in queries: search_result = await service.run_job("search", query=query, limit=5) answer = str(getattr(search_result, "answer", "") or "").strip() if answer: results.append(answer) used_queries.append(query) return self._merge_search_results(results), used_queries async def _sync_embeddings(self) -> dict[str, Any]: """更新已变化文件的向量,并清理已删除文件。""" if not embedding_client.configured: return {"enabled": False, "embedded": 0, "removed": 0} items = await self.memories() current_ids = {item.id for item in items} previous_hashes = await self.repo.embedding_hashes(self.id) # 文件不存在后,旧向量也要删除。 stale_ids = set(previous_hashes) - current_ids for memory_id in stale_ids: await self.repo.delete_embedding(self.id, memory_id) # 只处理新增或内容有变化的文件。 pending = [ item for item in items if previous_hashes.get(item.id) != embedding_fingerprint(item.content) ] if not pending: return {"enabled": True, "embedded": 0, "removed": len(stale_ids)} vectors = await embedding_client.embed([item.content for item in pending]) for item, vector in zip(pending, vectors): await self.repo.upsert_embedding( system=self.id, memory_id=item.id, content=item.content, source=item.source, content_hash=embedding_fingerprint(item.content), embedding=vector, ) return {"enabled": True, "embedded": len(pending), "removed": len(stale_ids)} async def _semantic_search(self, message: str) -> tuple[str, int]: """查找内容相近的文件,并过滤分数过低的结果。""" if not embedding_client.configured: return "", 0 try: query_vector = (await embedding_client.embed([message]))[0] matches = await self.repo.search_embeddings(self.id, query_vector, limit=5) except Exception: # 这一路失败时,主流程仍可使用 ReMe 文件搜索。 return "", 0 matches = [ match for match in matches if float(match.get("score", 0.0)) >= settings.embedding_min_score ] blocks = [ ( f"========== {match['source']} [vector_score={float(match['score']):.4f}] ==========" f"\n{match['content']}" ) for match in matches ] return self._merge_search_results(blocks), len(matches) async def _answer_from_memory(self, message: str, evidence: str) -> str: """根据搜索结果生成回答,模型不可用时直接返回原文。""" if not evidence: return "没有检索到与这个问题相关的历史记忆。" try: return await self.llm.chat( ( "你是一个使用文件记忆的 AI 助手。" "只能依据提供的记忆证据回答,不要编造证据中没有的事实。" "如果证据不足,要明确说明不确定。直接回答用户问题,简洁自然。" ), f"用户问题:{message}\n\n记忆证据:\n{evidence}", ) except Exception: # 保留原始内容比吞掉已经找到的结果更方便排查。 return f"根据检索到的记忆:\n\n{evidence}" async def chat(self, message: str) -> ChatResponse: """完成文件写入、索引更新、搜索、回答和操作记录。""" service = await self._ensure() write_candidate = self._should_memorize(message) memory_result: Any | None = None # 问题和临时请求不写文件,只处理值得长期保存的陈述。 if write_candidate: memory_result = await service.run_job( "auto_memory", messages=[{"name": "user", "role": "user", "content": message}], session_id=settings.workspace_id, memory_hint=( "只记录稳定偏好、项目事实、进度变化、可复用的故障经验和长期流程规则。" "问题、寒暄和临时请求不应写入长期记忆。" ), ) # 写入后立即更新索引,保证本轮搜索能看到刚保存的内容。 reindex_result = await service.run_job("reindex") try: embedding_sync = await self._sync_embeddings() except Exception as exc: # 向量服务异常时仍保留 ReMe 自带的文件检索。 embedding_sync = {"enabled": True, "embedded": 0, "removed": 0, "error": str(exc)} queries = await self._rewrite_query(message) # 文件搜索和相近内容搜索各跑一次,最后合并并去重。 lexical_search, used_queries = await self._search_memory(service, queries) semantic_search, semantic_count = await self._semantic_search(message) search = self._merge_search_results([lexical_search, semantic_search]) answer = await self._answer_from_memory(message, search) # 记录查询词和搜索摘要,出现误匹配时可回看整个过程。 await self._audit( "REME/auto_memory+search", details={ "write_candidate": write_candidate, "memory": str(getattr(memory_result, "answer", "skipped"))[:800], "reindex": str(reindex_result.metadata)[:800], "embedding_sync": embedding_sync, "queries": queries, "used_queries": used_queries, "semantic_results": semantic_count, "search": str(search)[:1000], }, ) return ChatResponse( system="reme", answer=answer, mode="real-sdk", memory_context=await self.memories(), memory_events=[ { "event": "file-memory", "write_candidate": write_candidate, "queries": queries, "retrieval": "hybrid", "semantic_results": semantic_count, "result": str(search)[:2000], } ], audit_events=await self.audit(), ) async def reset(self) -> None: """清空公共数据,关闭服务并删除 ReMe 工作目录。""" await super().reset() # 先关闭服务,避免删除目录后仍有后台句柄占用文件。 if self._service is not None: await self._service.close() self._service = None working_dir = settings.data_dir / "reme" if working_dir.exists(): import shutil shutil.rmtree(working_dir)