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): id = "reme" _PASSIVE_JOB_OVERRIDES = { # ReMe starts these as background/cron jobs by default. This local lab # keeps memory work request-driven to avoid silent file changes or LLM # cost; chat() performs an explicit reindex before every retrieval. "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}, } # These lightweight expansions strengthen ReMe's BM25 branch. The project # intentionally leaves ReMe's internal embedding_store disabled and uses # pgvector as the separate Embedding retrieval signal when configured. _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: 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]: """Expose ReMe's Markdown files in the common inspector contract.""" root = settings.data_dir / "reme" if not root.exists(): return [] items: list[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: 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: 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]: """Build short queries for ReMe's file search.""" normalized = message.strip() queries: list[str] = [normalized] if normalized else [] # Preserve explicit Latin/domain tokens such as RAG, FastAPI and # PostgreSQL; ReMe's tokenizer can match these reliably. 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)) # Add short Chinese chunks as a fallback, while excluding question # words that add noise to a lexical search. 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]: """Combine deterministic expansions with optional LLM query rewrite.""" 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: # ReMe remains usable when the optional rewrite call fails. The # deterministic expansions above cover the common engineering terms. pass return self._dedupe_queries(queries) @staticmethod def _merge_search_results(results: list[str]) -> str: """Merge duplicate snippets returned by multiple ReMe search queries.""" 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]]: 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]: """Embed changed ReMe files and remove vectors for deleted files.""" 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: 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: # A temporary Embedding outage should not take down file-memory # retrieval; ReMe's built-in file search remains available. 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: 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)