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>
This commit is contained in:
@@ -25,6 +25,7 @@ import logging
|
|||||||
import os
|
import os
|
||||||
import queue
|
import queue
|
||||||
import threading
|
import threading
|
||||||
|
import time
|
||||||
from typing import Any, Dict, List, Optional
|
from typing import Any, Dict, List, Optional
|
||||||
|
|
||||||
import httpx
|
import httpx
|
||||||
@@ -128,6 +129,14 @@ class MC2MemoryProvider(MemoryProvider):
|
|||||||
# Neue Fakten können gelandet sein → Recall-Cache invalidieren.
|
# Neue Fakten können gelandet sein → Recall-Cache invalidieren.
|
||||||
self._recall_cache.clear()
|
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:
|
def _run(self) -> None:
|
||||||
while not self._stop.is_set():
|
while not self._stop.is_set():
|
||||||
try:
|
try:
|
||||||
@@ -135,20 +144,32 @@ class MC2MemoryProvider(MemoryProvider):
|
|||||||
except queue.Empty:
|
except queue.Empty:
|
||||||
continue
|
continue
|
||||||
try:
|
try:
|
||||||
httpx.post(f"{MC_URL}/api/memory/learn",
|
self._post_learn(msgs)
|
||||||
json={"messages": msgs, "source": "hermes"},
|
|
||||||
headers=_headers(), timeout=120.0)
|
|
||||||
except Exception as exc:
|
|
||||||
log.debug("mc2-memory learn POST failed: %s", exc)
|
|
||||||
finally:
|
finally:
|
||||||
self._q.task_done()
|
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 --------------------
|
# -- Context-only: keine Agent-Tools → kein Tool-Loop --------------------
|
||||||
def get_tool_schemas(self) -> List[Dict[str, Any]]:
|
def get_tool_schemas(self) -> List[Dict[str, Any]]:
|
||||||
return []
|
return []
|
||||||
|
|
||||||
|
def on_session_end(self, messages: List[Dict[str, Any]]) -> None:
|
||||||
|
self._flush()
|
||||||
|
|
||||||
def shutdown(self) -> None:
|
def shutdown(self) -> None:
|
||||||
self._stop.set()
|
self._flush()
|
||||||
|
|
||||||
|
|
||||||
def register(ctx) -> None:
|
def register(ctx) -> None:
|
||||||
|
|||||||
Reference in New Issue
Block a user