Skip to content
Closed
15 changes: 15 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -239,3 +239,18 @@ outputs
evaluation/data/
test_add_pipeline.py
test_file_pipeline.py
data/
config.yaml

# ── Runtime / ephemeral (Violet addition 2026-05-27) ──
daemon/bridge.pid
memos-local/
bridge-status.json
apps/memos-local-plugin/daemon/
apps/memos-local-plugin/memos-local/
apps/memos-local-plugin/.memos-node-bin
apps/memos-local-plugin/bridge-status.json
apps/memos-local-plugin/prod_check.cjs
data.stale-backup/
scripts/trigger-scoring.py
data
113 changes: 84 additions & 29 deletions apps/memos-local-plugin/adapters/hermes/memos_provider/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -63,8 +63,8 @@
if str(_PLUGIN_DIR) not in sys.path:
sys.path.insert(0, str(_PLUGIN_DIR))

from bridge_client import BridgeError, MemosBridgeClient # noqa: E402
from daemon_manager import ensure_bridge_running, ensure_viewer_daemon # noqa: E402
from bridge_client import BridgeError, MemosBridgeClient, MemosHttpClient # noqa: E402
from daemon_manager import ensure_bridge_running, ensure_viewer_daemon, probe_viewer_status, kill_zombie_bridges # noqa: E402


try: # pragma: no cover — host-provided base class, absent in unit tests
Expand Down Expand Up @@ -243,7 +243,7 @@ class MemTensorProvider(MemoryProvider):
"""

def __init__(self) -> None:
self._bridge: MemosBridgeClient | None = None
self._bridge: MemosBridgeClient | MemosHttpClient | None = None
self._reconnect_lock = threading.Lock()
self._session_id: str = ""
self._episode_id: str = ""
Expand Down Expand Up @@ -314,34 +314,70 @@ def initialize(self, session_id: str, **kwargs: Any) -> None: # type: ignore[ov
except Exception as err:
logger.warning("MemOS: failed to start bridge — %s", err)
return

# Kill zombie bridges from previous sessions before deciding
# how to connect.
try:
ensure_viewer_daemon()
except Exception as err:
logger.warning("MemOS: viewer daemon check failed — %s", err)
new_bridge: MemosBridgeClient | None = None
try:
new_bridge = MemosBridgeClient()
# Register the fallback LLM handler BEFORE we open the
# session so it is available the very first time the
# plugin's facade asks for help (e.g. on the first
# `turn.start` retrieval call).
new_bridge.register_host_handler(
"host.llm.complete",
self._handle_host_llm_complete,
)
self._bridge = new_bridge
self._open_session(session_id)
logger.info(
"MemOS: bridge ready session=%s platform=%s (episode deferred)",
self._session_id,
self._platform,
)
except Exception as err:
logger.warning("MemOS: bridge init failed — %s", err)
if new_bridge is not None:
zombies = kill_zombie_bridges()
if zombies:
logger.info("MemOS: killed %d zombie bridge(s)", zombies)
except Exception:
pass

# If the daemon is already running on the viewer port, connect
# to it over HTTP instead of spawning a new stdio bridge. This
# eliminates zombie bridge accumulation.
viewer_status = probe_viewer_status()
if viewer_status == "running_memos":
try:
http_bridge: MemosHttpClient | None = MemosHttpClient()
http_bridge.register_host_handler(
"host.llm.complete",
self._handle_host_llm_complete,
)
self._bridge = http_bridge
self._open_session(session_id)
logger.info(
"MemOS: bridge ready (HTTP) session=%s platform=%s (episode deferred)",
self._session_id,
self._platform,
)
except Exception as err:
logger.warning("MemOS: HTTP bridge failed, falling back to stdio — %s", err)
with contextlib.suppress(Exception):
new_bridge.close()
self._bridge = None
http_bridge.close() # type: ignore[union-attr]
self._bridge = None
viewer_status = "free" # force stdio fallback below

if self._bridge is None:
try:
ensure_viewer_daemon()
except Exception as err:
logger.warning("MemOS: viewer daemon check failed — %s", err)
new_bridge: MemosBridgeClient | None = None
try:
new_bridge = MemosBridgeClient()
# Register the fallback LLM handler BEFORE we open the
# session so it is available the very first time the
# plugin's facade asks for help (e.g. on the first
# `turn.start` retrieval call).
new_bridge.register_host_handler(
"host.llm.complete",
self._handle_host_llm_complete,
)
self._bridge = new_bridge
self._open_session(session_id)
logger.info(
"MemOS: bridge ready (stdio) session=%s platform=%s (episode deferred)",
self._session_id,
self._platform,
)
except Exception as err:
logger.warning("MemOS: bridge init failed — %s", err)
if new_bridge is not None:
with contextlib.suppress(Exception):
new_bridge.close()
self._bridge = None
# Register a Hermes plugin hook to capture tool calls as they
# happen. The `post_tool_call` hook fires after every tool
# dispatch (write_file, terminal, search_files, etc.) with the
Expand Down Expand Up @@ -1746,6 +1782,25 @@ def _reconnect_bridge(self, session_id: str = "", *, timeout: float = 30.0) -> N
logger.info("MemOS: old bridge closed (pid=%s)", old_pid)

ensure_bridge_running()
# Try HTTP first if daemon is running
viewer_status = probe_viewer_status()
if viewer_status == "running_memos":
try:
http_bridge: MemosHttpClient | None = MemosHttpClient()
http_bridge.register_host_handler(
"host.llm.complete",
self._handle_host_llm_complete,
)
self._bridge = http_bridge
self._open_session(session_id, timeout=timeout)
logger.info("MemOS: reconnected via HTTP")
return
except Exception as err:
logger.warning("MemOS: HTTP reconnect failed, falling back to stdio — %s", err)
with contextlib.suppress(Exception):
http_bridge.close() # type: ignore[union-attr]
self._bridge = None

try:
ensure_viewer_daemon()
except Exception as err:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,9 @@
import shutil
import subprocess
import threading
import time
import urllib.error
import urllib.request

from pathlib import Path
from typing import TYPE_CHECKING, Any
Expand Down Expand Up @@ -419,4 +422,123 @@ def _resolve(self, msg: dict[str, Any]) -> None:
entry["error"] = msg["error"]
else:
entry["result"] = msg.get("result")
entry["event"].set()


class MemosHttpClient:
"""JSON-RPC 2.0 client that talks to the daemon bridge over HTTP.

