Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
128 changes: 128 additions & 0 deletions cron/bot_chat_delivery.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,128 @@
"""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 logging
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

logger = logging.getLogger(__name__)
_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 _records(root: Path) -> list[tuple[Path, dict]]:
records = []
for path in root.glob("*.json"):
try:
record = json.loads(path.read_text(encoding="utf-8"))
except (json.JSONDecodeError, UnicodeDecodeError) as exc:
# Keep damaged receipts as evidence; never replay them or block peers.
logger.error("Unreadable deferred Bot Chat receipt %s: %s", path, exc)
continue
records.append((path, record))
return records


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((record["sequence"] for _, record in _records(root)), 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(root: Path | None = None) -> None:
"""Serialize drains across processes without holding the producer lock."""
root = root if root is not None else _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

with _FileLock(root / ".lock"):
records = sorted(_records(root), key=lambda item: item[1]["sequence"])
for path, _ in records:
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)
job = record["job"]
job.pop("_bot_chat_delivery_receipts", None)
try:
error = _deliver_to_bot_chat(job, record["content"], record["profile"], deferred=record)
except Exception as exc:
# The claim survives uncertainty; one failed attempt must not stop peers.
error = f"{type(exc).__name__}: {exc}"
logger.exception("Deferred Bot Chat delivery %s failed", record["id"])
receipt = job.get("_bot_chat_delivery_receipts", {}).get(
f"bot-chat:{record['profile'] or '(own)'}")
status = "transferred" if receipt else "ambiguous" if error else "settled"
record.update(status=status, 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()
root = home / "cron" / "bot_chat_pending"
if not root.is_dir():
return
with _running_lock:
if home in _running:
return
_running.add(home)

def run():
try:
drain(root)
finally:
with _running_lock:
_running.discard(home)

threading.Thread(target=contextvars.copy_context().run, args=(run,), daemon=True,
name="cron-bot-chat-drain").start()
5 changes: 5 additions & 0 deletions cron/scheduler.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
41 changes: 31 additions & 10 deletions cron/scheduler_delivery.py
Original file line number Diff line number Diff line change
Expand Up @@ -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: Optional[dict] = None) -> 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
Expand All @@ -679,7 +679,11 @@ def _deliver_to_bot_chat(job: dict, content: str, profile: str) -> Optional[str]
)
try:
source_home = get_hermes_home().resolve()
home = (get_profile_dir(profile) if profile else source_home).resolve()
from pathlib import Path
home = (Path(deferred["home"]) if deferred is not None else
get_profile_dir(profile) if profile else source_home).resolve()
if deferred is not None and not (home / "state.db").is_file():
return f"bot-chat delivery target no longer exists: {home}; do not resend"
# run_one_job/claim_fire attach the durable execution id before delivery. The
# transient fallback supports direct helper callers, never deduping recurring
# runs by their (potentially identical) output or previous last_run timestamp.
Expand All @@ -690,9 +694,26 @@ def _deliver_to_bot_chat(job: dict, content: str, profile: str) -> Optional[str]
[str(source_home), job_id, str(run_id), str(home)],
ensure_ascii=False, separators=(",", ":"),
).encode("utf-8")).hexdigest()
if deferred is not None:
key = deferred["id"]
# 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:
if pending["content"] != content or pending["home"] != str(home):
raise ValueError("delivery id already belongs to a different payload")
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:
Expand Down Expand Up @@ -735,13 +756,13 @@ def _fail(msg: str, **log_kwargs) -> str:
from agent.delegation_context import delegated_child_subprocess_env
from tools.environments.local import strip_launch_profile_env
env = strip_launch_profile_env(delegated_child_subprocess_env(os.environ))
if profile:
argv += ["-p", profile]
# -p owns profile resolution; this scheduler's HERMES_HOME must not shadow it.
env.pop("HERMES_HOME", None)
else:
# Multiplex workers carry the profile in a ContextVar, not os.environ.
env["HERMES_HOME"] = str(source_home)
if not home.is_dir():
return _fail(f"bot-chat delivery target no longer exists: {home}; do not resend")
# Discovery (or deferred admission) owns the destination, not HOME or a
# subsequently changed active_profile. Do not resolve the name a second time.
env["HERMES_HOME"] = str(home)
if home.parent.name != "profiles":
argv += ["-p", "default"]

query_file = None
try:
Expand All @@ -761,7 +782,7 @@ def _fail(msg: str, **log_kwargs) -> str:
if result.returncode != 0:
tail = (result.stderr or result.stdout or "").strip()[-500:]
return _fail(
f"bot-chat delivery to profile '{profile_label}' failed (exit {result.returncode})"
f"bot-chat delivery to profile '{profile_label}' failed (exit {result.returncode}) at {home}"
+ (f": {tail}" if tail else ""))
logger.info("Job '%s': delivered to Bot Chat of profile '%s'", job_id, profile_label)
return None
Expand Down
107 changes: 107 additions & 0 deletions evals/botmode-dm-delivery/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,107 @@
# Native Bot Mode delivery probe

Real Electron, production Python backend and tool execution, disposable HOME/HERMES_HOME,
loopback scripted inference (no paid model). Linux seat fixture; run from repository root:

```sh
npm ci --no-audit --no-fund
cp evals/botmode-dm-delivery/probe-dm-delivery.spec.ts apps/desktop/e2e/
git apply evals/botmode-dm-delivery/mock-trigger.patch
(cd apps/desktop && npm run build)
(cd apps/desktop && DISPLAY=:0 XAUTHORITY=/run/user/1000/xauth_cnpsqU \
XDG_RUNTIME_DIR=/run/user/1000 VIRTUAL_ENV="$VIRTUAL_ENV" \
HERMES_DESKTOP_CDP_PORT=off npx playwright test e2e/probe-dm-delivery.spec.ts --reporter=list)
git apply -R evals/botmode-dm-delivery/mock-trigger.patch
rm apps/desktop/e2e/probe-dm-delivery.spec.ts
```

Use the current seat's actual Xauthority path and an existing runtime venv.
Artifacts default to `/tmp/botmode-dm-review/native` (override with `BOT_DM_EVIDENCE`); sandbox path is printed. The
fixture's generated hermes shim pins every child to this checkout, not an installed launcher.

## Verified results

Base `cf35e7351e770`: nested beta quiet CLI executes real `message_agent` to the named
alpha Desktop owner. Admission, execution and reply succeed, and the quiet sender
receives its completion. A reload then shows the attributed incoming message. The
live pre-reload view showed the reply but omitted the incoming row in this fixture;
that renderer-refresh behavior is not fixed by this cron change.

Base CLI-owner case: separate cron producer returns exact `SESSION_NOT_OWNED` refusal.
Release owner, run scheduler tick, open Beta in Desktop: no cron output, assertion red.
Fixed case: queued receipt; release owner; real synchronous scheduler tick drains;
Beta Desktop renders `CLI_OWNER_CRON_SENTINEL` and its reply once. Final two-case
native run: **2 passed (1.2m)**. Supported-owner nested case also passed independently
on base (**1 passed (40.0s)**).

This does not establish native Windows parity, general retry after unowned CLI
failure, or correction of the issue #105460 `--in ~` premise. CLI title resolution
is profile-DB-based, not workspace selection; the named live-owner route is positive.

The queue is deliberately at-most-once after claim. A crash before spawning but
after claiming remains inspectable as claimed; it is not retried automatically.

## Independent review follow-up

The probe now admits the named Beta destination from the default scheduler, then
changes the ticker's `HOME` to a different existing directory before drain.
Published head `c91dfcbfe810c` fails with `Profile 'beta' does not exist`, leaving
zero sentinel inputs in Beta. The follow-up pins the admitted home and ID; Beta
renders one input and reply. Set `BOT_DM_CORRUPT=1` to add one malformed JSON
record alongside the valid admission: before the follow-up the real tick raises
`JSONDecodeError`; afterward the damaged record stays on disk while Beta delivers.

Fresh built native run with both root change and corruption: **2 passed (1.4m)**,
including the existing nested `message_agent` control. Receipts/screenshots:
`/tmp/botmode-dm-review/{native-red2,corrupt-red,final-native}` and matching `.log`
files. The first follow-up run also exposed a fixture mistake (the changed HOME
was not created, so `--in ~` correctly refused); that failed receipt is retained
in `native-green.log`, and both source legs were rerun with an existing HOME.
Prior `/tmp/botmode-dm-recovery*` evidence remains untouched.

## Per-record delivery exception isolation

Set `BOT_DM_EXCEPTION=1` and select `-g "cron output"` for the controlled native
exception probe. `exception-tick.py` raises `PermissionError` at the actual
post-discovery target `Path.is_dir()` boundary, not at the delivery helper.
The same exception was first reproduced with real directory traversal permission
loss on Python 3.11/Linux; the retained test uses a portable controlled fault.
Before the guard, native tick exits 1, leaving the head claimed and sibling queued.
Afterward, the head is ambiguous, the sibling settles and renders once, and a
second real tick replays neither. Logs: `/tmp/botmode-dm-exception-{red,green}.log`;
receipts and screenshot: `/tmp/botmode-dm-exception/{red,green}/`.

The review's repeated-head starvation claim is not reachable: the claim commits
before delivery and later scans skip every non-queued record. One failed tick is
real; recurring replay of that same head is not. Indefinite queued/payload retention
is intentional, with no TTL or automatic ambiguous retry introduced here.

## Ordinary custom-root fallback (#104066 / #104055)

`probe-cron-root.spec.ts` adds the never-deferred sibling: copy it to
`apps/desktop/e2e/` and run with the same native fixture (no mock-trigger patch
needed for this case). It keeps default's real Desktop Bot Chat lease, submits
ordinary cron output to unowned Alpha from a separate Python producer under a
custom Hermes root, and holds the real quiet CLI child at loopback inference.
The child shim PID must match Alpha's real CLI lease; default's lease is unchanged.
After release, Alpha has exactly one input and Desktop renders the output.
The same case removes unused Beta and verifies delivery neither recreates Beta nor
creates a second `.hermes` root under HOME.

Both `origin/main`'s scheduler and pre-follow-up `c827ae179d67c` fail with
`Profile 'alpha' does not exist` before any recipient turn. Fixed native run:
**1 passed (46.3s)**. Two invariant tests exercise the actual CLI startup resolver
across named/default/own destinations with a changed active profile, and refusal
when the destination is missing initially or disappears during discovery:
**5 failed before, 5 passed after**. The old env-clearing test is replaced by these
behavior checks rather than retaining the broken expectation.

Evidence: `/tmp/botmode-cron-root/{before2,origin-main,after}.log`,
`after/{owners.json,children.log,rows.json,result.json,missing.json,ordinary-recipient.png}`.
The first fixture attempt (`before.log`) used the wrong default row label; the
actual Desktop label is Hermes. No production failure is claimed for that attempt.
Full cron directory: **1346 passed, 1 skipped across 116 files**; sibling mailbox,
DM, gateway consumer and profile tests: **124 passed, 3 skipped across 4 files**.
Credit @fangliquanflq's #104066 for the root-boundary diagnosis and anchoring fix;
this combined branch reuses its already-resolved destination instead of repeating
name resolution. No retry or receipt semantics change.
41 changes: 41 additions & 0 deletions evals/botmode-dm-delivery/exception-tick.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,41 @@
"""Controlled filesystem fault; run only against the native disposable sandbox."""
import json
import os
from pathlib import Path

from cron import scheduler_delivery as delivery
from cron.bot_chat_delivery import _root
from cron.scheduler import tick

root = _root()
original_which = delivery.shutil.which
original_is_dir = Path.is_dir
armed = False
raised = 0


def resolve_cli(*args, **kwargs):
global armed
armed = raised == 0
return original_which(*args, **kwargs)


def is_dir(path):
global armed, raised
if armed and path == Path(os.environ["HERMES_HOME"]) / "profiles" / "beta":
armed = False
raised += 1
raise PermissionError("controlled target traversal denied after discovery")
return original_is_dir(path)


delivery.shutil.which = resolve_cli
Path.is_dir = is_dir
try:
tick(verbose=False)
tick(verbose=False)
finally:
Path.is_dir = original_is_dir
delivery.shutil.which = original_which
records = [json.loads(p.read_text(encoding="utf-8")) for p in root.glob("*.json") if p.name != "broken.json"]
Path(os.environ["BOT_DM_EXCEPTION_RECEIPT"]).write_text(json.dumps({"raised": raised, "records": records}, indent=2), encoding="utf-8")
26 changes: 26 additions & 0 deletions evals/botmode-dm-delivery/mock-trigger.patch
Original file line number Diff line number Diff line change
@@ -0,0 +1,26 @@
diff --git a/tests-js/scripts/mock-server.ts b/tests-js/scripts/mock-server.ts
index 6dffd2cf0f388..09c6de07d81b2 100644
--- a/tests-js/scripts/mock-server.ts
+++ b/tests-js/scripts/mock-server.ts
@@ -500,6 +500,21 @@ export function startMockServer(options: MockServerOptions = {}): Promise<MockSe
_receivedUserTexts.push(userText)
}

+ const dmMatch = userText.match(/E2E_DM\(([^)]+)\)\[([\s\S]*)\]/)
+ const lastUserIndex = messages.lastIndexOf(lastUserMsg)
+ if (dmMatch && !messages.slice(lastUserIndex + 1).some(m => m?.role === 'tool')) {
+ const turn: ScriptedTurn = {
+ text: '',
+ toolCalls: [{ name: 'message_agent', args: { target: dmMatch[1], message: dmMatch[2] } }],
+ }
+ if (stream) {
+ streamScriptedTurn(res, model, turn)
+ } else {
+ nonStreamingScriptedTurn(res, model, turn)
+ }
+ return
+ }
+
const isInterimTrigger = userText.includes('E2E_INTERIM_TRIGGER')
const isSidebarTrigger = userText.includes('E2E_SIDEBAR_TRIGGER')
const isSidebarCrossTrigger = userText.includes('E2E_SIDEBAR_CROSS')
Loading
Loading