diff --git a/apps/desktop/src/plugins/hermes-bots/relay.test.ts b/apps/desktop/src/plugins/hermes-bots/relay.test.ts index ac46d5543e..99064c450a 100644 --- a/apps/desktop/src/plugins/hermes-bots/relay.test.ts +++ b/apps/desktop/src/plugins/hermes-bots/relay.test.ts @@ -115,19 +115,24 @@ async function loadRelay() { /** Fire the gateway's pending-envelope broadcast and let the debounced drain * run to completion. */ -async function pushAndSettle(times = 1) { - const listener = (hostMock.onEvent as ReturnType).mock.calls.at(-1)?.[1] as () => void +async function pushAndSettle(times = 1, event?: { connectionId?: string }) { + const listener = (hostMock.onEvent as ReturnType).mock.calls.at(-1)?.[1] as (payload?: { + connectionId?: string + }) => void for (let i = 0; i < times; i += 1) { - listener() + listener(event) } await vi.advanceTimersByTimeAsync(RELAY_PUSH_DEBOUNCE_MS + 10) } +const PLUGIN_DECISIONS_KEY = 'hermes.desktop.pluginDecisions.v2' + beforeEach(() => { vi.useFakeTimers() vi.clearAllMocks() + window.localStorage.removeItem(PLUGIN_DECISIONS_KEY) hostMock.onEvent = vi.fn(() => vi.fn()) hostMock.profileRoutes = vi.fn(async () => [route('a'), route('b')]) hostMock.requestProfile = vi.fn(async () => ({})) @@ -171,11 +176,11 @@ describe('push-notified drain (#93091)', () => { stopBotRelay() }) - it('keeps the interval poll as a BACKSTOP at the slow cadence', async () => { + it('keeps the interval poll as a BACKSTOP only when the shell cannot signal work', async () => { // The poll was 4s back when it WAS the delivery path — which (before // route retention) meant a fresh WebSocket dial + teardown per connection - // every 4s. Push carries envelope latency now; the poll only covers older - // backends and events that never reach the tap. + // every 4s. With the push door, an idle tick must not open a socket. + // Shells that cannot broadcast pending mail still poll. const calls = respondWith(() => ({ envelopes: [] })) const { startBotRelay, stopBotRelay } = await loadRelay() @@ -189,9 +194,22 @@ describe('push-notified drain (#93091)', () => { await vi.advanceTimersByTimeAsync(2000) - expect(calls.filter(call => call.method === 'bot_relay.outbox.drain')).toHaveLength(2) + expect(calls.filter(call => call.method === 'bot_relay.outbox.drain')).toHaveLength(0) stopBotRelay() + + hostMock.onEvent = undefined + calls.length = 0 + const legacy = await loadRelay() + + legacy.startBotRelay() + await vi.advanceTimersByTimeAsync(0) + calls.length = 0 + await vi.advanceTimersByTimeAsync(RELAY_DRAIN_INTERVAL_MS) + + expect(calls.filter(call => call.method === 'bot_relay.outbox.drain')).toHaveLength(2) + + legacy.stopBotRelay() }) it('re-schedules a push that raced an in-flight drain instead of dropping it', async () => { @@ -276,6 +294,53 @@ describe('push-notified drain (#93091)', () => { }) }) +describe('the 30s drain does not open a gateway socket with nothing to deliver (#118856)', () => { + it('does not dial when the push door is present and no route has outbox work', async () => { + const calls = respondWith(() => ({ envelopes: [] })) + const { startBotRelay, stopBotRelay } = await loadRelay() + + startBotRelay() + await vi.advanceTimersByTimeAsync(0) + calls.length = 0 + + await vi.advanceTimersByTimeAsync(RELAY_DRAIN_INTERVAL_MS) + + expect(calls.filter(call => call.method === 'bot_relay.outbox.drain')).toHaveLength(0) + + stopBotRelay() + }) + + it('does not dial a route that did not signal outbox work', async () => { + const calls = respondWith(() => ({ envelopes: [] })) + const { startBotRelay, stopBotRelay } = await loadRelay() + + startBotRelay() + await vi.advanceTimersByTimeAsync(0) + calls.length = 0 + + await pushAndSettle(1, { connectionId: 'a' }) + + expect(calls.filter(call => call.method === 'bot_relay.outbox.drain').map(call => call.connectionId)).toEqual(['a']) + + stopBotRelay() + }) + + it('does not dial when Bot Mode is off, even if an outbox event arrives', async () => { + window.localStorage.setItem(PLUGIN_DECISIONS_KEY, JSON.stringify({ 'hermes-bots': false })) + + const calls = respondWith(() => ({ envelopes: [{ id: 'env-1', target_connection: 'b' }] })) + const { startBotRelay, stopBotRelay } = await loadRelay() + + startBotRelay() + await pushAndSettle(1, { connectionId: 'a' }) + await vi.advanceTimersByTimeAsync(RELAY_DRAIN_INTERVAL_MS) + + expect(calls.filter(call => call.method === 'bot_relay.outbox.drain')).toHaveLength(0) + + stopBotRelay() + }) +}) + describe('relay-route socket retention (#93594)', () => { it('pins each connection ONCE across many drain ticks', async () => { const pins = trackRetention() diff --git a/apps/desktop/src/plugins/hermes-bots/relay.ts b/apps/desktop/src/plugins/hermes-bots/relay.ts index e159ea247d..ba3a0e6b80 100644 --- a/apps/desktop/src/plugins/hermes-bots/relay.ts +++ b/apps/desktop/src/plugins/hermes-bots/relay.ts @@ -9,7 +9,10 @@ import { host, LruCache } from '@hermes/plugin-sdk' +import { pluginActive } from '@/contrib/plugins-store' + import { botHandle, clearBotAttention, noteBotAttention } from './data' +import { ID } from './shared' import type { ProfileRoute, RosterRow } from './types' // ── cross-connection bot relay ──────────────────────────────────────────── @@ -28,8 +31,10 @@ import type { ProfileRoute, RosterRow } from './types' // relay degrades to whatever subset of connections supports it. const RELAY_ROSTER_INTERVAL_MS = 60_000 // Backstop cadence only (#93594): the push path below carries envelope latency, -// so the interval poll exists for older backends and missed events — 30s -// matches LIVE_SESSION_STATUS_BACKSTOP_INTERVAL_MS. It was 4s back when the +// so the interval poll exists for older backends that cannot signal +// `bot_relay.outbox.pending`. On a shell with that door, the tick must not +// open a gateway socket for a route that has not signaled outbox work — that +// dial/teardown is the idle ~30s cost (#118856). It was 4s back when the // poll WAS the delivery path, which (before route retention) also meant a // fresh WebSocket dial + teardown per registered connection every 4s. const RELAY_DRAIN_INTERVAL_MS = 30_000 @@ -106,6 +111,37 @@ const relay: RelayLifecycle = { // exemption). stopBotRelay releases everything. const relayRouteRetentions = new Map void>() +// Routes whose gateway has signaled `bot_relay.outbox.pending` since the last +// drain pass. `*` means the event carried no connection id (local/legacy +// primary, or an older payload) — every route may hold that envelope, so the +// pass still visits them. A route absent from this set has no outbox work; +// opening a socket to drain it is the idle dial. +const RELAY_OUTBOX_ANY = '*' +const routesWithOutboxWork = new Set() + +function relayBotModeOn(): boolean { + return pluginActive(ID) +} + +function noteRelayOutboxWork(event?: { connectionId?: string }) { + const id = String(event?.connectionId || '').trim() + + routesWithOutboxWork.add(id || RELAY_OUTBOX_ANY) +} + +function takeRelayOutboxWork(): { any: boolean; ids: Set } { + const any = routesWithOutboxWork.has(RELAY_OUTBOX_ANY) + const ids = new Set(routesWithOutboxWork) + + routesWithOutboxWork.clear() + + return { any, ids } +} + +function routeSignaledOutboxWork(connectionId: string, work: { any: boolean; ids: Set }): boolean { + return work.any || work.ids.has(connectionId) +} + /** One reachable gateway plus a representative route onto it. The route carries * identity only, so the human label comes from the registry (connectionLabels). */ interface RelayConnection { @@ -272,7 +308,7 @@ const relayAgentsCache = new LruCache(RELAY_AGENTS_CACH /** Push every gateway the union roster of agents on the OTHER connections. */ async function syncRelayRosters() { - if (relay.disposed || relay.rosterBusy) { + if (relay.disposed || relay.rosterBusy || !relayBotModeOn()) { return } @@ -372,7 +408,7 @@ async function syncRelayRosters() { * connection's own socket; the reply (or error) is posted back to the * sender gateway for its waiter. */ async function drainRelayOutboxes() { - if (relay.disposed) { + if (relay.disposed || !relayBotModeOn()) { return } @@ -398,6 +434,13 @@ async function drainRelayOutboxes() { return } + // A shell with the push door already knows which routes have mail. Do not + // open a throwaway socket to ask a route that has not signaled — Bot Mode + // off is the early return above; no outbox work is this skip (#118856). + // Older shells have no event, so the interval remains their only drain. + const pushDoor = typeof host.onEvent === 'function' + const work = pushDoor ? takeRelayOutboxWork() : { any: true, ids: new Set() } + const byId = new Map(connections.map(connection => [connection.id, connection])) // Phase 1 — claim every gateway's outbox before delivering anything. The @@ -410,6 +453,10 @@ async function drainRelayOutboxes() { const queued: RelayQueuedEnvelope[] = [] for (const sender of connections) { + if (pushDoor && !routeSignaledOutboxWork(sender.id, work)) { + continue + } + try { const res = await host.requestProfile<{ envelopes?: RelayEnvelope[] }>( sender.route, @@ -588,10 +635,15 @@ export function startBotRelay() { // Push path: the gateway change watcher broadcasts when an envelope hits // the outbox; drain immediately (debounced) instead of waiting the poll - // out. Feature-detected — older shells have no host.onEvent — and the 4s - // poll above stays untouched as the backstop either way. + // out. Feature-detected — older shells have no host.onEvent — and only + // those shells still open a socket from the interval. A signaled event + // names its connection when the renderer stamped one; an untagged event + // (local/legacy primary) may belong to any route. if (relay.pushUnsub === null && typeof host.onEvent === 'function') { - relay.pushUnsub = host.onEvent('bot_relay.outbox.pending', () => scheduleRelayPushDrain()) + relay.pushUnsub = host.onEvent('bot_relay.outbox.pending', event => { + noteRelayOutboxWork(event) + scheduleRelayPushDrain() + }) } } @@ -600,6 +652,7 @@ export function stopBotRelay() { // A rerun remembered mid-drain must not leak into the next start — // it would fire one stale drain after restart. relay.drainRerun = false + routesWithOutboxWork.clear() // Queued deliveries check `disposed` before they run; forget the lane tails // so a restart starts every target fresh instead of behind stale chains. relayLanes.clear()