| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411 |
- 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)
|