refactor(tui_gateway): split _apply_model_switch into phase helpers, compact controller/browser/relay/watcher/billing helpers
This commit is contained in:
@@ -1,28 +1,13 @@
|
||||
"""Bot-relay JSON-RPC handlers — the gateway side of cross-connection A2A.
|
||||
|
||||
Connections ARE the peer set: every gateway the Desktop holds a socket to
|
||||
(local, remote URL, SSH, Hermes Cloud, docker) must be able to find every
|
||||
other connection's agents and message them. The Desktop is the relay — it
|
||||
owns every socket — and these four methods are the door it uses on EACH
|
||||
connected gateway:
|
||||
|
||||
- ``bot_relay.roster.sync`` — Desktop pushes the union roster of agents on
|
||||
the OTHER connections into this gateway's ``bot_relay/roster.json``, so
|
||||
``message_agent`` can resolve cross-connection targets and Bot Chat
|
||||
prompts list them (capability-epoch refresh picks up changes).
|
||||
- ``bot_relay.outbox.drain`` — Desktop collects envelopes queued here by
|
||||
``message_agent`` for targets on other connections.
|
||||
- ``bot_relay.deliver`` — Desktop hands an envelope to the TARGET
|
||||
gateway; this method runs the same one-turn Bot Chat delivery local DMs
|
||||
use and returns the reply text.
|
||||
- ``bot_relay.reply`` — Desktop writes the reply (or a delivery
|
||||
error) back on the SENDER gateway; the waiter spawned at send time picks
|
||||
it up and wakes the sending agent via the standard completion path.
|
||||
|
||||
Storage/validation plumbing lives in ``tools/bot_relay.py``. Handlers are
|
||||
rebound onto server.py's globals at install time (see method_ctx.py) and may
|
||||
reference server module globals (``_ok``, ``_err``) not imported here; this
|
||||
module's own helpers reach them via keyword defaults.
|
||||
Connections ARE the peer set: the Desktop owns every gateway socket (local, remote,
|
||||
SSH, Cloud, docker) and relays between them through these four doors on EACH gateway:
|
||||
``roster.sync`` (push the union roster of OTHER connections' agents so ``message_agent``
|
||||
resolves them), ``outbox.drain`` (collect envelopes queued here for other connections),
|
||||
``deliver`` (run a one-turn Bot Chat delivery on the TARGET gateway, return the reply),
|
||||
``reply`` (write the reply/error back on the SENDER gateway for the waiter to pick up).
|
||||
Storage/validation plumbing lives in ``tools/bot_relay.py``. Handlers are rebound onto
|
||||
server.py's globals (method_ctx.py) and reference ``_ok``/``_err`` etc. bare.
|
||||
"""
|
||||
|
||||
import os
|
||||
@@ -45,12 +30,8 @@ def _run_delivery(profile: str, tmp: str) -> 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=600,
|
||||
local_delivery_command(profile, tmp), capture_output=True, text=True, encoding="utf-8",
|
||||
errors="replace", timeout=600,
|
||||
)
|
||||
|
||||
|
||||
@@ -58,9 +39,8 @@ def _run_delivery(profile: str, tmp: str) -> subprocess.CompletedProcess:
|
||||
def _(rid, params: dict, _root=_relay_root) -> dict:
|
||||
"""Replace this gateway's view of agents on OTHER connections.
|
||||
|
||||
Params: ``agents`` — list of rows ``{profile, handle, connection_id,
|
||||
connection_label?, title?, description?}``. Rows failing validation are
|
||||
dropped, not fatal. Result: ``{count}`` (accepted rows).
|
||||
Params: ``agents`` — rows ``{profile, handle, connection_id, connection_label?, title?,
|
||||
description?}``; rows failing validation are dropped, not fatal. Result: ``{count}``.
|
||||
"""
|
||||
try:
|
||||
from tools.bot_relay import write_remote_roster
|
||||
@@ -72,10 +52,9 @@ def _(rid, params: dict, _root=_relay_root) -> dict:
|
||||
|
||||
@method("bot_relay.outbox.drain")
|
||||
def _(rid, params: dict, _root=_relay_root) -> dict:
|
||||
"""Claim every pending cross-connection envelope queued on this gateway.
|
||||
"""Claim every pending cross-connection envelope queued on this gateway. Result: ``{envelopes}``.
|
||||
|
||||
Claimed envelopes move to ``claimed/`` atomically, so concurrent drains
|
||||
(two Desktop windows) can't double-deliver. Result: ``{envelopes}``.
|
||||
Claimed envelopes move to ``claimed/`` atomically, so concurrent drains can't double-deliver.
|
||||
"""
|
||||
try:
|
||||
from tools.bot_relay import claim_pending_envelopes
|
||||
@@ -89,12 +68,10 @@ def _(rid, params: dict, _root=_relay_root) -> dict:
|
||||
def _(rid, params: dict, _root=_relay_root, _run=_run_delivery) -> dict:
|
||||
"""Deliver a relayed DM into a profile's Bot Chat ON THIS GATEWAY.
|
||||
|
||||
Params: ``profile`` (target on this install), ``message`` (already
|
||||
attribution-prefixed by the sender gateway). Runs the same one-turn
|
||||
``hermes -p <profile> chat -c "Bot Chat"`` transport local DMs use and
|
||||
returns ``{reply}`` — the target agent's response text. Blocking by
|
||||
design (the Desktop calls it from its relay worker, off any UI path;
|
||||
the RPC pool keeps it off the WS reader thread).
|
||||
Params: ``profile`` (target on this install), ``message`` (already attribution-prefixed).
|
||||
Runs the same one-turn ``hermes -p <profile> chat -c "Bot Chat"`` transport local DMs
|
||||
use and returns ``{reply}``. Blocking by design (the Desktop calls it from its relay
|
||||
worker; the RPC pool keeps it off the WS reader thread).
|
||||
"""
|
||||
import os
|
||||
import subprocess
|
||||
@@ -110,7 +87,6 @@ def _(rid, params: dict, _root=_relay_root, _run=_run_delivery) -> dict:
|
||||
|
||||
if len(message) > MESSAGE_MAX_CHARS + 200: # + attribution headroom
|
||||
return _err(rid, 4091, "message too long")
|
||||
|
||||
root = _root()
|
||||
known = {"default"}
|
||||
profiles_dir = root / "profiles"
|
||||
@@ -120,13 +96,11 @@ def _(rid, params: dict, _root=_relay_root, _run=_run_delivery) -> dict:
|
||||
if resolved not in known:
|
||||
return _err(rid, 4092, f"no profile '{profile}' on this gateway")
|
||||
|
||||
# When THIS gateway already hosts the target's Bot Chat live (the
|
||||
# Desktop has it open), the subprocess transport is fenced out by the
|
||||
# single-owner lease and the payload dropped. Land the DM in the live
|
||||
# session as a normal user turn via prompt.submit instead — the
|
||||
# composer's choke point, so role alternation, persistence and
|
||||
# streaming behave as a typed message would. (Nested: needs server
|
||||
# globals via method_ctx rebinding.)
|
||||
# When THIS gateway already hosts the target's Bot Chat live, the subprocess
|
||||
# transport is fenced out by the single-owner lease and the payload dropped. Land
|
||||
# the DM in the live session via prompt.submit — the composer's choke point, so
|
||||
# role alternation, persistence and streaming behave as a typed message would.
|
||||
# (Nested: needs server globals via method_ctx rebinding.)
|
||||
def _live_bot_chat_sid(profile_name: str) -> str:
|
||||
from tools.bot_mode_probe import BOT_CHAT_TITLE
|
||||
|
||||
@@ -144,58 +118,46 @@ def _(rid, params: dict, _root=_relay_root, _run=_run_delivery) -> dict:
|
||||
|
||||
live_sid = _live_bot_chat_sid(resolved)
|
||||
if live_sid:
|
||||
# queued=True: a teammate's DM runs as the NEXT turn. It must never
|
||||
# interrupt or steer a turn already in flight (the default busy
|
||||
# mode does); hundreds of arrivals simply queue in arrival order.
|
||||
# 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})
|
||||
if "error" in submitted:
|
||||
return submitted
|
||||
return _ok(
|
||||
rid,
|
||||
{"reply": f"Delivered into @{resolved}'s open Bot Chat; the reply will appear there."},
|
||||
rid, {"reply": f"Delivered into @{resolved}'s open Bot Chat; the reply will appear there."}
|
||||
)
|
||||
|
||||
fd, tmp = tempfile.mkstemp(prefix="hermes-relay-dm-", suffix=".txt", text=True)
|
||||
try:
|
||||
with os.fdopen(fd, "w", encoding="utf-8") as f:
|
||||
f.write(message)
|
||||
# Per-profile turn lock serializes with any other delivery turn into
|
||||
# this profile (relay or local message_agent) and covers only the
|
||||
# turn execution window. Worst-case handler hold is lock wait
|
||||
# (bot_mode.turn_wait_seconds, default 120s) + the 600s turn timeout,
|
||||
# doubled when the retry policy grants one bounded re-run — callers
|
||||
# must tolerate ~1320s before assuming failure.
|
||||
# Per-profile turn lock serializes with any other delivery turn into this
|
||||
# profile and covers only the turn window. Worst-case hold is lock wait
|
||||
# (bot_mode.turn_wait_seconds, default 120s) + the 600s turn timeout, doubled
|
||||
# when the retry policy grants one re-run — callers must tolerate ~1320s.
|
||||
with acquire_turn_lock(root, resolved):
|
||||
proc = _run(resolved, tmp)
|
||||
if proc.returncode != 0:
|
||||
# Retry policy: transient classes re-run the SAME session
|
||||
# once; context_overflow also re-runs the same session — the
|
||||
# retried turn's pre-API compaction pass compacts the
|
||||
# over-threshold Bot Chat transcript first (the sanctioned
|
||||
# compression lever; no fresh session is ever minted).
|
||||
# Auth/quota/config classes never retry.
|
||||
# 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 transcript first (no fresh session is
|
||||
# ever minted). Auth/quota/config classes never retry.
|
||||
from tools.bot_failure_reasons import (
|
||||
RETRY_NONE,
|
||||
classify_agent_error,
|
||||
retry_action,
|
||||
RETRY_NONE, classify_agent_error, retry_action,
|
||||
)
|
||||
|
||||
first_detail = (proc.stderr or proc.stdout or "").strip()[-500:]
|
||||
if retry_action(classify_agent_error(first_detail)) != RETRY_NONE:
|
||||
proc = _run(resolved, tmp)
|
||||
finally:
|
||||
try:
|
||||
with contextlib.suppress(OSError):
|
||||
os.unlink(tmp)
|
||||
except OSError:
|
||||
pass
|
||||
if proc.returncode != 0:
|
||||
from tools.bot_failure_reasons import classify_agent_error
|
||||
|
||||
detail = (proc.stderr or proc.stdout or "").strip()[-500:]
|
||||
return _err(
|
||||
rid,
|
||||
5092,
|
||||
f"delivery turn failed: {detail or proc.returncode}",
|
||||
rid, 5092, f"delivery turn failed: {detail or proc.returncode}",
|
||||
data={"reason": classify_agent_error(detail)},
|
||||
)
|
||||
return _ok(rid, {"reply": (proc.stdout or "").strip()})
|
||||
@@ -212,8 +174,8 @@ def _(rid, params: dict, _root=_relay_root, _run=_run_delivery) -> dict:
|
||||
def _(rid, params: dict, _root=_relay_root) -> dict:
|
||||
"""Write a relayed reply (or delivery error) for a sender-side waiter.
|
||||
|
||||
Params: ``id`` (envelope id), ``reply`` and/or ``error``, optional
|
||||
``reason`` (typed failure code, see ``tools.bot_failure_reasons``).
|
||||
Params: ``id`` (envelope id), ``reply`` and/or ``error``, optional ``reason``
|
||||
(typed failure code, see ``tools.bot_failure_reasons``).
|
||||
"""
|
||||
envelope_id = str(params.get("id") or "").strip()
|
||||
if not envelope_id:
|
||||
@@ -222,11 +184,8 @@ def _(rid, params: dict, _root=_relay_root) -> dict:
|
||||
from tools.bot_relay import write_reply
|
||||
|
||||
write_reply(
|
||||
_root(),
|
||||
envelope_id,
|
||||
reply=str(params.get("reply") or ""),
|
||||
error=str(params.get("error") or ""),
|
||||
reason=str(params.get("reason") or ""),
|
||||
_root(), envelope_id, reply=str(params.get("reply") or ""),
|
||||
error=str(params.get("error") or ""), reason=str(params.get("reason") or ""),
|
||||
)
|
||||
return _ok(rid, {"ok": True})
|
||||
except ValueError as e:
|
||||
@@ -241,12 +200,8 @@ def register(server) -> None:
|
||||
|
||||
server._LONG_HANDLERS = server._LONG_HANDLERS | methods_groups.LONG_HANDLERS
|
||||
for name in (
|
||||
"get_hosted_room_service",
|
||||
"_WORKER_UNAVAILABLE",
|
||||
"_profile_name",
|
||||
"_requested_profile",
|
||||
"_api_server_key",
|
||||
"_room_link_run_storage_durable",
|
||||
"get_hosted_room_service", "_WORKER_UNAVAILABLE", "_profile_name", "_requested_profile",
|
||||
"_api_server_key", "_room_link_run_storage_durable",
|
||||
):
|
||||
setattr(server, name, getattr(methods_groups, name))
|
||||
methods_groups.bind_server(server)
|
||||
|
||||
Reference in New Issue
Block a user