fix(cron): keep deferred Bot Chat delivery bound to admission
Carry the original destination home and delivery ID into deferred drain and its child, rather than re-resolving a mutable profile/root. Missing destinations fail closed; supported-owner handoffs remain transferred, not ambiguous failures. Capture the producer root before the background thread starts, and retain/log malformed JSON without stopping healthy admissions or the whole cron tick. Two invariants reproduced failures on the published head. Real Electron root change and malformed-record cases are red before and green after; nested DM control remains passing. No automatic retry of claimed or uncertain turns.
This commit is contained in:
@@ -7,6 +7,7 @@ from __future__ import annotations
|
||||
|
||||
import contextvars
|
||||
import json
|
||||
import logging
|
||||
import threading
|
||||
from pathlib import Path
|
||||
|
||||
@@ -14,6 +15,7 @@ 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()
|
||||
|
||||
@@ -29,6 +31,19 @@ def read_pending(key: str) -> dict | None:
|
||||
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)
|
||||
@@ -38,17 +53,16 @@ def defer(key: str, job: dict, content: str, profile: str, home: Path) -> dict:
|
||||
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
|
||||
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() -> None:
|
||||
def drain(root: Path | None = None) -> None:
|
||||
"""Serialize drains across processes without holding the producer lock."""
|
||||
root = _root()
|
||||
root = root if root is not None else _root()
|
||||
if root.is_dir():
|
||||
with _FileLock(root / ".drain.lock"):
|
||||
_drain(root)
|
||||
@@ -59,9 +73,9 @@ def _drain(root: Path) -> None:
|
||||
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"):
|
||||
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":
|
||||
@@ -76,8 +90,13 @@ def _drain(root: Path) -> None:
|
||||
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)
|
||||
job = record["job"]
|
||||
job.pop("_bot_chat_delivery_receipts", None)
|
||||
error = _deliver_to_bot_chat(job, record["content"], record["profile"], deferred=record)
|
||||
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)
|
||||
|
||||
@@ -85,7 +104,8 @@ def _drain(root: Path) -> None:
|
||||
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():
|
||||
root = home / "cron" / "bot_chat_pending"
|
||||
if not root.is_dir():
|
||||
return
|
||||
with _running_lock:
|
||||
if home in _running:
|
||||
@@ -94,7 +114,7 @@ def drain_in_background() -> None:
|
||||
|
||||
def run():
|
||||
try:
|
||||
drain()
|
||||
drain(root)
|
||||
finally:
|
||||
with _running_lock:
|
||||
_running.discard(home)
|
||||
|
||||
@@ -653,7 +653,7 @@ def _get_bot_chat_delivery_timeout() -> int:
|
||||
return 600
|
||||
|
||||
|
||||
def _deliver_to_bot_chat(job: dict, content: str, profile: str, *, deferred: bool = False) -> 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
|
||||
@@ -679,7 +679,11 @@ def _deliver_to_bot_chat(job: dict, content: str, profile: str, *, deferred: boo
|
||||
)
|
||||
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.
|
||||
@@ -690,6 +694,8 @@ def _deliver_to_bot_chat(job: dict, content: str, profile: str, *, deferred: boo
|
||||
[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)
|
||||
@@ -750,7 +756,12 @@ def _deliver_to_bot_chat(job: dict, content: str, profile: str, *, deferred: boo
|
||||
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:
|
||||
if deferred is not None:
|
||||
# Admission owns the destination, not the current profile-name resolver.
|
||||
env["HERMES_HOME"] = str(home)
|
||||
if home.parent.name != "profiles":
|
||||
argv += ["-p", "default"] # Ignore a subsequently changed active_profile.
|
||||
elif profile:
|
||||
argv += ["-p", profile]
|
||||
# -p owns profile resolution; this scheduler's HERMES_HOME must not shadow it.
|
||||
env.pop("HERMES_HOME", None)
|
||||
|
||||
@@ -16,7 +16,7 @@ rm apps/desktop/e2e/probe-dm-delivery.spec.ts
|
||||
```
|
||||
|
||||
Use the current seat's actual Xauthority path and an existing runtime venv.
|
||||
Artifacts are retained in `/tmp/botmode-dm-recovery`; sandbox path is printed. The
|
||||
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
|
||||
@@ -40,3 +40,21 @@ is profile-DB-based, not workspace selection; the named live-owner route is posi
|
||||
|
||||
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.
|
||||
|
||||
@@ -9,7 +9,7 @@ const repo = path.resolve(import.meta.dirname, '../../..')
|
||||
const python = path.join(process.env.VIRTUAL_ENV || path.join(repo, '.venv'), 'bin', 'python')
|
||||
let fixture: MockBackendFixture
|
||||
let env: Record<string, string>
|
||||
const evidence = '/tmp/botmode-dm-recovery'
|
||||
const evidence = process.env.BOT_DM_EVIDENCE || '/tmp/botmode-dm-review/native'
|
||||
|
||||
test.beforeAll(async () => {
|
||||
fs.mkdirSync(evidence, { recursive: true })
|
||||
@@ -45,16 +45,18 @@ test('cron output waits for a CLI-only owner and arrives after owner release', a
|
||||
test.setTimeout(240_000)
|
||||
const output = fs.openSync(path.join(evidence, 'cli-owner.log'), 'w')
|
||||
const child = spawn(python, ['-m', 'hermes_cli.main', '-p', 'beta', 'chat', '--in', '~', '-c', 'Bot Chat', '--create-if-missing', '-Q', '-q', 'CLI_OWNER_HOLD'], { cwd: repo, env, stdio: ['ignore', output, output] })
|
||||
const cronEnv = { ...env, HERMES_HOME: path.join(fixture.sandbox.hermesHome, 'profiles', 'beta') }
|
||||
const cronEnv = { ...env, HERMES_HOME: fixture.sandbox.hermesHome }
|
||||
try {
|
||||
await fixture.mock.waitForHeldCompletion()
|
||||
const script = 'import json; from cron.scheduler_delivery import _deliver_to_bot_chat; j={"id":"cli-residual","name":"CLI residual","execution_id":"fixed-execution"}; result=_deliver_to_bot_chat(j,"CLI_OWNER_CRON_SENTINEL",""); print(json.dumps({"result":result,"job":j}))'
|
||||
const script = 'import json; from cron.scheduler_delivery import _deliver_to_bot_chat; j={"id":"cli-residual","name":"CLI residual","execution_id":"fixed-execution"}; result=_deliver_to_bot_chat(j,"CLI_OWNER_CRON_SENTINEL","beta"); print(json.dumps({"result":result,"job":j}))'
|
||||
const result = JSON.parse(execFileSync(python, ['-c', script], { env: cronEnv, cwd: repo, encoding: 'utf8', timeout: 30_000 }))
|
||||
console.log('CLI_OWNER_CRON_ADMISSION', JSON.stringify(result))
|
||||
fs.writeFileSync(path.join(evidence, 'cli-owner-admission.json'), JSON.stringify(result, null, 2))
|
||||
if (process.env.BOT_DM_CORRUPT) fs.writeFileSync(path.join(fixture.sandbox.hermesHome, 'cron', 'bot_chat_pending', 'broken.json'), '{')
|
||||
fixture.mock.releaseHeldStream()
|
||||
await expect.poll(() => child.exitCode, { timeout: 60_000 }).toBe(0)
|
||||
const ticker = spawn(python, ['-c', 'import time; from cron.scheduler import tick; from cron.bot_chat_delivery import _running; tick(verbose=False);\nwhile _running: time.sleep(0.1)'], { env: cronEnv, cwd: repo, stdio: ['ignore', output, output] })
|
||||
fs.mkdirSync(path.join(fixture.sandbox.root, 'changed-launch-home'), { recursive: true })
|
||||
const ticker = spawn(python, ['-c', 'import time; from cron.scheduler import tick; from cron.bot_chat_delivery import _running; tick(verbose=False);\nwhile _running: time.sleep(0.1)'], { env: { ...cronEnv, HOME: path.join(fixture.sandbox.root, 'changed-launch-home') }, cwd: repo, stdio: ['ignore', output, output] })
|
||||
await expect.poll(() => ticker.exitCode, { timeout: 90_000 }).toBe(0)
|
||||
await openBot('beta')
|
||||
expect(dbMessages('beta').filter(([role, text]) => role === 'user' && text.includes('CLI_OWNER_CRON_SENTINEL'))).toHaveLength(1)
|
||||
|
||||
83
tests/cron/test_bot_chat_pending_identity.py
Normal file
83
tests/cron/test_bot_chat_pending_identity.py
Normal file
@@ -0,0 +1,83 @@
|
||||
"""A deferred delivery keeps its original destination and admission identity."""
|
||||
import subprocess
|
||||
from pathlib import Path
|
||||
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
|
||||
from tools.bot_live_delivery import read_delivery_result
|
||||
|
||||
|
||||
@pytest.mark.parametrize("recipient", ["cli", "desktop", "renamed"])
|
||||
def test_deferred_destination_does_not_follow_root_changes(tmp_path, monkeypatch, recipient):
|
||||
source = tmp_path / "source"
|
||||
home = tmp_path / "original" / "profiles" / "beta"
|
||||
other = tmp_path / "other" / "profiles" / "beta"
|
||||
home.mkdir(parents=True)
|
||||
other.mkdir(parents=True)
|
||||
monkeypatch.setenv("HERMES_HOME", str(source))
|
||||
monkeypatch.setattr("hermes_cli.profiles.get_profile_dir", lambda _: home)
|
||||
db = SessionDB(db_path=home / "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=home)
|
||||
assert refusal is None
|
||||
job = {"id": "job", "execution_id": "execution"}
|
||||
try:
|
||||
assert "queued" in delivery._deliver_to_bot_chat(job, "output", "beta")
|
||||
key = job["_bot_chat_delivery_receipts"]["bot-chat:beta"]["delivery_id"]
|
||||
finally:
|
||||
lease.release()
|
||||
db.close()
|
||||
monkeypatch.setattr("hermes_cli.profiles.get_profile_dir", lambda _: other)
|
||||
run = Mock(return_value=subprocess.CompletedProcess([], 0, "", ""))
|
||||
monkeypatch.setattr(delivery.subprocess, "run", run)
|
||||
monkeypatch.setattr(delivery.shutil, "which", lambda _: "/bin/hermes")
|
||||
if recipient == "desktop":
|
||||
lease, refusal = try_acquire_active_session(
|
||||
session_id="chat", surface="desktop", config={}, registry_home=home,
|
||||
metadata={"bot_live_delivery_consumer": True, "live_session_id": "live"})
|
||||
assert refusal is None
|
||||
elif recipient == "renamed":
|
||||
home.rename(home.with_name("renamed"))
|
||||
try:
|
||||
with monkeypatch.context() as changed:
|
||||
changed.setenv("HERMES_HOME", str(tmp_path / "new-source"))
|
||||
queue.drain(source / "cron" / "bot_chat_pending")
|
||||
queue.drain(source / "cron" / "bot_chat_pending")
|
||||
if recipient == "cli":
|
||||
assert run.call_count == 1
|
||||
argv = run.call_args.args[0]
|
||||
assert "-p" not in argv
|
||||
assert Path(run.call_args.kwargs["env"]["HERMES_HOME"]) == home
|
||||
elif recipient == "desktop":
|
||||
run.assert_not_called()
|
||||
receipt = read_delivery_result(home, key)
|
||||
assert receipt is not None and receipt["status"] == "queued"
|
||||
assert queue.read_pending(key)["status"] == "transferred"
|
||||
else:
|
||||
run.assert_not_called()
|
||||
assert not home.exists()
|
||||
assert queue.read_pending(key)["status"] == "ambiguous"
|
||||
assert read_delivery_result(other, key) is None
|
||||
finally:
|
||||
lease.release()
|
||||
|
||||
|
||||
def test_corrupt_record_is_retained_without_blocking_other_admissions(tmp_path, monkeypatch):
|
||||
monkeypatch.setenv("HERMES_HOME", str(tmp_path))
|
||||
queue.defer("a" * 64, {"id": "job"}, "first", "", tmp_path)
|
||||
broken = tmp_path / "cron" / "bot_chat_pending" / "broken.json"
|
||||
broken.write_text("{", encoding="utf-8")
|
||||
seen = []
|
||||
monkeypatch.setattr(delivery, "_deliver_to_bot_chat", lambda j, c, p, **kw: seen.append(c))
|
||||
queue.drain()
|
||||
queue.defer("b" * 64, {"id": "next"}, "second", "", tmp_path)
|
||||
queue.drain()
|
||||
assert seen == ["first", "second"]
|
||||
assert broken.read_text(encoding="utf-8") == "{"
|
||||
@@ -536,7 +536,7 @@ error. A delivery failure does not count toward the job's `failure_streak`
|
||||
- `bot-chat:<profile>` 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. 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/<receipt-id>.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.
|
||||
- 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/<receipt-id>.json`. Later scheduler ticks deliver after that owner releases the chat, in admission order. Deferred work retains its admitted destination home and receipt ID even if the scheduler's launch root changes; a missing/renamed destination is not recreated or resolved to another profile. A `transferred` pending record points to the live-owner receipt, not a failed turn. Malformed JSON records are retained and logged without blocking other queued outputs. 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/<receipt-id>.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.
|
||||
|
||||
|
||||
Reference in New Issue
Block a user