feat(bot-mode): carry the sender through local and relay deliveries
a bot dm arrived as an ordinary user message. the only trace of the sender
was the "Message from" text prefix, which the model reads and nothing else
does. the recipient's memory provider saw its own configured user.
message_agent now passes the sender as {"id": "bot:<profile>", "name":
<handle>, "is_bot": true} to the delivery runner (--author <json>), which sets
HERMES_TURN_AUTHOR on the recipient one-shot only. the -Q turn reads it and
passes turn_author into run_conversation. the desktop relay forwards the
envelope's from_profile/from_handle to bot_relay.deliver, which sets the same
variable on its delivery turn. the runner drops any inherited author first so
a delivery without one stays unattributed. the text prefix is unchanged.
This commit is contained in:
@@ -125,6 +125,8 @@ interface RelayAgentRow {
|
||||
interface RelayEnvelope {
|
||||
id?: string
|
||||
message?: string
|
||||
from_profile?: string
|
||||
from_handle?: string
|
||||
target_connection?: string
|
||||
target_profile?: string
|
||||
}
|
||||
@@ -399,7 +401,10 @@ async function drainRelayOutboxes() {
|
||||
'bot_relay.deliver',
|
||||
{
|
||||
profile: String(envelope?.target_profile || ''),
|
||||
message: String(envelope?.message || '')
|
||||
message: String(envelope?.message || ''),
|
||||
// The target gateway sets HERMES_TURN_AUTHOR on the delivery turn from these two.
|
||||
from_profile: String(envelope?.from_profile || ''),
|
||||
from_handle: String(envelope?.from_handle || '')
|
||||
},
|
||||
RELAY_DELIVER_TIMEOUT_MS
|
||||
)
|
||||
|
||||
@@ -4055,9 +4055,15 @@ def _sync_cli_session_id_from_agent(cli) -> None:
|
||||
|
||||
|
||||
def _run_quiet_single_query(cli, effective_query):
|
||||
"""Quiet (-Q) one-shot turn: run, print the response (stderr for errors/session_id), then sys.exit with the automation exit code."""
|
||||
"""Quiet (-Q) one-shot turn: run, print the response (stderr for errors/session_id), then sys.exit with the automation exit code.
|
||||
The turn's author comes from HERMES_TURN_AUTHOR, which only a bot-to-bot dispatcher sets on this subprocess."""
|
||||
from agent.turn_author import turn_author_from_env
|
||||
|
||||
try:
|
||||
result = cli.agent.run_conversation(user_message=effective_query, conversation_history=cli.conversation_history)
|
||||
result = cli.agent.run_conversation(
|
||||
user_message=effective_query, conversation_history=cli.conversation_history,
|
||||
turn_author=turn_author_from_env(),
|
||||
)
|
||||
except KeyboardInterrupt:
|
||||
_emit_interrupted_session_end(cli, reason="keyboard_interrupt")
|
||||
print(f"\nsession_id: {cli.session_id}", file=sys.stderr)
|
||||
|
||||
@@ -0,0 +1,59 @@
|
||||
"""``hermes chat -Q`` passes the dispatcher's HERMES_TURN_AUTHOR to ``run_conversation`` as ``turn_author``.
|
||||
|
||||
A bot-to-bot delivery runs the recipient's turn as a ``-Q`` subprocess with that variable set;
|
||||
a human's ``-Q`` run has it unset and the turn stays unattributed.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
from types import SimpleNamespace
|
||||
|
||||
import pytest
|
||||
|
||||
import cli
|
||||
from agent.turn_author import TURN_AUTHOR_ENV
|
||||
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
def _plain_one_shot_env(monkeypatch):
|
||||
monkeypatch.delenv("HERMES_KANBAN_GOAL_MODE", raising=False)
|
||||
monkeypatch.delenv("HERMES_KANBAN_TASK", raising=False)
|
||||
|
||||
|
||||
def _fake_cli(recorded):
|
||||
def run_conversation(**kwargs):
|
||||
recorded.append(kwargs)
|
||||
return {"final_response": "ok"}
|
||||
|
||||
agent = SimpleNamespace(run_conversation=run_conversation, session_id="s-1")
|
||||
return SimpleNamespace(agent=agent, conversation_history=[], session_id="s-1")
|
||||
|
||||
|
||||
def _run(monkeypatch, env_value):
|
||||
if env_value is None:
|
||||
monkeypatch.delenv(TURN_AUTHOR_ENV, raising=False)
|
||||
else:
|
||||
monkeypatch.setenv(TURN_AUTHOR_ENV, env_value)
|
||||
recorded = []
|
||||
with pytest.raises(SystemExit) as exc:
|
||||
cli._run_quiet_single_query(_fake_cli(recorded), "hello")
|
||||
assert exc.value.code == 0
|
||||
assert len(recorded) == 1
|
||||
assert recorded[0]["user_message"] == "hello"
|
||||
return recorded[0]
|
||||
|
||||
|
||||
def test_quiet_one_shot_passes_turn_author_from_env(monkeypatch, capsys):
|
||||
author = {"id": "bot:coder", "name": "coder", "is_bot": True}
|
||||
kwargs = _run(monkeypatch, json.dumps(author))
|
||||
assert kwargs["turn_author"] == author
|
||||
assert capsys.readouterr().out.strip() == "ok"
|
||||
|
||||
|
||||
def test_quiet_one_shot_without_env_passes_none(monkeypatch):
|
||||
assert _run(monkeypatch, None)["turn_author"] is None
|
||||
|
||||
|
||||
def test_quiet_one_shot_junk_env_passes_none(monkeypatch):
|
||||
assert _run(monkeypatch, "not json")["turn_author"] is None
|
||||
@@ -212,14 +212,26 @@ def _capture_spawn(monkeypatch):
|
||||
|
||||
|
||||
def _runner_parts(command):
|
||||
"""(mode, dm_file, transport argv) of a runner command; the optional ``--author <json>`` pair is skipped."""
|
||||
parts = shlex.split(command)
|
||||
marker = parts.index("--run-delivery")
|
||||
if parts[marker + 1] == "--author":
|
||||
marker += 2
|
||||
argv = parts[marker + 3 :]
|
||||
if argv[:1] == ["--profile-home"]:
|
||||
argv = argv[2:]
|
||||
return parts[marker + 1], parts[marker + 2], argv
|
||||
|
||||
|
||||
def _runner_author(command):
|
||||
"""The parsed ``--author`` payload of a runner command, or None when the pair is absent."""
|
||||
parts = shlex.split(command)
|
||||
marker = parts.index("--run-delivery")
|
||||
if parts[marker + 1] != "--author":
|
||||
return None
|
||||
return json.loads(parts[marker + 2])
|
||||
|
||||
|
||||
def test_local_delivery_command_and_ack(tmp_path, monkeypatch):
|
||||
calls = _capture_spawn(monkeypatch)
|
||||
home = _managed_home(tmp_path, teammates=("researcher",))
|
||||
@@ -264,6 +276,8 @@ def test_local_delivery_command_and_ack(tmp_path, monkeypatch):
|
||||
# message body rides the temp file, never the command line
|
||||
assert "PAYLOAD_SENTINEL_7A91" not in command
|
||||
assert "$(" not in command
|
||||
# the sender rides the runner argv as a stable id plus display handle
|
||||
assert _runner_author(command) == {"id": "bot:default", "name": "hermes", "is_bot": True}
|
||||
|
||||
# attribution prefix applied server-side; body verbatim inside the file
|
||||
content = Path(dm_file).read_text(encoding="utf-8")
|
||||
@@ -313,6 +327,8 @@ def test_peer_delivery_command(tmp_path, monkeypatch):
|
||||
mode, _dm_file, transport_argv = _runner_parts(calls[0]["command"])
|
||||
assert mode == "stdin"
|
||||
assert transport_argv == ["hermes", "-p", "default", "peer", "dm", "spark/researcher"]
|
||||
# the peer child reads the author from its env and forwards it in the request body
|
||||
assert _runner_author(calls[0]["command"]) == {"id": "bot:default", "name": "hermes", "is_bot": True}
|
||||
|
||||
# bare peer name targets the peer's main agent
|
||||
result2 = json.loads(
|
||||
@@ -339,6 +355,32 @@ def test_named_profile_sender_prefix(tmp_path, monkeypatch):
|
||||
assert Path(dm_file).read_text(encoding="utf-8").startswith(
|
||||
"Message from 🤖 coder (@coder): "
|
||||
)
|
||||
assert _runner_author(calls[0]["command"]) == {"id": "bot:coder", "name": "coder", "is_bot": True}
|
||||
|
||||
|
||||
def test_delivery_command_author_json_survives_quoting_and_windows_slash_rewrite(tmp_path, monkeypatch):
|
||||
"""The author pair sits between ``--run-delivery`` and the mode; its JSON must come back
|
||||
byte-exact through shlex, and the Windows forward-slash rewrite must leave it alone."""
|
||||
author = {"id": "bot:default", "name": "hermes", "is_bot": True}
|
||||
command = bot_mode_dm._delivery_command(["hermes", "-p", "x"], str(tmp_path / "dm.txt"),
|
||||
stdin_file=False, author=author)
|
||||
parts = shlex.split(command)
|
||||
assert parts[2:4] == ["--run-delivery", "--author"]
|
||||
assert json.loads(parts[4]) == author
|
||||
assert parts[5] == "query-file"
|
||||
assert _runner_parts(command) == ("query-file", str(tmp_path / "dm.txt"), ["hermes", "-p", "x"])
|
||||
|
||||
monkeypatch.setattr(sys, "platform", "win32")
|
||||
command = bot_mode_dm._delivery_command(["hermes", "-p", "x"], "C:\\Users\\me\\dm.txt",
|
||||
stdin_file=False, author={"id": "bot:default", "name": 'q"q', "is_bot": True})
|
||||
parts = shlex.split(command)
|
||||
assert json.loads(parts[4]) == {"id": "bot:default", "name": 'q"q', "is_bot": True}
|
||||
assert parts[6] == "C:/Users/me/dm.txt"
|
||||
|
||||
# without an author the legacy shape is produced unchanged
|
||||
command = bot_mode_dm._delivery_command(["hermes"], "dm.txt", stdin_file=True)
|
||||
assert shlex.split(command)[2:4] == ["--run-delivery", "stdin"]
|
||||
assert _runner_author(command) is None
|
||||
|
||||
|
||||
def test_spawn_failure_reports_error(tmp_path, monkeypatch):
|
||||
@@ -526,9 +568,96 @@ def test_query_file_delivery_closes_stdin_for_initial_attempt_and_retry(
|
||||
assert not dm_file.exists()
|
||||
|
||||
|
||||
@pytest.mark.parametrize("args", [[], ["--run-delivery"], ["--run-delivery", "bad", "x"]])
|
||||
def test_delivery_main_rejects_invalid_cli(args):
|
||||
@pytest.mark.parametrize(
|
||||
"args",
|
||||
[
|
||||
[],
|
||||
["--run-delivery"],
|
||||
["--run-delivery", "bad", "x"],
|
||||
["--run-delivery", "--author"],
|
||||
["--run-delivery", "--author", '{"id":"bot:x","is_bot":true}'],
|
||||
["--run-delivery", "--author", "not json", "query-file", "x", "hermes"],
|
||||
["--run-delivery", "--author", "[1]", "query-file", "x", "hermes"],
|
||||
],
|
||||
)
|
||||
def test_delivery_main_rejects_invalid_cli(args, monkeypatch):
|
||||
spawned = []
|
||||
monkeypatch.setattr(subprocess, "run", lambda *a, **k: spawned.append(a))
|
||||
assert bot_mode_dm._delivery_main(args) == 2
|
||||
assert not spawned
|
||||
|
||||
|
||||
@pytest.mark.parametrize("mode", ["stdin", "query-file"])
|
||||
def test_delivery_main_sets_turn_author_env_on_child(tmp_path, monkeypatch, mode):
|
||||
from agent.turn_author import TURN_AUTHOR_ENV
|
||||
|
||||
dm_file = tmp_path / "message.txt"
|
||||
dm_file.write_text("secret", encoding="utf-8")
|
||||
calls = []
|
||||
|
||||
def fake_run(argv, **kwargs):
|
||||
calls.append((argv, kwargs))
|
||||
return subprocess.CompletedProcess(argv, 0, stdout="", stderr="")
|
||||
|
||||
monkeypatch.setattr(subprocess, "run", fake_run)
|
||||
monkeypatch.setenv("HERMES_DM_TEST_MARKER", "kept")
|
||||
author = {"id": "bot:coder", "name": "coder", "is_bot": True}
|
||||
|
||||
returncode = bot_mode_dm._delivery_main(
|
||||
["--run-delivery", "--author", json.dumps(author), mode, str(dm_file), "hermes", "-p", "researcher"]
|
||||
)
|
||||
|
||||
assert returncode == 0
|
||||
assert len(calls) == 1
|
||||
argv, kwargs = calls[0]
|
||||
assert argv[:3] == ["hermes", "-p", "researcher"]
|
||||
assert json.loads(kwargs["env"][TURN_AUTHOR_ENV]) == author
|
||||
assert kwargs["env"]["HERMES_DM_TEST_MARKER"] == "kept"
|
||||
assert not dm_file.exists()
|
||||
|
||||
|
||||
def test_delivery_main_without_author_drops_inherited_turn_author(tmp_path, monkeypatch):
|
||||
"""A runner spawned from inside a bot's own turn inherits that turn's HERMES_TURN_AUTHOR;
|
||||
the legacy no-author shape must not pass it on as the recipient's author."""
|
||||
from agent.turn_author import TURN_AUTHOR_ENV
|
||||
|
||||
dm_file = tmp_path / "message.txt"
|
||||
dm_file.write_text("secret", encoding="utf-8")
|
||||
calls = []
|
||||
|
||||
def fake_run(argv, **kwargs):
|
||||
calls.append(kwargs)
|
||||
return subprocess.CompletedProcess(argv, 0, stdout="", stderr="")
|
||||
|
||||
monkeypatch.setattr(subprocess, "run", fake_run)
|
||||
monkeypatch.setenv(TURN_AUTHOR_ENV, json.dumps({"id": "bot:previous", "name": "previous", "is_bot": True}))
|
||||
|
||||
assert bot_mode_dm._delivery_main(["--run-delivery", "query-file", str(dm_file), "hermes"]) == 0
|
||||
assert TURN_AUTHOR_ENV not in calls[0]["env"]
|
||||
|
||||
|
||||
def test_real_delivery_command_round_trip_carries_author(tmp_path):
|
||||
"""End to end through a real subprocess: the runner argv built by ``_delivery_command`` puts
|
||||
HERMES_TURN_AUTHOR into the transport child's environment."""
|
||||
dm_file = tmp_path / "message.txt"
|
||||
dm_file.write_text("secret", encoding="utf-8")
|
||||
observed = tmp_path / "observed.txt"
|
||||
child = tmp_path / "child.py"
|
||||
child.write_text(
|
||||
"import os, pathlib, sys\n"
|
||||
"pathlib.Path(sys.argv[1]).write_text(os.environ.get('HERMES_TURN_AUTHOR', 'unset'), encoding='utf-8')\n",
|
||||
encoding="utf-8",
|
||||
)
|
||||
author = {"id": "bot:default", "name": "hermes", "is_bot": True}
|
||||
command = bot_mode_dm._delivery_command(
|
||||
[sys.executable, str(child), str(observed)], str(dm_file), stdin_file=False, author=author
|
||||
)
|
||||
|
||||
result = subprocess.run(shlex.split(command), check=False)
|
||||
|
||||
assert result.returncode == 0
|
||||
assert json.loads(observed.read_text(encoding="utf-8")) == author
|
||||
assert not dm_file.exists()
|
||||
|
||||
|
||||
@pytest.mark.parametrize("mode", ["stdin", "query-file"])
|
||||
|
||||
@@ -584,3 +584,28 @@ def test_message_agent_surfaces_runtime_offline_refusal(tmp_path, monkeypatch):
|
||||
assert "offline" in out.get("error", "")
|
||||
# fail-fast means no envelope was queued
|
||||
assert bot_relay.claim_pending_envelopes(home) == []
|
||||
|
||||
|
||||
# ── delivery turn author (HERMES_TURN_AUTHOR on the recipient turn) ──────────
|
||||
|
||||
|
||||
def test_delivery_turn_author_from_envelope_sender_fields():
|
||||
author = bot_relay.delivery_turn_author("ops", "ops-bot")
|
||||
assert author == {"id": "bot:ops", "name": "ops-bot", "is_bot": True}
|
||||
# The display name falls back to the profile; the id never comes from the handle.
|
||||
assert bot_relay.delivery_turn_author("ops", "") == {"id": "bot:ops", "name": "ops", "is_bot": True}
|
||||
assert bot_relay.delivery_turn_author("", "ops-bot") is None
|
||||
assert bot_relay.delivery_turn_author(None, None) is None
|
||||
|
||||
|
||||
def test_delivery_env_carries_only_the_given_author(monkeypatch):
|
||||
"""The dispatcher's own HERMES_TURN_AUTHOR never reaches the child: dropped without an author, replaced with one."""
|
||||
from agent.turn_author import TURN_AUTHOR_ENV
|
||||
|
||||
monkeypatch.setenv("HERMES_RELAY_TEST_MARKER", "kept")
|
||||
monkeypatch.setenv(TURN_AUTHOR_ENV, json.dumps({"id": "bot:previous", "name": "previous", "is_bot": True}))
|
||||
|
||||
assert TURN_AUTHOR_ENV not in bot_relay.delivery_env(None)
|
||||
env = bot_relay.delivery_env(bot_relay.delivery_turn_author("ops", "ops"))
|
||||
assert json.loads(env[TURN_AUTHOR_ENV]) == {"id": "bot:ops", "name": "ops", "is_bot": True}
|
||||
assert env["HERMES_RELAY_TEST_MARKER"] == "kept"
|
||||
|
||||
@@ -193,3 +193,43 @@ def test_deliver_write_failure_still_removes_tempfile(home, monkeypatch, tmp_pat
|
||||
assert "error" in err
|
||||
assert made, "mkstemp was never reached"
|
||||
assert not glob.glob(str(tmp_path / "hermes-relay-dm-*")), "tempfile leaked"
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def fake_runs(monkeypatch):
|
||||
"""Fake ``subprocess.run`` that records each call's kwargs; ``outcomes`` holds (returncode, stderr) per call."""
|
||||
calls, outcomes = [], []
|
||||
|
||||
def _fake_run(argv, **kwargs):
|
||||
calls.append(kwargs)
|
||||
code, err = outcomes.pop(0) if outcomes else (0, "")
|
||||
|
||||
class _Proc:
|
||||
returncode, stdout, stderr = code, "ok" if code == 0 else "", err
|
||||
|
||||
return _Proc()
|
||||
|
||||
monkeypatch.setattr("subprocess.run", _fake_run)
|
||||
return calls, outcomes
|
||||
|
||||
|
||||
@pytest.mark.parametrize("sender, expected", [
|
||||
({"from_profile": "scout", "from_handle": "scout"}, {"id": "bot:scout", "name": "scout", "is_bot": True}),
|
||||
({}, None),
|
||||
], ids=["sender fields", "no sender fields"])
|
||||
def test_deliver_child_env_carries_the_envelope_sender_on_every_attempt(home, monkeypatch, fake_runs, sender, expected):
|
||||
"""HERMES_TURN_AUTHOR on the child comes from the envelope's sender fields alone: the retry gets the same
|
||||
author, and without sender fields a stale author on the gateway's own environment never reaches the child."""
|
||||
from agent.turn_author import TURN_AUTHOR_ENV
|
||||
|
||||
calls, outcomes = fake_runs
|
||||
outcomes.extend([(1, "HTTP 429 rate limit"), (0, "")])
|
||||
monkeypatch.setenv("HERMES_RELAY_TEST_MARKER", "kept")
|
||||
monkeypatch.setenv(TURN_AUTHOR_ENV, json.dumps({"id": "bot:stale", "name": "stale", "is_bot": True}))
|
||||
|
||||
_result(srv._methods["bot_relay.deliver"](1, {"profile": "ops", "message": "ping", **sender}))
|
||||
|
||||
envs = [c["env"] for c in calls]
|
||||
assert len(envs) == 2
|
||||
assert [json.loads(e[TURN_AUTHOR_ENV]) if TURN_AUTHOR_ENV in e else None for e in envs] == [expected, expected]
|
||||
assert all(e["HERMES_RELAY_TEST_MARKER"] == "kept" for e in envs)
|
||||
|
||||
+40
-15
@@ -214,6 +214,8 @@ def message_agent_tool(target: str = "", message: str = "", task_id: Optional[st
|
||||
return _roster_err("target is required.")
|
||||
content = f"Message from 🤖 {_handle(me)} (@{_handle(me)}): " + body
|
||||
delivery = dict(task_id=task_id, agent=agent)
|
||||
# Attribution for the recipient's memory hooks; the text prefix above stays the human-facing signature.
|
||||
author = {"id": f"bot:{me}", "name": _handle(me), "is_bot": True}
|
||||
|
||||
# Peer target: '<peer>/<agent>' or a bare registered peer name.
|
||||
peer_match = _PEER_TARGET_RE.match(raw_target)
|
||||
@@ -226,7 +228,8 @@ def message_agent_tool(target: str = "", message: str = "", task_id: Optional[st
|
||||
# load_config(), while the roster above reads the machine-root config — the CLI must run
|
||||
# in that same profile or a secondary-profile bot sees an empty registry.
|
||||
return _start_delivery(["hermes", "-p", _self_profile_name(root), "peer", "dm", dm_target], content,
|
||||
f"@{peer_profile or peer_name} on peer '{peer_name}'", stdin_file=True, **delivery)
|
||||
f"@{peer_profile or peer_name} on peer '{peer_name}'", stdin_file=True,
|
||||
author=author, **delivery)
|
||||
|
||||
# Local teammate.
|
||||
is_local_shape = bool(_LOCAL_TARGET_RE.match(raw_target))
|
||||
@@ -246,7 +249,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, profile_home=roster_homes[resolved], **delivery)
|
||||
stdin_file=False, profile_home=roster_homes[resolved], author=author, **delivery)
|
||||
|
||||
|
||||
def _try_relay_delivery(root: Path, raw_target: str, content: str, me: str, *,
|
||||
@@ -355,7 +358,7 @@ def _delivery_lock(argv: list[str], *, stdin_file: bool):
|
||||
return acquire_turn_lock(_hermes_root(Path(_default_home())), argv[2])
|
||||
|
||||
|
||||
def _run_local_turn(argv: list[str], dm_file: str) -> int:
|
||||
def _run_local_turn(argv: list[str], dm_file: str, *, env: Optional[dict[str, str]] = None) -> int:
|
||||
"""One Bot Chat turn via ``--query-file`` (plus one policy-gated retry); re-emits
|
||||
the transport's streams and returns its exit code. Transient failures re-run the
|
||||
same session; a context_overflow re-run lets the retried turn's pre-API compaction
|
||||
@@ -363,7 +366,7 @@ def _run_local_turn(argv: list[str], dm_file: str) -> int:
|
||||
|
||||
def _turn():
|
||||
return subprocess.run([*argv, "--query-file", dm_file], check=False, stdin=subprocess.DEVNULL,
|
||||
capture_output=True, text=True)
|
||||
capture_output=True, text=True, env=env)
|
||||
|
||||
proc = _turn()
|
||||
if proc.returncode != 0:
|
||||
@@ -462,11 +465,13 @@ def _local_delivery_home(argv: list[str]) -> Path | None:
|
||||
|
||||
|
||||
def _run_delivery(argv: list[str], dm_file: str, *, stdin_file: bool,
|
||||
profile_home: Path | None = None) -> int:
|
||||
profile_home: Path | None = None, author: Optional[dict] = 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.
|
||||
``author`` rides to the child as HERMES_TURN_AUTHOR: the local turn reads it directly and
|
||||
``hermes peer dm`` forwards it in the request body.
|
||||
|
||||
Local (query-file) turns get one policy-gated retry (#93091 item 5): transient failures re-run the same
|
||||
session; a context_overflow re-run lets the retried turn's pre-API compaction pass compact the Bot Chat
|
||||
@@ -489,20 +494,24 @@ def _run_delivery(argv: list[str], dm_file: str, *, stdin_file: bool,
|
||||
if record is not None:
|
||||
return _wait_live_dm(record["profile_home"], record["delivery_id"])
|
||||
try:
|
||||
from tools.bot_relay import delivery_env
|
||||
|
||||
env = delivery_env(author)
|
||||
with _delivery_lock(argv, stdin_file=stdin_file):
|
||||
if not stdin_file:
|
||||
return _run_local_turn(argv, dm_file)
|
||||
return _run_local_turn(argv, dm_file, env=env)
|
||||
# Keep the file open until the transport exits; cleanup occurs
|
||||
# after subprocess.run returns, not merely after stdin reaches EOF.
|
||||
with open(dm_file, "r", encoding="utf-8") as stream:
|
||||
return subprocess.run(argv, stdin=stream, check=False).returncode
|
||||
return subprocess.run(argv, stdin=stream, check=False, env=env).returncode
|
||||
finally:
|
||||
_unlink_dm_file(dm_file)
|
||||
|
||||
|
||||
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."""
|
||||
profile_home: Path | None = None, author: Optional[dict] = None) -> str:
|
||||
"""Build an argv-safe command for the cleanup-owning background runner:
|
||||
``--run-delivery [--author <json>] <mode> <dm_file> [--profile-home <path>] <argv...>``."""
|
||||
runner_argv = [sys.executable, str(Path(__file__).resolve()), "--run-delivery",
|
||||
"stdin" if stdin_file else "query-file", dm_file]
|
||||
if profile_home is not None:
|
||||
@@ -512,11 +521,15 @@ def _delivery_command(argv: list[str], dm_file: str, *, stdin_file: bool,
|
||||
# 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).
|
||||
runner_argv = [part.replace("\\", "/") for part in runner_argv]
|
||||
if author:
|
||||
# Inserted after the slash rewrite: JSON escapes are backslashes too.
|
||||
runner_argv[3:3] = ["--author", json.dumps(author, separators=(",", ":"))]
|
||||
return shlex.join(runner_argv)
|
||||
|
||||
|
||||
def _start_delivery(argv: list[str], content: str, label: str, *, stdin_file: bool,
|
||||
task_id: Optional[str], agent: Any, profile_home: Path | None = None) -> str:
|
||||
task_id: Optional[str], agent: Any, profile_home: Path | None = None,
|
||||
author: Optional[dict] = 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:
|
||||
@@ -528,7 +541,7 @@ def _start_delivery(argv: list[str], content: str, label: str, *, stdin_file: bo
|
||||
"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)
|
||||
command = _delivery_command(argv, dm_file, stdin_file=False, profile_home=profile_home, author=author)
|
||||
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.")
|
||||
@@ -538,7 +551,7 @@ def _start_delivery(argv: list[str], content: str, label: str, *, stdin_file: bo
|
||||
result["process_id"] = notification["process_id"]
|
||||
return json.dumps(result)
|
||||
try:
|
||||
command = _delivery_command(argv, dm_file, stdin_file=stdin_file, profile_home=profile_home)
|
||||
command = _delivery_command(argv, dm_file, stdin_file=stdin_file, profile_home=profile_home, author=author)
|
||||
except BaseException:
|
||||
_unlink_dm_file(dm_file)
|
||||
raise
|
||||
@@ -585,13 +598,25 @@ def _spawn_delivery(command: str, label: str, *, dm_file: Optional[str] = None,
|
||||
|
||||
|
||||
def _delivery_main(args: list[str]) -> int:
|
||||
if len(args) < 3 or args[0] != "--run-delivery" or args[1] not in ("stdin", "query-file"):
|
||||
"""Runner entry: ``--run-delivery [--author <json>] <mode> <dm_file> [--profile-home <path>] <argv...>``.
|
||||
Malformed argv exits 2 without touching the DM file; only ``_delivery_command`` builds this argv."""
|
||||
if not args or args[0] != "--run-delivery":
|
||||
return 2
|
||||
rest, author = args[1:], None
|
||||
if rest[:1] == ["--author"]:
|
||||
from agent.turn_author import parse_turn_author
|
||||
|
||||
author = parse_turn_author(rest[1]) if len(rest) > 1 else None
|
||||
if author is None:
|
||||
return 2
|
||||
rest = rest[2:]
|
||||
if len(rest) < 2 or rest[0] not in ("stdin", "query-file"):
|
||||
return 2
|
||||
try:
|
||||
argv, profile_home = args[3:], None
|
||||
argv, profile_home = rest[2:], 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)
|
||||
return _run_delivery(argv, rest[1], stdin_file=rest[0] == "stdin", profile_home=profile_home, author=author)
|
||||
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.
|
||||
|
||||
@@ -384,6 +384,27 @@ def local_delivery_command(profile: str, query_file: str) -> list[str]:
|
||||
return [_hermes_cli(), "-p", profile, *BOT_CHAT_TURN_ARGS, "--query-file", query_file]
|
||||
|
||||
|
||||
def delivery_turn_author(from_profile: Any, from_handle: Any) -> Optional[dict]:
|
||||
"""The author of a relayed DM's recipient turn, built from the envelope's sender fields.
|
||||
None when the envelope names no sender, so an unattributed delivery stays unattributed."""
|
||||
profile = str(from_profile or "").strip()
|
||||
if not profile:
|
||||
return None
|
||||
return {"id": f"bot:{profile}", "name": str(from_handle or "").strip() or profile, "is_bot": True}
|
||||
|
||||
|
||||
def delivery_env(author: Optional[dict]) -> dict[str, str]:
|
||||
"""Environment for one delivery turn's ``hermes`` child. The dispatcher's own HERMES_TURN_AUTHOR is
|
||||
dropped first so a delivery without an author never inherits the author of the turn that sent it."""
|
||||
from agent.turn_author import TURN_AUTHOR_ENV, turn_author_env
|
||||
|
||||
env = dict(os.environ)
|
||||
env.pop(TURN_AUTHOR_ENV, None)
|
||||
if author:
|
||||
env.update(turn_author_env(author))
|
||||
return env
|
||||
|
||||
|
||||
# Two deliveries into the SAME profile must never run Bot Chat turns concurrently.
|
||||
# Deliveries are separate ``hermes`` subprocesses, so the lock is a per-profile
|
||||
# lockfile under ``<root>/bot_relay/locks/`` held with ``fcntl.flock`` for exactly
|
||||
|
||||
@@ -25,11 +25,11 @@ def _relay_root() -> Path:
|
||||
return home.parent.parent if home.parent.name == "profiles" else home
|
||||
|
||||
|
||||
def _run_delivery(profile: str, tmp: str) -> subprocess.CompletedProcess:
|
||||
def _run_delivery(profile: str, tmp: str, env: dict | None = None) -> subprocess.CompletedProcess:
|
||||
from tools.bot_relay import local_delivery_command
|
||||
return subprocess.run(
|
||||
local_delivery_command(profile, tmp), capture_output=True, text=True, encoding="utf-8",
|
||||
errors="replace", timeout=TURN_ATTEMPT_TIMEOUT_SECONDS)
|
||||
errors="replace", timeout=TURN_ATTEMPT_TIMEOUT_SECONDS, env=env)
|
||||
|
||||
|
||||
@method("bot_relay.roster.sync")
|
||||
@@ -102,6 +102,11 @@ def _(rid, params: dict, _root=_relay_root, _run=_run_delivery) -> dict:
|
||||
def _detail(p) -> str:
|
||||
return (p.stderr or p.stdout or "").strip()[-500:]
|
||||
|
||||
# The Desktop forwards the envelope's sender fields; the turn's author labels memory only
|
||||
# and grants nothing. Absent fields leave the turn unattributed, as before.
|
||||
from tools.bot_relay import delivery_env, delivery_turn_author
|
||||
turn_env = delivery_env(delivery_turn_author(params.get("from_profile"), params.get("from_handle")))
|
||||
|
||||
fd, tmp = tempfile.mkstemp(prefix="hermes-relay-dm-", suffix=".txt", text=True)
|
||||
try:
|
||||
with os.fdopen(fd, "w", encoding="utf-8") as f:
|
||||
@@ -113,7 +118,7 @@ def _(rid, params: dict, _root=_relay_root, _run=_run_delivery) -> dict:
|
||||
# turn timeout below — doubled when the retry policy grants one bounded re-run — so clients
|
||||
# calling bot_relay.deliver must tolerate ~1320s before assuming failure. See #93091.
|
||||
with acquire_turn_lock(root, resolved):
|
||||
proc = _run(resolved, tmp)
|
||||
proc = _run(resolved, tmp, turn_env)
|
||||
if proc.returncode != 0:
|
||||
# Retry policy: transient classes re-run the SAME session once; context_overflow
|
||||
# too — the retried turn's pre-API compaction pass compacts the over-threshold
|
||||
@@ -122,7 +127,7 @@ def _(rid, params: dict, _root=_relay_root, _run=_run_delivery) -> dict:
|
||||
from tools.bot_failure_reasons import (
|
||||
RETRY_NONE, classify_agent_error, retry_action)
|
||||
if retry_action(classify_agent_error(_detail(proc))) != RETRY_NONE:
|
||||
proc = _run(resolved, tmp)
|
||||
proc = _run(resolved, tmp, turn_env)
|
||||
finally:
|
||||
with contextlib.suppress(OSError):
|
||||
os.unlink(tmp)
|
||||
|
||||
Reference in New Issue
Block a user