Drop-in replacement for ``MemosBridgeClient`` when a daemon is already
running on the viewer port. Instead of spawning a new subprocess, this
client POSTs JSON-RPC envelopes to ``/api/v1/rpc`` on the daemon's HTTP
server. This eliminates zombie bridge accumulation.

Limitations vs. stdio client:
- No reverse-direction RPC (``host.llm.complete``). The daemon's own
stdio bridge handles host LLM fallback internally.
- No ``notify()`` (notifications have no response; use ``request()``
for all calls — the daemon handles both).
"""

def __init__(
self,
*,
port: int = 18800,
host: str = "127.0.0.1",
api_key: str | None = None,
) -> None:
self._base_url = f"http://{host}:{port}/api/v1/rpc"
self._lock = threading.Lock()
self._next_id = 1
self._closed = False
self._api_key = api_key
self._host_handlers: dict[str, Callable[[dict[str, Any]], Any]] = {}

@property
def pid(self) -> int:
"""Return 0 — HTTP client has no subprocess."""
return 0

# ─── Public API (matches MemosBridgeClient) ──

def request(
self,
method: str,
params: Any = None,
*,
timeout: float = 30.0,
) -> dict[str, Any]:
if self._closed:
raise BridgeError("transport_closed", "HTTP client is closed")
with self._lock:
rpc_id = self._next_id
self._next_id += 1

envelope = {
"jsonrpc": "2.0",
"id": rpc_id,
"method": method,
"params": params,
}
payload = json.dumps(envelope, ensure_ascii=False).encode("utf-8")

req = urllib.request.Request(
self._base_url,
data=payload,
headers={"Content-Type": "application/json"},
method="POST",
)
if self._api_key:
req.add_header("Authorization", f"Bearer {self._api_key}")

try:
with urllib.request.urlopen(req, timeout=timeout) as resp:
body = resp.read(4_194_304) # 4 MiB safety cap
except urllib.error.HTTPError as exc:
# Try to read the JSON-RPC error body
try:
err_body = exc.read(4_194_304)
err_json = json.loads(err_body.decode("utf-8", errors="replace"))
err_obj = err_json.get("error") or {}
raise BridgeError(
err_obj.get("data", {}).get("code") or str(err_obj.get("code", "internal")),
err_obj.get("message", f"HTTP {exc.code}"),
err_obj.get("data"),
) from exc
except (json.JSONDecodeError, UnicodeDecodeError):
raise BridgeError("internal", f"HTTP {exc.code}: {exc.reason}") from exc
except urllib.error.URLError as exc:
raise BridgeError("transport_closed", str(exc)) from exc

result = json.loads(body.decode("utf-8", errors="replace"))
if "error" in result:
e = result["error"]
raise BridgeError(
(e.get("data") or {}).get("code") or str(e.get("code", "internal")),
e.get("message", "unknown error"),
e.get("data"),
)
return result.get("result") or {}

def notify(self, method: str, params: Any = None) -> None:
"""No-op for HTTP — use request() for all calls."""
try:
self.request(method, params, timeout=5.0)
except Exception:
pass

def on_event(self, cb: Callable[[dict[str, Any]], None]) -> None:
"""No-op — SSE events not supported over HTTP transport."""

def on_log(self, cb: Callable[[dict[str, Any]], None]) -> None:
"""No-op — SSE logs not supported over HTTP transport."""

def register_host_handler(
self,
method: str,
handler: Callable[[dict[str, Any]], Any],
) -> None:
"""Store the handler but it won't be called (daemon handles host LLM internally)."""
self._host_handlers[method] = handler

