fix(api-server): a peer DM into a Bot Chat open in Desktop is answered by that chat, not a second writer

`hermes peer dm` posts POST /api/sessions/{id}/chat. When the target is the
profile's canonical Bot Chat and a Desktop holds it live, that Desktop session
owns the chat's single-writer lease, and every other writer is refused
SESSION_NOT_OWNED — per-session exclusivity is correctness, enforced
unconditionally (hermes_cli/active_sessions.py). The API server neither takes
that lease nor checks it, so the turn ran beside the owner: the open chat never
showed the message or the reply, the live session's context never learned of
them, and two writers appended to one transcript in state.db.

Hand the message to the owner's mailbox instead, as local DMs
(tools/bot_mode_dm.py) and relayed DMs (tui_gateway/methods_bot_relay.py,
budget so the peer still gets the reply on the same call. A turn still running
at that deadline answers 202 with the delivery id, and `peer dm` reports the
message as queued in that chat instead of printing "(no reply)".

Only the canonical Bot Chat's own compression lineage is handed off: a peer turn
into any other session, or into a Bot Chat nobody holds, runs here as before.

Tests: the four-row table is the whole discriminator (owner answers, owner still
running, another session, nobody holding the chat) and the client row pins that a
queued answer reads as delivered.
This commit is contained in:
John Paul Soliva
2026-09-18 17:49:52 +09:00
committed by Teknium
parent 3c7943fbcc
commit b4b34178b7
4 changed files with 195 additions and 1 deletions

View File

@@ -3187,6 +3187,55 @@ class APIServerAdapter(OpenAICompatRoutesMixin, BasePlatformAdapter):
return ""
return "confirmed" if runtime else "accepted"
async def _answer_through_live_bot_chat(self, ctx: Dict[str, Any]) -> Optional["web.Response"]:
"""Hand a turn aimed at a canonical Bot Chat that a Desktop holds live to that owner.
This is the ``hermes peer dm`` transport. Running the turn here would make this process a
second writer beside the lease holder: the open chat never shows the message or the reply,
its live context never learns of them, and the two transcripts interleave in state.db.
Local and relayed DMs already hand such a message to the owner's mailbox
(``tools/bot_mode_dm.py``, ``tui_gateway/methods_bot_relay.py``). This waits for the owner's
receipt on the same budget as the local path, so the peer still gets the reply on this call.
"""
message = ctx["user_message"]
if not isinstance(message, str):
return None
db = await self._ensure_session_db_async()
if db is None:
return None
home = Path(db.db_path).parent
session_id = ctx["session_id"]
from tools.bot_live_delivery import deliver_to_live_owner, find_canonical_live_owner, read_delivery_result
from tools.bot_mode_dm import _LIVE_WAIT_SECONDS
def _admit() -> Optional[Dict[str, Any]]:
owner = find_canonical_live_owner(home)
# Only the canonical Bot Chat's own lineage: a peer turn into any other session runs here.
if owner is None or db.get_compression_tip(session_id) != owner["session_id"]:
return None
return deliver_to_live_owner(home, owner, message, author=ctx["run_kwargs"]["turn_author"])
record = await asyncio.to_thread(_admit)
if record is None:
return None
delivery_id = record["delivery_id"]
deadline = time.monotonic() + _LIVE_WAIT_SECONDS
while record["status"] in ("queued", "claimed") and time.monotonic() < deadline:
await asyncio.sleep(0.5)
record = await asyncio.to_thread(read_delivery_result, home, delivery_id) or record
headers = self._session_headers(session_id, ctx["gateway_session_key"])
if record["status"] == "settled":
return web.json_response(
{"object": "hermes.session.chat.completion", "session_id": session_id,
"message": {"role": "assistant", "content": record.get("reply") or ""},
"usage": {}, "runtime": {}, "delivery_id": delivery_id}, headers=headers)
if record["status"] in ("queued", "claimed"):
return web.json_response(
{"object": "hermes.session.chat.queued", "session_id": session_id,
"status": record["status"], "delivery_id": delivery_id}, status=202, headers=headers)
return _error_response(record.get("error") or f"Bot Chat delivery {record['status']}", 502,
code=record.get("reason") or record["status"], headers=headers)
@_admit_api_agent_request
async def _handle_session_chat(self, request: "web.Request") -> "web.Response":
"""POST /api/sessions/{session_id}/chat — one synchronous agent turn (plus the delivery lanes'
@@ -3201,6 +3250,9 @@ class APIServerAdapter(OpenAICompatRoutesMixin, BasePlatformAdapter):
ctx, err = await self._prepare_session_chat(request)
if err is not None:
return err
handed_off = await self._answer_through_live_bot_chat(ctx)
if handed_off is not None:
return handed_off
gateway_session_key = ctx["gateway_session_key"]
session_id = ctx["session_id"]
history = await self._conversation_history_for_session(session_id)

View File

@@ -365,6 +365,14 @@ def _peer_dm(args, message: str, peer_name: str, profile: str | None, base: str,
file=sys.stderr)
return 1
return _peer_failure(peer_name, exc)
if result.get("object") == "hermes.session.chat.queued":
# The peer's Bot Chat is open in its Desktop and that turn outlasted the peer's wait: the
# message is in the open chat and is answered there, so a resend would run it twice.
queued_in = result.get("session_id") or session_id
return _emit(args, {"peer": peer_name, "profile": profile, "session_id": queued_in,
"status": result.get("status") or "queued", "delivery_id": result.get("delivery_id")},
[f"Peer '{peer_name}' has its Bot Chat open, so the message went into that chat (session "
f"{queued_in}) and is answered there. The reply cannot come back on this call. Do NOT resend."])
msg = result.get("message")
reply = str(msg.get("content") or "") if isinstance(msg, dict) else ""
# A successful bare silence marker is a delivery decision, not a message:

View File

@@ -0,0 +1,130 @@
"""A peer DM into a Bot Chat that a Desktop holds open is answered BY that open chat.
``hermes peer dm`` posts to ``/api/sessions/{id}/chat`` on the peer. When the peer's canonical Bot
Chat is open in its Desktop, the Desktop session holds the chat's single-writer lease; running the
turn in the API server beside it made a second writer the open chat never saw. The message now goes
through the owner's mailbox, like local and relayed DMs, and the owner's receipt carries the reply.
"""
from __future__ import annotations
import json
import threading
import time
from types import SimpleNamespace
from unittest.mock import AsyncMock, patch
import pytest
from aiohttp import web
from aiohttp.test_utils import TestClient, TestServer
from gateway.config import PlatformConfig
from gateway.platforms.api_server import APIServerAdapter
from hermes_cli.subcommands import peer as peer_mod
from hermes_state import SessionDB
from tools import bot_live_delivery as mailbox
AUTHOR = {"id": "bot:cto", "name": "cto", "is_bot": True}
def _app(adapter):
app = web.Application()
app.router.add_post("/api/sessions/{session_id}/chat", adapter._handle_session_chat)
return app
def _owner_answers(home, reply):
"""What the Desktop's live session does: claim the delivery, run it as its next turn, settle it."""
def _run():
owner = mailbox.find_canonical_live_owner(home)
for _ in range(200):
claimed = mailbox.claim_pending_delivery(home, owner)
if claimed is not None:
mailbox.complete_delivery(home, claimed["delivery_id"], status="settled", reply=reply)
return
time.sleep(0.02)
thread = threading.Thread(target=_run, daemon=True)
thread.start()
return thread
@pytest.mark.asyncio
@pytest.mark.parametrize(
("target", "open_in_desktop", "owner_replies", "status", "content", "turn_ran_here"),
[
("bot-chat", True, True, 200, "pong", False),
("bot-chat", True, False, 202, None, False),
("scratch", True, True, 200, "ran here", True),
("bot-chat", False, False, 200, "ran here", True),
],
ids=["open-bot-chat-answers", "open-bot-chat-still-running", "other-session", "bot-chat-not-open"],
)
async def test_a_peer_turn_into_an_open_bot_chat_is_answered_by_its_live_owner(
tmp_path, monkeypatch, target, open_in_desktop, owner_replies, status, content, turn_ran_here
):
"""Only the canonical Bot Chat's live owner takes the turn; every other session, and a Bot Chat
nobody holds, still runs here. A turn the owner has not finished inside the wait is reported as
queued in that chat, never as a failure the sender would resend."""
home = tmp_path.resolve()
monkeypatch.setenv("HERMES_HOME", str(home))
monkeypatch.setattr("tools.bot_mode_dm._LIVE_WAIT_SECONDS", 1.0)
db = SessionDB(home / "state.db")
db.create_session("bot-chat", "desktop")
db.set_session_title("bot-chat", "Bot Chat")
db.create_session("scratch", "api_server")
lease = None
if open_in_desktop:
from hermes_cli.active_sessions import try_acquire_active_session
lease, refusal = try_acquire_active_session(
session_id="bot-chat", surface="desktop", config={}, registry_home=home, track_liveness=True,
metadata={"live_session_id": "live-1", "bot_live_delivery_consumer": True})
assert lease is not None and refusal is None
owner = _owner_answers(home, "pong") if owner_replies and target == "bot-chat" else None
adapter = APIServerAdapter(PlatformConfig(enabled=True))
adapter._session_db = db
try:
with patch.object(adapter, "_run_agent", AsyncMock(return_value=({"final_response": "ran here"}, {}))) as run:
async with TestClient(TestServer(_app(adapter))) as cli:
resp = await cli.post(f"/api/sessions/{target}/chat", json={"message": "ping", "author": AUTHOR})
body = await resp.json()
if owner is not None:
owner.join(5)
assert resp.status == status, body
assert run.called is turn_ran_here
admitted = sorted((home / "runtime" / "bot_live_delivery").glob("*.json"))
if turn_ran_here:
assert body["message"]["content"] == content
assert not admitted
return
[record] = [json.loads(path.read_text()) for path in admitted]
assert (record["message"], record["author"], record["owner"]["live_session_id"]) == ("ping", AUTHOR, "live-1")
assert body["delivery_id"] == record["delivery_id"]
if status == 200:
assert body["message"]["content"] == content
else:
assert (body["object"], body["status"]) == ("hermes.session.chat.queued", "queued")
finally:
if lease is not None:
lease.release()
db.close()
@pytest.mark.parametrize("as_json", [False, True], ids=["text", "json"])
def test_peer_dm_reports_a_turn_queued_in_the_open_bot_chat_as_delivered(monkeypatch, capsys, as_json):
"""The queued answer means the message IS in the peer's open Bot Chat: say so, succeed, and tell
the sender not to resend, instead of printing ``(no reply)`` as if the turn were empty."""
monkeypatch.setattr(peer_mod, "_ensure_bot_chat", lambda base, key: "bot-chat")
monkeypatch.setattr(peer_mod, "_request", lambda url, key, **kw: {
"object": "hermes.session.chat.queued", "session_id": "bot-chat", "status": "claimed", "delivery_id": "d" * 32})
code = peer_mod._peer_dm(SimpleNamespace(json=as_json), "hello", "mini", None, "http://peer:8642", "key")
out = capsys.readouterr().out
assert code == 0
assert "(no reply)" not in out
if as_json:
assert json.loads(out) == {"peer": "mini", "profile": None, "session_id": "bot-chat",
"status": "claimed", "delivery_id": "d" * 32}
else:
assert "went into that chat (session bot-chat)" in out and "Do NOT resend" in out

View File

@@ -247,7 +247,11 @@ Use `peer dm` only for short queries and receipts because it holds one HTTP
connection until the turn finishes. If the peer takes the message but the turn outlasts that
connection, the message is already in the peer's Bot Chat and the turn keeps running there, so the
command says exactly that instead of reporting the peer unreachable — resending would run the turn
twice. A timeout while connecting still reports the peer unreachable. For a long turn, `peer run` returns a
twice. A timeout while connecting still reports the peer unreachable. When the peer's Bot Chat is
open in its Desktop, the message is handed to that open chat and runs there as its next turn —
whoever is watching sees it, and the reply still comes back on the call; if that turn is still going
after five minutes, `peer dm` reports the message as queued in that chat rather than lost. For a long
turn, `peer run` returns a
`run_id` immediately; poll it with `peer status`. The run inherits the
canonical Bot Chat transcript, and a stable `--idempotency-key` makes a retry
return the original run instead of starting duplicate work. Use `peer stop`