feat: complete MC2 Kahlschlag (remove governor, memory, eigenleben, obsolete deploy scripts and frontend widgets)
Ampel / ampel (push) Successful in 24s
Ampel / ampel (push) Successful in 24s
This commit is contained in:
@@ -25,15 +25,12 @@ from routers import (
|
||||
chronik,
|
||||
connect,
|
||||
console,
|
||||
eigenleben,
|
||||
events,
|
||||
gateway_proxy,
|
||||
governor,
|
||||
health,
|
||||
hermes_ui,
|
||||
ideen,
|
||||
maintenance,
|
||||
memory,
|
||||
models,
|
||||
routing,
|
||||
system,
|
||||
@@ -43,7 +40,6 @@ from routers import (
|
||||
)
|
||||
from routers import reminders as reminders_router
|
||||
from services import ketten_digest, metrics_history, reminders, sentry, warmer
|
||||
from services import memory as memory_svc
|
||||
from starlette.requests import Request
|
||||
|
||||
# Zentrales Logging — Level via MC_LOG_LEVEL (INFO default). Eine Konfiguration
|
||||
@@ -74,13 +70,6 @@ async def lifespan(app: FastAPI) -> AsyncGenerator[None, None]:
|
||||
tasks.append(asyncio.create_task(ketten_digest.digest_loop()))
|
||||
# 24-h-Metrik-Verlauf (Cockpit-Zeitachse): 10-s-Sampler, Ringpuffer, persistiert.
|
||||
tasks.append(asyncio.create_task(metrics_history.sampler_loop()))
|
||||
if memory_svc.AUTO_DEDUPE_ENABLED:
|
||||
tasks.append(asyncio.create_task(memory_svc.auto_dedupe_loop()))
|
||||
log.info(
|
||||
"Mem0-Auto-Dedupe aktiv (alle %ss, Schwelle %s)",
|
||||
memory_svc.AUTO_DEDUPE_INTERVAL,
|
||||
memory_svc.AUTO_DEDUPE_THRESHOLD,
|
||||
)
|
||||
# Geteilter HTTP-Client zur lokalen Engine: Keep-Alive/Connection-Pooling statt neuer Client
|
||||
# pro /v1-Anfrage (spart Sockets/TIME_WAIT unter parallelen Agent-Strömen von Zed/Kilo).
|
||||
app.state.gw_client = httpx.AsyncClient(
|
||||
@@ -121,7 +110,6 @@ app.include_router(models.router)
|
||||
app.include_router(routing.router)
|
||||
app.include_router(system.router)
|
||||
app.include_router(connect.router)
|
||||
app.include_router(memory.router)
|
||||
app.include_router(agent.router)
|
||||
app.include_router(
|
||||
voice.router
|
||||
@@ -142,8 +130,6 @@ app.include_router(maintenance.router)
|
||||
app.include_router(auftragsbuch.router) # Vorschlags-Inbox (Mensch-Gate als Klick)
|
||||
app.include_router(ideen.router) # Ideen-Queue (natives Hermes-Kanban) — Tür der Zentrale
|
||||
app.include_router(chronik.router) # Timeline der autonomen Taten (Announce-Store)
|
||||
app.include_router(governor.router) # Token-Wächter (:8100) — Zählerstand fürs Cockpit
|
||||
app.include_router(eigenleben.router) # „Von allein": Skills + Vorschlags-Bilanz der Box
|
||||
app.include_router(events.router) # SSE-Eventstrom /api/events (P3a) — Invalidation-Bus
|
||||
app.include_router(wissen.router) # Wissens-Vault (Traum-Notizen) read-only
|
||||
app.include_router(zeitmaschine.router) # Snapshots ansehen + Ein-Klick-Restore (detached)
|
||||
|
||||
@@ -1,54 +0,0 @@
|
||||
"""Eigenleben-Endpoints — „Von allein“-Ansicht (Skills + Vorschlags-Bilanz), read-only."""
|
||||
|
||||
from fastapi import APIRouter, HTTPException
|
||||
from pydantic import BaseModel
|
||||
from services import eigenleben
|
||||
|
||||
router = APIRouter(prefix="/api")
|
||||
|
||||
|
||||
class SkillItem(BaseModel):
|
||||
id: str
|
||||
name: str
|
||||
kategorie: str | None
|
||||
aktiv: bool | None
|
||||
beschreibung: str
|
||||
version: str
|
||||
autor: str
|
||||
herkunft: str
|
||||
herkunft_label: str
|
||||
genutzt: int
|
||||
zuletzt_genutzt: str | None
|
||||
erstellt: str | None
|
||||
status: str
|
||||
geaendert: float
|
||||
|
||||
|
||||
class EigenlebenOverviewResponse(BaseModel):
|
||||
available: bool
|
||||
skills: list[SkillItem]
|
||||
selbst_erstellt: int
|
||||
bilanz: dict
|
||||
|
||||
|
||||
class SkillTextResponse(BaseModel):
|
||||
id: str
|
||||
text: str
|
||||
|
||||
|
||||
@router.get("/eigenleben", response_model=EigenlebenOverviewResponse)
|
||||
def get_overview() -> dict:
|
||||
return eigenleben.overview()
|
||||
|
||||
|
||||
# :path — Kategorie-Skills haben einen Schrägstrich in der ID (autonomous-ai-agents/
|
||||
# hermes-agent). Als normales {skill_id}-Segment matcht die Route dann nicht, der Request
|
||||
# fällt in den SPA-Fallback (index.html, Status 200) und die UI zeigt „Konnte den
|
||||
# Skill-Text nicht laden". Die Pfad-Validierung (max. eine Ebene, keine . /..) macht
|
||||
# services.eigenleben.skill_text.
|
||||
@router.get("/eigenleben/skill/{skill_id:path}", response_model=SkillTextResponse)
|
||||
def get_skill(skill_id: str) -> dict:
|
||||
text = eigenleben.skill_text(skill_id)
|
||||
if text is None:
|
||||
raise HTTPException(404, "Skill nicht gefunden.")
|
||||
return {"id": skill_id, "text": text}
|
||||
@@ -1,40 +0,0 @@
|
||||
"""Governor — Fenster auf den Token-Waechter (:8100).
|
||||
|
||||
Der Governor ist ein eigenstaendiger, absichtlich winziger Proxy ohne Datenbank: er
|
||||
sitzt zwischen den Coding-Agenten und diesem Gateway, zaehlt ehrlich mit (echte
|
||||
`usage.prompt_tokens` aus jeder Antwort) und zieht bei ueberlangen Sitzungen die
|
||||
Notbremse. Hier wird nichts Neues erhoben — nur sein Status-Endpunkt gleichursprünglich
|
||||
fuer die Oberflaeche verfuegbar gemacht, damit das Frontend nicht per CORS auf einen
|
||||
zweiten Port ausweichen muss.
|
||||
|
||||
Faellt der Governor aus, liefert dieser Router `ok: false` statt eines Fehlers: die
|
||||
Kachel zeigt dann „nicht erreichbar" und das Cockpit bleibt heil.
|
||||
"""
|
||||
|
||||
import os
|
||||
|
||||
import httpx
|
||||
from fastapi import APIRouter
|
||||
|
||||
router = APIRouter(prefix="/api")
|
||||
|
||||
GOVERNOR_URL = os.environ.get("MC_GOVERNOR_URL", "http://127.0.0.1:8100")
|
||||
|
||||
|
||||
@router.get("/governor")
|
||||
async def governor_status() -> dict:
|
||||
"""Momentaufnahme des Token-Waechters. Nie werfen — die Kachel darf nie das Cockpit reissen."""
|
||||
try:
|
||||
async with httpx.AsyncClient(timeout=3.0) as c:
|
||||
r = await c.get(f"{GOVERNOR_URL}/governor/status")
|
||||
r.raise_for_status()
|
||||
data = r.json()
|
||||
data["reachable"] = True
|
||||
return data
|
||||
except Exception as exc:
|
||||
return {
|
||||
"ok": False,
|
||||
"reachable": False,
|
||||
"error": f"{type(exc).__name__}: {exc}",
|
||||
"url": GOVERNOR_URL,
|
||||
}
|
||||
@@ -1,92 +0,0 @@
|
||||
"""Memory-Endpoints (geteiltes Gedächtnis). LAN-only, kein Token in 2.0-Phase 3."""
|
||||
|
||||
import time
|
||||
|
||||
from fastapi import APIRouter, HTTPException
|
||||
from pydantic import BaseModel
|
||||
from services import memory
|
||||
from services.voice_metrics import park, record_stage
|
||||
|
||||
router = APIRouter(prefix="/api")
|
||||
|
||||
|
||||
class MemIn(BaseModel):
|
||||
content: str
|
||||
category: str = "knowledge"
|
||||
source: str = "manual"
|
||||
|
||||
|
||||
class MemUp(BaseModel):
|
||||
content: str | None = None
|
||||
category: str | None = None
|
||||
|
||||
|
||||
class DedupeIn(BaseModel):
|
||||
apply: bool = False
|
||||
threshold: float = 0.85
|
||||
|
||||
|
||||
class LearnIn(BaseModel):
|
||||
text: str | None = None
|
||||
messages: list[dict] | None = None
|
||||
source: str = "auto"
|
||||
category: str = "knowledge"
|
||||
|
||||
|
||||
@router.get("/memory/export")
|
||||
def export() -> dict:
|
||||
return memory.export_text()
|
||||
|
||||
|
||||
@router.get("/memory/graph")
|
||||
def graph(min_score: float = 0.45, top_k: int = 3) -> dict:
|
||||
"""Fakten als Ähnlichkeits-Graph (Knoten + semantische Kanten) für die Visualisierung."""
|
||||
return memory.graph(min_score=min_score, top_k=top_k)
|
||||
|
||||
|
||||
@router.post("/memory/learn", status_code=201)
|
||||
def learn(body: LearnIn) -> dict:
|
||||
"""Auto-Lernen: Gesprächs-Turns/Text durchreichen → Mem0 extrahiert Fakten selbst."""
|
||||
return memory.learn(text=body.text, messages=body.messages,
|
||||
source=body.source, category=body.category)
|
||||
|
||||
|
||||
@router.post("/memory/dedupe")
|
||||
def dedupe(body: DedupeIn) -> dict:
|
||||
return memory.dedupe(apply=body.apply, threshold=body.threshold)
|
||||
|
||||
|
||||
@router.get("/memory")
|
||||
def list_mem(q: str = "", category: str = "") -> list[dict]:
|
||||
# Semantischer Retrieve (q gesetzt) = u.a. Hermes' Mem0-Prefetch VOR jedem Turn. Dauer messen
|
||||
# + parken, damit der laufende Voice-Chat-Turn sie als Unter-Detail seiner Hirn-Zeit einsammelt.
|
||||
if q:
|
||||
_t0 = time.perf_counter()
|
||||
res = memory.list_memories(q=q, category=category)
|
||||
_ms = (time.perf_counter() - _t0) * 1000.0
|
||||
record_stage("memory_retrieve", _ms)
|
||||
park("retrieve", _ms)
|
||||
return res
|
||||
return memory.list_memories(q=q, category=category)
|
||||
|
||||
|
||||
@router.post("/memory", status_code=201)
|
||||
def add(body: MemIn) -> dict:
|
||||
if body.category not in memory.CATEGORIES:
|
||||
raise HTTPException(400, f"Kategorie '{body.category}' unbekannt.")
|
||||
return memory.add_memory(body.content, body.category, body.source)
|
||||
|
||||
|
||||
@router.put("/memory/{mid}")
|
||||
def update(mid: str, body: MemUp) -> dict:
|
||||
res = memory.update_memory(mid, content=body.content, category=body.category)
|
||||
if not res:
|
||||
raise HTTPException(404, "Eintrag nicht gefunden")
|
||||
return res
|
||||
|
||||
|
||||
@router.delete("/memory/{mid}")
|
||||
def delete(mid: str) -> dict:
|
||||
if not memory.delete_memory(mid):
|
||||
raise HTTPException(404, "Eintrag nicht gefunden")
|
||||
return {"ok": True}
|
||||
@@ -1,219 +0,0 @@
|
||||
"""Eigenleben — was die Box von allein gebaut und vorgeschlagen hat.
|
||||
|
||||
Beantwortet die User-Frage vom 15.07.2026: „Es gibt keinen Viewpoint, welche Skills
|
||||
erstellt wurden, was die Box vorgeschlagen hat, welche Skills Lucy hat und wie die
|
||||
funktionieren." Reine LESE-Schicht, keine Schreib-Operationen:
|
||||
|
||||
- Skills aus ~/.hermes/skills (= Lucys Skill-Satz; die Worker-Profile werkstatt/betrieb
|
||||
tragen Kopien desselben Satzes). Herkunft dreistufig klassifiziert:
|
||||
selbst → weder mitgeliefert noch aus unserem Repo = die Box/Lucy hat ihn erzeugt
|
||||
repo → liegt in ~/mission-control-v2/deploy/skills (von uns gebaut + deployt)
|
||||
bundled → liegt in ~/.hermes/hermes-agent/skills (kam mit Hermes mit)
|
||||
- Nutzungszahlen aus ~/.hermes/skills/.usage.json (Hermes' eigener Zähler).
|
||||
- Bilanz aus Auftragsbuch-Chronik (/srv/models/mc2-auftragsbuch.json, angenommene
|
||||
Branches) + Ablehnungs-Lern-Journal (/srv/models/mc2-ablehnungen.jsonl).
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import re
|
||||
import subprocess
|
||||
import time
|
||||
from pathlib import Path
|
||||
|
||||
SKILLS_DIR = Path.home() / ".hermes" / "skills"
|
||||
BUNDLED_DIR = Path.home() / ".hermes" / "hermes-agent" / "skills"
|
||||
REPO_DIR = Path.home() / "mission-control-v2" / "deploy" / "skills"
|
||||
HERMES_SRC = Path.home() / ".hermes" / "hermes-agent"
|
||||
HERMES_PY = HERMES_SRC / "venv" / "bin" / "python"
|
||||
USAGE_FILE = SKILLS_DIR / ".usage.json"
|
||||
AUFTRAGSBUCH_FILE = Path("/srv/models/mc2-auftragsbuch.json")
|
||||
ABLEHNUNGEN_FILE = Path("/srv/models/mc2-ablehnungen.jsonl")
|
||||
|
||||
_HERKUNFT_LABEL = {
|
||||
"selbst": "Selbst erstellt",
|
||||
"repo": "Von uns gebaut",
|
||||
"bundled": "Mit Hermes mitgeliefert",
|
||||
}
|
||||
|
||||
|
||||
def _frontmatter(text: str) -> dict:
|
||||
"""Sehr kleiner Frontmatter-Leser: nur die flachen key: value-Zeilen des ersten
|
||||
---Blocks (name/description/version/author reichen hier; kein YAML-Import nötig)."""
|
||||
out: dict = {}
|
||||
if not text.startswith("---"):
|
||||
return out
|
||||
for line in text.split("\n", 1)[1].split("\n"):
|
||||
if line.strip() == "---":
|
||||
break
|
||||
m = re.match(r"^([A-Za-z_][\w-]*):\s*(.*)$", line)
|
||||
if not m:
|
||||
continue
|
||||
val = m.group(2).strip().strip('"').strip("'")
|
||||
if val:
|
||||
out[m.group(1)] = val
|
||||
return out
|
||||
|
||||
|
||||
def _norm(s: str) -> str:
|
||||
return re.sub(r"[\s_-]+", "", (s or "").lower())
|
||||
|
||||
|
||||
# aktiv/inaktiv: dieselbe Wahrheit wie Hermes selbst (und sein eingebautes WebUI) —
|
||||
# hermes_cli.banner.get_available_skills() prüft Plattform/Voraussetzungen/Config und
|
||||
# liefert, was der Agent WIRKLICH laden kann. Wir rufen sie im Hermes-venv auf (MC2
|
||||
# läuft in eigenem Python) und cachen großzügig — der Skill-Satz ändert sich selten.
|
||||
_aktiv_cache: dict = {"ts": 0.0, "namen": None}
|
||||
_AKTIV_EVERY = 600.0 # s
|
||||
|
||||
|
||||
def _aktive_namen() -> set[str] | None:
|
||||
"""Normalisierte Namen aller ladbaren Skills; None = nicht ermittelbar (Dev-Modus)."""
|
||||
if not HERMES_PY.is_file():
|
||||
return None
|
||||
now = time.time()
|
||||
if _aktiv_cache["namen"] is not None and now - _aktiv_cache["ts"] < _AKTIV_EVERY:
|
||||
return _aktiv_cache["namen"]
|
||||
code = ("from hermes_cli.banner import get_available_skills;import json;"
|
||||
"print(json.dumps(get_available_skills()))")
|
||||
try:
|
||||
r = subprocess.run([str(HERMES_PY), "-c", code], capture_output=True,
|
||||
text=True, timeout=60, cwd=str(HERMES_SRC))
|
||||
idx = r.stdout.find("{")
|
||||
kategorien = json.loads(r.stdout[idx:]) if idx >= 0 else {}
|
||||
namen = {_norm(n) for liste in kategorien.values() for n in liste}
|
||||
except Exception:
|
||||
return _aktiv_cache["namen"] # alter Stand ist besser als flackerndes None
|
||||
_aktiv_cache.update(ts=now, namen=namen)
|
||||
return namen
|
||||
|
||||
|
||||
def list_skills() -> list[dict]:
|
||||
"""Alle Skills mit Beschreibung, Herkunft und Nutzungszahlen (leer im Dev-Modus)."""
|
||||
if not SKILLS_DIR.is_dir():
|
||||
return []
|
||||
try:
|
||||
usage_raw = json.loads(USAGE_FILE.read_text(encoding="utf-8"))
|
||||
except Exception:
|
||||
usage_raw = {}
|
||||
usage = {_norm(k): v for k, v in usage_raw.items()}
|
||||
aktive = _aktive_namen()
|
||||
|
||||
def eintrag(md: Path, rel: str, kategorie: str | None) -> dict:
|
||||
try:
|
||||
fm = _frontmatter(md.read_text(encoding="utf-8", errors="replace"))
|
||||
except Exception:
|
||||
fm = {}
|
||||
if (BUNDLED_DIR / rel).is_dir():
|
||||
herkunft = "bundled"
|
||||
elif (REPO_DIR / rel).is_dir():
|
||||
herkunft = "repo"
|
||||
else:
|
||||
herkunft = "selbst"
|
||||
kurz = rel.split("/")[-1]
|
||||
u = usage.get(_norm(kurz)) or usage.get(_norm(fm.get("name", ""))) or {}
|
||||
aktiv = None
|
||||
if aktive is not None:
|
||||
aktiv = _norm(kurz) in aktive or _norm(fm.get("name", "")) in aktive
|
||||
return {
|
||||
"id": rel,
|
||||
"name": fm.get("name") or kurz,
|
||||
"kategorie": kategorie,
|
||||
"aktiv": aktiv,
|
||||
"beschreibung": fm.get("description") or "",
|
||||
"version": fm.get("version") or "",
|
||||
"autor": fm.get("author") or "",
|
||||
"herkunft": herkunft,
|
||||
"herkunft_label": _HERKUNFT_LABEL[herkunft],
|
||||
"genutzt": int(u.get("use_count") or 0),
|
||||
"zuletzt_genutzt": u.get("last_used_at") or None,
|
||||
"erstellt": u.get("created_at") or None,
|
||||
"status": u.get("state") or "active",
|
||||
"geaendert": md.stat().st_mtime,
|
||||
}
|
||||
|
||||
skills: list[dict] = []
|
||||
for d in sorted(SKILLS_DIR.iterdir()):
|
||||
if not d.is_dir() or d.name.startswith("."):
|
||||
continue # .archive (Curator) & Co. sind keine aktiven Skills
|
||||
md = d / "SKILL.md"
|
||||
if md.is_file():
|
||||
skills.append(eintrag(md, d.name, None))
|
||||
continue
|
||||
# Kategorie-Ordner (mitgelieferte Sammlungen wie research/, apple/): eine Ebene
|
||||
# tiefer liegen die echten Skills (research/blogwatcher/SKILL.md) — je einzeln zeigen.
|
||||
for sub in sorted(d.iterdir()):
|
||||
smd = sub / "SKILL.md"
|
||||
if sub.is_dir() and not sub.name.startswith(".") and smd.is_file():
|
||||
skills.append(eintrag(smd, f"{d.name}/{sub.name}", d.name))
|
||||
# Selbst erstellte zuerst — das ist der Star der Ansicht.
|
||||
order = {"selbst": 0, "repo": 1, "bundled": 2}
|
||||
skills.sort(key=lambda s: (order[s["herkunft"]], -s["genutzt"], s["name"]))
|
||||
return skills
|
||||
|
||||
|
||||
def skill_text(skill_id: str) -> str | None:
|
||||
"""Voller SKILL.md-Text („wie funktioniert der Skill") — max. Kategorie/Name, keine Pfad-Tricks."""
|
||||
if not re.fullmatch(r"[\w.-]+(/[\w.-]+)?", skill_id or ""):
|
||||
return None
|
||||
# Das Zeichen-Set oben erlaubt „." und „.." als Segment — seit die Route ein
|
||||
# :path-Parameter ist, hier explizit raus (kein Klettern über SKILLS_DIR hinaus).
|
||||
if any(seg in (".", "..") for seg in skill_id.split("/")):
|
||||
return None
|
||||
md = SKILLS_DIR / skill_id / "SKILL.md"
|
||||
if not md.is_file():
|
||||
return None
|
||||
try:
|
||||
return md.read_text(encoding="utf-8", errors="replace")
|
||||
except Exception:
|
||||
return None
|
||||
|
||||
|
||||
def bilanz() -> dict:
|
||||
"""Vorschlags-Bilanz: was die Box vorschlug und was daraus wurde (eine Timeline)."""
|
||||
eintraege: list[dict] = []
|
||||
try:
|
||||
branches = json.loads(AUFTRAGSBUCH_FILE.read_text(encoding="utf-8")).get("branches") or {}
|
||||
for branch, info in branches.items():
|
||||
eintraege.append({
|
||||
"art": "angenommen" if info.get("state") == "eingespielt" else info.get("state", "?"),
|
||||
"titel": branch,
|
||||
"detail": info.get("detail") or "",
|
||||
"ts": float(info.get("ts") or 0),
|
||||
})
|
||||
except Exception:
|
||||
pass
|
||||
try:
|
||||
for line in ABLEHNUNGEN_FILE.read_text(encoding="utf-8").splitlines():
|
||||
if not line.strip():
|
||||
continue
|
||||
try:
|
||||
j = json.loads(line)
|
||||
except Exception:
|
||||
continue
|
||||
eintraege.append({
|
||||
"art": "abgelehnt",
|
||||
"titel": j.get("subject") or j.get("branch") or "?",
|
||||
"detail": j.get("grund") or "",
|
||||
"ts": float(j.get("ts") or 0),
|
||||
})
|
||||
except Exception:
|
||||
pass
|
||||
eintraege.sort(key=lambda e: -e["ts"])
|
||||
return {
|
||||
"eintraege": eintraege,
|
||||
"angenommen": sum(1 for e in eintraege if e["art"] == "angenommen"),
|
||||
"abgelehnt": sum(1 for e in eintraege if e["art"] == "abgelehnt"),
|
||||
}
|
||||
|
||||
|
||||
def overview() -> dict:
|
||||
skills = list_skills()
|
||||
b = bilanz()
|
||||
return {
|
||||
"available": SKILLS_DIR.is_dir(),
|
||||
"skills": skills,
|
||||
"selbst_erstellt": sum(1 for s in skills if s["herkunft"] == "selbst"),
|
||||
"bilanz": b,
|
||||
}
|
||||
@@ -1,176 +0,0 @@
|
||||
"""
|
||||
Geteiltes Gedächtnis (die „Verfassung") — jetzt auto-lernend & semantisch über Mem0.
|
||||
|
||||
Dieser Service ist nur noch ein dünner HTTP-Client auf den Mem0-Sidecar (mem0_service/app.py,
|
||||
läuft im ~/.mem0/venv unter Python 3.12). Die `/api/memory`-API-Form bleibt unverändert, damit
|
||||
UI und MCP-Server kompatibel bleiben. Neu gegenüber der alten flachen SQLite:
|
||||
|
||||
- search (q gesetzt) ist SEMANTISCH (Vektor/Embeddings) statt LIKE-Textsuche, mit Relevanz-Score.
|
||||
- learn() reicht Gesprächs-Turns durch → Mem0 EXTRAHIERT Fakten selbst (Auto-Lernen).
|
||||
- Dedup macht Mem0 beim Auto-Lernen selbst; der manuelle Kurator unten bleibt als Komfort.
|
||||
|
||||
5 Kategorien (user · instruction · stable · versioned · ephemeral) bleiben als Metadaten erhalten.
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import logging
|
||||
import os
|
||||
import re
|
||||
from difflib import SequenceMatcher
|
||||
|
||||
import httpx
|
||||
from config import MEM0_SERVICE_URL
|
||||
|
||||
log = logging.getLogger(__name__)
|
||||
|
||||
CATEGORIES = ("identity", "knowledge", "rules", "events")
|
||||
|
||||
_TIMEOUT = httpx.Timeout(60.0, connect=5.0) # LLM-Extraktion kann ein paar Sekunden dauern
|
||||
|
||||
|
||||
def _get(path: str, **params) -> list | dict:
|
||||
r = httpx.get(f"{MEM0_SERVICE_URL}{path}", params=params, timeout=_TIMEOUT)
|
||||
r.raise_for_status()
|
||||
return r.json()
|
||||
|
||||
|
||||
def _post(path: str, data: dict) -> dict:
|
||||
r = httpx.post(f"{MEM0_SERVICE_URL}{path}", json=data, timeout=_TIMEOUT)
|
||||
r.raise_for_status()
|
||||
return r.json()
|
||||
|
||||
|
||||
def _put(path: str, data: dict) -> dict:
|
||||
r = httpx.put(f"{MEM0_SERVICE_URL}{path}", json=data, timeout=_TIMEOUT)
|
||||
r.raise_for_status()
|
||||
return r.json()
|
||||
|
||||
|
||||
def _delete(path: str) -> dict:
|
||||
r = httpx.delete(f"{MEM0_SERVICE_URL}{path}", timeout=_TIMEOUT)
|
||||
r.raise_for_status()
|
||||
return r.json()
|
||||
|
||||
|
||||
def list_memories(q: str = "", category: str = "") -> list[dict]:
|
||||
"""Alle Fakten oder — wenn q gesetzt — die semantisch ähnlichsten (mit `score`)."""
|
||||
return _get("/memory", **{k: v for k, v in (("q", q), ("category", category)) if v})
|
||||
|
||||
|
||||
def add_memory(content: str, category: str = "stable", source: str = "manual") -> dict:
|
||||
"""Einen Fakt VERBATIM speichern (keine LLM-Umformung). Auto-Lernen → learn()."""
|
||||
return _post("/memory", {"content": content.strip(), "category": category, "source": source})
|
||||
|
||||
|
||||
def update_memory(mid: str, content: str | None = None, category: str | None = None) -> dict | None:
|
||||
try:
|
||||
return _put(f"/memory/{mid}", {"content": content, "category": category})
|
||||
except httpx.HTTPStatusError as exc:
|
||||
if exc.response.status_code == 404:
|
||||
return None
|
||||
raise
|
||||
|
||||
|
||||
def delete_memory(mid: str) -> bool:
|
||||
try:
|
||||
_delete(f"/memory/{mid}")
|
||||
return True
|
||||
except httpx.HTTPStatusError as exc:
|
||||
if exc.response.status_code == 404:
|
||||
return False
|
||||
raise
|
||||
|
||||
|
||||
def learn(text: str | None = None, messages: list[dict] | None = None,
|
||||
source: str = "auto", category: str = "stable") -> dict:
|
||||
"""Auto-Lernen: Text/Gesprächs-Turns durchreichen → Mem0 extrahiert die Fakten selbst."""
|
||||
return _post("/learn", {"text": text, "messages": messages,
|
||||
"source": source, "category": category})
|
||||
|
||||
|
||||
def graph(min_score: float = 0.45, top_k: int = 3) -> dict:
|
||||
"""Fakten als Ähnlichkeits-Graph (Knoten + semantische Kanten) für die UI-Visualisierung."""
|
||||
return _get("/graph", min_score=min_score, top_k=top_k)
|
||||
|
||||
|
||||
def export_text() -> dict:
|
||||
rows = sorted(list_memories(), key=lambda r: (r.get("category", ""), r.get("updated_at", "")))
|
||||
lines = ["# Mission Control — Gedächtnis\n"]
|
||||
current = ""
|
||||
for r in rows:
|
||||
if r.get("category") != current:
|
||||
current = r.get("category", "")
|
||||
lines.append(f"\n## {current}\n")
|
||||
src = r.get("source", "")
|
||||
when = (r.get("updated_at") or "")[:10]
|
||||
lines.append(f"- {r.get('content', '')} _(Quelle: {src}, {when})_")
|
||||
return {"text": "\n".join(lines), "count": len(rows)}
|
||||
|
||||
|
||||
# --- Manueller Kurator (deterministisch, kein LLM) ---------------------------
|
||||
def _norm(s: str) -> str:
|
||||
s = re.sub(r"[^\w\s]", " ", s.lower(), flags=re.UNICODE)
|
||||
return re.sub(r"\s+", " ", s).strip()
|
||||
|
||||
|
||||
def dedupe(apply: bool = False, threshold: float = 0.85) -> dict:
|
||||
"""Findet Dubletten (exakt/enthalten/ähnlich) je Kategorie, behält den längsten
|
||||
Eintrag. Mem0 dedupliziert beim Auto-Lernen schon semantisch — das hier ist der
|
||||
manuelle Komfort-Knopf fürs UI (z.B. nach vielen Verbatim-Importen)."""
|
||||
rows = sorted(list_memories(), key=lambda r: (-len(r.get("content", "")), r.get("created_at", "")))
|
||||
used: set[str] = set()
|
||||
groups: list[dict] = []
|
||||
for i, a in enumerate(rows):
|
||||
if a["id"] in used:
|
||||
continue
|
||||
na = _norm(a.get("content", ""))
|
||||
if not na:
|
||||
continue
|
||||
dups = []
|
||||
for b in rows[i + 1:]:
|
||||
if b["id"] in used or b.get("category") != a.get("category"):
|
||||
continue
|
||||
nb = _norm(b.get("content", ""))
|
||||
if not nb:
|
||||
continue
|
||||
if nb in na or na in nb or SequenceMatcher(None, na, nb).ratio() >= threshold:
|
||||
dups.append(b); used.add(b["id"])
|
||||
if dups:
|
||||
used.add(a["id"])
|
||||
groups.append({
|
||||
"keep": {"id": a["id"], "content": a.get("content"), "category": a.get("category")},
|
||||
"remove": [{"id": d["id"], "content": d.get("content")} for d in dups],
|
||||
})
|
||||
dup_count = sum(len(g["remove"]) for g in groups)
|
||||
removed = 0
|
||||
if apply:
|
||||
for g in groups:
|
||||
for d in g["remove"]:
|
||||
if delete_memory(d["id"]):
|
||||
removed += 1
|
||||
return {"groups": groups, "duplicate_count": dup_count, "removed": removed, "applied": apply}
|
||||
|
||||
|
||||
# ── Auto-Dedupe (D16e): Dubletten von selbst wegräumen, kein Knopf-Zwang ─────────────
|
||||
# Deterministisch in UNSERER Schicht (wie reminders, Verdikt 11c) statt am Hermes-cron —
|
||||
# läuft als Hintergrund-Task. Höhere Schwelle als der manuelle Knopf (0.9 statt 0.85), weil
|
||||
# OHNE Mensch-Review angewandt wird: lieber eine Dublette stehen lassen als etwas Distinktes
|
||||
# löschen. Mem0 dedupliziert beim Auto-Lernen schon semantisch — das hier fängt den Rest
|
||||
# (Verbatim-Importe, wiederholte Fakten) ab, damit das Gedächtnis nicht ohne Zutun aufbläht.
|
||||
AUTO_DEDUPE_ENABLED = os.environ.get("MC_MEM_DEDUPE_ENABLED", "1") != "0"
|
||||
AUTO_DEDUPE_INTERVAL = int(os.environ.get("MC_MEM_DEDUPE_INTERVAL", str(24 * 3600))) # täglich
|
||||
AUTO_DEDUPE_START_DELAY = int(os.environ.get("MC_MEM_DEDUPE_START_DELAY", "300")) # 5 min nach Start
|
||||
AUTO_DEDUPE_THRESHOLD = float(os.environ.get("MC_MEM_DEDUPE_THRESHOLD", "0.75"))
|
||||
|
||||
|
||||
async def auto_dedupe_loop() -> None:
|
||||
"""Hintergrund-Task: räumt periodisch Gedächtnis-Dubletten weg (apply=True)."""
|
||||
await asyncio.sleep(AUTO_DEDUPE_START_DELAY) # Mem0-Sidecar nach Start hochkommen lassen
|
||||
while True:
|
||||
try:
|
||||
res = await asyncio.to_thread(dedupe, True, AUTO_DEDUPE_THRESHOLD) # sync HTTP → Thread
|
||||
if res.get("removed"):
|
||||
log.info("mem-dedupe: %d Dublette(n) automatisch entfernt", res["removed"])
|
||||
except Exception:
|
||||
log.debug("mem-dedupe: Lauf fehlgeschlagen", exc_info=True)
|
||||
await asyncio.sleep(AUTO_DEDUPE_INTERVAL)
|
||||
Reference in New Issue
Block a user