fix(bot-mode): carry the relay sender into a live Bot Chat turn
When the target Bot Chat is already open on this gateway, the relay handler delivers through `prompt.submit` with `queued: true`, and that branch dropped the envelope's sender. The model still saw the text prefix, but the turn reached the agent unattributed, the exact case the subprocess branch fixes. The relay handler now stamps the author on the submit as a `DeliveryAuthor`, an in-process object a JSON client cannot build, so `prompt.submit` accepts it the way it accepts a hosted-room callback and refuses a dict with error 4124. The busy queue keeps an authored envelope in its own slot, the drain hands the author to the turn runner, and the runner passes it to an agent that declares the keyword. A plain prompt after an authored dm carries no author. Local deliveries to a desktop-owned Bot Chat take the live-owner mailbox instead. The admission intent and the mailbox record now carry the author, a retry under the same id with a different author is refused, and the owner gateway hands the author to the turn it runs. Isolated compute turns still run unattributed, because the compute-host frame has no author field.
This commit is contained in:
@@ -93,3 +93,16 @@ def test_only_canonical_capable_owner_receives_across_compression(tmp_path, capa
|
||||
finally:
|
||||
lease.release()
|
||||
db.close()
|
||||
|
||||
|
||||
def test_delivery_keeps_the_sender_and_refuses_a_different_one_under_the_same_id(tmp_path):
|
||||
from tools import bot_live_delivery as mailbox
|
||||
|
||||
owner = dict(profile_home=str(tmp_path.resolve()), session_id="chat", lease_id="lease", live_session_id="live")
|
||||
author = {"id": "bot:coder", "name": "coder", "is_bot": True}
|
||||
queued = mailbox.deliver_to_live_owner(tmp_path, owner, "hello", delivery_id="b" * 32, author=author)
|
||||
assert queued["author"] == author
|
||||
assert mailbox.deliver_to_live_owner(tmp_path, owner, "hello", delivery_id="b" * 32, author=author) == queued
|
||||
with pytest.raises(ValueError):
|
||||
mailbox.deliver_to_live_owner(tmp_path, owner, "hello", delivery_id="b" * 32, author={**author, "id": "bot:other"})
|
||||
assert "author" not in mailbox.deliver_to_live_owner(tmp_path, owner, "no sender", delivery_id="c" * 32)
|
||||
|
||||
@@ -359,8 +359,7 @@ def test_named_profile_sender_prefix(tmp_path, monkeypatch):
|
||||
|
||||
|
||||
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."""
|
||||
"""The author JSON sits between ``--run-delivery`` and the mode, survives shlex, and the Windows slash rewrite skips it."""
|
||||
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)
|
||||
@@ -418,6 +417,7 @@ def test_live_dm_admitted_before_waiter_failure(tmp_path, monkeypatch):
|
||||
assert record is not None
|
||||
assert record["owner"] == owner
|
||||
assert record["message"] == "Message from 🤖 hermes (@hermes): hello"
|
||||
assert record["author"] == {"id": "bot:default", "name": "hermes", "is_bot": True}
|
||||
assert "notification_error" in result
|
||||
|
||||
|
||||
@@ -587,8 +587,14 @@ def test_delivery_main_rejects_invalid_cli(args, monkeypatch):
|
||||
assert not spawned
|
||||
|
||||
|
||||
@pytest.mark.parametrize("mode", ["stdin", "query-file"])
|
||||
def test_delivery_main_sets_turn_author_env_on_child(tmp_path, monkeypatch, mode):
|
||||
@pytest.mark.parametrize("mode, author", [
|
||||
("stdin", {"id": "bot:coder", "name": "coder", "is_bot": True}),
|
||||
("query-file", {"id": "bot:coder", "name": "coder", "is_bot": True}),
|
||||
("query-file", None),
|
||||
], ids=["stdin", "query-file", "no author"])
|
||||
def test_delivery_main_child_env_carries_only_the_argv_author(tmp_path, monkeypatch, mode, author):
|
||||
"""The ``--author`` payload becomes HERMES_TURN_AUTHOR on the child. Without it the runner drops the
|
||||
variable it inherited from the sending bot's own turn instead of passing it on as the recipient's author."""
|
||||
from agent.turn_author import TURN_AUTHOR_ENV
|
||||
|
||||
dm_file = tmp_path / "message.txt"
|
||||
@@ -601,44 +607,22 @@ def test_delivery_main_sets_turn_author_env_on_child(tmp_path, monkeypatch, mode
|
||||
|
||||
monkeypatch.setattr(subprocess, "run", fake_run)
|
||||
monkeypatch.setenv("HERMES_DM_TEST_MARKER", "kept")
|
||||
author = {"id": "bot:coder", "name": "coder", "is_bot": True}
|
||||
monkeypatch.setenv(TURN_AUTHOR_ENV, json.dumps({"id": "bot:previous", "name": "previous", "is_bot": True}))
|
||||
author_args = ["--author", json.dumps(author)] if author else []
|
||||
|
||||
returncode = bot_mode_dm._delivery_main(
|
||||
["--run-delivery", "--author", json.dumps(author), mode, str(dm_file), "hermes", "-p", "researcher"]
|
||||
)
|
||||
["--run-delivery", *author_args, mode, str(dm_file), "hermes", "-p", "researcher"])
|
||||
|
||||
assert returncode == 0
|
||||
assert len(calls) == 1
|
||||
argv, kwargs = calls[0]
|
||||
[(argv, kwargs)] = calls
|
||||
assert argv[:3] == ["hermes", "-p", "researcher"]
|
||||
assert json.loads(kwargs["env"][TURN_AUTHOR_ENV]) == author
|
||||
assert kwargs["env"]["HERMES_DM_TEST_MARKER"] == "kept"
|
||||
assert (json.loads(kwargs["env"][TURN_AUTHOR_ENV]) if TURN_AUTHOR_ENV in kwargs["env"] else None) == author
|
||||
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."""
|
||||
"""Through a real subprocess, the runner argv built by ``_delivery_command`` sets HERMES_TURN_AUTHOR on the child."""
|
||||
dm_file = tmp_path / "message.txt"
|
||||
dm_file.write_text("secret", encoding="utf-8")
|
||||
observed = tmp_path / "observed.txt"
|
||||
|
||||
@@ -62,14 +62,15 @@ def test_imported_crash_marker_never_autocontinues(tmp_path):
|
||||
def test_local_work_blocks_mailbox_claim_without_consuming_envelope(monkeypatch, tmp_path):
|
||||
import tools.bot_live_delivery as mailbox
|
||||
owner = {"lease_id": "lease", "live_session_id": "live", "session_id": "chat"}
|
||||
pending = [{"id": "receipt", "message": "imported"}]
|
||||
author = {"id": "bot:coder", "name": "coder", "is_bot": True}
|
||||
pending = [{"id": "receipt", "message": "imported", "author": author}]
|
||||
monkeypatch.setattr(mailbox, "find_canonical_live_owner", lambda home: owner)
|
||||
monkeypatch.setattr(mailbox, "claim_pending_delivery", lambda home, pinned: pending.pop(0))
|
||||
receipts = []
|
||||
monkeypatch.setattr(mailbox, "complete_delivery", lambda *args, **kwargs: receipts.append((args, kwargs)))
|
||||
submitted = []
|
||||
def submit(rid, sid, session, text, **kwargs):
|
||||
submitted.append(text)
|
||||
submitted.append((text, kwargs.get("turn_author")))
|
||||
kwargs["terminal_callback"]({"status": "settled", "text": "reply"})
|
||||
return True
|
||||
poll = rebind(session_notifications._poll_bot_live_delivery_once, {
|
||||
@@ -87,6 +88,6 @@ def test_local_work_blocks_mailbox_claim_without_consuming_envelope(monkeypatch,
|
||||
assert poll("other-live", session) is False
|
||||
assert pending
|
||||
assert poll("live", session) is True
|
||||
assert submitted == ["imported"] and not pending
|
||||
assert submitted == [("imported", author)] and not pending
|
||||
assert receipts[0][0][1] == "receipt"
|
||||
assert receipts[0][1]["reply"] == "reply"
|
||||
|
||||
@@ -0,0 +1,124 @@
|
||||
"""A relayed dm into a live Bot Chat keeps its sender all the way into ``run_conversation``.
|
||||
|
||||
The relay handler stamps the author as an in-process ``DeliveryAuthor``; ``prompt.submit`` accepts only
|
||||
that object, so a dashboard client cannot claim a bot identity through the same RPC.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import threading
|
||||
import types
|
||||
|
||||
import pytest
|
||||
|
||||
from tools.bot_relay import DeliveryAuthor
|
||||
from tui_gateway import server as srv
|
||||
|
||||
AUTHOR = {"id": "bot:coder", "name": "coder", "is_bot": True}
|
||||
OTHER = {"id": "bot:writer", "name": "writer", "is_bot": True}
|
||||
|
||||
|
||||
def _result(resp):
|
||||
return resp["result"] if "result" in resp else resp
|
||||
|
||||
|
||||
def _session(agent=None, **extra):
|
||||
return {
|
||||
"agent": agent if agent is not None else types.SimpleNamespace(),
|
||||
"session_key": "gw-session-key", "history": [], "history_lock": threading.Lock(),
|
||||
"history_version": 0, "running": False, "attached_images": [], "image_counter": 0, "cols": 80,
|
||||
"slash_worker": None, "show_reasoning": False, "tool_progress_mode": "all", "inflight_turn": None,
|
||||
"transport": None, **extra,
|
||||
}
|
||||
|
||||
|
||||
class _InlineThread:
|
||||
def __init__(self, target=None, daemon=None, args=(), kwargs=None):
|
||||
self._target, self._args, self._kwargs = target, args, kwargs or {}
|
||||
|
||||
def start(self):
|
||||
if self._target is not None:
|
||||
self._target(*self._args, **self._kwargs)
|
||||
|
||||
def is_alive(self):
|
||||
return False
|
||||
|
||||
def join(self, timeout=None):
|
||||
return None
|
||||
|
||||
|
||||
@pytest.fixture()
|
||||
def turn_env(monkeypatch, tmp_path):
|
||||
monkeypatch.setattr(srv.threading, "Thread", _InlineThread)
|
||||
for name in ("_emit", "_wire_callbacks", "_sync_agent_model_with_config", "_register_session_cwd",
|
||||
"_tts_stream_begin", "_sync_session_key_after_compress"):
|
||||
monkeypatch.setattr(srv, name, lambda *a, **k: None)
|
||||
monkeypatch.setattr(srv, "_session_cwd", lambda session: str(tmp_path))
|
||||
monkeypatch.setattr(srv, "_get_usage", lambda agent: {})
|
||||
|
||||
|
||||
def test_live_relay_stamps_the_sender_as_a_delivery_author(tmp_path, monkeypatch):
|
||||
home = tmp_path / ".hermes"
|
||||
(home / "profiles" / "ops").mkdir(parents=True)
|
||||
monkeypatch.setenv("HERMES_HOME", str(home))
|
||||
submitted = []
|
||||
monkeypatch.setitem(srv._methods, "prompt.submit", lambda rid, p: submitted.append(p) or srv._ok(rid, {"status": "streaming"}))
|
||||
monkeypatch.setattr(srv, "_profile_home", lambda name: home / "profiles" / name)
|
||||
monkeypatch.setitem(srv._sessions, "live-ops",
|
||||
{"profile_home": str(home / "profiles" / "ops"), "pending_title": "Bot Chat", "history": []})
|
||||
|
||||
_result(srv._methods["bot_relay.deliver"](1, {"profile": "ops", "message": "ping", "from_profile": "coder", "from_handle": "coder"}))
|
||||
|
||||
assert submitted == [{"session_id": "live-ops", "text": "ping", "queued": True, "_turn_author": DeliveryAuthor(AUTHOR)}]
|
||||
|
||||
|
||||
def test_prompt_submit_refuses_a_client_supplied_author(monkeypatch):
|
||||
srv._sessions["sid"] = _session()
|
||||
try:
|
||||
resp = srv._methods["prompt.submit"]("r", {"session_id": "sid", "text": "hi", "_turn_author": dict(AUTHOR)})
|
||||
finally:
|
||||
srv._sessions.pop("sid", None)
|
||||
assert resp["error"]["code"] == 4124
|
||||
|
||||
|
||||
def test_busy_relay_dms_queue_with_their_authors_and_drain_with_them(monkeypatch):
|
||||
dispatched = []
|
||||
monkeypatch.setattr(srv, "_run_prompt_submit", lambda rid, sid, _s, text, **kw: dispatched.append((text, kw)))
|
||||
monkeypatch.setattr(srv, "_interrupt_busy_session", lambda *a, **k: None)
|
||||
session = _session(agent=types.SimpleNamespace(interrupt=lambda: None), running=True)
|
||||
srv._sessions["sid"] = session
|
||||
try:
|
||||
for text, author in (("ping", AUTHOR), ("hello", OTHER), ("human note", None)):
|
||||
params = {"session_id": "sid", "text": text, "queued": True}
|
||||
if author:
|
||||
params["_turn_author"] = DeliveryAuthor(author)
|
||||
assert _result(srv._methods["prompt.submit"]("r", params)) == {"status": "queued"}
|
||||
# Authored envelopes never merge with each other or with the human's text.
|
||||
assert session["queued_prompt"]["turn_author"] == AUTHOR
|
||||
assert [e["text"] for e in session["queued_prompts"]] == ["hello", "human note"]
|
||||
assert session["queued_prompts"][0]["turn_author"] == OTHER
|
||||
assert "turn_author" not in session["queued_prompts"][1]
|
||||
for _ in range(3):
|
||||
session["running"] = False
|
||||
assert srv._drain_queued_prompt("d", "sid", session) is True
|
||||
finally:
|
||||
srv._sessions.pop("sid", None)
|
||||
assert [(t, kw.get("turn_author")) for t, kw in dispatched] == [("ping", AUTHOR), ("hello", OTHER), ("human note", None)]
|
||||
|
||||
|
||||
def test_turn_runner_passes_the_author_only_when_set_and_only_to_an_agent_that_declares_it(turn_env):
|
||||
seen = []
|
||||
|
||||
def accepting(user_message, *, turn_author="not passed", **kwargs):
|
||||
seen.append(turn_author)
|
||||
return {"final_response": "ok"}
|
||||
|
||||
def legacy(user_message, conversation_history=None, stream_callback=None, persist_user_message=None, task_id=None):
|
||||
seen.append("legacy called")
|
||||
return {"final_response": "ok"}
|
||||
|
||||
for fn, author in ((accepting, AUTHOR), (accepting, None), (legacy, AUTHOR)):
|
||||
agent = types.SimpleNamespace(session_id="a", run_conversation=fn, clear_interrupt=lambda: None)
|
||||
srv._run_prompt_submit("rid", "ui-sid", _session(agent=agent, running=True), "ping", turn_author=author)
|
||||
|
||||
assert seen == [AUTHOR, "not passed", "legacy called"]
|
||||
@@ -121,7 +121,7 @@ def _write(path: Path, record: dict[str, Any]) -> None:
|
||||
|
||||
def deliver_to_live_owner(
|
||||
profile_home: Path | str, owner: dict[str, Any], message: str,
|
||||
*, delivery_id: str | None = None,
|
||||
*, delivery_id: str | None = None, author: dict[str, Any] | None = None,
|
||||
) -> dict[str, Any]:
|
||||
"""Return durable admission immediately, without waiting for the owner.
|
||||
|
||||
@@ -136,7 +136,7 @@ def deliver_to_live_owner(
|
||||
path = root / f"{key}.json"
|
||||
existing = _read(path)
|
||||
if existing is not None:
|
||||
if existing["owner"] != pinned or existing["message"] != message:
|
||||
if existing["owner"] != pinned or existing["message"] != message or existing.get("author") != author:
|
||||
raise ValueError("delivery id already belongs to a different payload")
|
||||
return existing
|
||||
# Wall time can roll back. Permanent receipts retain the admission
|
||||
@@ -146,7 +146,7 @@ def deliver_to_live_owner(
|
||||
if (record := _read(candidate)) is not None), default=0) + 1
|
||||
record = dict(delivery_id=key, id=key, owner=pinned, **pinned,
|
||||
message=message, status="queued", created_at=time.time_ns(),
|
||||
sequence=sequence)
|
||||
sequence=sequence, **({"author": dict(author)} if author else {}))
|
||||
_write(path, record)
|
||||
return record
|
||||
|
||||
|
||||
@@ -402,7 +402,7 @@ def _run_local_turn(argv: list[str], dm_file: str, *, env: Optional[dict[str, st
|
||||
return proc.returncode
|
||||
|
||||
|
||||
def _admit_live_dm(profile_home: Path | None, dm_file: str) -> dict | None:
|
||||
def _admit_live_dm(profile_home: Path | None, dm_file: str, author: Optional[dict] = None) -> 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,
|
||||
@@ -418,7 +418,8 @@ def _admit_live_dm(profile_home: Path | None, dm_file: str) -> dict | None:
|
||||
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())
|
||||
delivery_id=hashlib.sha256(str(Path(dm_file).resolve()).encode()).hexdigest(),
|
||||
**({"author": author} if author else {}))
|
||||
try:
|
||||
fd = os.open(intent_path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
|
||||
except FileExistsError:
|
||||
@@ -433,7 +434,7 @@ def _admit_live_dm(profile_home: Path | None, dm_file: str) -> dict | None:
|
||||
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"])
|
||||
delivery_id=intent["delivery_id"], author=intent.get("author"))
|
||||
return record
|
||||
|
||||
|
||||
@@ -470,8 +471,7 @@ def _run_delivery(argv: list[str], dm_file: str, *, stdin_file: bool,
|
||||
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.
|
||||
``author`` rides to the child as HERMES_TURN_AUTHOR; ``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
|
||||
@@ -484,7 +484,7 @@ def _run_delivery(argv: list[str], dm_file: str, *, stdin_file: bool,
|
||||
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)
|
||||
record = _admit_live_dm(home, dm_file, author)
|
||||
except Exception as exc:
|
||||
print(json.dumps({"status": "ambiguous", "delivery_id": hashlib.sha256(
|
||||
str(Path(dm_file).resolve()).encode()).hexdigest(),
|
||||
@@ -534,7 +534,7 @@ def _start_delivery(argv: list[str], content: str, label: str, *, stdin_file: bo
|
||||
dm_file = _write_dm_file(content)
|
||||
if profile_home is not None:
|
||||
try:
|
||||
record = _admit_live_dm(profile_home, dm_file)
|
||||
record = _admit_live_dm(profile_home, dm_file, author)
|
||||
except Exception as exc:
|
||||
return json.dumps({"status": "ambiguous", "delivery_id": hashlib.sha256(
|
||||
str(Path(dm_file).resolve()).encode()).hexdigest(),
|
||||
@@ -598,8 +598,7 @@ def _spawn_delivery(command: str, label: str, *, dm_file: Optional[str] = None,
|
||||
|
||||
|
||||
def _delivery_main(args: list[str]) -> int:
|
||||
"""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."""
|
||||
"""Runner entry for the argv ``_delivery_command`` builds. Malformed argv exits 2 without touching the DM file."""
|
||||
if not args or args[0] != "--run-delivery":
|
||||
return 2
|
||||
rest, author = args[1:], None
|
||||
|
||||
@@ -384,6 +384,22 @@ def local_delivery_command(profile: str, query_file: str) -> list[str]:
|
||||
return [_hermes_cli(), "-p", profile, *BOT_CHAT_TURN_ARGS, "--query-file", query_file]
|
||||
|
||||
|
||||
class DeliveryAuthor:
|
||||
"""A relayed turn's author as an in-process object. A JSON client cannot build one, so
|
||||
``prompt.submit`` trusts it the way it trusts a hosted-room callback."""
|
||||
|
||||
__slots__ = ("author",)
|
||||
|
||||
def __init__(self, author: dict) -> None:
|
||||
self.author = dict(author)
|
||||
|
||||
def __eq__(self, other: object) -> bool:
|
||||
return isinstance(other, DeliveryAuthor) and other.author == self.author
|
||||
|
||||
def __repr__(self) -> str:
|
||||
return f"DeliveryAuthor({self.author!r})"
|
||||
|
||||
|
||||
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."""
|
||||
|
||||
@@ -90,10 +90,16 @@ def _(rid, params: dict, _root=_relay_root, _run=_run_delivery) -> dict:
|
||||
if isinstance(record, dict) and (record.get("profile_home") or None) == want_home
|
||||
and _session_live_title(
|
||||
record, _session_lookup_key(record, fallback=live_sid)) == BOT_CHAT_TITLE), "")
|
||||
# The Desktop forwards the envelope's sender. The author labels memory only and grants nothing.
|
||||
from tools.bot_relay import DeliveryAuthor, delivery_env, delivery_turn_author
|
||||
author = delivery_turn_author(params.get("from_profile"), params.get("from_handle"))
|
||||
if live_sid:
|
||||
# queued=True: a teammate's DM runs as the NEXT turn and never interrupts or steers a
|
||||
# turn in flight (the default busy mode does); arrivals queue in order.
|
||||
submitted = _methods["prompt.submit"](rid, {"session_id": live_sid, "text": message, "queued": True})
|
||||
submit_params: dict = {"session_id": live_sid, "text": message, "queued": True}
|
||||
if author:
|
||||
submit_params["_turn_author"] = DeliveryAuthor(author)
|
||||
submitted = _methods["prompt.submit"](rid, submit_params)
|
||||
if "error" in submitted:
|
||||
return submitted
|
||||
reply = f"Delivered into @{resolved}'s open Bot Chat; the reply will appear there."
|
||||
@@ -102,9 +108,7 @@ 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. The author labels memory only and grants nothing.
|
||||
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")))
|
||||
turn_env = delivery_env(author)
|
||||
|
||||
fd, tmp = tempfile.mkstemp(prefix="hermes-relay-dm-", suffix=".txt", text=True)
|
||||
try:
|
||||
|
||||
@@ -464,7 +464,7 @@ def _persist_session_row_for_submit(rid, session):
|
||||
return None
|
||||
|
||||
|
||||
def _run_after_agent_ready(rid, sid, session, text, display_kind, hosted_terminal_callback):
|
||||
def _run_after_agent_ready(rid, sid, session, text, display_kind, hosted_terminal_callback, turn_author=None):
|
||||
"""Turn thread body: patient wait for a deferred build (a slow build must not eat the
|
||||
accepted in-flight message), then run."""
|
||||
# The wait delivers the prompt when the still-running build completes, honors a cancel promptly, notices
|
||||
@@ -494,7 +494,7 @@ def _run_after_agent_ready(rid, sid, session, text, display_kind, hosted_termina
|
||||
return
|
||||
_run_prompt_submit(
|
||||
rid, sid, session, text, display_kind=display_kind,
|
||||
terminal_callback=hosted_terminal_callback)
|
||||
terminal_callback=hosted_terminal_callback, turn_author=turn_author)
|
||||
|
||||
|
||||
_TRUNCATION_PARAMS = (
|
||||
@@ -548,6 +548,13 @@ def _(rid, params: dict) -> dict:
|
||||
session, err = _sess_nowait(params, rid)
|
||||
if err:
|
||||
return err
|
||||
from tools.bot_relay import DeliveryAuthor
|
||||
|
||||
# Only the relay handler can build a DeliveryAuthor. A dict here is a client claiming a sender.
|
||||
raw_author = params.get("_turn_author")
|
||||
if raw_author is not None and not isinstance(raw_author, DeliveryAuthor):
|
||||
return _err(rid, 4124, "turn author is stamped by the gateway, never by a client")
|
||||
turn_author = raw_author.author if raw_author is not None else None
|
||||
hosted_task = params.get("_hosted_task")
|
||||
hosted_terminal_callback = params.get("_hosted_terminal_callback")
|
||||
internal_hosted_submit = hosted_task is not None or hosted_terminal_callback is not None
|
||||
@@ -592,7 +599,7 @@ def _(rid, params: dict) -> dict:
|
||||
return _err(rid, 4091, "hosted room member session is busy")
|
||||
busy_transport = t or session.get("transport")
|
||||
busy_response = _handle_busy_submit(
|
||||
rid, sid, session, text, busy_transport, queued=bool(params.get("queued")))
|
||||
rid, sid, session, text, busy_transport, queued=bool(params.get("queued")), turn_author=turn_author)
|
||||
if busy_response is not None:
|
||||
return busy_response
|
||||
raw_rebind_ids = params.get("rebind_survivor_row_ids")
|
||||
@@ -604,6 +611,9 @@ def _(rid, params: dict) -> dict:
|
||||
if err is not None:
|
||||
return err
|
||||
if turn_isolation:
|
||||
if turn_author:
|
||||
logger.debug("isolated compute turns carry no author yet; the turn from %s runs unattributed",
|
||||
turn_author.get("id"))
|
||||
isolated_response = _submit_prompt_to_compute_host(
|
||||
rid, sid, session, text, display_kind=display_kind)
|
||||
if not isolated_response.get("error"):
|
||||
@@ -628,7 +638,7 @@ def _(rid, params: dict) -> dict:
|
||||
_start_agent_build(sid, session)
|
||||
run_thread = threading.Thread(
|
||||
target=lambda: _run_after_agent_ready(
|
||||
rid, sid, session, text, display_kind, hosted_terminal_callback),
|
||||
rid, sid, session, text, display_kind, hosted_terminal_callback, turn_author),
|
||||
daemon=True)
|
||||
# Handle lets session.interrupt tell a live turn from a stuck `running` flag.
|
||||
session["_run_thread"] = run_thread
|
||||
|
||||
@@ -503,7 +503,8 @@ def _prepare_turn_input(sid: str, session: dict, st: _TurnRun, text: Any, images
|
||||
|
||||
def _invoke_agent(
|
||||
sid: str, session: dict, st: _TurnRun, prompt: Any, run_message: Any, streamer,
|
||||
images: list[str], display_kind: str | None, display_metadata: dict | None) -> None:
|
||||
images: list[str], display_kind: str | None, display_metadata: dict | None,
|
||||
turn_author: dict | None = None) -> None:
|
||||
"""Wire the streaming callbacks and run the conversation into ``st.result``."""
|
||||
agent = st.agent
|
||||
|
||||
@@ -539,6 +540,8 @@ def _invoke_agent(
|
||||
if display_kind and "persist_user_display_kind" in run_params:
|
||||
run_kwargs["persist_user_display_kind"] = display_kind
|
||||
run_kwargs["persist_user_display_metadata"] = display_metadata
|
||||
if turn_author and "turn_author" in run_params:
|
||||
run_kwargs["turn_author"] = turn_author
|
||||
# Live-rename hook: auto-titling fires inside the turn prologue.
|
||||
_title_key = session.get("session_key") or sid
|
||||
agent._on_session_title = lambda t, _src, _k=_title_key: _emit(
|
||||
@@ -782,7 +785,8 @@ def _run_prompt_submit(
|
||||
rid, sid: str, session: dict, text: Any, *, display_kind: str | None = None,
|
||||
display_metadata: dict | None = None, image_paths: list[str] | None = None,
|
||||
queued_prompt_generation: int | None = None,
|
||||
terminal_callback: Callable[[dict[str, Any]], None] | None = None) -> bool:
|
||||
terminal_callback: Callable[[dict[str, Any]], None] | None = None,
|
||||
turn_author: dict | None = None) -> bool:
|
||||
admitted = _admit_prompt_turn(sid, session, text, image_paths, queued_prompt_generation)
|
||||
if admitted is None:
|
||||
return False
|
||||
@@ -825,7 +829,7 @@ def _run_prompt_submit(
|
||||
prompt, run_message, cols, streamer = prepared
|
||||
_invoke_agent(
|
||||
sid, session, st, prompt, run_message, streamer, images, display_kind,
|
||||
display_metadata)
|
||||
display_metadata, turn_author)
|
||||
status_note = _absorb_turn_result(
|
||||
sid, session, st, text, display_kind, display_metadata)
|
||||
payload, raw, status = _complete_turn_payload(session, st, status_note, cols)
|
||||
|
||||
@@ -125,10 +125,12 @@ def _ac_inflight_original(session: dict) -> str:
|
||||
return str(turn.get("user") or "").strip() if isinstance(turn, dict) else ""
|
||||
|
||||
|
||||
def _enqueue_prompt(session: dict, text: Any, transport: Any, image_paths: list[str] | None = None) -> None:
|
||||
def _enqueue_prompt(session: dict, text: Any, transport: Any, image_paths: list[str] | None = None,
|
||||
turn_author: dict | None = None) -> None:
|
||||
"""Queue a message for the next turn. Text-only arrivals share a slot and merge losslessly (like the
|
||||
consecutive-user merge in ``repair_message_sequence``); image-bearing ones stay separate envelopes so attachment
|
||||
chronology survives. ``transport`` is pinned so the drained turn streams to its sender."""
|
||||
consecutive-user merge in ``repair_message_sequence``); image-bearing and authored ones stay separate
|
||||
envelopes so attachment chronology and the sender survive. ``transport`` is pinned so the drained turn
|
||||
streams to its sender."""
|
||||
image_paths = list(image_paths or [])
|
||||
# Scrub live-turn self-duplicates first so the text merge below can't glue "{original}\n\n{later}" and re-fire the
|
||||
# original after a correction settles.
|
||||
@@ -138,10 +140,12 @@ def _enqueue_prompt(session: dict, text: Any, transport: Any, image_paths: list[
|
||||
# Never queue a text-only self-copy of the live prompt: draining it would restart it.
|
||||
if text_only and text.strip() == _ac_inflight_original(session) != "":
|
||||
return
|
||||
queued = {"text": text, "transport": transport, **({"image_paths": image_paths} if image_paths else {})}
|
||||
queued = {"text": text, "transport": transport, **({"image_paths": image_paths} if image_paths else {}),
|
||||
**({"turn_author": turn_author} if turn_author else {})}
|
||||
existing = session.get("queued_prompt")
|
||||
if (existing and text_only and isinstance(existing.get("text"), str)
|
||||
and not existing.get("image_paths") and not session.get("queued_prompts")):
|
||||
if (existing and text_only and not turn_author and isinstance(existing.get("text"), str)
|
||||
and not existing.get("image_paths") and not existing.get("turn_author")
|
||||
and not session.get("queued_prompts")):
|
||||
prev = existing["text"]
|
||||
existing["text"] = f"{prev}\n\n{text}" if prev and text else (prev or text)
|
||||
elif existing:
|
||||
@@ -233,7 +237,8 @@ def _ac_try_correction(rid, session: dict, agent: Any, method: str, plain_text:
|
||||
return _ok(rid, {"status": status})
|
||||
|
||||
|
||||
def _handle_busy_submit(rid, sid: str, session: dict, text: Any, transport: Any, queued: bool = False) -> dict | None:
|
||||
def _handle_busy_submit(rid, sid: str, session: dict, text: Any, transport: Any, queued: bool = False,
|
||||
turn_author: dict | None = None) -> dict | None:
|
||||
"""Apply ``display.busy_input_mode`` to a mid-turn prompt instead of rejecting it (rejection made clients busy-retry
|
||||
and drop sends): ``interrupt`` (default) → redirect, falling back to hard interrupt + queue; ``queue`` → queue only;
|
||||
``steer`` → inject after the current atomic action. ``queued=True`` (client queue drain) forces queue mode: a "run
|
||||
@@ -264,7 +269,7 @@ def _handle_busy_submit(rid, sid: str, session: dict, text: Any, transport: Any,
|
||||
if image_paths:
|
||||
session["attached_images"] = image_paths + list(session.get("attached_images", []))
|
||||
return None
|
||||
_enqueue_prompt(session, text, transport, image_paths=image_paths)
|
||||
_enqueue_prompt(session, text, transport, image_paths=image_paths, turn_author=turn_author)
|
||||
session["last_active"] = time.time()
|
||||
# Attachments need their own model invocation: queue without cancelling so the user gets both results in order.
|
||||
# ``steer`` must NEVER escalate to a hard interrupt: it would kill the live turn AND drop ``AIAgent._pending_steer``
|
||||
@@ -307,10 +312,12 @@ def _drain_queued_prompt(rid, sid: str, session: dict) -> bool:
|
||||
kwargs: dict = {"queued_prompt_generation": queue_generation}
|
||||
if queued.get("image_paths"):
|
||||
kwargs["image_paths"] = queued["image_paths"]
|
||||
# The compute-host frame has no author field, so only the inline runner receives it.
|
||||
author_kwargs = {"turn_author": queued["turn_author"]} if queued.get("turn_author") else {}
|
||||
dispatch_failed = False
|
||||
try:
|
||||
if not use_compute_host:
|
||||
_run_prompt_submit(rid, sid, session, queued["text"], **kwargs)
|
||||
_run_prompt_submit(rid, sid, session, queued["text"], **kwargs, **author_kwargs)
|
||||
elif (resp := _submit_prompt_to_compute_host(rid, sid, session, queued["text"], **kwargs)).get("error"):
|
||||
with session["history_lock"]:
|
||||
session["running"] = False
|
||||
|
||||
@@ -529,7 +529,8 @@ def _poll_bot_live_delivery_once(sid: str, session: dict) -> bool:
|
||||
|
||||
try:
|
||||
started = _run_prompt_submit(f"__bot_dm__{delivery_id}", sid, session, claimed["message"],
|
||||
image_paths=[], terminal_callback=terminal_receipt)
|
||||
image_paths=[], terminal_callback=terminal_receipt,
|
||||
turn_author=claimed.get("author") or None)
|
||||
except Exception as exc:
|
||||
_notif_release_turn(session)
|
||||
terminal_receipt({"status": "failed", "error": str(exc)})
|
||||
|
||||
Reference in New Issue
Block a user