fix(desktop): do not dial the bot relay drain with nothing to deliver
The 30s drain opened a throwaway gateway socket even when Bot Mode was off or the route had no outbox work. With the push door, only a route that signaled bot_relay.outbox.pending is drained. Older shells without that door still poll. The background-scope reset from #119836 is unchanged.
This commit is contained in:
@@ -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<typeof vi.fn>).mock.calls.at(-1)?.[1] as () => void
|
||||
async function pushAndSettle(times = 1, event?: { connectionId?: string }) {
|
||||
const listener = (hostMock.onEvent as ReturnType<typeof vi.fn>).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()
|
||||
|
||||
@@ -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<string, () => 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<string>()
|
||||
|
||||
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<string> } {
|
||||
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<string> }): 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<string, RelayAgentRow[]>(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<string>() }
|
||||
|
||||
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()
|
||||
|
||||
Reference in New Issue
Block a user