"""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 ""`, 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