Un al doilea mesaj trimis cât Claude încă lucra aștepta până se termina turul 1 — corecția „stai, nu în master" ajungea după ce greșeala era gata. Verificat în producție înainte de commit: mesajul 2 stătea 25s blocat în lock, apoi pornea ca tur separat. Acum canalele de chat pot ține un proces `claude` viu per canal, cu stdin deschis, și al doilea mesaj intră în ACELAȘI tur. - `src/claude_runner.py` — ClaudeProcess (steering, respawn cu --resume, drenare stderr, respawn la comutarea OpenRouter) + RunnerRegistry (max_live, reaper pe inactivitate, stop_all la shutdown) - `src/stream_json.py` — parser stream-json partajat cu `_run_claude`; pur, nu aruncă niciodată pe is_error (PlanningSession retrimite pe error_max_turns și depinde de asta) - `src/sentinels.py` — un singur loc pentru __AUDIO__/__STEERED__, în loc de 4 verificări copiate; repară și bug-ul preexistent prin care WhatsApp posta literal `__AUDIO__:/cale` - dispecer în `send_message`: lock.acquire(blocking=False) — eșecul de a lua lock-ul ESTE „rulează un tur", ceea ce elimină flagul inflight din decizie și cursa TOCTOU odată cu el - `/stop` oprește turul, nu sesiunea — active.json rămâne valid - rate limit prin proces persistent vine ca result.is_error, nu ca exit code; convertit înapoi în același RuntimeError, altfel fallback-ul local nu s-ar mai declanșa niciodată, în tăcere Steering-ul nu face niciodată cross-adapter (un mesaj text nu intră într-un tur voice: împart același channel_id). Mesajele steered dintr-un tur care pică sunt re-livrate, nu pierdute. Testat live cu CLI-ul real: corecție la secunda 10 dintr-un tur de 24s, un singur result, num_turns=2. Notă: mesajele steered sunt împachetate în [EXTERNAL CONTENT], deci o corecție formulată ca override agresiv poate fi refuzată ca prompt injection — pentru oprire folosește /stop. Suită: 1199 passed, 12 failed (toate pre-existente pe HEAD curat). Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01SiJGsZVSEGjRHZEJiXaxCC
549 lines
22 KiB
Python
549 lines
22 KiB
Python
"""Persistent Claude CLI processes ("steering") for Echo-Core.
|
|
|
|
Why two paths exist
|
|
--------------------
|
|
Echo-Core has ALWAYS run the Claude CLI one-shot: `_run_claude` in
|
|
`claude_session.py` spawns `claude -p "<prompt>"`, reads the response, and
|
|
the process exits. That is simple and fine for a request/response turn, but
|
|
it means a second message sent while turn 1 is still running can't reach
|
|
Claude until turn 1 finishes — by then it's too late to say "wait, don't do
|
|
that".
|
|
|
|
`ClaudeProcess` below keeps one `claude` subprocess ALIVE per channel, with
|
|
stdin held open, so a second message can be written into the SAME turn
|
|
(`steer()`) instead of queueing behind it. `RunnerRegistry` owns the set of
|
|
live processes: a cap (`max_live`), an idle reaper, and shutdown.
|
|
|
|
`heartbeat.py` and `planning_session.py` deliberately stay on the one-shot
|
|
path (`_run_claude` / `start_session` / `resume_session`, unchanged) — there
|
|
is no human waiting to correct a 3am cron job or a planning conversation
|
|
mid-turn, so a live process there would only add RAM (measured 292-541 MB
|
|
per process) for a capability nobody uses. Only interactive chat channels
|
|
(Discord/Telegram/WhatsApp) route through this module, and only when
|
|
`steering.enabled` is on (that flag and the dispatch decision live in
|
|
`claude_session.py`/`router.py` — this module doesn't read config itself).
|
|
|
|
Manual repro recipe (T15): to exercise steering by hand, ask Echo for
|
|
something with a 30s+ `sleep` in Bash ("run `sleep 40 && echo done` then
|
|
tell me the weather"), then send a second message on the same channel
|
|
while it's still running. Watch for the "steered N chars" log line below.
|
|
|
|
Threading model
|
|
----------------
|
|
Thread-based, NOT asyncio: `send_message` (in `claude_session.py`) is sync
|
|
code called via `asyncio.to_thread` from the async adapters, so a second
|
|
event loop here would be unused complexity. Three kinds of threads touch a
|
|
single `ClaudeProcess`: the caller's own thread running `run_turn()`
|
|
(blocking read of stdout), a caller's thread calling `steer()` concurrently,
|
|
and the shared reaper thread. `_stdin_lock` is the single lock serializing
|
|
all of them around `self.inflight` and stdin writes — see `steer()`'s
|
|
docstring for the race it closes and the one it deliberately does not (and
|
|
cannot: it is an inter-process race, not a Python one).
|
|
"""
|
|
|
|
import atexit
|
|
import dataclasses
|
|
import enum
|
|
import json
|
|
import logging
|
|
import shutil
|
|
import subprocess
|
|
import threading
|
|
import time
|
|
from collections import deque
|
|
from pathlib import Path
|
|
from typing import Callable
|
|
|
|
from src.claude_session import (
|
|
CLAUDE_BIN,
|
|
DEFAULT_MODEL,
|
|
DEFAULT_TIMEOUT,
|
|
PROJECT_ROOT,
|
|
_safe_env,
|
|
build_system_prompt,
|
|
is_rate_limit_error,
|
|
)
|
|
from src.stream_json import consume_stream
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Constants
|
|
# ---------------------------------------------------------------------------
|
|
|
|
MAX_LIVE_DEFAULT = 2 # measured 292-541 MB RSS per live `claude` process (8 GB host)
|
|
IDLE_MINUTES_DEFAULT = 20
|
|
|
|
# A2: `_safe_env()` checks this at spawn to decide whether to route through
|
|
# OpenRouter. On a one-shot process that's re-evaluated every turn; on a
|
|
# persistent one it would freeze at spawn time. Compared against its state
|
|
# each turn (see `_respawn_if_needed`) — the ONLY justified respawn-on-change
|
|
# (respawning on personality/ edits was explicitly rejected, see plan
|
|
# "Riscul #2").
|
|
OPENROUTER_SEMAPHORE = PROJECT_ROOT / ".use_openrouter"
|
|
|
|
|
|
def _wrap_external_content(text: str) -> str:
|
|
"""Same injection-protection wrapping `start_session`/`resume_session`
|
|
use — T1: a steered message MUST NOT bypass it just because it goes in
|
|
over stdin instead of argv (Section 3 finding S1, a security regression
|
|
otherwise)."""
|
|
return f"[EXTERNAL CONTENT]\n{text}\n[END EXTERNAL CONTENT]"
|
|
|
|
|
|
def _user_turn_line(text: str) -> str:
|
|
wrapped = _wrap_external_content(text)
|
|
return json.dumps(
|
|
{"type": "user", "message": {"role": "user", "content": wrapped}},
|
|
ensure_ascii=False,
|
|
)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# steer() outcome — X2: enum, not bool, so the fallback logic can't be
|
|
# skipped by accident.
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
class SteerStatus(enum.Enum):
|
|
STEERED = "steered"
|
|
# DANGER, not the safe case (C1): no turn was in flight when the write
|
|
# landed, so the CLI treated it as a brand-new turn N+1 whose stdout
|
|
# nobody else is reading. `SteerResult.turn` already holds that turn's
|
|
# full result — decision D4 ("consume the stream you just started").
|
|
RAN_AS_TURN = "ran_as_turn"
|
|
# T7: the stdin write itself raised (process was dead/dying). The text
|
|
# was re-dispatched as a fresh `run_turn()` (respawning if needed) —
|
|
# `SteerResult.turn` holds ITS result. Never lost, per CLAUDE.md's
|
|
# "turnul nu se pierde niciodată".
|
|
PROCESS_DEAD = "process_dead"
|
|
|
|
|
|
@dataclasses.dataclass
|
|
class SteerResult:
|
|
"""Return value of `ClaudeProcess.steer()`.
|
|
|
|
`turn` is populated for every status except STEERED, and IS the
|
|
response for that text — bundling it here (instead of a bare enum)
|
|
makes it structurally impossible for a caller to see a non-STEERED
|
|
status and call `run_turn(text)` again "just in case": that would
|
|
double-send the message and bill it twice.
|
|
"""
|
|
|
|
status: SteerStatus
|
|
turn: dict | None = None
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# ClaudeProcess
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
class ClaudeProcess:
|
|
"""One live `claude` subprocess for a single channel."""
|
|
|
|
def __init__(self, channel_id: str, model: str = DEFAULT_MODEL,
|
|
session_id: str | None = None, cwd: Path | str | None = None):
|
|
self.channel_id = channel_id
|
|
self.model = model
|
|
self.session_id = session_id
|
|
self.cwd = cwd or PROJECT_ROOT
|
|
self.proc: subprocess.Popen | None = None
|
|
self.inflight = False
|
|
self.last_active = time.monotonic()
|
|
|
|
self._stdin_lock = threading.Lock()
|
|
self._lifecycle_lock = threading.Lock()
|
|
self._pending_steers: list[str] = []
|
|
self._stderr_buf: deque[str] = deque(maxlen=50)
|
|
self._openrouter_at_spawn = False
|
|
|
|
# -- lifecycle ---------------------------------------------------------
|
|
|
|
def alive(self) -> bool:
|
|
return self.proc is not None and self.proc.poll() is None
|
|
|
|
def _build_cmd(self) -> list[str]:
|
|
cmd = [
|
|
CLAUDE_BIN, "-p",
|
|
"--input-format", "stream-json",
|
|
"--output-format", "stream-json", "--verbose",
|
|
"--model", self.model,
|
|
"--system-prompt", build_system_prompt(),
|
|
"--dangerously-skip-permissions",
|
|
"--autocompact", "auto",
|
|
]
|
|
if self.session_id:
|
|
cmd += ["--resume", self.session_id]
|
|
return cmd
|
|
|
|
def _spawn(self) -> None:
|
|
if not shutil.which(CLAUDE_BIN):
|
|
raise FileNotFoundError(
|
|
"Claude CLI not found. "
|
|
"Install: https://docs.anthropic.com/en/docs/claude-code"
|
|
)
|
|
self._openrouter_at_spawn = OPENROUTER_SEMAPHORE.exists()
|
|
self.proc = subprocess.Popen(
|
|
self._build_cmd(),
|
|
stdin=subprocess.PIPE,
|
|
stdout=subprocess.PIPE,
|
|
stderr=subprocess.PIPE,
|
|
text=True,
|
|
bufsize=1, # line-buffered: a steer() write must reach the CLI promptly
|
|
env=_safe_env(),
|
|
cwd=str(self.cwd),
|
|
)
|
|
self.inflight = False
|
|
self.last_active = time.monotonic()
|
|
self._stderr_buf.clear()
|
|
threading.Thread(target=self._drain_stderr, args=(self.proc,), daemon=True).start()
|
|
logger.info("channel=%s: spawned claude pid=%s resume=%s",
|
|
self.channel_id, self.proc.pid, bool(self.session_id))
|
|
|
|
def _drain_stderr(self, proc: subprocess.Popen) -> None:
|
|
"""T3: MANDATORY. An undrained stderr pipe fills (~64KB) and
|
|
deadlocks the child mid-turn with no exception on our side —
|
|
this thread is the only thing preventing that."""
|
|
try:
|
|
for line in proc.stderr:
|
|
self._stderr_buf.append(line.rstrip("\n"))
|
|
except (ValueError, OSError):
|
|
pass # pipe torn down under us — process is exiting, nothing left to drain
|
|
|
|
def _respawn_if_needed(self) -> None:
|
|
"""Call with `_lifecycle_lock` held."""
|
|
if self.alive() and OPENROUTER_SEMAPHORE.exists() != self._openrouter_at_spawn:
|
|
logger.info("channel=%s: .use_openrouter changed since spawn — respawning (A2)",
|
|
self.channel_id)
|
|
self.stop()
|
|
if not self.alive():
|
|
self._spawn()
|
|
|
|
def stop(self) -> None:
|
|
"""Terminate the process. Safe to call repeatedly / when already dead.
|
|
|
|
Takes `_stdin_lock` (E-C2): a concurrent `steer()` must never write
|
|
into a pipe we're closing underneath it.
|
|
"""
|
|
with self._stdin_lock:
|
|
proc = self.proc
|
|
if proc is None:
|
|
return
|
|
try:
|
|
proc.stdin.close()
|
|
except (BrokenPipeError, OSError, ValueError):
|
|
pass
|
|
try:
|
|
proc.terminate()
|
|
proc.wait(timeout=3)
|
|
except subprocess.TimeoutExpired:
|
|
try:
|
|
proc.kill()
|
|
proc.wait(timeout=3)
|
|
except OSError:
|
|
pass
|
|
except OSError:
|
|
pass
|
|
logger.info("channel=%s: process stopped", self.channel_id)
|
|
|
|
# -- turns ---------------------------------------------------------
|
|
|
|
def run_turn(self, text: str, on_text: Callable[[str], None] | None = None,
|
|
timeout: int = DEFAULT_TIMEOUT) -> dict:
|
|
"""Run one turn. Returns exactly the dict `_run_claude` returns, and
|
|
raises the exact same exceptions for the same reasons
|
|
(FileNotFoundError, TimeoutError, `RuntimeError("Claude CLI error
|
|
(exit 1): ...")` on a rate limit) so `router.py`'s existing handling
|
|
— including the local-fallback trigger — keeps working unchanged.
|
|
"""
|
|
with self._lifecycle_lock:
|
|
self._respawn_if_needed()
|
|
|
|
line = _user_turn_line(text)
|
|
with self._stdin_lock:
|
|
self.inflight = True
|
|
self._pending_steers = []
|
|
try:
|
|
self.proc.stdin.write(line + "\n")
|
|
self.proc.stdin.flush()
|
|
except (BrokenPipeError, ValueError, OSError) as exc:
|
|
self.inflight = False
|
|
raise RuntimeError(f"Claude CLI error: failed to start turn: {exc}") from exc
|
|
|
|
return self._read_until_result(on_text=on_text, timeout=timeout)
|
|
|
|
def steer(self, text: str) -> SteerResult:
|
|
"""Write *text* into the current in-flight turn's stdin, without
|
|
opening a new turn — the common case.
|
|
|
|
C1/E-C1: the only race a lock can close is between OUR threads —
|
|
`steer()` checks `self.inflight` under the same `_stdin_lock` that
|
|
`run_turn()`/`_read_until_result()` flips it under, so nothing here
|
|
can observe a stale `True`. What no lock can close is the CLI's own
|
|
internal turn boundary: it may have already finished turn N and
|
|
gone idle before our write physically reaches its stdin pipe, in
|
|
which case the CLI treats the write as opening turn N+1 regardless
|
|
of what we believed. That is `RAN_AS_TURN` — see decision D4 on
|
|
`SteerStatus` above. (A permanent stdout-owning reader thread would
|
|
close this residual race too; deferred as a documented upgrade, not
|
|
required for the accepted fix.)
|
|
"""
|
|
line = _user_turn_line(text)
|
|
wrote = False
|
|
became_new_turn = False
|
|
with self._stdin_lock:
|
|
was_inflight = self.inflight
|
|
if self.proc is not None:
|
|
try:
|
|
self.proc.stdin.write(line + "\n")
|
|
self.proc.stdin.flush()
|
|
wrote = True
|
|
except (BrokenPipeError, ValueError, OSError):
|
|
wrote = False
|
|
if wrote and was_inflight:
|
|
self._pending_steers.append(text)
|
|
elif wrote:
|
|
became_new_turn = True
|
|
self.inflight = True # we now own reading this orphaned turn
|
|
self._pending_steers = []
|
|
|
|
if wrote and was_inflight:
|
|
logger.info("channel=%s: steered %d chars into in-flight turn",
|
|
self.channel_id, len(text))
|
|
return SteerResult(SteerStatus.STEERED)
|
|
|
|
if became_new_turn:
|
|
logger.info(
|
|
"channel=%s: steer() found no in-flight turn — the write became "
|
|
"turn N+1, consuming its stream now (D4)", self.channel_id,
|
|
)
|
|
turn = self._read_until_result(on_text=None, timeout=DEFAULT_TIMEOUT)
|
|
return SteerResult(SteerStatus.RAN_AS_TURN, turn=turn)
|
|
|
|
# T7: the write itself failed — process dead/dying. Never lose the
|
|
# message: re-dispatch as a fresh (respawning) turn.
|
|
logger.info("channel=%s: steer() hit a dead process — falling back to run_turn (T7)",
|
|
self.channel_id)
|
|
turn = self.run_turn(text)
|
|
return SteerResult(SteerStatus.PROCESS_DEAD, turn=turn)
|
|
|
|
def pop_pending_steers(self) -> list[str]:
|
|
"""C3: texts steered into the turn that just failed (rate limit,
|
|
timeout, crash) — the caller MUST re-dispatch them ("turnul nu se
|
|
pierde niciodată"). Empty once the turn they belonged to succeeds."""
|
|
pending, self._pending_steers = self._pending_steers, []
|
|
return pending
|
|
|
|
def _read_until_result(self, on_text: Callable[[str], None] | None,
|
|
timeout: int) -> dict:
|
|
proc = self.proc
|
|
timed_out = threading.Event()
|
|
|
|
def _watchdog():
|
|
try:
|
|
proc.wait(timeout=timeout)
|
|
except subprocess.TimeoutExpired:
|
|
timed_out.set()
|
|
try:
|
|
proc.kill() # M2: next run_turn() respawns with --resume
|
|
except OSError:
|
|
pass
|
|
|
|
watchdog = threading.Thread(target=_watchdog, daemon=True)
|
|
watchdog.start()
|
|
|
|
def _on_init(sid: str) -> None:
|
|
self.session_id = sid
|
|
|
|
try:
|
|
result = consume_stream(proc.stdout, on_text=on_text, on_init=_on_init)
|
|
finally:
|
|
with self._stdin_lock:
|
|
self.inflight = False
|
|
self.last_active = time.monotonic()
|
|
|
|
if timed_out.is_set():
|
|
raise TimeoutError(f"Claude CLI timed out after {timeout}s")
|
|
|
|
if result is None:
|
|
stderr_tail = "\n".join(self._stderr_buf)[-500:]
|
|
raise RuntimeError(
|
|
f"Claude CLI error: no result line in stream. stderr: {stderr_tail}"
|
|
)
|
|
|
|
if result.get("session_id"):
|
|
# E-T4: refresh every turn, not just the first — a respawn after
|
|
# turn 5 must --resume turn 5's session, not turn 1's.
|
|
self.session_id = result["session_id"]
|
|
|
|
if result.get("is_error"):
|
|
detail = result.get("result", "")
|
|
if is_rate_limit_error(detail):
|
|
# T2/C4: persistent process surfaces a rate limit as
|
|
# `result.is_error`, never a nonzero exit code. Must raise
|
|
# this EXACT phrasing or `is_rate_limit_error()` in
|
|
# router.py/scheduler.py won't recognize it and the local
|
|
# fallback never fires.
|
|
raise RuntimeError(f"Claude CLI error (exit 1): {detail}")
|
|
# any other is_error: return it, don't raise (C4 — parser and
|
|
# this wrapper both stay non-throwing for non-rate-limit errors)
|
|
|
|
self._pending_steers = []
|
|
return result
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# RunnerRegistry
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
class RunnerRegistry:
|
|
"""Owns the set of live `ClaudeProcess` instances, one per channel.
|
|
|
|
Module-level global state (like `claude_session._session_locks`) — see
|
|
`get_registry()` below. X6: no LRU eviction; at `max_live` capacity a
|
|
new channel simply degrades to the one-shot path (T10/M1) rather than
|
|
evicting someone else's live process.
|
|
"""
|
|
|
|
def __init__(self, max_live: int = MAX_LIVE_DEFAULT,
|
|
idle_minutes: int = IDLE_MINUTES_DEFAULT,
|
|
reap_interval_s: float = 60.0):
|
|
if max_live <= 0:
|
|
# X7: don't fail silently — a mis-set config value should be
|
|
# loud in the logs, even though the *behavior* (always degrade)
|
|
# is already correct without special-casing it below.
|
|
logger.warning(
|
|
"steering.max_live=%d <= 0 — steering effectively disabled, "
|
|
"every channel degrades to one-shot", max_live,
|
|
)
|
|
self.max_live = max_live
|
|
self.idle_minutes = max(1, idle_minutes) # X7: clamp, never "reap instantly"
|
|
self._procs: dict[str, ClaudeProcess] = {}
|
|
self._lock = threading.Lock()
|
|
self._reap_interval_s = reap_interval_s
|
|
self._stopping = threading.Event()
|
|
self._reaper = threading.Thread(target=self._reap_loop, daemon=True)
|
|
self._reaper.start()
|
|
atexit.register(self.stop_all)
|
|
|
|
def get(self, channel_id: str, model: str = DEFAULT_MODEL,
|
|
session_id: str | None = None, cwd: Path | str | None = None) -> ClaudeProcess | None:
|
|
"""Return the channel's live process, creating one if under
|
|
`max_live`. Returns None (T10) when at capacity and this channel
|
|
doesn't already have one — the caller must degrade to one-shot.
|
|
|
|
An existing entry for *channel_id* is always returned as-is,
|
|
regardless of capacity — which is also what rules out the M1
|
|
"two writers on one session_id" scenario: a channel already
|
|
holding a live process can never be told to fall back to one-shot
|
|
by this method.
|
|
"""
|
|
with self._lock:
|
|
proc = self._procs.get(channel_id)
|
|
if proc is not None:
|
|
return proc
|
|
if len(self._procs) >= self.max_live:
|
|
logger.info(
|
|
"channel=%s: max_live=%d reached, degrading to one-shot",
|
|
channel_id, self.max_live,
|
|
)
|
|
return None
|
|
proc = ClaudeProcess(channel_id, model=model, session_id=session_id, cwd=cwd)
|
|
self._procs[channel_id] = proc
|
|
return proc
|
|
|
|
def stop(self, channel_id: str) -> bool:
|
|
with self._lock:
|
|
proc = self._procs.pop(channel_id, None)
|
|
if proc is None:
|
|
return False
|
|
proc.stop()
|
|
return True
|
|
|
|
def stop_all(self) -> None:
|
|
"""Shutdown path (T4/A1) — zero orphaned `claude` processes after a
|
|
`systemctl restart`. Registered with `atexit` in `__init__`; a
|
|
caller in `main.py` should also invoke this explicitly on SIGTERM
|
|
for a synchronous shutdown (atexit is the fallback, not the primary
|
|
path)."""
|
|
self._stopping.set()
|
|
with self._lock:
|
|
procs = list(self._procs.values())
|
|
self._procs.clear()
|
|
for proc in procs:
|
|
try:
|
|
proc.stop()
|
|
except Exception:
|
|
logger.exception("channel=%s: error stopping during stop_all", proc.channel_id)
|
|
|
|
def live_count(self) -> int:
|
|
with self._lock:
|
|
return len(self._procs)
|
|
|
|
# -- reaper ---------------------------------------------------------
|
|
|
|
def _reap_loop(self) -> None:
|
|
while not self._stopping.wait(self._reap_interval_s):
|
|
self._reap_once()
|
|
|
|
def _reap_once(self) -> None:
|
|
"""T11: resilient to any single process's `stop()` raising — the
|
|
loop (and the rest of the batch) must keep going, and every reap
|
|
must be logged (a silently-dead reaper is an invisible RAM leak)."""
|
|
cutoff = time.monotonic() - self.idle_minutes * 60
|
|
with self._lock:
|
|
candidates = [
|
|
cid for cid, p in self._procs.items()
|
|
if not p.inflight and p.last_active < cutoff
|
|
]
|
|
for cid in candidates:
|
|
try:
|
|
with self._lock:
|
|
proc = self._procs.get(cid)
|
|
# Re-check under lock: it may have gone inflight, been
|
|
# stopped, or been replaced since the snapshot above.
|
|
if proc is None or proc.inflight or proc.last_active >= cutoff:
|
|
continue
|
|
del self._procs[cid]
|
|
proc.stop()
|
|
logger.info("reaper: stopped idle channel=%s (idle >= %d min)",
|
|
cid, self.idle_minutes)
|
|
except Exception:
|
|
logger.exception("reaper: error stopping channel=%s — continuing", cid)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Module-level singleton — mirrors `claude_session._session_locks`.
|
|
# ---------------------------------------------------------------------------
|
|
|
|
_registry: RunnerRegistry | None = None
|
|
_registry_lock = threading.Lock()
|
|
|
|
|
|
def get_registry(max_live: int = MAX_LIVE_DEFAULT,
|
|
idle_minutes: int = IDLE_MINUTES_DEFAULT) -> RunnerRegistry:
|
|
"""Lazy singleton. *max_live*/*idle_minutes* only apply on first call —
|
|
later calls just return the existing registry (config is read by
|
|
`router.py` per call and passed in here; this function doesn't re-read
|
|
config.json itself)."""
|
|
global _registry
|
|
if _registry is not None:
|
|
return _registry
|
|
with _registry_lock:
|
|
if _registry is None:
|
|
_registry = RunnerRegistry(max_live=max_live, idle_minutes=idle_minutes)
|
|
return _registry
|
|
|
|
|
|
def reset_registry_for_tests() -> None:
|
|
"""H2: tests must not leak live (fake) processes into each other. Stops
|
|
and drops the current singleton so the next `get_registry()` call
|
|
builds a fresh one."""
|
|
global _registry
|
|
with _registry_lock:
|
|
if _registry is not None:
|
|
_registry.stop_all()
|
|
_registry = None
|