fix: cron and local DMs reach an open Desktop Bot Chat
Route local producers to durable owner ingress before attempting the unowned CLI lane. Preserve per-run/per-message IDs and receipt-first retry handling; never fall back after ambiguous admission. Report cron admission as queued, not completed or failed, in job status, the execution ledger and CLI/tool UX. Native isolated Electron validation reproduces SESSION_NOT_OWNED on main for both idle and busy owners. Fixed owner consumes idle cron, busy cron, local DM and mounted-chat cron exactly once, keeps its lease, yields to queued human input, and preserves the prior model-request prefix and tool schema. Inference alone used a deterministic loopback wire stub; no paid model call.
This commit is contained in:
@@ -23,7 +23,7 @@ from hermes_time import now as _now
|
||||
logger = logging.getLogger(__name__)
|
||||
_KNOWN_STATUSES = {"claimed", "running", "completed", "failed", "unknown"}
|
||||
_KNOWN_SOURCES = {"builtin", "direct", "external"}
|
||||
_KNOWN_DELIVERY_OUTCOMES = {"delivered", "failed", "suppressed", "suppressed_acked", "not_configured"}
|
||||
_KNOWN_DELIVERY_OUTCOMES = {"queued", "delivered", "failed", "suppressed", "suppressed_acked", "not_configured"}
|
||||
_TERMINAL_STATUSES = {"completed", "failed", "unknown"}
|
||||
|
||||
|
||||
|
||||
+3
-2
@@ -74,6 +74,7 @@ def _initialize_schema(conn: sqlite3.Connection) -> None:
|
||||
"CREATE INDEX IF NOT EXISTS idx_executions_status_claimed "
|
||||
"ON executions(status, claimed_at DESC, id DESC)"
|
||||
)
|
||||
add_column_if_missing(conn, "executions", "delivery_outcome", "delivery_outcome TEXT")
|
||||
add_column_if_missing(conn, "executions", "scheduled_instant", "scheduled_instant TEXT")
|
||||
conn.execute(
|
||||
"CREATE INDEX IF NOT EXISTS idx_executions_occurrence "
|
||||
@@ -247,10 +248,10 @@ def finish_execution(
|
||||
cur = conn.execute(
|
||||
"""UPDATE executions
|
||||
SET status=?, finished_at=?, error=?, handoff_pending=0,
|
||||
handoff_started_at=NULL
|
||||
handoff_started_at=NULL, delivery_outcome=?
|
||||
WHERE id=? AND status IN ('claimed','running')
|
||||
AND process_id=? AND pid=?""",
|
||||
(status, now, detail, execution_id, _PROCESS_ID, os.getpid()),
|
||||
(status, now, detail, delivery_outcome, execution_id, _PROCESS_ID, os.getpid()),
|
||||
)
|
||||
if cur.rowcount != 1:
|
||||
return None
|
||||
|
||||
+12
-1
@@ -2599,9 +2599,12 @@ def _record_fire_ownership_lost(job_id: str, fire_owner: Optional[str], executio
|
||||
def _classify_delivery_outcome(
|
||||
*, delivery_error, should_deliver: bool, unresolved_origin: bool,
|
||||
normalized_deliver: str, incident_acked: bool, success: bool,
|
||||
delivery_queued=None,
|
||||
) -> str:
|
||||
if delivery_error:
|
||||
return "failed"
|
||||
if should_deliver and delivery_queued:
|
||||
return "queued"
|
||||
if should_deliver and unresolved_origin:
|
||||
return "not_configured"
|
||||
if should_deliver and normalized_deliver != "local":
|
||||
@@ -2803,7 +2806,13 @@ def _finish_interrupted_run(job: dict, execution_id: str, delivery_error: Option
|
||||
def _finish_completed_run(d: _RunDelivery, fire_owner: Optional[str], execution_id: str) -> bool:
|
||||
"""mark_job_run (owner-fenced) + execution ledger row for a run that reached delivery."""
|
||||
job = d.job
|
||||
if not d.should_deliver and job.get("last_delivery_queued"):
|
||||
from cron.jobs import update_job
|
||||
update_job(job["id"], {"last_delivery_queued": None})
|
||||
job["last_delivery_queued"] = None
|
||||
mark_kwargs = {"delivery_error": d.delivery_error}
|
||||
if d.success and not d.delivery_error and d.should_deliver and job.get("last_delivery_queued"):
|
||||
mark_kwargs["status"] = "delivery_queued"
|
||||
if fire_owner is not None:
|
||||
mark_kwargs["expected_fire_owner"] = fire_owner
|
||||
if d.blocked_config:
|
||||
@@ -2816,6 +2825,7 @@ def _finish_completed_run(d: _RunDelivery, fire_owner: Optional[str], execution_
|
||||
return True
|
||||
delivery_outcome = _classify_delivery_outcome(
|
||||
delivery_error=d.delivery_error,
|
||||
delivery_queued=job.get("last_delivery_queued"),
|
||||
should_deliver=d.should_deliver,
|
||||
unresolved_origin=d.unresolved_origin,
|
||||
# Read the lane the notice was actually routed through (failure_deliver on failure).
|
||||
@@ -2861,7 +2871,8 @@ def _deliver_crash_failure(
|
||||
)
|
||||
delivery_outcome = _classify_delivery_outcome(
|
||||
delivery_error=delivery_error, should_deliver=True, unresolved_origin=unresolved_origin,
|
||||
normalized_deliver=normalized_deliver, incident_acked=False, success=False)
|
||||
normalized_deliver=normalized_deliver, incident_acked=False, success=False,
|
||||
delivery_queued=job.get("last_delivery_queued"))
|
||||
if delivery_outcome in ("delivered", "not_configured"):
|
||||
_mark_incident_alerted(failure_incident_id)
|
||||
return delivery_error, delivery_outcome
|
||||
|
||||
+89
-20
@@ -642,12 +642,67 @@ def _get_bot_chat_delivery_timeout() -> int:
|
||||
|
||||
|
||||
def _deliver_to_bot_chat(job: dict, content: str, profile: str) -> Optional[str]:
|
||||
"""Deliver job output into a profile's canonical Bot Chat as a real inbound user turn, via
|
||||
``hermes [-p <profile>] chat --in ~ -c "Bot Chat" --create-if-missing -Q --query-file`` — the
|
||||
Bot Mode agent-to-agent lane, so canonical-session rules apply and it is alternation-safe.
|
||||
``profile`` is ``""`` for the job's own profile. None on success, else an error string."""
|
||||
"""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
|
||||
string so existing Optional[str] callers cannot misreport admission as delivery.
|
||||
``profile`` is ``""`` for the job's own profile.
|
||||
"""
|
||||
import hashlib
|
||||
import json
|
||||
import tempfile
|
||||
import uuid
|
||||
from hermes_constants import get_hermes_home
|
||||
from hermes_cli.profiles import get_profile_dir
|
||||
from tools.bot_live_delivery import (
|
||||
deliver_to_live_owner, find_canonical_live_owner, read_delivery_result,
|
||||
)
|
||||
|
||||
job_id = job.get("id", "?")
|
||||
profile_label = profile or "(own)"
|
||||
message = (
|
||||
f'[Cronjob "{job.get("name", job_id)}" output — scheduled job, not the user. '
|
||||
f"Review it, act on anything that needs action, and summarize "
|
||||
f"for the chat.]\n\n{content}"
|
||||
)
|
||||
try:
|
||||
source_home = get_hermes_home().resolve()
|
||||
home = (get_profile_dir(profile) if profile else source_home).resolve()
|
||||
# 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.
|
||||
run_id = job.get("execution_id")
|
||||
if not run_id:
|
||||
run_id = job.setdefault("_bot_chat_run_id", uuid.uuid4().hex)
|
||||
key = hashlib.sha256(json.dumps(
|
||||
[str(source_home), job_id, str(run_id), str(home)],
|
||||
ensure_ascii=False, separators=(",", ":"),
|
||||
).encode("utf-8")).hexdigest()
|
||||
# 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:
|
||||
owner = find_canonical_live_owner(home)
|
||||
if owner is not None:
|
||||
receipt = deliver_to_live_owner(home, owner, message, delivery_id=key)
|
||||
if receipt is not None:
|
||||
if receipt["message"] != message:
|
||||
raise ValueError("delivery id already belongs to a different payload")
|
||||
status = receipt["status"]
|
||||
target = f"bot-chat:{profile_label}"
|
||||
receipts = job.setdefault("_bot_chat_delivery_receipts", {})
|
||||
receipts[target] = {"status": status, "delivery_id": key}
|
||||
logger.info("Job '%s': Bot Chat %s receipt=%s status=%s",
|
||||
job_id, profile_label, key, status)
|
||||
if status == "settled":
|
||||
return None
|
||||
detail = ("completion unverified; do not resend" if status in ("queued", "claimed")
|
||||
else receipt.get("error") or receipt.get("reason") or "not completed")
|
||||
return f"{target} {status} (receipt {key}): {detail}"
|
||||
except Exception as exc:
|
||||
# Discovery/admission uncertainty must never open a second-writer fallback.
|
||||
return f"bot-chat delivery to profile '{profile_label}' unverified: {exc}"
|
||||
|
||||
hermes_bin = shutil.which("hermes")
|
||||
if hermes_bin:
|
||||
argv = [hermes_bin]
|
||||
@@ -671,14 +726,9 @@ def _deliver_to_bot_chat(job: dict, content: str, profile: str) -> Optional[str]
|
||||
argv += ["-p", profile]
|
||||
# -p owns profile resolution; this scheduler's HERMES_HOME must not shadow it.
|
||||
env.pop("HERMES_HOME", None)
|
||||
|
||||
# Prefix marks this as scheduled output, not the human (Bot Mode sender-attribution).
|
||||
message = (
|
||||
f'[Cronjob "{job.get("name", job_id)}" output — scheduled job, not the user. '
|
||||
f"Review it, act on anything that needs action, and summarize "
|
||||
f"for the chat.]\n\n{content}"
|
||||
)
|
||||
profile_label = profile or "(own)"
|
||||
else:
|
||||
# Multiplex workers carry the profile in a ContextVar, not os.environ.
|
||||
env["HERMES_HOME"] = str(source_home)
|
||||
|
||||
query_file = None
|
||||
try:
|
||||
@@ -998,14 +1048,21 @@ def _cron_delivery_notify_enabled(cfg: Optional[dict]) -> bool:
|
||||
|
||||
def _record_delivery_verification(job: dict, unverified_targets: list) -> None:
|
||||
"""Persist ``last_delivery_unverified``: list of ``platform:chat_id`` targets acked with no
|
||||
evidence, or None. Skips the write when unchanged; never raises (bookkeeping must not fail a
|
||||
evidence, or None, alongside queued Bot Chat receipts. Never raises (bookkeeping must not fail a
|
||||
delivery)."""
|
||||
new_value = list(unverified_targets) or None
|
||||
if (job.get("last_delivery_unverified") or None) == new_value:
|
||||
queued = {target: receipt for target, receipt in
|
||||
job.get("_bot_chat_delivery_receipts", {}).items()
|
||||
if receipt["status"] in ("queued", "claimed")} or None
|
||||
values = {key: value for key, value in {
|
||||
"last_delivery_unverified": new_value, "last_delivery_queued": queued,
|
||||
}.items() if (job.get(key) or None) != value}
|
||||
if not values:
|
||||
return
|
||||
job.update(values)
|
||||
try:
|
||||
from cron.jobs import update_job
|
||||
update_job(job["id"], {"last_delivery_unverified": new_value})
|
||||
update_job(job["id"], values)
|
||||
except Exception as exc: # pragma: no cover - defensive
|
||||
logger.debug("Job '%s': could not record delivery verification: %s", job.get("id"), exc)
|
||||
|
||||
@@ -1599,8 +1656,10 @@ def _deliver_result(
|
||||
running) the live adapter is tried first (E2EE rooms can't use the standalone HTTP path), then
|
||||
standalone fallback. ``for_failure=True`` routes failure-category notices through the job's
|
||||
``failure_deliver`` override when present (NS-788). Returns None on success, else an error."""
|
||||
job.pop("_bot_chat_delivery_receipts", None)
|
||||
targets = _resolve_delivery_targets(job, for_failure=for_failure)
|
||||
if not targets:
|
||||
_record_delivery_verification(job, [])
|
||||
return _unresolved_delivery_outcome(job, for_failure)
|
||||
|
||||
# Restart-safe workers have no live gateway adapters: hand the send back through a durable
|
||||
@@ -1610,10 +1669,16 @@ def _deliver_result(
|
||||
# and that nested delivery must not be keyed under the outer execution id.
|
||||
external_execution = os.environ.get("_HERMES_CRON_EXTERNAL_WORKER", "")
|
||||
if (external_execution and adapters is None
|
||||
and external_execution == str(job.get("execution_id") or "")):
|
||||
and external_execution == str(job.get("execution_id") or "")
|
||||
and any(target["platform"] != BOT_CHAT_PLATFORM for target in targets)):
|
||||
from cron.delivery_queue import enqueue_and_wait
|
||||
|
||||
return enqueue_and_wait(external_execution, job, content, for_failure=for_failure)
|
||||
_record_delivery_verification(job, [])
|
||||
error = enqueue_and_wait(external_execution, job, content, for_failure=for_failure)
|
||||
from cron.jobs import get_job
|
||||
refreshed = get_job(job["id"]) or {}
|
||||
job["last_delivery_queued"] = refreshed.get("last_delivery_queued")
|
||||
return error
|
||||
|
||||
from gateway.config import load_gateway_config
|
||||
|
||||
@@ -1677,12 +1742,16 @@ def _deliver_result(
|
||||
|
||||
delivery_errors = []
|
||||
for target in targets:
|
||||
# bot-chat targets bypass gateway adapters: output becomes an inbound turn in the target
|
||||
# profile's Bot Chat via the chat CLI lane. Must precede the Platform enum, which lacks it.
|
||||
# Bot Chat owns admission; never concurrently resume a live owner's transcript.
|
||||
if target["platform"] == BOT_CHAT_PLATFORM:
|
||||
bot_chat_error = _deliver_to_bot_chat(job, content, target["chat_id"])
|
||||
if bot_chat_error:
|
||||
delivery_errors.append(bot_chat_error)
|
||||
receipt_target = f"bot-chat:{target['chat_id'] or '(own)'}"
|
||||
receipt = job.get("_bot_chat_delivery_receipts", {}).get(receipt_target)
|
||||
if not receipt or receipt["status"] not in ("queued", "claimed"):
|
||||
delivery_errors.append(bot_chat_error)
|
||||
if receipt and receipt["status"] == "ambiguous":
|
||||
unverified_targets.append(bot_chat_error)
|
||||
continue
|
||||
|
||||
t = _prepare_target_delivery(
|
||||
|
||||
+5
-1
@@ -167,6 +167,8 @@ def _last_run_display(job: Dict[str, Any]) -> str:
|
||||
last_status = job["last_status"]
|
||||
if last_status == "ok":
|
||||
return color("ok", Colors.GREEN)
|
||||
if last_status == "delivery_queued":
|
||||
return color("delivery_queued: completion unverified; do not resend", Colors.YELLOW)
|
||||
if last_status == "delivery_failed":
|
||||
# Agent succeeded but the result never reached the user — not green; last_error is None.
|
||||
return color(f"delivery_failed: {job.get('last_delivery_error') or '?'}", Colors.YELLOW)
|
||||
@@ -216,6 +218,8 @@ def _job_rows(job: Dict[str, Any]) -> List[tuple[str, str]]:
|
||||
def _job_warnings(job: Dict[str, Any]) -> List[str]:
|
||||
"""Delivery / fire warning lines for one job in ``cron list``."""
|
||||
lines = []
|
||||
if queued := job.get("last_delivery_queued"):
|
||||
lines.append(f"Delivery queued (completion unverified; do not resend): {queued}")
|
||||
if job.get("last_delivery_error"):
|
||||
lines.append(f"{color('⚠ Delivery failed:', Colors.YELLOW)} {job['last_delivery_error']}")
|
||||
# A live adapter acked the last send but returned no message_id / raw_response
|
||||
@@ -491,7 +495,7 @@ def _cron_doctor_issues_for_job(job: Dict[str, Any]) -> List[str]:
|
||||
issues: List[str] = []
|
||||
last_status = str(job.get("last_status") or "").strip().lower()
|
||||
# "delivery_failed" = the agent run succeeded; the delivery issue below reports it.
|
||||
if last_status and last_status not in {"ok", "delivery_failed"}:
|
||||
if last_status and last_status not in {"ok", "delivery_failed", "delivery_queued"}:
|
||||
issues.append(f"last run failed: {str(job.get('last_error') or 'unknown error').strip()}")
|
||||
if delivery_err := str(job.get("last_delivery_error") or "").strip():
|
||||
issues.append(f"last delivery failed: {delivery_err}")
|
||||
|
||||
@@ -0,0 +1,70 @@
|
||||
"""Cron admission must not become a second writer or a completed-delivery claim."""
|
||||
from pathlib import Path
|
||||
from unittest.mock import Mock
|
||||
|
||||
from cron import scheduler_delivery as delivery
|
||||
from tools import bot_live_delivery as mailbox
|
||||
|
||||
|
||||
def test_live_delivery_retry_keeps_receipt_across_owner_loss(tmp_path, monkeypatch):
|
||||
monkeypatch.setattr(Path, "home", lambda: tmp_path)
|
||||
source = tmp_path / "custom-home"
|
||||
monkeypatch.setenv("HERMES_HOME", str(source))
|
||||
subprocess_run = Mock(side_effect=AssertionError("live owner must not spawn CLI"))
|
||||
monkeypatch.setattr(delivery.subprocess, "run", subprocess_run)
|
||||
from hermes_cli.profiles import get_profile_dir
|
||||
|
||||
for profile, home in [("", source), ("research", get_profile_dir("research"))]:
|
||||
owner = dict(profile_home=str(home.resolve()), session_id="bot", lease_id="lease",
|
||||
live_session_id="live")
|
||||
discovery = Mock(return_value=owner)
|
||||
monkeypatch.setattr(mailbox, "find_canonical_live_owner", discovery)
|
||||
job = dict(id="digest", name="Digest", execution_id="first-run")
|
||||
pending = delivery._deliver_to_bot_chat(job, "payload", profile)
|
||||
assert pending and "queued" in pending
|
||||
records = list((home / "runtime/bot_live_delivery").glob("*.json"))
|
||||
assert len(records) == 1
|
||||
key = records[0].stem
|
||||
record = mailbox.read_delivery_result(home, key)
|
||||
assert record and record["message"].endswith("payload")
|
||||
discovery.side_effect = AssertionError("receipt must precede discovery")
|
||||
assert delivery._deliver_to_bot_chat(dict(job), "payload", profile) == pending
|
||||
mailbox.claim_pending_delivery(home, owner)
|
||||
mailbox.complete_delivery(home, key, status="ambiguous", error="owner died")
|
||||
outcome = delivery._deliver_to_bot_chat(dict(job), "payload", profile)
|
||||
assert outcome and "ambiguous" in outcome
|
||||
discovery.side_effect = None
|
||||
next_job = dict(job, execution_id="next-run")
|
||||
outcome = delivery._deliver_to_bot_chat(next_job, "payload", profile)
|
||||
assert outcome and "queued" in outcome
|
||||
assert len(list((home / "runtime/bot_live_delivery").glob("*.json"))) == 2
|
||||
subprocess_run.assert_not_called()
|
||||
|
||||
|
||||
def test_result_records_pending_until_terminal_receipt(tmp_path, monkeypatch):
|
||||
from cron import jobs
|
||||
from gateway import config
|
||||
|
||||
monkeypatch.setenv("HERMES_HOME", str(tmp_path))
|
||||
monkeypatch.delenv("_HERMES_CRON_EXTERNAL_WORKER", raising=False)
|
||||
owner = dict(profile_home=str(tmp_path.resolve()), session_id="bot", lease_id="lease",
|
||||
live_session_id="live")
|
||||
monkeypatch.setattr(mailbox, "find_canonical_live_owner", lambda home: owner)
|
||||
monkeypatch.setattr(delivery._sched, "load_config", lambda: {})
|
||||
monkeypatch.setattr(config, "load_gateway_config", lambda: None)
|
||||
monkeypatch.setattr(delivery.subprocess, "run", Mock(side_effect=AssertionError("CLI")))
|
||||
updates = []
|
||||
monkeypatch.setattr(jobs, "update_job", lambda key, values: updates.append(values))
|
||||
job = dict(id="digest", execution_id="run", deliver="bot-chat")
|
||||
error = delivery._deliver_result(job, "payload")
|
||||
assert error is None
|
||||
queued = updates[-1]["last_delivery_queued"]
|
||||
assert queued and next(iter(queued.values()))["status"] == "queued"
|
||||
assert delivery._sched._classify_delivery_outcome(
|
||||
delivery_error=error, delivery_queued=queued, should_deliver=True, unresolved_origin=False,
|
||||
normalized_deliver="bot-chat", incident_acked=False, success=True) == "queued"
|
||||
record = mailbox.claim_pending_delivery(tmp_path, owner)
|
||||
assert record is not None
|
||||
mailbox.complete_delivery(tmp_path, record["delivery_id"], status="settled", reply="done")
|
||||
assert delivery._deliver_result(job, "payload") is None
|
||||
assert updates[-1]["last_delivery_queued"] is None
|
||||
@@ -214,7 +214,10 @@ def _capture_spawn(monkeypatch):
|
||||
def _runner_parts(command):
|
||||
parts = shlex.split(command)
|
||||
marker = parts.index("--run-delivery")
|
||||
return parts[marker + 1], parts[marker + 2], parts[marker + 3 :]
|
||||
argv = parts[marker + 3 :]
|
||||
if argv[:1] == ["--profile-home"]:
|
||||
argv = argv[2:]
|
||||
return parts[marker + 1], parts[marker + 2], argv
|
||||
|
||||
|
||||
def test_local_delivery_command_and_ack(tmp_path, monkeypatch):
|
||||
@@ -355,6 +358,54 @@ def test_spawn_failure_reports_error(tmp_path, monkeypatch):
|
||||
assert "could not be started" in result["error"]
|
||||
|
||||
|
||||
def test_live_dm_admitted_before_waiter_failure(tmp_path, monkeypatch):
|
||||
from tools import bot_live_delivery as live
|
||||
|
||||
home = _managed_home(tmp_path)
|
||||
target = home / "profiles" / "researcher"
|
||||
owner = dict(profile_home=str(target), session_id="bot", lease_id="lease", live_session_id="live")
|
||||
monkeypatch.setattr(live, "find_canonical_live_owner", lambda h: owner if Path(h) == target else None)
|
||||
monkeypatch.setattr(bot_mode_dm, "_dm_dir", lambda: tmp_path)
|
||||
monkeypatch.setenv("HERMES_HOME", str(tmp_path / "wrong-home"))
|
||||
import tools.terminal_tool as terminal
|
||||
monkeypatch.setattr(terminal, "terminal_tool", lambda *a, **k: json.dumps({"error": "spawn failed"}))
|
||||
|
||||
result = json.loads(bot_mode_dm.message_agent_tool("researcher", "hello", agent=_FakeAgent(home)))
|
||||
assert result["status"] == "queued"
|
||||
record = live.read_delivery_result(target, result["delivery_id"])
|
||||
assert record is not None
|
||||
assert record["owner"] == owner
|
||||
assert record["message"] == "Message from 🤖 hermes (@hermes): hello"
|
||||
assert "notification_error" in result
|
||||
|
||||
|
||||
def test_live_dm_runner_retry_never_reexecutes_failed_claim(tmp_path, monkeypatch, capsys):
|
||||
from tools import bot_live_delivery as live
|
||||
|
||||
home = _managed_home(tmp_path)
|
||||
target = home / "profiles" / "researcher"
|
||||
monkeypatch.setenv("HERMES_HOME", str(home))
|
||||
owner = dict(profile_home=str(target), session_id="bot", lease_id="lease", live_session_id="live")
|
||||
monkeypatch.setattr(live, "find_canonical_live_owner", lambda h: owner)
|
||||
monkeypatch.setattr(bot_mode_dm, "_LIVE_WAIT_SECONDS", 0)
|
||||
monkeypatch.setattr(subprocess, "run", lambda *a, **k: pytest.fail("must not launch a model turn"))
|
||||
dm_file = tmp_path / "message.txt"
|
||||
dm_file.write_text("hello", encoding="utf-8")
|
||||
argv = ["hermes", "-p", "researcher"]
|
||||
assert bot_mode_dm._run_delivery(argv, str(dm_file), stdin_file=False) == 0
|
||||
queued = json.loads(capsys.readouterr().out)
|
||||
assert queued["status"] == "queued"
|
||||
claimed = live.claim_pending_delivery(target, owner)
|
||||
assert claimed is not None
|
||||
live.complete_delivery(target, claimed["delivery_id"], status="failed", error="HTTP 429 rate limit")
|
||||
monkeypatch.setattr(live, "find_canonical_live_owner", lambda h: None)
|
||||
assert bot_mode_dm._run_delivery(argv, str(dm_file), stdin_file=False) == 1
|
||||
failed = json.loads(capsys.readouterr().out)
|
||||
assert failed["status"] == "failed"
|
||||
assert failed["delivery_id"] == queued["delivery_id"]
|
||||
assert dm_file.read_text(encoding="utf-8") == "hello"
|
||||
|
||||
|
||||
# ── plaintext tempfile lifecycle ─────────────────────────────────────────────
|
||||
|
||||
|
||||
|
||||
+117
-10
@@ -16,6 +16,7 @@ local → ``hermes -p <name> chat --in ~ -c "Bot Chat" --create-if-missing -Q
|
||||
from __future__ import annotations
|
||||
|
||||
import contextlib
|
||||
import hashlib
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
@@ -43,6 +44,7 @@ MESSAGE_MAX_CHARS = 16000
|
||||
# the machine dies between spawn ack and the runner's finally.
|
||||
_DM_DIR_NAME = "hermes-dm"
|
||||
_DM_STALE_SECONDS = 24 * 60 * 60
|
||||
_LIVE_WAIT_SECONDS = 300
|
||||
|
||||
# '<peer>/<agent>' — peer names are lowercase (``hermes peer`` normalizes them).
|
||||
_PEER_TARGET_RE = re.compile(r"^([a-z0-9][a-z0-9_-]{0,63})/([a-zA-Z0-9][a-zA-Z0-9_-]{0,63})$")
|
||||
@@ -192,7 +194,8 @@ def message_agent_tool(target: str = "", message: str = "", task_id: Optional[st
|
||||
return _err(f"Bot Mode gate check failed: {exc}")
|
||||
|
||||
root, me = _hermes_root(Path(home)), _self_profile_name(Path(home))
|
||||
roster = [name for name, _dir in _roster(root)]
|
||||
roster_homes = dict(_roster(root))
|
||||
roster = list(roster_homes)
|
||||
peers = _peers(root)
|
||||
teammates = [_handle(n) for n in roster if n != me]
|
||||
|
||||
@@ -243,7 +246,7 @@ def message_agent_tool(target: str = "", message: str = "", task_id: Optional[st
|
||||
"machine, or on a registered peer. Pick a name from the roster "
|
||||
"(roles are listed in your system prompt).")
|
||||
return _start_delivery(["hermes", "-p", resolved, *BOT_CHAT_TURN_ARGS], content, f"@{_handle(resolved)}",
|
||||
stdin_file=False, **delivery)
|
||||
stdin_file=False, profile_home=roster_homes[resolved], **delivery)
|
||||
|
||||
|
||||
def _try_relay_delivery(root: Path, raw_target: str, content: str, me: str, *,
|
||||
@@ -396,9 +399,73 @@ def _run_local_turn(argv: list[str], dm_file: str) -> int:
|
||||
return proc.returncode
|
||||
|
||||
|
||||
def _run_delivery(argv: list[str], dm_file: str, *, stdin_file: bool) -> int:
|
||||
"""Run one DM transport and remove its plaintext file after consumption. The turn
|
||||
window (not the enqueue) holds the target profile's cross-process lock, so two
|
||||
def _admit_live_dm(profile_home: Path | None, dm_file: str) -> dict | None:
|
||||
"""Pin intent before admission; retries may inspect, never change transport."""
|
||||
from tools.bot_live_delivery import (
|
||||
_fsync_dir, deliver_to_live_owner, find_canonical_live_owner, read_delivery_result,
|
||||
)
|
||||
|
||||
intent: dict[str, Any]
|
||||
intent_path = Path(dm_file + ".live.json")
|
||||
if intent_path.exists():
|
||||
intent = json.loads(intent_path.read_text(encoding="utf-8"))
|
||||
else:
|
||||
assert profile_home is not None
|
||||
owner = find_canonical_live_owner(profile_home)
|
||||
if owner is None:
|
||||
return None
|
||||
intent = dict(owner=owner, message=Path(dm_file).read_text(encoding="utf-8"),
|
||||
delivery_id=hashlib.sha256(str(Path(dm_file).resolve()).encode()).hexdigest())
|
||||
try:
|
||||
fd = os.open(intent_path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
|
||||
except FileExistsError:
|
||||
intent = json.loads(intent_path.read_text(encoding="utf-8"))
|
||||
else:
|
||||
with os.fdopen(fd, "w", encoding="utf-8") as stream:
|
||||
json.dump(intent, stream)
|
||||
stream.flush()
|
||||
os.fsync(stream.fileno())
|
||||
_fsync_dir(intent_path.parent)
|
||||
home = intent["owner"]["profile_home"]
|
||||
record = read_delivery_result(home, intent["delivery_id"])
|
||||
if record is None:
|
||||
record = deliver_to_live_owner(home, intent["owner"], intent["message"],
|
||||
delivery_id=intent["delivery_id"])
|
||||
return record
|
||||
|
||||
|
||||
def _wait_live_dm(home: str, delivery_id: str) -> int:
|
||||
from tools.bot_live_delivery import read_delivery_result
|
||||
|
||||
deadline = time.monotonic() + _LIVE_WAIT_SECONDS
|
||||
while True:
|
||||
record = read_delivery_result(home, delivery_id)
|
||||
status = record["status"] if record else "ambiguous"
|
||||
if status not in ("queued", "claimed") or time.monotonic() >= deadline:
|
||||
break
|
||||
time.sleep(min(0.5, max(0, deadline - time.monotonic())))
|
||||
payload = {key: record[key] for key in ("reply", "error", "reason") if record and record.get(key)}
|
||||
payload.update(status=status, delivery_id=delivery_id)
|
||||
if status in ("queued", "claimed", "ambiguous"):
|
||||
payload["detail"] = "Delivery remains pending or its outcome is unknown. Do not resend; receipt is retained."
|
||||
print(json.dumps(payload))
|
||||
return 0 if status in ("settled", "queued", "claimed") else 1
|
||||
|
||||
|
||||
def _local_delivery_home(argv: list[str]) -> Path | None:
|
||||
cli = (argv[0] if argv else "").rsplit("\\", 1)[-1].rsplit("/", 1)[-1]
|
||||
if len(argv) < 3 or cli not in ("hermes", "hermes.exe") or argv[1] != "-p":
|
||||
return None
|
||||
from tools.bot_mode_probe import _hermes_root, _roster
|
||||
|
||||
return dict(_roster(_hermes_root(Path(_default_home())))).get(argv[2])
|
||||
|
||||
|
||||
def _run_delivery(argv: list[str], dm_file: str, *, stdin_file: bool,
|
||||
profile_home: Path | None = None) -> int:
|
||||
"""Route to the live owner before attempting a CLI transport. Live deliveries
|
||||
retain their intent/payload and immutable receipt; only CLI/peer payloads are
|
||||
removed after consumption. The CLI turn window holds the profile lock, so two
|
||||
deliveries into one profile queue; a bounded wait ends in a 'target_busy' refusal.
|
||||
|
||||
Local (query-file) turns get one policy-gated retry (#93091 item 5): transient failures re-run the same
|
||||
@@ -407,6 +474,20 @@ def _run_delivery(argv: list[str], dm_file: str, *, stdin_file: bool) -> int:
|
||||
ever minted. Auth/quota/config failures never retry. Peer transports (stdin mode) retry on their own
|
||||
gateway's deliver path, not here.
|
||||
"""
|
||||
# The live consumer owns turn admission; never compete for its CLI lease.
|
||||
if not stdin_file:
|
||||
home = profile_home or _local_delivery_home(argv)
|
||||
if home is not None or Path(dm_file + ".live.json").exists():
|
||||
try:
|
||||
record = _admit_live_dm(home, dm_file)
|
||||
except Exception as exc:
|
||||
print(json.dumps({"status": "ambiguous", "delivery_id": hashlib.sha256(
|
||||
str(Path(dm_file).resolve()).encode()).hexdigest(),
|
||||
"error": f"Live admission outcome unknown: {exc}. Do not resend.",
|
||||
"evidence_file": dm_file}))
|
||||
return 1
|
||||
if record is not None:
|
||||
return _wait_live_dm(record["profile_home"], record["delivery_id"])
|
||||
try:
|
||||
with _delivery_lock(argv, stdin_file=stdin_file):
|
||||
if not stdin_file:
|
||||
@@ -419,10 +500,14 @@ def _run_delivery(argv: list[str], dm_file: str, *, stdin_file: bool) -> int:
|
||||
_unlink_dm_file(dm_file)
|
||||
|
||||
|
||||
def _delivery_command(argv: list[str], dm_file: str, *, stdin_file: bool) -> str:
|
||||
def _delivery_command(argv: list[str], dm_file: str, *, stdin_file: bool,
|
||||
profile_home: Path | None = None) -> str:
|
||||
"""Build an argv-safe command for the cleanup-owning background runner."""
|
||||
runner_argv = [sys.executable, str(Path(__file__).resolve()), "--run-delivery",
|
||||
"stdin" if stdin_file else "query-file", dm_file, *argv]
|
||||
"stdin" if stdin_file else "query-file", dm_file]
|
||||
if profile_home is not None:
|
||||
runner_argv.extend(["--profile-home", str(Path(profile_home).resolve())])
|
||||
runner_argv.extend(argv)
|
||||
if sys.platform == "win32":
|
||||
# The tracked local backend uses Git Bash on native Windows: forward slashes keep drive
|
||||
# paths executable there; backslash paths are parsed as command names (exit 127).
|
||||
@@ -431,11 +516,29 @@ def _delivery_command(argv: list[str], dm_file: str, *, stdin_file: bool) -> str
|
||||
|
||||
|
||||
def _start_delivery(argv: list[str], content: str, label: str, *, stdin_file: bool,
|
||||
task_id: Optional[str], agent: Any) -> str:
|
||||
task_id: Optional[str], agent: Any, profile_home: Path | None = None) -> str:
|
||||
"""Create a DM file and transfer its cleanup ownership to the runner."""
|
||||
dm_file = _write_dm_file(content)
|
||||
if profile_home is not None:
|
||||
try:
|
||||
record = _admit_live_dm(profile_home, dm_file)
|
||||
except Exception as exc:
|
||||
return json.dumps({"status": "ambiguous", "delivery_id": hashlib.sha256(
|
||||
str(Path(dm_file).resolve()).encode()).hexdigest(),
|
||||
"error": f"Live delivery admission could not be confirmed: {exc}. Do not resend.",
|
||||
"evidence_file": dm_file})
|
||||
if record is not None:
|
||||
command = _delivery_command(argv, dm_file, stdin_file=False, profile_home=profile_home)
|
||||
notification = json.loads(_spawn_delivery(command, label, task_id=task_id, agent=agent))
|
||||
result = dict(status=record["status"], delivery_id=record["delivery_id"], to=label,
|
||||
detail="Durably queued for the live Bot Chat owner. Do NOT wait or resend; finish your turn.")
|
||||
if notification.get("error"):
|
||||
result["notification_error"] = notification["error"]
|
||||
elif notification.get("process_id"):
|
||||
result["process_id"] = notification["process_id"]
|
||||
return json.dumps(result)
|
||||
try:
|
||||
command = _delivery_command(argv, dm_file, stdin_file=stdin_file)
|
||||
command = _delivery_command(argv, dm_file, stdin_file=stdin_file, profile_home=profile_home)
|
||||
except BaseException:
|
||||
_unlink_dm_file(dm_file)
|
||||
raise
|
||||
@@ -485,7 +588,10 @@ def _delivery_main(args: list[str]) -> int:
|
||||
if len(args) < 3 or args[0] != "--run-delivery" or args[1] not in ("stdin", "query-file"):
|
||||
return 2
|
||||
try:
|
||||
return _run_delivery(args[3:], args[2], stdin_file=args[1] == "stdin")
|
||||
argv, profile_home = args[3:], None
|
||||
if len(argv) >= 2 and argv[0] == "--profile-home":
|
||||
profile_home, argv = Path(argv[1]), argv[2:]
|
||||
return _run_delivery(argv, args[2], stdin_file=args[1] == "stdin", profile_home=profile_home)
|
||||
except Exception as exc:
|
||||
# 'target_busy': the queued delivery gave up after its bounded wait — surface the
|
||||
# structured payload on stdout so the completion notification carries it back.
|
||||
@@ -521,4 +627,5 @@ def _session_title(agent: Any) -> str:
|
||||
|
||||
|
||||
if __name__ == "__main__": # pragma: no cover - exercised as a background process
|
||||
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
|
||||
raise SystemExit(_delivery_main(sys.argv[1:]))
|
||||
|
||||
@@ -168,6 +168,8 @@ def _manual_run_delivery_note(deliver: str, refreshed: Dict[str, Any]) -> str:
|
||||
return " (output saved locally only)"
|
||||
err = str(refreshed.get("last_delivery_error") or "").strip()
|
||||
if not err:
|
||||
if refreshed.get("last_delivery_queued"):
|
||||
return " (output queued for Bot Chat; completion unverified, do not resend)"
|
||||
return " (output was delivered there by the job itself)"
|
||||
return f" (⚠ delivery FAILED: {err[:200]})"
|
||||
|
||||
@@ -317,7 +319,7 @@ def _run_claimed_job(job: Dict[str, Any], extra_prompt: Optional[str] = None) ->
|
||||
# That is NOT a success for the caller — the calling agent relays this result — so report it as
|
||||
# failed and surface the delivery error, which lives in last_delivery_error (last_error is None for
|
||||
# these runs, and a bare success=False with error=None reads as an unexplained failure). See #83993.
|
||||
ok = last_status == "ok"
|
||||
ok = last_status in {"ok", "delivery_queued"}
|
||||
if execution is not None and execution.get("status") != "completed":
|
||||
ok = False
|
||||
run_error = execution.get("error") or f"execution ended in {execution.get('status') or 'unknown'} state"
|
||||
|
||||
@@ -286,7 +286,7 @@ Platforms in the first group have explicit, validated target syntax — named ch
|
||||
|
||||
For **Telegram topics**, use `telegram:<chat_id>:<thread_id>` (e.g., `telegram:-1001234567890:17585`). For **Slack threads**, the third segment is the parent message's `thread_ts` (e.g., `slack:C0123ABCD45:1700000000.000100`), so it only applies when replying under an existing message.
|
||||
|
||||
**Bot Chat** (`bot-chat`, `bot-chat:<profile>`) is a machine-local pseudo-platform, not a gateway adapter: the scheduler delivers by running `hermes [-p <profile>] chat --in ~ -c "Bot Chat" --create-if-missing -Q --query-file <tmp>` — the same lane Bot Mode agent-to-agent messages use — so the output arrives as a real inbound turn in the profile's canonical Bot Chat and the bot runs a full agent turn on it (alternation-safe by construction; this is the chat command lane, not a transcript mirror). The bare token targets the job's own profile; the named form is validated against `~/.hermes/profiles/` at create time and again at fire time, and never resolves across machines. Bot-chat targets are excluded from the `all` routing token and from delivery preflight (no gateway credentials involved). The per-delivery subprocess timeout is `cron.bot_chat_delivery_timeout_seconds` (default 600).
|
||||
**Bot Chat** (`bot-chat`, `bot-chat:<profile>`) is a machine-local pseudo-platform, not a gateway adapter. A mailbox-capable canonical live owner receives durable admission immediately (idle or busy); only that owner executes the incoming turn. `scheduler_delivery._deliver_to_bot_chat` resolves the target with `get_profile_dir` or the job's current `get_hermes_home`, derives the receipt ID from the source home, job ID, durable `execution_id`, and target home, and checks the receipt before discovering an owner. An existing receipt never permits CLI fallback. Without a mailbox owner it retains `hermes [-p <profile>] chat --in ~ -c "Bot Chat" --create-if-missing -Q --query-file <tmp>` and normal ownership fencing. Both lanes deliver a real inbound turn, not a transcript mirror. Queued/claimed receipts populate `last_delivery_queued` with receipt IDs. The delivery aggregator excludes admission notices from genuine errors and records execution `delivery_outcome=queued`; successful jobs use `last_status=delivery_queued`. Genuine errors on mixed targets take precedence as failed while retaining queued receipt metadata. The target profile’s durable receipt is authoritative for terminal completion. Queued is the historical admission outcome, not proof of delivery. Historical cron status does not automatically track later receipt completion. Bot-chat targets are excluded from `all` and credential preflight. Bot-chat-only external workers bypass the gateway delivery queue; mixed external-worker targets retain gateway handoff. `cron.bot_chat_delivery_timeout_seconds` (default 600) bounds only the legacy subprocess lane.
|
||||
|
||||
### Response Wrapping
|
||||
|
||||
|
||||
@@ -109,6 +109,8 @@ Bots message each other with attribution, and you can hand work off from any cha
|
||||
- **Renamed Bots keep their tags in sync** — give a Bot a friendly name (the pencil in its chat header, or `hermes profile rename`) and it becomes taggable by that name: a Bot titled *Research Buddy* answers to `@research-buddy` (and `@researchbuddy`), in regular chats and in group rooms alike. The composer's `@` autocomplete offers the renamed tag and also matches when you type the old profile name, which keeps resolving too.
|
||||
- **Direct messages** — every Bot Chat carries the `message_agent` tool: a Bot messages a teammate by calling `message_agent(target="researcher", message="…")`. The tool validates the target against the live roster, prefixes the sender's `Message from 🤖 <sender> (@<sender>):` attribution automatically, and delivers into the teammate's canonical Bot Chat. Delivery is **fire-and-forget**: the sender gets an acknowledgement, finishes its turn, and the reply arrives later as a background completion notification. The message travels as a real parameter (nothing shell-interpreted — quotes, `$(...)`, and backticks arrive verbatim), and the Bot composes its own message rather than forwarding your words. The teammate roster — names **and roles** from each profile's title/description — is part of every Bot Chat's system prompt, so Bots know who does what before choosing a recipient. The tool exists **only** in canonical Bot Chat sessions on Bot-Mode-managed installs; regular chats, group-room member sessions, and CLI sessions never see it.
|
||||
|
||||
Local messages also reach a Bot Chat that stays open in Desktop or the TUI. The receiving backend keeps ownership: it reads durable ingress on its existing notification poller, admits immediately when idle, or waits until the running turn and already queued human prompts finish. A `queued` acknowledgement confirms durable admission, **not** a completed reply. The target profile retains the delivery ID and receipt under `runtime/bot_live_delivery/`; `settled` confirms completion. A crashed or cancelled imported turn is not automatically replayed, and pending work pinned to a departed owner remains inspectable rather than being silently rerun. Do not resend a delivery whose outcome is unknown. Older backends without live-delivery capability retain the existing ownership refusal; restart that backend after upgrading.
|
||||
|
||||
The backend teaches each Bot's canonical Bot Chat session the messaging protocol automatically at prompt-build time — including when a teammate opens it headlessly from the CLI. Only the canonical Bot Chat gets the protocol section; your regular sessions and your SOUL.md stay untouched. This is controlled by `agent.bot_mode_protocol` in `config.yaml` (default: on):
|
||||
|
||||
```yaml
|
||||
|
||||
@@ -496,6 +496,9 @@ 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. Without a live mailbox owner, the existing `hermes chat -c "Bot Chat" --create-if-missing` lane remains available (normal session ownership checks still apply).
|
||||
- **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.
|
||||
|
||||
### Routing intent (`all`)
|
||||
|
||||
|
||||
Reference in New Issue
Block a user