Files
mission-control-v2/hermes/plugins/mc2-memory/__init__.py
T
Hitonabi 5c3f50dfa5 Fix: mc2-memory Provider flusht offene Turns bei Session-Ende
shutdown() + on_session_end() schreiben verbliebene Queue-Eintraege synchron raus
(_flush), damit kurzlebige Prozesse den Hintergrund-Worker nicht mitten im /learn-POST
killen. Live verifiziert: hands-off Auto-Lernen via Gateway (source=hermes) + Recall.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-27 20:12:13 +02:00

177 lines
6.8 KiB
Python

"""
MC2 Memory — Hermes-Memory-Provider, der ans geteilte **Mem0-Gedächtnis** hängt (über MC2 :9001,
das den Mem0-Sidecar :8765 proxyt). Macht das Gedächtnis für Hermes **hands-off**:
- `sync_turn` → nach jedem Turn werden die Turn-Messages an `/api/memory/learn` geschickt;
Mem0 EXTRAHIERT dauerhafte Fakten selbst (infer=True), dedupliziert semantisch. Nicht-blockierend
(Hintergrund-Worker), damit der Agent nie auf die LLM-Extraktion wartet.
- `prefetch` → vor jedem Turn semantische Suche gegen die Nutzer-Message; relevante Fakten werden
als Kontext eingeblendet (automatischer Recall).
**Context-only:** `get_tool_schemas()` liefert `[]` → der Agent bekommt KEINE Memory-Tools. Damit
entfällt das Tool-Loop-Risiko, wegen dem das native Memory-Toolset abgeschaltet ist
(`agent.disabled_toolsets: [..., memory]` bleibt). Single Source of Truth bleibt der MC2-Sidecar
(NoThink-Fix + alleiniger Chroma-Besitzer) — dieser Provider ist nur ein dünner HTTP-Client.
Aktivierung (in ~/.hermes/config.yaml):
memory:
memory_enabled: true
provider: mc2-memory
"""
from __future__ import annotations
import logging
import os
import queue
import threading
import time
from typing import Any, Dict, List, Optional
import httpx
from agent.memory_provider import MemoryProvider
log = logging.getLogger(__name__)
MC_URL = os.environ.get("MC_URL", "http://127.0.0.1:9001").rstrip("/")
MC_TOKEN = os.environ.get("MC_TOKEN", "")
RECALL_LIMIT = int(os.environ.get("MC2_MEMORY_RECALL_LIMIT", "5"))
RECALL_MIN_SCORE = float(os.environ.get("MC2_MEMORY_RECALL_MIN_SCORE", "0.3"))
MIN_USER_LEN = int(os.environ.get("MC2_MEMORY_MIN_USER_LEN", "12")) # zu Kurzes nicht lernen
def _headers() -> dict:
return {"X-MC-Token": MC_TOKEN} if MC_TOKEN else {}
class MC2MemoryProvider(MemoryProvider):
def __init__(self) -> None:
self._session_id = ""
self._agent_context = "primary"
self._q: "queue.Queue[list]" = queue.Queue(maxsize=200)
self._worker: Optional[threading.Thread] = None
self._stop = threading.Event()
self._recall_cache: Dict[str, str] = {}
# -- Identität / Verfügbarkeit -------------------------------------------
def name(self) -> str:
return "mc2-memory"
def is_available(self) -> bool:
try:
r = httpx.get(f"{MC_URL}/api/health", timeout=3.0)
return r.status_code == 200
except Exception:
return False
def initialize(self, session_id: str, **kwargs) -> None:
self._session_id = session_id
# Nur im primären Agent-Kontext lernen (cron/subagent-Systemprompts nicht einlernen).
self._agent_context = kwargs.get("agent_context", "primary")
if self._worker is None or not self._worker.is_alive():
self._stop.clear()
self._worker = threading.Thread(target=self._run, name="mc2-memory-learn", daemon=True)
self._worker.start()
def system_prompt_block(self) -> str:
return (
"Du hast ein dauerhaftes, geteiltes Langzeitgedächtnis (MC2/Mem0). Relevante Fakten "
"werden dir vor einem Turn automatisch eingeblendet; neue dauerhafte Fakten über den "
"Nutzer und das Projekt werden nach dem Turn automatisch gelernt — du musst dafür "
"nichts tun."
)
# -- Recall (vor dem Turn) -----------------------------------------------
def prefetch(self, query: str, *, session_id: str = "") -> str:
q = (query or "").strip()
if len(q) < 3:
return ""
if q in self._recall_cache:
return self._recall_cache[q]
try:
r = httpx.get(f"{MC_URL}/api/memory", params={"q": q}, headers=_headers(), timeout=5.0)
r.raise_for_status()
items = r.json()
except Exception as exc: # nie den Turn blockieren/abbrechen
log.debug("mc2-memory prefetch failed: %s", exc)
return ""
facts = [
i for i in items
if i.get("score") is None or float(i.get("score", 0)) >= RECALL_MIN_SCORE
][:RECALL_LIMIT]
out = ""
if facts:
lines = "\n".join(f"- {f.get('content', '')}" for f in facts)
out = f"Relevante Fakten aus dem Langzeitgedächtnis:\n{lines}"
# Cache klein halten
if len(self._recall_cache) > 64:
self._recall_cache.clear()
self._recall_cache[q] = out
return out
# -- Lernen (nach dem Turn) ----------------------------------------------
def sync_turn(self, user_content: str, assistant_content: str, *,
session_id: str = "", messages: Optional[List[Dict[str, Any]]] = None) -> None:
if self._agent_context != "primary":
return
u = (user_content or "").strip()
if len(u) < MIN_USER_LEN:
return
msgs: List[Dict[str, str]] = [{"role": "user", "content": u}]
a = (assistant_content or "").strip()
if a:
msgs.append({"role": "assistant", "content": a[:4000]})
try:
self._q.put_nowait(msgs)
except queue.Full:
log.debug("mc2-memory learn queue full — Turn übersprungen")
# Neue Fakten können gelandet sein → Recall-Cache invalidieren.
self._recall_cache.clear()
def _post_learn(self, msgs: list) -> None:
try:
httpx.post(f"{MC_URL}/api/memory/learn",
json={"messages": msgs, "source": "hermes"},
headers=_headers(), timeout=90.0)
except Exception as exc:
log.debug("mc2-memory learn POST failed: %s", exc)
def _run(self) -> None:
while not self._stop.is_set():
try:
msgs = self._q.get(timeout=1.0)
except queue.Empty:
continue
try:
self._post_learn(msgs)
finally:
self._q.task_done()
def _flush(self) -> None:
"""Offene Turns garantiert rausschreiben — bei Session-Ende/CLI-Exit, wo der
Prozess sofort beendet wird (daemon-Worker würde sonst mitten im POST sterben)."""
self._stop.set()
if self._worker and self._worker.is_alive():
self._worker.join(timeout=95.0) # laufenden Worker-POST zu Ende lassen
while True: # vom Worker nicht mehr abgeholte Turns
try:
msgs = self._q.get_nowait()
except queue.Empty:
break
self._post_learn(msgs)
# -- Context-only: keine Agent-Tools → kein Tool-Loop --------------------
def get_tool_schemas(self) -> List[Dict[str, Any]]:
return []
def on_session_end(self, messages: List[Dict[str, Any]]) -> None:
self._flush()
def shutdown(self) -> None:
self._flush()
def register(ctx) -> None:
ctx.register_memory_provider(MC2MemoryProvider())