diff --git a/apps/desktop/src/plugins/hermes-bots/plugin.js b/apps/desktop/src/plugins/hermes-bots/plugin.js index 68c9f3c2b9aa..7b6b6690a596 100644 --- a/apps/desktop/src/plugins/hermes-bots/plugin.js +++ b/apps/desktop/src/plugins/hermes-bots/plugin.js @@ -1190,6 +1190,222 @@ function handleSessionsGatewayTransition() { .then(() => scheduleGroupChatServerSync($groupChats.get())) } +// ── cross-connection bot relay ──────────────────────────────────────────── +// Connections ARE the peer set: every gateway this 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 via message_agent. The +// Desktop is the relay — it owns every socket. Two loops: +// - roster loop: pushes each gateway the union roster of agents on the +// OTHER connections (bot_relay.roster.sync), so message_agent resolves +// cross-connection targets and Bot Chat prompts list them; +// - drain loop: collects queued envelopes from every gateway +// (bot_relay.outbox.drain), delivers each on the target connection's +// own socket (bot_relay.deliver), and posts the reply back to the +// sender gateway (bot_relay.reply) where a waiter wakes the sender. +// Older backends without the RPCs fail per-call and are skipped — the +// relay degrades to whatever subset of connections supports it. +const RELAY_ROSTER_INTERVAL_MS = 60_000 +const RELAY_DRAIN_INTERVAL_MS = 4_000 +let relayDisposed = false +let relayRosterTimer = null +let relayDrainTimer = null +let relayRosterBusy = false +let relayDrainBusy = false + +/** One representative route per reachable connection id. */ +async function relayConnections() { + if (typeof host.profileRoutes !== 'function' || typeof host.requestProfile !== 'function') { + return [] + } + + try { + const routes = await host.profileRoutes() + const byConnection = new Map() + + for (const route of Array.isArray(routes) ? routes : []) { + const id = String(route?.connectionId || '') + + if (id && !byConnection.has(id)) { + byConnection.set(id, route) + } + } + + return [...byConnection.entries()].map(([id, route]) => ({ id, route })) + } catch { + return [] + } +} + +/** The agents living on one connection, as relay roster rows. */ +async function relayAgentsOn(connection) { + try { + const res = await host.requestProfile(connection.route, 'profiles.list', { include_sessions: false }) + const profiles = Array.isArray(res?.profiles) ? res.profiles : [] + const label = String( + connection.route?.connectionLabel || connection.route?.label || connection.id + ) + + return profiles + .map(profile => ({ + profile: String(profile?.name || ''), + handle: botHandle(profile?.name, profile), + connection_id: connection.id, + connection_label: label, + title: String(profile?.ui_meta?.['hermes-bots']?.title || profile?.display_name || ''), + description: String(profile?.description || '') + })) + .filter(row => row.profile) + } catch { + return [] + } +} + +/** Push every gateway the union roster of agents on the OTHER connections. */ +async function syncRelayRosters() { + if (relayDisposed || relayRosterBusy) { + return + } + + relayRosterBusy = true + + try { + const connections = await relayConnections() + + if (connections.length < 2) { + return + } + + const agentsByConnection = new Map() + await Promise.all( + connections.map(async connection => { + agentsByConnection.set(connection.id, await relayAgentsOn(connection)) + }) + ) + + await Promise.all( + connections.map(async connection => { + const others = [] + + for (const [id, agents] of agentsByConnection) { + if (id !== connection.id) { + others.push(...agents) + } + } + + try { + await host.requestProfile(connection.route, 'bot_relay.roster.sync', { agents: others }) + } catch { + // Older backend without the relay RPCs — skip this connection. + } + }) + ) + } finally { + relayRosterBusy = false + } +} + +/** Drain every gateway's outbox and deliver each envelope on the target + * connection's own socket; the reply (or error) is posted back to the + * sender gateway for its waiter. */ +async function drainRelayOutboxes() { + if (relayDisposed || relayDrainBusy) { + return + } + + relayDrainBusy = true + + try { + const connections = await relayConnections() + + if (connections.length < 2) { + return + } + + const byId = new Map(connections.map(connection => [connection.id, connection])) + + for (const sender of connections) { + let envelopes = [] + + try { + const res = await host.requestProfile(sender.route, 'bot_relay.outbox.drain', {}) + envelopes = Array.isArray(res?.envelopes) ? res.envelopes : [] + } catch { + continue + } + + for (const envelope of envelopes) { + if (relayDisposed) { + return + } + + const envelopeId = String(envelope?.id || '') + const target = byId.get(String(envelope?.target_connection || '')) + const postReply = async payload => { + try { + await host.requestProfile(sender.route, 'bot_relay.reply', { id: envelopeId, ...payload }) + } catch { + // Sender gateway unreachable — its waiter times out with guidance. + } + } + + if (!envelopeId) { + continue + } + + if (!target) { + await postReply({ error: `connection '${envelope?.target_connection}' is not connected to this Desktop right now` }) + continue + } + + try { + const res = await host.requestProfile(target.route, 'bot_relay.deliver', { + profile: String(envelope?.target_profile || ''), + message: String(envelope?.message || '') + }) + await postReply({ reply: String(res?.reply || '') }) + } catch (error) { + await postReply({ error: String(error?.message || error || 'delivery failed') }) + } + } + } + } finally { + relayDrainBusy = false + } +} + +function startBotRelay() { + relayDisposed = false + + // Source-shape test harnesses evaluate plugin.js without DOM timers — + // the relay only runs where a real event loop exists. + if (typeof setInterval !== 'function' || typeof clearInterval !== 'function') { + return + } + + if (relayRosterTimer === null) { + relayRosterTimer = setInterval(() => void syncRelayRosters(), RELAY_ROSTER_INTERVAL_MS) + void syncRelayRosters() + } + + if (relayDrainTimer === null) { + relayDrainTimer = setInterval(() => void drainRelayOutboxes(), RELAY_DRAIN_INTERVAL_MS) + } +} + +function stopBotRelay() { + relayDisposed = true + + if (relayRosterTimer !== null) { + clearInterval(relayRosterTimer) + relayRosterTimer = null + } + + if (relayDrainTimer !== null) { + clearInterval(relayDrainTimer) + relayDrainTimer = null + } +} + /** Per-bot appearance + display meta, persisted via ctx.storage: * { [botName]: { shape, color, title } } */ const $botMeta = atom({}) @@ -11954,10 +12170,14 @@ export default { pluginCtx = ctx groupChatSyncDisposed = false startFaceClock() + // The cross-connection relay rides every gateway socket this Desktop + // holds: roster sync + envelope drain/deliver/reply loops. + startBotRelay() // Disabling the plugin (or a hot reload) must actually stop the clock — // before this, the rAF loop + 1Hz document scan ran until app restart. if (typeof ctx.onDispose === 'function') { ctx.onDispose(stopFaceClock) + ctx.onDispose(stopBotRelay) } // @-mention autocomplete: typing "@rese…" in ANY composer offers the @@ -12306,20 +12526,22 @@ export default { // Identification only. Each line names the agent the user's tag // resolves to (friendly title + device for cross-connection rows), - // so the agent knows exactly who "@research-buddy" is without the - // renderer ever acting on the user's behalf. + // so the agent knows exactly who "@research-buddy" is. Cross- + // connection targets carry the '@connection' suffix message_agent + // resolves against the Desktop-synced relay roster. const lines = mentionedBots.map(bot => { const handle = botHandle(bot.name, bot) const title = String(botRosterMeta(bot, $botMeta.get())?.title || bot.ui_meta?.['hermes-bots']?.title || bot.title || '').trim() + const target = bot.remoteSource && bot.connectionId ? `${handle}@${bot.connectionId}` : handle const where = bot.remoteSource - ? ` — on ${bot.connectionLabel || bot.connectionId}` + ? ` — on ${bot.connectionLabel || bot.connectionId} (message_agent target: "${target}")` : '' return `@${handle} = agent profile "${bot.name}"${title ? ` ("${title}")` : ''}${where}` }) const note = '\n\n[@mentions resolved from the Bot Mode roster — the user is referring to: ' + lines.join('; ') + - '. If they want one of these agents contacted, compose your own message and send it with your message_agent tool; never forward the user\u2019s text verbatim. If this session has no message_agent tool, agent messaging is unavailable here — say so.]' + '. If they want one of these agents contacted, compose your own message and send it with your message_agent tool (agents on other connected machines are reachable too — the Desktop relays it); never forward the user\u2019s text verbatim. If this session has no message_agent tool, agent messaging is unavailable here — say so.]' return { ...draft, text: text + note } } } diff --git a/apps/desktop/src/plugins/hermes-bots/tests/bot-relay.test.mjs b/apps/desktop/src/plugins/hermes-bots/tests/bot-relay.test.mjs new file mode 100644 index 000000000000..86fb0d46eefb --- /dev/null +++ b/apps/desktop/src/plugins/hermes-bots/tests/bot-relay.test.mjs @@ -0,0 +1,55 @@ +import assert from 'node:assert/strict' +import { readFileSync } from 'node:fs' +import test from 'node:test' + +// The cross-connection bot relay (Aug 2026 ruling: connections ARE the peer +// set). Source-shape contracts: +// - both relay loops exist, start in register(), and stop via ctx.onDispose; +// - the roster loop pushes bot_relay.roster.sync with agents from OTHER +// connections only; +// - the drain loop wires drain → deliver → reply, and posts an error reply +// when the target connection is gone (waiter must never dangle); +// - the remote-row toast no longer tells users messaging is device-local; +// - the middleware note carries the message_agent target for cross- +// connection rows instead of implying they're unreachable. + +const pluginSource = readFileSync(new URL('../plugin.js', import.meta.url), 'utf8') + +test('relay loops start in register() and stop on dispose', () => { + assert.match(pluginSource, /startBotRelay\(\)/) + assert.match(pluginSource, /ctx\.onDispose\(stopBotRelay\)/) + // teardown really clears both timers + const stop = pluginSource.slice(pluginSource.indexOf('function stopBotRelay')) + assert.match(stop.slice(0, 500), /clearInterval\(relayRosterTimer\)/) + assert.match(stop.slice(0, 500), /clearInterval\(relayDrainTimer\)/) +}) + +test('roster loop syncs OTHER connections agents to each gateway', () => { + const sync = pluginSource.slice( + pluginSource.indexOf('async function syncRelayRosters'), + pluginSource.indexOf('async function drainRelayOutboxes') + ) + assert.match(sync, /bot_relay\.roster\.sync/) + assert.match(sync, /id !== connection\.id/) +}) + +test('drain loop wires drain → deliver → reply with error fallback', () => { + const drain = pluginSource.slice( + pluginSource.indexOf('async function drainRelayOutboxes'), + pluginSource.indexOf('function startBotRelay') + ) + assert.match(drain, /bot_relay\.outbox\.drain/) + assert.match(drain, /bot_relay\.deliver/) + assert.match(drain, /bot_relay\.reply/) + // a missing target connection still posts a reply (error) for the waiter + assert.match(drain, /is not connected to this Desktop right now/) +}) + +test('remote-row dead-end toast is gone (rows open; relay carries DMs)', () => { + assert.doesNotMatch(pluginSource, /Gateway stays on this device/) +}) + +test('middleware note names the cross-connection message_agent target', () => { + assert.match(pluginSource, /message_agent target: "\$\{target\}"/) + assert.match(pluginSource, /agents on other connected machines are reachable too/) +}) diff --git a/tests/tools/test_bot_relay.py b/tests/tools/test_bot_relay.py new file mode 100644 index 000000000000..8692cd663666 --- /dev/null +++ b/tests/tools/test_bot_relay.py @@ -0,0 +1,290 @@ +"""Tests: cross-connection bot relay (tools/bot_relay.py + message_agent route). + +Connections ARE the peer set: every Desktop-connected gateway must be +message_agent-reachable. These tests pin the gateway-side plumbing — +roster validation, target resolution (incl. ambiguity), outbox claim +atomicity, reply write validation — and the two behavior contracts the +relay adds to message_agent: + +- a target resolving against the Desktop-synced relay roster is queued as + an envelope and acknowledged like any DM (fire-and-forget, waiter spawned); +- the legacy-SOUL dedupe (empty protocol section) NO LONGER strips the tool: + the injection/execution gates key on managed-install, not section text. +""" + +import json +import re +from pathlib import Path + +import pytest + +from tools import bot_relay +from tools.bot_mode_dm import ( + MESSAGE_AGENT_TOOL_NAME, + ensure_message_agent_tool, + message_agent_tool, +) + + +@pytest.fixture() +def root(tmp_path): + return tmp_path + + +def _rows(): + return [ + { + "profile": "default", + "handle": "hermes", + "connection_id": "cloud-1", + "connection_label": "Hermes Cloud", + "title": "Moxie", + "description": "Main cloud agent", + }, + { + "profile": "researcher", + "handle": "researcher", + "connection_id": "ssh-vps", + "connection_label": "VPS", + }, + ] + + +# ── roster ─────────────────────────────────────────────────────────────────── + + +def test_roster_roundtrip_and_validation(root): + rows = _rows() + [ + {"profile": "", "handle": "x", "connection_id": "c"}, # no profile + {"profile": "bad name!", "connection_id": "c"}, # bad charset + "not-a-dict", + {"profile": "default", "handle": "hermes", "connection_id": "cloud-1"}, # dupe + ] + count = bot_relay.write_remote_roster(root, rows) + assert count == 2 + back = bot_relay.read_remote_roster(root) + assert [r["profile"] for r in back] == ["default", "researcher"] + assert back[0]["title"] == "Moxie" + + +def test_roster_read_missing_and_corrupt(root): + assert bot_relay.read_remote_roster(root) == [] + base = bot_relay.relay_root(root) + base.mkdir(parents=True) + (base / bot_relay.ROSTER_FILE).write_text("{corrupt", encoding="utf-8") + assert bot_relay.read_remote_roster(root) == [] + + +def test_resolve_remote_target_forms(root): + bot_relay.write_remote_roster(root, _rows()) + roster = bot_relay.read_remote_roster(root) + assert bot_relay.resolve_remote_target("researcher", roster)["connection_id"] == "ssh-vps" + assert bot_relay.resolve_remote_target("@hermes", roster)["profile"] == "default" + # profile name resolves too + assert bot_relay.resolve_remote_target("default", roster)["connection_id"] == "cloud-1" + # exact connection-qualified form + assert bot_relay.resolve_remote_target("hermes@cloud-1", roster)["profile"] == "default" + assert bot_relay.resolve_remote_target("hermes@nope", roster) is None + assert bot_relay.resolve_remote_target("ghost", roster) is None + + +def test_resolve_ambiguous_handle_across_connections(root): + rows = _rows() + [ + {"profile": "researcher", "handle": "researcher", "connection_id": "cloud-1"} + ] + bot_relay.write_remote_roster(root, rows) + roster = bot_relay.read_remote_roster(root) + assert bot_relay.resolve_remote_target("researcher", roster) == "ambiguous" + match = bot_relay.resolve_remote_target("researcher@ssh-vps", roster) + assert match["connection_id"] == "ssh-vps" + forms = bot_relay.remote_target_forms(roster) + assert "researcher@ssh-vps" in forms and "researcher@cloud-1" in forms + assert "hermes" in forms # unique handle stays bare + + +# ── outbox / replies ───────────────────────────────────────────────────────── + + +def test_enqueue_claim_is_atomic_and_single_shot(root): + bot_relay.write_remote_roster(root, _rows()) + roster = bot_relay.read_remote_roster(root) + target = bot_relay.resolve_remote_target("researcher", roster) + env = bot_relay.enqueue_envelope( + root, target=target, message="hi", sender_profile="work", sender_handle="work" + ) + assert re.match(r"^[0-9a-f]{32}$", env["id"]) + claimed = bot_relay.claim_pending_envelopes(root) + assert [e["id"] for e in claimed] == [env["id"]] + assert claimed[0]["target_connection"] == "ssh-vps" + assert claimed[0]["message"] == "hi" + # second drain: nothing (no double delivery) + assert bot_relay.claim_pending_envelopes(root) == [] + + +def test_write_reply_validates_envelope_id(root): + with pytest.raises(ValueError): + bot_relay.write_reply(root, "../../etc/passwd", reply="x") + path = bot_relay.write_reply(root, "a" * 32, reply="pong") + data = json.loads(Path(path).read_text(encoding="utf-8")) + assert data["reply"] == "pong" and not data["error"] + + +def test_waiter_command_quotes_and_targets_reply_file(root): + env = {"id": "b" * 32, "target_handle": "researcher", "target_connection": "ssh-vps"} + cmd = bot_relay.waiter_command(root, env) + assert ("b" * 32) in cmd and "-c" in cmd + assert "rm -rf" not in cmd # sanity: single quoted -c payload + + +# ── message_agent integration: relay route + legacy-SOUL gate fix ─────────── + +import textwrap + + +def _managed_home(tmp_path, *, legacy_soul=False): + home = tmp_path / ".hermes" + home.mkdir(exist_ok=True) + d = home / "profiles" / "researcher" + d.mkdir(parents=True, exist_ok=True) + (d / "profile.yaml").write_text( + textwrap.dedent( + """\ + description: teammate for tests + ui_meta: + hermes-bots: + shape: cloud + """ + ), + encoding="utf-8", + ) + if legacy_soul: + (home / "SOUL.md").write_text( + "# Soul\n\n## Messaging other agents\nold shellout protocol\n", + encoding="utf-8", + ) + return home + + +class _FakeDB: + def __init__(self, home, title): + self.db_path = str(home / "state.db") + self._title = title + + def get_session_title(self, _sid): + return self._title + + +class _FakeAgent: + def __init__(self, home, title="Bot Chat"): + self._session_db = _FakeDB(home, title) + self.session_id = "sess-1" + self._session_title_hint = None + self._bot_mode_protocol = True + self.tools: list = [] + self.valid_tool_names: set = set() + + +@pytest.fixture(autouse=True) +def _fresh_probe_cache(): + from tools import bot_mode_probe + + bot_mode_probe._reset_cache_for_tests() + yield + bot_mode_probe._reset_cache_for_tests() + + +def test_tool_injects_despite_legacy_soul_protocol(tmp_path): + """The legacy-SOUL dedupe empties the SECTION, never the TOOL. + + Regression: upgraded installs whose SOUL.md still carries the old + plugin-appended protocol silently lost message_agent because the gate + keyed on section non-emptiness. + """ + from tools import bot_mode_probe + + home = _managed_home(tmp_path, legacy_soul=True) + # Premise: the dedupe really does empty the section for this profile... + assert bot_mode_probe.get_bot_mode_protocol_section(home) == "" + # ...but the install is managed, so the tool must still inject. + agent = _FakeAgent(home) + assert ensure_message_agent_tool(agent) is True + assert [t["function"]["name"] for t in agent.tools] == [MESSAGE_AGENT_TOOL_NAME] + + +def test_relay_route_queues_envelope_and_spawns_waiter(tmp_path, monkeypatch): + home = _managed_home(tmp_path) + bot_relay.write_remote_roster(home, [ + {"profile": "default", "handle": "hermes", "connection_id": "cloud-1", + "connection_label": "Hermes Cloud", "title": "Moxie"}, + ]) + + spawned = {} + + def _fake_spawn(command, label, *, task_id, agent): + spawned["command"] = command + spawned["label"] = label + return json.dumps({"status": "sent", "to": label}) + + monkeypatch.setattr("tools.bot_mode_dm._spawn_delivery", _fake_spawn) + agent = _FakeAgent(home) + out = json.loads(message_agent_tool(target="hermes", message="ping", agent=agent)) + assert out.get("status") == "sent" + assert "Hermes Cloud" in spawned["label"] + # envelope landed in the outbox with attribution prefixed + pending = bot_relay.claim_pending_envelopes(home) + assert len(pending) == 1 + assert pending[0]["target_connection"] == "cloud-1" + assert pending[0]["target_profile"] == "default" + assert pending[0]["message"].startswith("Message from 🤖 hermes (@hermes): ping") + # waiter watches this envelope's reply file + assert pending[0]["id"] in spawned["command"] + + +def test_relay_route_ambiguous_target_errors_with_forms(tmp_path, monkeypatch): + home = _managed_home(tmp_path) + bot_relay.write_remote_roster(home, [ + {"profile": "scout", "handle": "scout", "connection_id": "cloud-1"}, + {"profile": "scout", "handle": "scout", "connection_id": "ssh-vps"}, + ]) + monkeypatch.setattr( + "tools.bot_mode_dm._spawn_delivery", + lambda *a, **k: json.dumps({"status": "sent"}), + ) + agent = _FakeAgent(home) + out = json.loads(message_agent_tool(target="scout", message="hi", agent=agent)) + assert "scout@cloud-1" in out.get("error", "") and "scout@ssh-vps" in out["error"] + # connection-qualified form goes through + out2 = json.loads(message_agent_tool(target="scout@ssh-vps", message="hi", agent=agent)) + assert out2.get("status") == "sent" + + +def test_unknown_target_error_mentions_connected_machines(tmp_path): + home = _managed_home(tmp_path) + agent = _FakeAgent(home) + out = json.loads(message_agent_tool(target="ghost", message="hi", agent=agent)) + assert "connected machine" in out.get("error", "") + + +def test_protocol_section_lists_remote_teammates(tmp_path): + from tools import bot_mode_probe + + home = _managed_home(tmp_path) + bot_relay.write_remote_roster(home, [ + {"profile": "default", "handle": "hermes", "connection_id": "cloud-1", + "connection_label": "Hermes Cloud", "title": "Moxie"}, + ]) + section = bot_mode_probe.get_bot_mode_protocol_section(home, force_refresh=True) + assert "OTHER connected machines" in section + assert "`@hermes` — on Hermes Cloud — Moxie" in section + + +def test_capability_fingerprint_changes_with_relay_roster(tmp_path): + from tools import bot_mode_probe + + home = _managed_home(tmp_path) + before = bot_mode_probe.capability_fingerprint(home) + bot_relay.write_remote_roster(home, [ + {"profile": "default", "handle": "hermes", "connection_id": "cloud-1"}, + ]) + after = bot_mode_probe.capability_fingerprint(home) + assert before != after # eternal Bot Chats refresh once on roster change diff --git a/tests/tui_gateway/test_bot_relay_methods.py b/tests/tui_gateway/test_bot_relay_methods.py new file mode 100644 index 000000000000..747c031f1026 --- /dev/null +++ b/tests/tui_gateway/test_bot_relay_methods.py @@ -0,0 +1,107 @@ +"""Tests: bot_relay.* JSON-RPC handlers (tui_gateway/methods_bot_relay.py). + +The Desktop's relay door on each connected gateway. Contracts: +- roster.sync persists validated rows and reports the accepted count; +- outbox.drain returns queued envelopes exactly once; +- deliver validates the target profile against THIS install and runs the + one-turn Bot Chat transport (subprocess is faked here — the argv contract + is what's pinned); +- reply writes the waiter's file and rejects malformed envelope ids. +""" + +from __future__ import annotations + +import json + +import pytest + +import tui_gateway.server as srv +from tools import bot_relay + + +@pytest.fixture +def home(tmp_path, monkeypatch): + h = tmp_path / ".hermes" + (h / "profiles" / "ops").mkdir(parents=True) + monkeypatch.setenv("HERMES_HOME", str(h)) + return h + + +def _result(envelope): + assert "error" not in envelope, envelope + return envelope["result"] + + +def test_roster_sync_persists_and_counts(home): + out = _result( + srv._methods["bot_relay.roster.sync"]( + 1, + { + "agents": [ + {"profile": "scout", "handle": "scout", "connection_id": "cloud-1"}, + {"profile": "", "connection_id": "cloud-1"}, # dropped + ] + }, + ) + ) + assert out["count"] == 1 + assert [r["profile"] for r in bot_relay.read_remote_roster(home)] == ["scout"] + + +def test_outbox_drain_returns_each_envelope_once(home): + target = {"profile": "scout", "handle": "scout", "connection_id": "cloud-1", + "connection_label": "", "title": "", "description": ""} + env = bot_relay.enqueue_envelope( + home, target=target, message="m", sender_profile="default", sender_handle="hermes" + ) + first = _result(srv._methods["bot_relay.outbox.drain"](1, {})) + assert [e["id"] for e in first["envelopes"]] == [env["id"]] + second = _result(srv._methods["bot_relay.outbox.drain"](2, {})) + assert second["envelopes"] == [] + + +def test_deliver_validates_profile_and_runs_transport(home, monkeypatch): + calls = {} + + class _Proc: + returncode = 0 + stdout = "pong from ops" + stderr = "" + + def _fake_run(argv, **kwargs): + calls["argv"] = argv + return _Proc() + + monkeypatch.setattr("subprocess.run", _fake_run) + out = _result( + srv._methods["bot_relay.deliver"](1, {"profile": "ops", "message": "ping"}) + ) + assert out["reply"] == "pong from ops" + argv = calls["argv"] + assert argv[:3] == ["hermes", "-p", "ops"] + assert "Bot Chat" in argv and "--query-file" in argv + + # 'hermes' alias resolves to default + _result(srv._methods["bot_relay.deliver"](2, {"profile": "hermes", "message": "x"})) + assert calls["argv"][:3] == ["hermes", "-p", "default"] + + # unknown profile refuses without spawning + calls.clear() + err = srv._methods["bot_relay.deliver"](3, {"profile": "ghost", "message": "x"}) + assert "error" in err and "ghost" in err["error"]["message"] + assert not calls + + +def test_deliver_requires_params(home): + err = srv._methods["bot_relay.deliver"](1, {"profile": "", "message": ""}) + assert "error" in err + + +def test_reply_roundtrip_and_id_validation(home): + envelope_id = "c" * 32 + _result(srv._methods["bot_relay.reply"](1, {"id": envelope_id, "reply": "hi"})) + path = bot_relay.relay_root(home) / bot_relay.REPLIES_DIR / f"{envelope_id}.json" + assert json.loads(path.read_text(encoding="utf-8"))["reply"] == "hi" + + err = srv._methods["bot_relay.reply"](2, {"id": "../evil"}) + assert "error" in err diff --git a/tools/bot_mode_dm.py b/tools/bot_mode_dm.py index c286984c5f14..25785a98a7d8 100644 --- a/tools/bot_mode_dm.py +++ b/tools/bot_mode_dm.py @@ -84,9 +84,11 @@ def message_agent_tool_schema() -> dict: "clearly relevant teammate when it genuinely helps the user's goal; " "don't fan out to several agents unless the user explicitly asked. " "Use the teammate roster in your system prompt (names + roles) to pick " - "the right recipient; targets: a teammate name (e.g. 'researcher'), or " + "the right recipient; targets: a teammate name (e.g. 'researcher'), " "'/' for an agent on a registered peer gateway " - "(e.g. 'spark/researcher', or just '' for the peer's main agent)." + "(e.g. 'spark/researcher', or just '' for the peer's main agent), " + "or an agent on another connected machine from your roster (use " + "'@' if the same handle exists on several)." ), "parameters": { "type": "object", @@ -135,11 +137,15 @@ def ensure_message_agent_tool(agent: Any) -> bool: and tool.get("function", {}).get("name") == MESSAGE_AGENT_TOOL_NAME ): return True - from tools.bot_mode_probe import BOT_CHAT_TITLE, get_bot_mode_protocol_section + from tools.bot_mode_probe import BOT_CHAT_TITLE, is_bot_mode_managed if _session_title(agent) != BOT_CHAT_TITLE: return False - if not get_bot_mode_protocol_section(_agent_home(agent)): + # Managed-install check, NOT section non-emptiness: a profile whose + # SOUL.md carries the legacy plugin-appended protocol text gets an + # empty section (dedupe) but must still receive the tool — otherwise + # upgraded installs silently lose A2A messaging (Aug 2026). + if not is_bot_mode_managed(_agent_home(agent)): return False if agent.tools is None: agent.tools = [] @@ -235,7 +241,7 @@ def message_agent_tool( # ── defense-in-depth gate: only a canonical Bot Chat may deliver ── home = _agent_home(agent) try: - from tools.bot_mode_probe import BOT_CHAT_TITLE, get_bot_mode_protocol_section + from tools.bot_mode_probe import BOT_CHAT_TITLE, is_bot_mode_managed title = _session_title(agent) if title != BOT_CHAT_TITLE: @@ -243,7 +249,7 @@ def message_agent_tool( "message_agent is only available in a Bot Mode 'Bot Chat' session. " "This session is not one; do not retry." ) - if not get_bot_mode_protocol_section(home): + if not is_bot_mode_managed(home): return _err( "This install is not Bot-Mode-managed (no bot roster); " "message_agent is unavailable. Do not retry." @@ -289,17 +295,35 @@ def message_agent_tool( return _spawn_delivery(command, label, task_id=task_id, agent=agent) # ── local teammate ── - if not _LOCAL_TARGET_RE.match(raw_target): + if not _LOCAL_TARGET_RE.match(raw_target) and "@" not in raw_target: return _err(f"Invalid target: {raw_target!r}.", roster=teammates, peers=peers) - resolved = _resolve_local_name(raw_target, roster) + resolved = _resolve_local_name(raw_target, roster) if _LOCAL_TARGET_RE.match(raw_target) else None if resolved is None: + # ── cross-connection teammate (Desktop relay) ── + # Every gateway connected to the user's Desktop is reachable: the + # relay roster lists agents on the other connections; delivery rides + # the Desktop's own persistent socket to that gateway. + relayed = _try_relay_delivery( + root, raw_target, body, me, sender_handle, task_id=task_id, agent=agent + ) + if relayed is not None: + return relayed return _err( - f"No teammate named '{raw_target}' on this install. " - "Pick a name from the roster (roles are listed in your system prompt).", + f"No teammate named '{raw_target}' on this install, on a connected " + "machine, or on a registered peer. Pick a name from the roster " + "(roles are listed in your system prompt).", roster=teammates, peers=peers, ) if resolved == me: + # Same-name target on ANOTHER connection (e.g. this gateway's + # 'default' messaging the cloud 'default') — try the relay before + # calling it a self-message. + relayed = _try_relay_delivery( + root, raw_target, body, me, sender_handle, task_id=task_id, agent=agent + ) + if relayed is not None: + return relayed return _err("You can't message yourself. Pick a teammate from the roster.") dm_file = _write_dm_file(prefix + body) @@ -310,6 +334,64 @@ def message_agent_tool( return _spawn_delivery(command, f"@{_handle(resolved)}", task_id=task_id, agent=agent) +def _try_relay_delivery( + root: Path, + raw_target: str, + body: str, + me: str, + sender_handle: str, + *, + task_id: Optional[str], + agent: Any, +) -> Optional[str]: + """Cross-connection delivery via the Desktop relay, or None if the + target doesn't resolve against the relay roster. + + The envelope is queued on disk; the Desktop drains it over RPC and + delivers on the target connection's own socket. A background waiter is + spawned immediately so the relayed reply wakes the sender through the + standard completion-notification path — identical UX to a local DM. + """ + try: + from tools.bot_relay import ( + enqueue_envelope, + read_remote_roster, + resolve_remote_target, + waiter_command, + ) + + roster = read_remote_roster(root) + if not roster: + return None + match = resolve_remote_target(raw_target, roster) + if match is None: + return None + if match == "ambiguous": + forms = ", ".join( + f"{r['handle']}@{r['connection_id']}" + for r in roster + if r["handle"].lower() == raw_target.strip().lstrip("@").lower() + ) + return _err( + f"'{raw_target}' exists on several connected machines — " + f"disambiguate with one of: {forms}." + ) + envelope = enqueue_envelope( + root, + target=match, + message=f"Message from 🤖 {sender_handle} (@{sender_handle}): {body}", + sender_profile=me, + sender_handle=sender_handle, + ) + label = f"@{match['handle']} on {match['connection_label'] or match['connection_id']}" + return _spawn_delivery( + waiter_command(root, envelope), label, task_id=task_id, agent=agent + ) + except Exception: + logger.debug("relay delivery attempt failed", exc_info=True) + return None + + def _write_dm_file(content: str) -> str: """The message rides a temp file — never inline shell text.""" fd, path = tempfile.mkstemp(prefix="hermes-dm-", suffix=".txt", text=True) diff --git a/tools/bot_mode_probe.py b/tools/bot_mode_probe.py index 66935f8c7eb4..ca92c6f26a41 100644 --- a/tools/bot_mode_probe.py +++ b/tools/bot_mode_probe.py @@ -92,6 +92,24 @@ def _roster(root: Path) -> list[tuple[str, Path]]: return entries +def is_bot_mode_managed(home: str | os.PathLike | None = None) -> bool: + """True when ANY profile on this install is Bot-Mode-managed. + + The tool-injection gate for ``message_agent`` — deliberately independent + of :func:`get_bot_mode_protocol_section`'s emptiness: a profile whose + SOUL.md carries the legacy plugin-appended protocol gets an empty + section (text dedupe) but must still get the tool. Never raises. + """ + try: + resolved = Path( + str(home) if home else (os.getenv("HERMES_HOME") or os.path.expanduser("~/.hermes")) + ) + root = _hermes_root(resolved) + return any(_is_bot_managed(d) for _n, d in _roster(root)) + except Exception: + return False + + def _soul_has_protocol(profile_dir: Path) -> bool: try: soul = profile_dir / "SOUL.md" @@ -174,6 +192,38 @@ def _peers(root: Path) -> list[str]: return [] +def _remote_paragraph(root: Path) -> str: + """Protocol addendum for agents on OTHER connected machines. + + Fed by the Desktop relay roster (``tools/bot_relay.py``) — every gateway + connected to the user's Desktop (local, remote URL, SSH, Hermes Cloud, + docker) syncs its agents here, so bots can DM across machines with the + same message_agent tool. Only rendered when the relay roster is + non-empty. + """ + try: + from tools.bot_relay import read_remote_roster, remote_target_forms + + roster = read_remote_roster(root) + except Exception: + return "" + if not roster: + return "" + lines = [] + for row, form in zip(roster, remote_target_forms(roster)): + where = row["connection_label"] or row["connection_id"] + role = " — ".join(p for p in (row["title"], row["description"]) if p) + lines.append( + f"- `@{form}` — on {where}" + (f" — {role}" if role else "") + ) + return ( + "\n\nTeammates on OTHER connected machines (reachable through the " + "Desktop relay — message them with message_agent exactly like local " + "teammates; replies arrive as completion notifications the same " + "way):\n" + "\n".join(lines) + ) + + def _peer_paragraph(root: Path) -> str: """Protocol addendum for cross-machine DMs — only when peers exist.""" peers = _peers(root) @@ -230,6 +280,7 @@ def _build_section(home: Path) -> str: f"You are `@{handle}`. Your teammates (live roster; roles from their " "profiles):\n" f"{roster_block}" + + _remote_paragraph(root) + _peer_paragraph(root) ) @@ -340,6 +391,18 @@ def capability_fingerprint(home: str | os.PathLike | None = None) -> str: surface["peers"] = _peers(_hermes_root(resolved)) except Exception: surface["peers"] = [] + try: + # The Desktop relay roster is part of the messaging surface too: + # connecting/disconnecting a machine, or agents appearing on one, + # must refresh eternal Bot Chat prompts the same way. + from tools.bot_relay import read_remote_roster + + surface["remote_roster"] = sorted( + f"{r['connection_id']}:{r['profile']}:{r['title']}" + for r in read_remote_roster(_hermes_root(resolved)) + ) + except Exception: + surface["remote_roster"] = [] try: blob = json.dumps(surface, sort_keys=True).encode("utf-8") return hashlib.sha256(blob).hexdigest()[:12] diff --git a/tools/bot_relay.py b/tools/bot_relay.py new file mode 100644 index 000000000000..2e7e20a85c41 --- /dev/null +++ b/tools/bot_relay.py @@ -0,0 +1,339 @@ +"""Bot Mode cross-connection relay — connections ARE the peer set. + +Every gateway connected to the user's Desktop (local, remote URL, SSH, +Hermes Cloud, docker) is a persistent line. This module is the gateway-side +half of the relay that rides those lines so agents on ANY connected gateway +can find and message agents on ANY other, with `message_agent` as the one +send path (Teknium ruling, Aug 2026 — the peers-vs-connections split was +itself the bug). + +How the relay works (three files under ``/bot_relay/``): + +- ``roster.json`` — the union roster of agents on OTHER connections, pushed + by the Desktop over each connection's WebSocket (``bot_relay.roster.sync``). + ``tools/bot_mode_probe.py`` folds it into the Bot Chat protocol section so + every bot knows every reachable teammate, and ``message_agent`` resolves + cross-connection targets against it. +- ``outbox/`` — envelopes queued by ``message_agent`` for targets that live + on another connection. The Desktop drains them (``bot_relay.outbox.drain``) + and delivers each to the target connection (``bot_relay.deliver``). +- ``replies/`` — one JSON per envelope, written when the Desktop relays the + target agent's reply back (``bot_relay.reply``). A background waiter + spawned at send time watches for it, so the reply wakes the sender through + the exact same completion-notification path local DMs already use. + +The gateway never holds another connection's credentials; the Desktop owns +every socket and does all cross-connection I/O. Everything here is plain +file plumbing on the gateway's own HERMES root — no network, never raises +out of the public helpers. +""" + +from __future__ import annotations + +import json +import logging +import os +import re +import shlex +import sys +import tempfile +import time +import uuid +from pathlib import Path +from typing import Any, Optional + +logger = logging.getLogger(__name__) + +RELAY_DIR_NAME = "bot_relay" +ROSTER_FILE = "roster.json" +OUTBOX_DIR = "outbox" +CLAIMED_DIR = "claimed" +REPLIES_DIR = "replies" + +# A reply must arrive before the waiter gives up. Cross-connection turns can +# be slow (remote model, cold gateway) — generous, but bounded. +REPLY_WAIT_SECONDS = 900 + +# Envelopes and replies older than this are stale artifacts (Desktop was +# closed, connection died) and are swept opportunistically. +STALE_AFTER_SECONDS = 6 * 3600 + +_HANDLE_RE = re.compile(r"^[a-zA-Z0-9][a-zA-Z0-9_-]{0,63}$") + + +def relay_root(root: Path | str) -> Path: + return Path(root) / RELAY_DIR_NAME + + +def _ensure_dirs(root: Path | str) -> Path: + base = relay_root(root) + for sub in (OUTBOX_DIR, CLAIMED_DIR, REPLIES_DIR): + (base / sub).mkdir(parents=True, exist_ok=True) + return base + + +# ── remote roster ──────────────────────────────────────────────────────────── + + +def _normalize_roster_row(row: Any) -> Optional[dict]: + """Validated, minimal roster row or None. + + Rows come from the Desktop over RPC — treat as untrusted input. A row + names an agent on another connection: profile name, taggable handle, + the connection id/label of the gateway that owns it, and optional + friendly title/description for the protocol section. + """ + if not isinstance(row, dict): + return None + profile = str(row.get("profile") or "").strip() + handle = str(row.get("handle") or "").strip().lstrip("@") + connection_id = str(row.get("connection_id") or "").strip() + if not profile or not connection_id: + return None + if not handle: + handle = "hermes" if profile == "default" else profile + if not _HANDLE_RE.match(handle) or not _HANDLE_RE.match(profile): + return None + out = { + "profile": profile, + "handle": handle, + "connection_id": connection_id, + "connection_label": str(row.get("connection_label") or "").strip()[:80], + "title": str(row.get("title") or "").strip()[:120], + "description": " ".join(str(row.get("description") or "").split())[:160], + } + return out + + +def write_remote_roster(root: Path | str, rows: Any) -> int: + """Atomically persist the Desktop-pushed remote roster. Returns count.""" + base = _ensure_dirs(root) + cleaned: list[dict] = [] + seen: set[tuple[str, str]] = set() + for row in rows if isinstance(rows, list) else []: + norm = _normalize_roster_row(row) + if not norm: + continue + key = (norm["connection_id"], norm["profile"]) + if key in seen: + continue + seen.add(key) + cleaned.append(norm) + cleaned.sort(key=lambda r: (r["connection_id"], r["profile"])) + payload = {"updated_at": int(time.time()), "agents": cleaned} + target = base / ROSTER_FILE + fd, tmp = tempfile.mkstemp(dir=str(base), prefix=".roster-", suffix=".tmp") + try: + with os.fdopen(fd, "w", encoding="utf-8") as f: + json.dump(payload, f, ensure_ascii=False, sort_keys=True) + os.replace(tmp, target) + except Exception: + try: + os.unlink(tmp) + except OSError: + pass + raise + return len(cleaned) + + +def read_remote_roster(root: Path | str) -> list[dict]: + """The current remote roster (possibly empty). Never raises.""" + try: + raw = (relay_root(root) / ROSTER_FILE).read_text(encoding="utf-8") + data = json.loads(raw) + agents = data.get("agents") if isinstance(data, dict) else None + if not isinstance(agents, list): + return [] + return [r for r in (_normalize_roster_row(a) for a in agents) if r] + except FileNotFoundError: + return [] + except Exception: + logger.debug("bot_relay roster read failed", exc_info=True) + return [] + + +def resolve_remote_target(raw_target: str, roster: list[dict]) -> Any: + """Resolve ``raw_target`` against the remote roster. + + Accepted forms: + - bare handle/profile (``moxie``) — must be unique across connections; + - ``@`` / ``@`` — exact. + + Returns the matched row, the string ``"ambiguous"`` when a bare form + matches agents on several connections, or None for no match. + """ + want = str(raw_target or "").strip().lstrip("@") + if not want: + return None + conn: Optional[str] = None + if "@" in want: + want, _, conn = want.partition("@") + want = want.strip() + conn = conn.strip() + if not want or not conn: + return None + matches = [] + for row in roster: + if want.lower() not in (row["handle"].lower(), row["profile"].lower()): + continue + if conn and row["connection_id"].lower() != conn.lower(): + continue + matches.append(row) + if not matches: + return None + if len(matches) > 1: + return "ambiguous" + return matches[0] + + +def remote_target_forms(roster: list[dict]) -> list[str]: + """Human/agent-facing target strings, ambiguity-aware.""" + by_handle: dict[str, int] = {} + for row in roster: + by_handle[row["handle"].lower()] = by_handle.get(row["handle"].lower(), 0) + 1 + forms = [] + for row in roster: + if by_handle[row["handle"].lower()] > 1: + forms.append(f"{row['handle']}@{row['connection_id']}") + else: + forms.append(row["handle"]) + return forms + + +# ── outbox / replies ───────────────────────────────────────────────────────── + + +def enqueue_envelope( + root: Path | str, + *, + target: dict, + message: str, + sender_profile: str, + sender_handle: str, +) -> dict: + """Queue a cross-connection DM for the Desktop relay. Returns envelope.""" + base = _ensure_dirs(root) + envelope = { + "id": uuid.uuid4().hex, + "created_at": int(time.time()), + "from_profile": sender_profile, + "from_handle": sender_handle, + "target_connection": target["connection_id"], + "target_profile": target["profile"], + "target_handle": target["handle"], + "message": message, + } + path = base / OUTBOX_DIR / f"{envelope['id']}.json" + fd, tmp = tempfile.mkstemp(dir=str(base / OUTBOX_DIR), prefix=".env-", suffix=".tmp") + with os.fdopen(fd, "w", encoding="utf-8") as f: + json.dump(envelope, f, ensure_ascii=False) + os.replace(tmp, path) + return envelope + + +def claim_pending_envelopes(root: Path | str) -> list[dict]: + """Drain the outbox (rename → claimed/, so a second drain can't double- + deliver). Sweeps stale claimed/reply artifacts opportunistically.""" + base = _ensure_dirs(root) + _sweep_stale(base) + out: list[dict] = [] + outbox = base / OUTBOX_DIR + for path in sorted(outbox.glob("*.json")): + claimed = base / CLAIMED_DIR / path.name + try: + os.replace(path, claimed) # atomic claim + out.append(json.loads(claimed.read_text(encoding="utf-8"))) + except (OSError, ValueError): + continue + return out + + +def write_reply( + root: Path | str, envelope_id: str, *, reply: str = "", error: str = "" +) -> Path: + """Persist the relayed reply (or delivery error) for the waiter.""" + base = _ensure_dirs(root) + safe = str(envelope_id or "").strip() + if not re.match(r"^[0-9a-f]{32}$", safe): + raise ValueError(f"invalid envelope id: {envelope_id!r}") + path = base / REPLIES_DIR / f"{safe}.json" + payload = { + "id": safe, + "at": int(time.time()), + "reply": str(reply or ""), + "error": str(error or ""), + } + fd, tmp = tempfile.mkstemp(dir=str(base / REPLIES_DIR), prefix=".rep-", suffix=".tmp") + with os.fdopen(fd, "w", encoding="utf-8") as f: + json.dump(payload, f, ensure_ascii=False) + os.replace(tmp, path) + return path + + +def _sweep_stale(base: Path) -> None: + cutoff = time.time() - STALE_AFTER_SECONDS + for sub in (CLAIMED_DIR, REPLIES_DIR, OUTBOX_DIR): + try: + for path in (base / sub).glob("*.json"): + try: + if path.stat().st_mtime < cutoff: + path.unlink() + except OSError: + continue + except OSError: + continue + + +# ── waiter (runs on the sender gateway via terminal background process) ───── + + +def waiter_command(root: Path | str, envelope: dict) -> str: + """Shell command that blocks until the reply file appears, then prints it. + + Spawned with ``terminal_tool(background=True, notify_on_complete=True)`` + so its stdout — the teammate's reply — arrives as the same completion + notification local DMs use. Stdlib-only; runs under the sender gateway's + interpreter. + """ + reply_path = str(relay_root(root) / REPLIES_DIR / f"{envelope['id']}.json") + label = f"@{envelope['target_handle']} on {envelope['target_connection']}" + code = ( + "import json,os,sys,time\n" + f"p = {reply_path!r}\n" + f"deadline = time.time() + {REPLY_WAIT_SECONDS}\n" + "while time.time() < deadline:\n" + " if os.path.exists(p):\n" + " d = json.load(open(p, encoding='utf-8'))\n" + " if d.get('error'):\n" + f" print('Delivery to {label} failed: ' + d['error'])\n" + " sys.exit(1)\n" + f" print('Reply from {label}:')\n" + " print(d.get('reply') or '(empty reply)')\n" + " sys.exit(0)\n" + " time.sleep(2)\n" + f"print('No reply from {label} within {REPLY_WAIT_SECONDS}s. The message may " + "still be delivered when the Desktop reconnects; do not resend blindly.')\n" + "sys.exit(1)\n" + ) + return f"{shlex.quote(sys.executable or 'python3')} -c {shlex.quote(code)}" + + +# ── delivery command (used by the deliver RPC on the TARGET gateway) ──────── + + +def local_delivery_command(profile: str, query_file: str) -> list[str]: + """argv that delivers a DM into ``profile``'s Bot Chat on THIS gateway.""" + return [ + "hermes", + "-p", + profile, + "chat", + "--in", + "~", + "-c", + "Bot Chat", + "--create-if-missing", + "-Q", + "--query-file", + query_file, + ] diff --git a/tui_gateway/methods_bot_relay.py b/tui_gateway/methods_bot_relay.py new file mode 100644 index 000000000000..e337fca3c9a4 --- /dev/null +++ b/tui_gateway/methods_bot_relay.py @@ -0,0 +1,168 @@ +"""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. +""" + +from .method_ctx import HandlerRegistry + +_registry = HandlerRegistry() +method = _registry.method + + +@method("bot_relay.roster.sync") +def _(rid, params: dict) -> 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). + """ + try: + import os + from pathlib import Path + + from tools.bot_relay import write_remote_roster + + home = Path(os.getenv("HERMES_HOME") or os.path.expanduser("~/.hermes")) + root = home.parent.parent if home.parent.name == "profiles" else home + count = write_remote_roster(root, params.get("agents")) + return _ok(rid, {"count": count}) + except Exception as e: + return _err(rid, 5090, str(e)) + + +@method("bot_relay.outbox.drain") +def _(rid, params: dict) -> dict: + """Claim every pending cross-connection envelope queued on this gateway. + + Claimed envelopes move to ``claimed/`` atomically, so concurrent drains + (two Desktop windows) can't double-deliver. Result: ``{envelopes}``. + """ + try: + import os + from pathlib import Path + + from tools.bot_relay import claim_pending_envelopes + + home = Path(os.getenv("HERMES_HOME") or os.path.expanduser("~/.hermes")) + root = home.parent.parent if home.parent.name == "profiles" else home + return _ok(rid, {"envelopes": claim_pending_envelopes(root)}) + except Exception as e: + return _err(rid, 5091, str(e)) + + +@method("bot_relay.deliver") +def _(rid, params: dict) -> 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 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). + """ + import os + import subprocess + import tempfile + from pathlib import Path + + profile = str(params.get("profile") or "").strip() + message = str(params.get("message") or "").strip() + if not profile or not message: + return _err(rid, 4090, "profile and message required") + try: + from tools.bot_mode_dm import MESSAGE_MAX_CHARS + from tools.bot_relay import local_delivery_command + + if len(message) > MESSAGE_MAX_CHARS + 200: # + attribution headroom + return _err(rid, 4091, "message too long") + + home = Path(os.getenv("HERMES_HOME") or os.path.expanduser("~/.hermes")) + root = home.parent.parent if home.parent.name == "profiles" else home + known = {"default"} + profiles_dir = root / "profiles" + if profiles_dir.is_dir(): + known.update(c.name for c in profiles_dir.iterdir() if c.is_dir()) + resolved = "default" if profile.lower() == "hermes" else profile + if resolved not in known: + return _err(rid, 4092, f"no profile '{profile}' on this gateway") + + fd, tmp = tempfile.mkstemp(prefix="hermes-relay-dm-", suffix=".txt", text=True) + with os.fdopen(fd, "w", encoding="utf-8") as f: + f.write(message) + try: + proc = subprocess.run( + local_delivery_command(resolved, tmp), + capture_output=True, + text=True, + timeout=600, + ) + finally: + try: + os.unlink(tmp) + except OSError: + pass + if proc.returncode != 0: + detail = (proc.stderr or proc.stdout or "").strip()[-500:] + return _err(rid, 5092, f"delivery turn failed: {detail or proc.returncode}") + return _ok(rid, {"reply": (proc.stdout or "").strip()}) + except subprocess.TimeoutExpired: + return _err(rid, 5093, "delivery turn timed out") + except Exception as e: + return _err(rid, 5094, str(e)) + + +@method("bot_relay.reply") +def _(rid, params: dict) -> dict: + """Write a relayed reply (or delivery error) for a sender-side waiter. + + Params: ``id`` (envelope id), ``reply`` and/or ``error``. + """ + envelope_id = str(params.get("id") or "").strip() + if not envelope_id: + return _err(rid, 4093, "id required") + try: + import os + from pathlib import Path + + from tools.bot_relay import write_reply + + home = Path(os.getenv("HERMES_HOME") or os.path.expanduser("~/.hermes")) + root = home.parent.parent if home.parent.name == "profiles" else home + write_reply( + root, + envelope_id, + reply=str(params.get("reply") or ""), + error=str(params.get("error") or ""), + ) + return _ok(rid, {"ok": True}) + except ValueError as e: + return _err(rid, 4094, str(e)) + except Exception as e: + return _err(rid, 5095, str(e)) + + +def register(server) -> None: + _registry.install(server) diff --git a/tui_gateway/server.py b/tui_gateway/server.py index 4fb51eeec420..ef1b98b5f9b2 100644 --- a/tui_gateway/server.py +++ b/tui_gateway/server.py @@ -275,6 +275,14 @@ def _thread_panic_hook(args): "profiles.get_asset", "profiles.list", "profiles.set_asset", + # Bot-relay RPCs: roster.sync/outbox.drain/reply are cheap file I/O, + # but bot_relay.deliver runs a FULL one-turn agent conversation + # (subprocess, up to 600s) — all four stay off the WS reader thread + # so a slow relay delivery can never block prompt.submit. + "bot_relay.roster.sync", + "bot_relay.outbox.drain", + "bot_relay.deliver", + "bot_relay.reply", # image.generate is a multi-second remote API round-trip. "image.generate", "projects.discover_repos", @@ -15703,6 +15711,7 @@ def _mcp_summarize_server(name, cfg): # noqa: E402 # over already exists; register() rebinds them onto this namespace. from . import ( # noqa: E402 methods_browser_control as _methods_browser_control, + methods_bot_relay as _methods_bot_relay, methods_complete as _methods_complete, methods_config as _methods_config, methods_images as _methods_images, @@ -15721,6 +15730,7 @@ def _mcp_summarize_server(name, cfg): # noqa: E402 _methods_tools, _methods_profiles, _methods_images, + _methods_bot_relay, ): _m.register(sys.modules[__name__]) del _m diff --git a/website/docs/user-guide/bot-mode.md b/website/docs/user-guide/bot-mode.md index 8a9eb0c9a785..7a9e198cb530 100644 --- a/website/docs/user-guide/bot-mode.md +++ b/website/docs/user-guide/bot-mode.md @@ -94,7 +94,7 @@ Groups are standalone rows in the same activity-ordered roster as Bot DMs. A Bot Bots message each other with attribution, and you can hand work off from any chat: -- **@mentions** — type `@researcher have a look at this` in any chat and the composer's `@` autocomplete helps you pick the right Bot; on send, the mention is resolved against the live roster and the active Bot is told exactly who you mean (profile, friendly name, and device for cross-connection Bots). The Bot then composes its own message and sends it with `message_agent` — your text is never forwarded verbatim, and the reply comes back attributed to that agent. An email address or an unknown `@` passes through untouched. Reaching a Bot on another machine goes through a registered peer gateway (see `hermes peer` below). +- **@mentions** — type `@researcher have a look at this` in any chat and the composer's `@` autocomplete helps you pick the right Bot; on send, the mention is resolved against the live roster and the active Bot is told exactly who you mean (profile, friendly name, and device for cross-connection Bots). The Bot then composes its own message and sends it with `message_agent` — your text is never forwarded verbatim, and the reply comes back attributed to that agent. An email address or an unknown `@` passes through untouched. Bots on other connected machines are reachable the same way: the Desktop relays the message over that connection's own socket (see *Bots across machines* below). - **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 🤖 (@):` 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. @@ -109,6 +109,14 @@ agent: Bot-to-bot delivery is per-invocation: the receiving Bot picks the message up when it next runs. Live interrupt of a Bot mid-conversation is future work. ::: +### Messaging across connected machines (the Desktop relay) + +Every gateway you register in **Settings → Connections** — local, remote URL, SSH, Hermes Cloud, docker — is a persistent line the Desktop holds open, and Bot Mode uses those lines for messaging automatically. No extra setup: + +- **Rosters propagate on their own.** While the Desktop runs, it periodically tells each connected gateway which agents live on the *other* connections. Every Bot Chat's teammate roster then lists them ("Teammates on OTHER connected machines"), with names, roles, and which machine they're on — and the roster refreshes when agents appear, disappear, or get renamed (capability epoch). +- **`message_agent` reaches them directly.** A Bot on your laptop messages the cloud agent with `message_agent(target="moxie", …)` exactly like a local teammate. If the same handle exists on several machines, disambiguate with `target="moxie@"` (the tool's error tells the Bot the exact forms). Delivery rides the Desktop: the sending gateway queues the message, the Desktop relays it to the target connection's own gateway, the target Bot runs a turn in its canonical Bot Chat, and the reply comes back to the sender as the same background completion notification local DMs use. +- **The Desktop is the courier.** Cross-connection delivery works while a Desktop that knows both connections is running (it holds the sockets and the credentials — gateways never see each other's auth). If the Desktop is closed mid-delivery, the sender's Bot is told the reply didn't arrive rather than left hanging. For always-on machine-to-machine messaging with no Desktop in the loop, register a peer (`hermes peer`, below) — the two routes coexist. + ### Bot-initiated DMs across machines (`hermes peer`) Bots on one machine can message Bots on **another machine's gateway** without any desktop in the loop. Register the other gateway as a *peer* (its API server URL + `API_SERVER_KEY`): @@ -130,7 +138,7 @@ Requirements: the peer machine runs the `api_server` gateway platform with a str When you register several backends in **Settings → Connections** — the local runtime, remote gateways, SSH hosts, Hermes Cloud instances — the roster shows the Bots from **every** connected source, persistently: SSH sources are inventoried without spawning anything on the remote box, and machines that are momentarily unreachable keep their last-known rows instead of vanishing. When the same profile name exists on several sources, handles disambiguate as `@name-device` (for example `@research-homelab`). A Bot's chats, sessions, memory, and routines live on the machine that owns the profile. -Clicking a Connections Bot does **not** hop your window onto that machine — stay in your chat and `@mention` it, seat it in a group chat, or create new agents on it directly with the **Create on** picker. Cloud and local agents share one roster this way: register your Hermes Cloud instance and your desktop (say, over Tailscale or SSH) and their Bots can message each other and sit in the same rooms, with each agent's work running on its own machine. +Clicking a Connections Bot does **not** hop your window onto that machine — stay in your chat and `@mention` it, seat it in a group chat, or create new agents on it directly with the **Create on** picker. Cloud and local agents share one roster this way: register your Hermes Cloud instance and your desktop (say, over Tailscale or SSH) and their Bots can message each other and sit in the same rooms, with each agent's work running on its own machine. Bot-to-bot DMs across those machines go through the Desktop relay automatically (see *Messaging across connected machines* above). See [Connecting Desktop to Many Hermes Instances](./multi-connection-desktop.md) for the full multi-connection guide.