def close(self) -> None:
self._closed = True
Original file line number Diff line number Diff line change
Expand Up @@ -336,6 +336,77 @@ def shutdown_bridge() -> None:
_bridge_ok = None


def probe_viewer_status() -> str:
"""Return the current viewer daemon status without side effects.

Returns one of: ``"running_memos"``, ``"free"``, ``"blocked"``.
This is a cheap, lock-free probe suitable for deciding whether to
spawn a new stdio bridge or connect to the existing daemon over HTTP.
"""
return _probe_viewer()


def kill_zombie_bridges() -> int:
"""Kill all bridge.cjs processes that are NOT the daemon on port 18800.

Returns the number of zombies killed. The daemon (the process that
owns port 18800) is left alone. This should be called early in the
provider's lifecycle to clean up leftovers from crashed sessions.
"""
import subprocess as sp

# Find the PID that owns port 18800 (the real daemon)
daemon_pid: int | None = None
try:
ss_out = sp.check_output(
["ss", "-tlnp"],
timeout=2.0,
text=True,
)
for line in ss_out.splitlines():
if ":18800" in line:
# ss output: users:(("node",pid=21246,fd=24))
import re
m = re.search(r"pid=(\d+)", line)
if m:
daemon_pid = int(m.group(1))
break
except Exception:
pass

# Find all bridge.cjs processes
killed = 0
try:
ps_out = sp.check_output(
["ps", "aux"],
timeout=2.0,
text=True,
)
for line in ps_out.splitlines():
if "bridge.cjs" not in line or "grep" in line:
continue
parts = line.split()
if len(parts) < 2:
continue
try:
pid = int(parts[1])
except ValueError:
continue
if daemon_pid is not None and pid == daemon_pid:
continue
# This is a zombie — kill it
try:
os.kill(pid, signal.SIGTERM)
killed += 1
logger.info("MemOS: killed zombie bridge pid=%d", pid)
except (OSError, ProcessLookupError):
pass
except Exception as err:
logger.debug("MemOS: zombie scan failed — %s", err)

return killed


def wait_for_process_exit(pid: int, timeout: float = 5.0) -> bool:
"""Wait for a process to exit.

Expand Down
7 changes: 5 additions & 2 deletions apps/memos-local-plugin/server/routes/auth.ts
Original file line number Diff line number Diff line change
Expand Up @@ -410,10 +410,13 @@ export function requireSession(
agent?: string | null,
): boolean {
// Public: auth endpoints + health (so the viewer can tell whether
// the backend is up BEFORE unlocking). Other API routes, including
// ping, fall through and are only open when no password is configured.
// the backend is up BEFORE unlocking). RPC endpoint is exempt
// because it's used by the local Python adapter (same machine, no
// browser session). Other API routes, including ping, fall through
// and are only open when no password is configured.
if (pathname.startsWith("/api/v1/auth/")) return true;
if (pathname === "/api/v1/health") return true;
if (pathname === "/api/v1/rpc") return true;

const state = readAuthState(homeDir);
if (!state) return true; // password protection off → open
Expand Down
Loading