diff --git a/cron/bot_chat_delivery.py b/cron/bot_chat_delivery.py new file mode 100644 index 0000000000..0448c42418 --- /dev/null +++ b/cron/bot_chat_delivery.py @@ -0,0 +1,103 @@ +"""Defer never-started cron outputs behind unsupported Bot Chat owners. + +Inspired by 686f6c61's queue proposal (#100319). Unlike retrying failed CLI +turns, only pending requests are eligible: a persisted claim never expires. +""" +from __future__ import annotations + +import contextvars +import json +import threading +from pathlib import Path + +from hermes_cli.active_sessions import _FileLock +from hermes_constants import get_hermes_home +from utils import atomic_json_write + +_running: set[Path] = set() +_running_lock = threading.Lock() + + +def _root() -> Path: + return get_hermes_home().resolve() / "cron" / "bot_chat_pending" + + +def read_pending(key: str) -> dict | None: + try: + return json.loads((_root() / f"{key}.json").read_text(encoding="utf-8")) + except FileNotFoundError: + return None + + +def defer(key: str, job: dict, content: str, profile: str, home: Path) -> dict: + root = _root() + root.mkdir(parents=True, exist_ok=True, mode=0o700) + with _FileLock(root / ".lock"): + record = read_pending(key) + if record is not None: + if record["content"] != content or record["home"] != str(home): + raise ValueError("delivery id already belongs to a different payload") + return record + sequence = max((json.loads(p.read_text(encoding="utf-8"))["sequence"] + for p in root.glob("*.json")), default=0) + 1 + record = dict(id=key, status="queued", job=job, content=content, + profile=profile, home=str(home), sequence=sequence) + atomic_json_write(root / f"{key}.json", record, fsync_dir=True, mode=0o600) + return record + + +def drain() -> None: + """Serialize drains across processes without holding the producer lock.""" + root = _root() + if root.is_dir(): + with _FileLock(root / ".drain.lock"): + _drain(root) + + +def _drain(root: Path) -> None: + """Claim before execution. Errors/interruptions never authorize another turn.""" + from cron.scheduler_delivery import _deliver_to_bot_chat + from tools.bot_live_delivery import find_canonical_live_owner, find_canonical_owner + + paths = sorted(root.glob("*.json"), + key=lambda p: json.loads(p.read_text(encoding="utf-8"))["sequence"]) + for path in paths: + with _FileLock(root / ".lock"): + record = json.loads(path.read_text(encoding="utf-8")) + if record["status"] != "queued": + continue + home = Path(record["home"]) + try: + owner = find_canonical_owner(home) + if owner is not None and find_canonical_live_owner(home) is None: + continue + except Exception: + # Discovery uncertainty is not permission to launch. + continue + record["status"] = "claimed" + atomic_json_write(path, record, fsync_dir=True, mode=0o600) + error = _deliver_to_bot_chat(record["job"], record["content"], record["profile"], deferred=True) + record.update(status="ambiguous" if error else "settled", error=error) + # A transferred live-owner receipt remains authoritative, including queued. + atomic_json_write(path, record, fsync_dir=True, mode=0o600) + + +def drain_in_background() -> None: + """Do not hold up unrelated cron ticks while the eventual Bot Chat turn runs.""" + home = get_hermes_home().resolve() + if not _root().is_dir(): + return + with _running_lock: + if home in _running: + return + _running.add(home) + + def run(): + try: + drain() + finally: + with _running_lock: + _running.discard(home) + + threading.Thread(target=contextvars.copy_context().run, args=(run,), daemon=True, + name="cron-bot-chat-drain").start() diff --git a/cron/scheduler.py b/cron/scheduler.py index d91f9ef374..99d8b0485f 100644 --- a/cron/scheduler.py +++ b/cron/scheduler.py @@ -3789,6 +3789,11 @@ def tick( logger.debug("Cron dispatch paused while gateway drains existing work") return 0 + from cron.bot_chat_delivery import drain, drain_in_background + if sync: + drain() + else: + drain_in_background() _maybe_reap_dead_owners() # Periodic worktree GC (6h, threaded) — the only sweep gateway-only boxes get. try: diff --git a/cron/scheduler_delivery.py b/cron/scheduler_delivery.py index 90c02b331c..c419f8eaa3 100644 --- a/cron/scheduler_delivery.py +++ b/cron/scheduler_delivery.py @@ -653,7 +653,7 @@ def _get_bot_chat_delivery_timeout() -> int: return 600 -def _deliver_to_bot_chat(job: dict, content: str, profile: str) -> Optional[str]: +def _deliver_to_bot_chat(job: dict, content: str, profile: str, *, deferred: bool = False) -> Optional[str]: """Hand output to the live Bot Chat owner, or use the legacy unowned CLI lane. None means completed; a queued/claimed receipt returns an explicit unverified status @@ -693,6 +693,19 @@ def _deliver_to_bot_chat(job: dict, content: str, profile: str) -> Optional[str] # Read BEFORE discovery: the previous owner may have exited after accepting. # No receipt state, including ambiguous/failed, authorizes a CLI replay. receipt = read_delivery_result(home, key) + if receipt is None and not deferred: + from cron.bot_chat_delivery import defer, read_pending + from tools.bot_live_delivery import find_canonical_owner + + pending = read_pending(key) + if pending is None and find_canonical_live_owner(home) is None and find_canonical_owner(home): + pending = defer(key, dict(job), content, profile, home) + if pending is not None: + status = pending["status"] + target = f"bot-chat:{profile_label}" + job.setdefault("_bot_chat_delivery_receipts", {})[target] = { + "status": status, "delivery_id": key} + return None if status == "settled" else f"{target} {status} (receipt {key}): completion unverified; do not resend" if receipt is None: owner = find_canonical_live_owner(home) if owner is not None: diff --git a/tests/cron/test_bot_chat_pending.py b/tests/cron/test_bot_chat_pending.py new file mode 100644 index 0000000000..ad811451fa --- /dev/null +++ b/tests/cron/test_bot_chat_pending.py @@ -0,0 +1,62 @@ +"""Only never-started cron delivery may wait for a CLI owner's release.""" +import subprocess +from unittest.mock import Mock + +import pytest + +from cron import bot_chat_delivery as queue +from cron import scheduler_delivery as delivery +from hermes_cli.active_sessions import try_acquire_active_session +from hermes_state import SessionDB + + +@pytest.mark.parametrize("error", [None, subprocess.TimeoutExpired("hermes", 1)]) +def test_cli_owner_deferral_and_attempt_fence(tmp_path, monkeypatch, error): + monkeypatch.setenv("HERMES_HOME", str(tmp_path)) + db = SessionDB(db_path=tmp_path / "state.db") + db.create_session(session_id="chat", source="cli") + db.set_session_title("chat", "Bot Chat") + lease, refusal = try_acquire_active_session(session_id="chat", surface="cli", config={}, registry_home=tmp_path) + assert refusal is None and lease is not None + run = Mock(side_effect=error, return_value=subprocess.CompletedProcess([], 0, "", "")) + monkeypatch.setattr(delivery.subprocess, "run", run) + monkeypatch.setattr(delivery.shutil, "which", lambda _: "/bin/hermes") + job = {"id": "job", "execution_id": "execution"} + try: + assert "queued" in delivery._deliver_to_bot_chat(job, "output", "") + key = job["_bot_chat_delivery_receipts"]["bot-chat:(own)"]["delivery_id"] + queue.drain() + run.assert_not_called() + lease.release() + queue.drain() + assert run.call_count == 1 + expected = "ambiguous" if error else "settled" + assert queue.read_pending(key)["status"] == expected + queue.drain() + delivery._deliver_to_bot_chat(job, "output", "") + assert run.call_count == 1 + finally: + lease.release() + db.close() + + +def test_pending_queue_uses_admission_order_and_keeps_claims(tmp_path, monkeypatch): + monkeypatch.setenv("HERMES_HOME", str(tmp_path)) + job = {"id": "job"} + queue.defer("f" * 64, job, "older", "", tmp_path) + queue.defer("a" * 64, job, "newer", "", tmp_path) + seen = [] + + def interrupted(job, content, profile, **kwargs): + seen.append(content) + raise KeyboardInterrupt + + monkeypatch.setattr(delivery, "_deliver_to_bot_chat", interrupted) + with pytest.raises(KeyboardInterrupt): + queue.drain() + assert seen == ["older"] + assert queue.read_pending("f" * 64)["status"] == "claimed" + monkeypatch.setattr(delivery, "_deliver_to_bot_chat", lambda j, c, p, **kw: seen.append(c)) + queue.drain() + queue.drain() + assert seen == ["older", "newer"] diff --git a/tools/bot_live_delivery.py b/tools/bot_live_delivery.py index 54d71c4d2b..ebec4efeec 100644 --- a/tools/bot_live_delivery.py +++ b/tools/bot_live_delivery.py @@ -25,12 +25,8 @@ _OWNER_KEYS = ("profile_home", "session_id", "lease_id", "live_session_id") _TERMINAL = frozenset({"settled", "failed", "cancelled", "ambiguous"}) -def find_canonical_live_owner(profile_home: Path | str) -> dict[str, Any] | None: - """Resolve exact Bot Chat's compression tip without creating/migrating its DB. - - Capability advertisement is mandatory; old Desktop/TUI processes must not - receive work they cannot consume. Registry errors propagate, failing closed. - """ +def find_canonical_owner(profile_home: Path | str) -> dict[str, Any] | None: + """Return the exact Bot Chat tip's lease, including unsupported CLI owners.""" from hermes_cli.active_sessions import active_session_registry_snapshot from hermes_state import SessionDB @@ -46,12 +42,18 @@ def find_canonical_live_owner(profile_home: Path | str) -> dict[str, Any] | None if not session_id: return None for entry in active_session_registry_snapshot(registry_home=home): - meta = entry.get("metadata") or {} - if (entry["session_id"] == session_id - and meta.get("bot_live_delivery_consumer") is True - and meta.get("live_session_id")): - return dict(profile_home=str(home), session_id=session_id, - lease_id=entry["lease_id"], live_session_id=meta["live_session_id"]) + if entry["session_id"] == session_id: + return {**entry, "profile_home": str(home)} + return None + + +def find_canonical_live_owner(profile_home: Path | str) -> dict[str, Any] | None: + """Only advertised consumers may receive owner-pinned mailbox deliveries.""" + entry = find_canonical_owner(profile_home) + meta = (entry or {}).get("metadata") or {} + if entry and meta.get("bot_live_delivery_consumer") is True and meta.get("live_session_id"): + return {key: entry[key] for key in ("profile_home", "session_id", "lease_id")} | { + "live_session_id": meta["live_session_id"]} return None diff --git a/website/docs/user-guide/features/cron.md b/website/docs/user-guide/features/cron.md index 73d96cdd54..28dd92ef6c 100644 --- a/website/docs/user-guide/features/cron.md +++ b/website/docs/user-guide/features/cron.md @@ -536,7 +536,7 @@ error. A delivery failure does not count toward the job's `failure_streak` - `bot-chat:` targets another profile **on the same machine**. Names are validated against `hermes profile list` when the job is created; profiles on other gateways or machines can never be targeted, so same-named profiles across machines are unambiguous. - Each delivery costs the target bot one full agent turn — mind the schedule frequency. - Composes with other targets (`bot-chat,telegram`) but is never included in `all`. -- If the canonical chat is open in a mailbox-capable Desktop/TUI backend, delivery is **durably queued immediately**, whether the bot is idle or busy. Only that live owner runs the incoming turn; cron does not start a competing CLI writer. Without a live mailbox owner, the existing `hermes chat -c "Bot Chat" --create-if-missing` lane remains available (normal session ownership checks still apply). +- If the canonical chat is open in a mailbox-capable Desktop/TUI backend, delivery is **durably queued immediately**, whether the bot is idle or busy. Only that live owner runs the incoming turn; cron does not start a competing CLI writer. If a CLI-only or older unsupported owner holds the chat, cron retains the never-started output under the sending profile's `cron/bot_chat_pending/.json`. Later scheduler ticks deliver after that owner releases the chat, in admission order. With no owner, the existing `hermes chat -c "Bot Chat" --create-if-missing` lane remains available (normal session ownership checks still apply). A deferred request is claimed before launching that lane; interruption or an uncertain subprocess result never causes an automatic resend. - **Queued is not completed.** Cron records receipt IDs and `queued`/`claimed` statuses in `last_delivery_queued`, with delivery outcome `queued` (neither delivered nor failed). A successful job shows `delivery_queued`; genuine errors on other targets still take precedence as delivery failures. The bot may complete later. The durable receipt in the target profile's `runtime/bot_live_delivery/.json` is authoritative; cron's historical status is not automatically refreshed. - Rechecking the same execution inspects its existing receipt, even if the owner has disappeared. It never falls back to another writer after acceptance. `failed`, `cancelled`, or `ambiguous` receipts are not automatically replayed; inspect the chat and receipt before intentionally starting new work. Each new cron execution has a distinct delivery ID.