From ac3f9a2dc4fe20fc24c1cfbdabf500a9d2b41ef3 Mon Sep 17 00:00:00 2001 From: kshitijk4poor <82637225+kshitijk4poor@users.noreply.github.com> Date: Mon, 24 Aug 2026 00:27:57 +0530 Subject: [PATCH] =?UTF-8?q?feat(bot-mode):=20per-profile=20turn=20lock=20?= =?UTF-8?q?=E2=80=94=20concurrent=20deliveries=20queue=20instead=20of=20ra?= =?UTF-8?q?cing=20(#93091)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- hermes_cli/config_defaults.py | 11 ++ tests/tools/test_bot_turn_lock.py | 304 ++++++++++++++++++++++++++++++ tools/bot_mode_dm.py | 52 ++++- tools/bot_relay.py | 106 ++++++++++- tui_gateway/methods_bot_relay.py | 21 ++- 5 files changed, 476 insertions(+), 18 deletions(-) create mode 100644 tests/tools/test_bot_turn_lock.py diff --git a/hermes_cli/config_defaults.py b/hermes_cli/config_defaults.py index a84cb5a93b..845c0f5388 100644 --- a/hermes_cli/config_defaults.py +++ b/hermes_cli/config_defaults.py @@ -2317,6 +2317,17 @@ DEFAULT_CONFIG = { }, }, + # Bot Mode (desktop-managed multi-agent) settings. + # Section defined here for the per-profile turn lock (#93091); other + # Bot Mode knobs may merge into this section from parallel work. + "bot_mode": { + # How long a second delivery into an already-busy target profile + # queues behind the current turn before failing with a structured + # 'target_busy' error. Deliveries are serialized per profile with a + # cross-process file lock so two turns never race one Bot Chat. + "turn_wait_seconds": 120, + }, + # Mattermost platform settings (gateway mode) "mattermost": { "require_mention": True, # Require @mention to respond in channels diff --git a/tests/tools/test_bot_turn_lock.py b/tests/tools/test_bot_turn_lock.py new file mode 100644 index 0000000000..bbf0fa437a --- /dev/null +++ b/tests/tools/test_bot_turn_lock.py @@ -0,0 +1,304 @@ +"""Tests: per-profile bot turn lock (#93091 — tools/bot_relay.py). + +Two deliveries into the same target profile must serialize on a +cross-process flock; the queued one waits a bounded budget and then fails +with a structured 'target_busy' refusal. Real flock on real (short) +tmp_path lockfiles — flock contends between separate fds even within one +process, so threads exercise the true kernel-lock semantics. +""" + +from __future__ import annotations + +import fcntl +import json +import os +import threading +import time + +import pytest + +from tools import bot_mode_dm, bot_relay +from tools.bot_relay import TurnBusyError, acquire_turn_lock, turn_lock_path + + +@pytest.fixture +def root(tmp_path): + # Keep the lockfile path SHORT (macOS-safe). + r = tmp_path / "r" + r.mkdir() + return r + + +def _hold_flock(path, hold_event, release_event): + """Grab the profile lock on a separate fd, signal, hold until told.""" + path.parent.mkdir(parents=True, exist_ok=True) + fd = os.open(path, os.O_RDWR | os.O_CREAT, 0o600) + fcntl.flock(fd, fcntl.LOCK_EX) + hold_event.set() + release_event.wait(timeout=10) + os.close(fd) # close releases the flock — process-death semantics + + +def test_second_delivery_waits_then_succeeds(root): + held = threading.Event() + release = threading.Event() + t = threading.Thread( + target=_hold_flock, args=(turn_lock_path(root, "ops"), held, release) + ) + t.start() + assert held.wait(timeout=5) + + # Release shortly after the waiter starts probing. + threading.Timer(0.3, release.set).start() + start = time.monotonic() + with acquire_turn_lock(root, "ops", timeout_seconds=5): + waited = time.monotonic() - start + t.join(timeout=5) + assert waited >= 0.2, "second delivery should have queued behind the holder" + + +def test_timeout_is_structured_target_busy(root): + held = threading.Event() + release = threading.Event() + t = threading.Thread( + target=_hold_flock, args=(turn_lock_path(root, "ops"), held, release) + ) + t.start() + assert held.wait(timeout=5) + try: + with pytest.raises(TurnBusyError) as excinfo: + with acquire_turn_lock(root, "ops", timeout_seconds=0.3): + pass # pragma: no cover — must not acquire + err = excinfo.value + assert err.reason == "target_busy" + assert err.profile == "ops" + assert err.waited_seconds >= 0.3 + assert "target_busy" in str(err) + assert "0s" in str(err) or "1s" in str(err) # rough wait duration surfaced + finally: + release.set() + t.join(timeout=5) + + +def test_different_profiles_do_not_contend(root): + held = threading.Event() + release = threading.Event() + t = threading.Thread( + target=_hold_flock, args=(turn_lock_path(root, "ops"), held, release) + ) + t.start() + assert held.wait(timeout=5) + try: + start = time.monotonic() + with acquire_turn_lock(root, "scout", timeout_seconds=2): + pass + assert time.monotonic() - start < 1.0 + finally: + release.set() + t.join(timeout=5) + + +def test_lock_released_when_holder_fd_closes(root): + """flock dies with the holder's fd — a crashed turn can't wedge the profile.""" + path = turn_lock_path(root, "ops") + path.parent.mkdir(parents=True, exist_ok=True) + fd = os.open(path, os.O_RDWR | os.O_CREAT, 0o600) + fcntl.flock(fd, fcntl.LOCK_EX) + os.close(fd) # simulate holder process death (kernel releases the lock) + with acquire_turn_lock(root, "ops", timeout_seconds=0.5): + pass # acquires immediately — no TurnBusyError + + +def test_reentry_after_clean_release(root): + with acquire_turn_lock(root, "ops", timeout_seconds=1): + pass + with acquire_turn_lock(root, "ops", timeout_seconds=1): + pass + + +def test_lock_path_is_short_and_sanitized(root): + p = turn_lock_path(root, "we/ird nam√©" + "x" * 200) + assert p.parent == bot_relay.relay_root(root) / bot_relay.LOCKS_DIR + assert len(p.name) <= 70 + assert "/" not in p.name.replace(".lock", "") + + +def test_turn_wait_seconds_falls_back_to_module_constant(monkeypatch): + def _boom(): + raise RuntimeError("no config") + + monkeypatch.setattr("hermes_cli.config.load_config", _boom) + assert bot_relay.turn_wait_seconds() == float(bot_relay.TURN_WAIT_SECONDS_FALLBACK) + + +def test_turn_wait_seconds_reads_config(monkeypatch): + monkeypatch.setattr( + "hermes_cli.config.load_config", + lambda: {"bot_mode": {"turn_wait_seconds": 7}}, + ) + assert bot_relay.turn_wait_seconds() == 7.0 + + +# ── wiring: local teammate delivery (tools/bot_mode_dm.py) ────────────────── + + +def test_run_delivery_holds_profile_lock_during_turn(root, tmp_path, monkeypatch): + """The local `hermes -p ` turn runs UNDER the profile lock.""" + home = root / ".hermes" + home.mkdir() + monkeypatch.setenv("HERMES_HOME", str(home)) + dm = tmp_path / "dm.txt" + dm.write_text("hi", encoding="utf-8") + observed = {} + + def _fake_run(argv, **kwargs): + # While the turn runs, a second acquire on the same profile must fail. + with pytest.raises(TurnBusyError): + with acquire_turn_lock(home, "ops", timeout_seconds=0.15): + pass # pragma: no cover + observed["argv"] = argv + + class _P: + returncode = 0 + + return _P() + + monkeypatch.setattr(bot_mode_dm.subprocess, "run", _fake_run) + rc = bot_mode_dm._run_delivery( + ["hermes", "-p", "ops", "chat"], str(dm), stdin_file=False + ) + assert rc == 0 + assert observed["argv"][:3] == ["hermes", "-p", "ops"] + # …and after the turn, the lock is free again. + with acquire_turn_lock(home, "ops", timeout_seconds=0.5): + pass + + +def test_delivery_main_reports_target_busy_json(root, tmp_path, monkeypatch, capsys): + """A queued delivery that exceeds its budget surfaces the structured error.""" + home = root / ".hermes" + home.mkdir() + monkeypatch.setenv("HERMES_HOME", str(home)) + monkeypatch.setattr(bot_relay, "turn_wait_seconds", lambda: 0.2) + dm = tmp_path / "dm.txt" + dm.write_text("hi", encoding="utf-8") + + held = threading.Event() + release = threading.Event() + t = threading.Thread( + target=_hold_flock, args=(turn_lock_path(home, "ops"), held, release) + ) + t.start() + assert held.wait(timeout=5) + try: + rc = bot_mode_dm._delivery_main( + ["--run-delivery", "query-file", str(dm), "hermes", "-p", "ops", "chat"] + ) + assert rc == 1 + payload = json.loads(capsys.readouterr().out.strip()) + assert payload["reason"] == "target_busy" # #93091 item-1 enum extension + assert "ops" in payload["error"] + finally: + release.set() + t.join(timeout=5) + assert not dm.exists(), "DM plaintext must be reclaimed even on refusal" + + +def test_peer_stdin_delivery_skips_local_lock(root, tmp_path, monkeypatch): + """Peer transports run their turn on the remote gateway — no local lock.""" + home = root / ".hermes" + home.mkdir() + monkeypatch.setenv("HERMES_HOME", str(home)) + dm = tmp_path / "dm.txt" + dm.write_text("hi", encoding="utf-8") + + held = threading.Event() + release = threading.Event() + t = threading.Thread( + target=_hold_flock, args=(turn_lock_path(home, "ops"), held, release) + ) + t.start() + assert held.wait(timeout=5) + try: + + def _fake_run(argv, **kwargs): + class _P: + returncode = 0 + + return _P() + + monkeypatch.setattr(bot_mode_dm.subprocess, "run", _fake_run) + rc = bot_mode_dm._run_delivery( + ["hermes", "peer", "dm", "spark/ops"], str(dm), stdin_file=True + ) + assert rc == 0 # did not contend with the held 'ops' lock + finally: + release.set() + t.join(timeout=5) + + +# ── wiring: relay deliver RPC (tui_gateway/methods_bot_relay.py) ───────────── + + +def test_relay_deliver_returns_target_busy_error(tmp_path, monkeypatch): + import tui_gateway.server as srv + + h = tmp_path / "h" + (h / "profiles" / "ops").mkdir(parents=True) + monkeypatch.setenv("HERMES_HOME", str(h)) + monkeypatch.setattr(bot_relay, "turn_wait_seconds", lambda: 0.2) + + spawned = {} + + def _fake_run(argv, **kwargs): # pragma: no cover — must not be reached + spawned["argv"] = argv + + monkeypatch.setattr("subprocess.run", _fake_run) + + held = threading.Event() + release = threading.Event() + t = threading.Thread( + target=_hold_flock, args=(turn_lock_path(h, "ops"), held, release) + ) + t.start() + assert held.wait(timeout=5) + try: + out = srv._methods["bot_relay.deliver"](1, {"profile": "ops", "message": "x"}) + assert "error" in out + assert out["error"]["code"] == 5096 + assert "target_busy" in out["error"]["message"] + assert not spawned, "turn must not spawn while the profile is busy" + finally: + release.set() + t.join(timeout=5) + + +def test_relay_deliver_serializes_then_succeeds(tmp_path, monkeypatch): + import tui_gateway.server as srv + + h = tmp_path / "h" + (h / "profiles" / "ops").mkdir(parents=True) + monkeypatch.setenv("HERMES_HOME", str(h)) + monkeypatch.setattr(bot_relay, "turn_wait_seconds", lambda: 5.0) + + class _Proc: + returncode = 0 + stdout = "pong" + stderr = "" + + monkeypatch.setattr("subprocess.run", lambda *a, **k: _Proc()) + + held = threading.Event() + release = threading.Event() + t = threading.Thread( + target=_hold_flock, args=(turn_lock_path(h, "ops"), held, release) + ) + t.start() + assert held.wait(timeout=5) + threading.Timer(0.3, release.set).start() + start = time.monotonic() + out = srv._methods["bot_relay.deliver"](1, {"profile": "ops", "message": "x"}) + t.join(timeout=5) + assert "error" not in out, out + assert out["result"]["reply"] == "pong" + assert time.monotonic() - start >= 0.2, "deliver should have queued" diff --git a/tools/bot_mode_dm.py b/tools/bot_mode_dm.py index 72191dd92a..43fa53c9ed 100644 --- a/tools/bot_mode_dm.py +++ b/tools/bot_mode_dm.py @@ -40,6 +40,7 @@ the same wake shape every Bot Mode agent already knows. from __future__ import annotations +import contextlib import json import logging import os @@ -517,18 +518,41 @@ def _unlink_dm_file(path: str) -> None: pass +def _delivery_lock(argv: list[str], *, stdin_file: bool): + """Per-profile turn lock context for a LOCAL teammate delivery (#93091). + + Local deliveries (``hermes -p chat …``) collide with relay + deliveries into the same profile — both run a Bot Chat turn on this + install — so the turn window is serialized on the shared cross-process + lock in ``tools.bot_relay``. Peer transports (stdin mode) run on the + remote gateway; their turn is locked THERE by its own deliver path. + """ + if stdin_file or len(argv) < 3 or argv[0] != "hermes" or argv[1] != "-p": + return contextlib.nullcontext() + from tools.bot_relay import acquire_turn_lock + + home = Path(os.getenv("HERMES_HOME") or os.path.expanduser("~/.hermes")) + return acquire_turn_lock(_hermes_root(home), argv[2]) + + def _run_delivery(argv: list[str], dm_file: str, *, stdin_file: bool) -> int: - """Run one DM transport and remove its plaintext file after consumption.""" + """Run one DM transport and remove its plaintext file after consumption. + + The turn execution window (not the enqueue) holds the target profile's + cross-process lock, so two deliveries into one profile queue instead of + racing; a bounded wait ends in a structured 'target_busy' refusal. + """ try: - if stdin_file: - # Keep the file open until the transport exits; cleanup occurs - # after subprocess.run returns, not merely after stdin reaches EOF. - with open(dm_file, "r", encoding="utf-8") as stream: - return subprocess.run(argv, stdin=stream, check=False).returncode - return subprocess.run( - [*argv, "--query-file", dm_file], - check=False, - ).returncode + with _delivery_lock(argv, stdin_file=stdin_file): + if stdin_file: + # Keep the file open until the transport exits; cleanup occurs + # after subprocess.run returns, not merely after stdin reaches EOF. + with open(dm_file, "r", encoding="utf-8") as stream: + return subprocess.run(argv, stdin=stream, check=False).returncode + return subprocess.run( + [*argv, "--query-file", dm_file], + check=False, + ).returncode finally: _unlink_dm_file(dm_file) @@ -639,6 +663,14 @@ def _delivery_main(args: list[str]) -> int: try: return _run_delivery(args[3:], dm_file, stdin_file=stdin_file) except Exception as exc: + # 'target_busy' extends the #93091 item-1 structured refusal enum: + # the queued delivery gave up after its bounded wait — surface the + # structured payload on stdout so the completion notification carries + # it back to the sending agent. + reason = getattr(exc, "reason", "") + if reason == "target_busy": + print(json.dumps({"error": str(exc), "reason": "target_busy"})) + return 1 print( f"message_agent delivery failed: {type(exc).__name__}: {exc}", file=sys.stderr, diff --git a/tools/bot_relay.py b/tools/bot_relay.py index 4f020f8eee..25e6104951 100644 --- a/tools/bot_relay.py +++ b/tools/bot_relay.py @@ -32,6 +32,7 @@ the sender fails fast instead of queueing a DM nobody will drain (#93091). from __future__ import annotations +import contextlib import json import logging import os @@ -42,7 +43,7 @@ import tempfile import time import uuid from pathlib import Path -from typing import Any, Optional +from typing import Any, Iterator, Optional logger = logging.getLogger(__name__) @@ -51,6 +52,11 @@ ROSTER_FILE = "roster.json" OUTBOX_DIR = "outbox" CLAIMED_DIR = "claimed" REPLIES_DIR = "replies" +LOCKS_DIR = "locks" + +# Fallback wait budget for a queued delivery turn when config is unreadable. +# The real knob is ``bot_mode.turn_wait_seconds`` in config.yaml. +TURN_WAIT_SECONDS_FALLBACK = 120 # A reply must arrive before the waiter gives up. Cross-connection turns can # be slow (remote model, cold gateway) — generous, but bounded. @@ -521,3 +527,101 @@ def local_delivery_command(profile: str, query_file: str) -> list[str]: "--query-file", query_file, ] + + +# ── per-profile turn lock (#93091) ─────────────────────────────────────────── +# +# Two deliveries into the SAME target profile must never run their Bot Chat +# turns concurrently: deliveries spawn separate ``hermes`` subprocesses, so +# an in-memory mutex is useless — the lock is a per-profile lockfile under +# ``/bot_relay/locks/`` held with ``fcntl.flock`` for exactly the turn +# execution window. flock is released by the kernel when the holder's fd +# closes (including process death), so a crashed turn can never wedge the +# profile. A queued delivery waits up to ``bot_mode.turn_wait_seconds`` and +# then fails with a structured 'target_busy' refusal instead of blocking +# forever. + + +class TurnBusyError(RuntimeError): + """A delivery turn is already running for the target profile. + + ``reason`` is 'target_busy' — extends the #93091 item-1 structured + refusal enum. ``waited_seconds`` is roughly how long the caller queued + behind the current turn before giving up. + """ + + reason = "target_busy" + + def __init__(self, profile: str, waited_seconds: float): + self.profile = profile + self.waited_seconds = waited_seconds + super().__init__( + f"target_busy: another delivery turn is already running for " + f"profile '{profile}' — queued behind it for ~{int(round(waited_seconds))}s " + "without it finishing. The message was NOT delivered; retry shortly." + ) + + +def turn_wait_seconds() -> float: + """Wait budget for a queued delivery turn (config, lazily read).""" + try: + from hermes_cli.config import cfg_get, load_config + + val = cfg_get(load_config(), "bot_mode", "turn_wait_seconds", default=None) + if val is not None: + return max(0.0, float(val)) + except Exception: + logger.debug("bot_mode.turn_wait_seconds read failed", exc_info=True) + return float(TURN_WAIT_SECONDS_FALLBACK) + + +def turn_lock_path(root: Path | str, profile: str) -> Path: + """Per-profile lockfile path (short — safe on macOS temp roots).""" + safe = re.sub(r"[^a-zA-Z0-9_-]", "_", str(profile or ""))[:64] or "_" + return relay_root(root) / LOCKS_DIR / f"{safe}.lock" + + +@contextlib.contextmanager +def acquire_turn_lock( + root: Path | str, profile: str, timeout_seconds: float | None = None +) -> Iterator[Path]: + """Hold ``profile``'s cross-process turn lock for the ``with`` body. + + Non-blocking flock probe + short-sleep retry loop up to the budget + (``bot_mode.turn_wait_seconds`` unless ``timeout_seconds`` is given), so + waiters acquire roughly in arrival order without ever deadlocking. + Raises :class:`TurnBusyError` when the budget is exhausted. On platforms + without ``fcntl`` (Windows) the lock degrades to a no-op — those + installs never had this race path in production. + """ + try: + import fcntl + except ImportError: # pragma: no cover — Windows + yield turn_lock_path(root, profile) + return + + budget = turn_wait_seconds() if timeout_seconds is None else max(0.0, float(timeout_seconds)) + path = turn_lock_path(root, profile) + path.parent.mkdir(parents=True, exist_ok=True) + fd = os.open(path, os.O_RDWR | os.O_CREAT, 0o600) + try: + start = time.monotonic() + deadline = start + budget + while True: + try: + fcntl.flock(fd, fcntl.LOCK_EX | fcntl.LOCK_NB) + break + except OSError: + now = time.monotonic() + if now >= deadline: + raise TurnBusyError(profile, now - start) + time.sleep(min(0.1, max(0.005, deadline - now))) + try: + yield path + finally: + try: + fcntl.flock(fd, fcntl.LOCK_UN) + except OSError: # pragma: no cover — kernel releases on close anyway + pass + finally: + os.close(fd) diff --git a/tui_gateway/methods_bot_relay.py b/tui_gateway/methods_bot_relay.py index ae266c6731..68fb77961e 100644 --- a/tui_gateway/methods_bot_relay.py +++ b/tui_gateway/methods_bot_relay.py @@ -94,7 +94,7 @@ def _(rid, params: dict) -> dict: return _err(rid, 4090, "profile and message required") try: from tools.bot_mode_dm import MESSAGE_MAX_CHARS - from tools.bot_relay import local_delivery_command + from tools.bot_relay import acquire_turn_lock, local_delivery_command if len(message) > MESSAGE_MAX_CHARS + 200: # + attribution headroom return _err(rid, 4091, "message too long") @@ -113,12 +113,16 @@ def _(rid, params: dict) -> dict: try: with os.fdopen(fd, "w", encoding="utf-8") as f: f.write(message) - proc = subprocess.run( - local_delivery_command(resolved, tmp), - capture_output=True, - text=True, - timeout=600, - ) + # Per-profile turn lock (#93091): serialize with any other + # delivery turn into this profile (relay or local message_agent). + # The lock covers only the turn execution window. + with acquire_turn_lock(root, resolved): + proc = subprocess.run( + local_delivery_command(resolved, tmp), + capture_output=True, + text=True, + timeout=600, + ) finally: try: os.unlink(tmp) @@ -131,6 +135,9 @@ def _(rid, params: dict) -> dict: except subprocess.TimeoutExpired: return _err(rid, 5093, "delivery turn timed out") except Exception as e: + # 'target_busy' extends the #93091 item-1 structured refusal enum. + if getattr(e, "reason", "") == "target_busy": + return _err(rid, 5096, str(e)) return _err(rid, 5094, str(e))