fix(desktop): key fan-out event consumption by (connectionId, profile)
Secondary-gateway events were tagged with connectionId (store/gateway fan-out) but no consumer read it: working/attention tracking, the pruneSecondaryGateways keep-set, and the profile-scoped event gates (skin.changed / change-watcher broadcasts / approval-mode reconcile) all keyed by session id + bare profile name. Every registered source exposes a 'default' profile (the roster force-unshifts it), so two connected gateways collided — gateway B's 'default' activity was attributed to gateway A's 'default', keeping the wrong socket alive and applying the wrong source's config/skin/cron changes. Thread connectionId through consumption using the existing composite backendScopeKey helper: - session-states records each registry-tagged event's (connectionId, profile) scope per runtime session; liveSessionScopes() projects the busy/needs-input ones as composite keys for the gateway keep-set. - recomputeKeptGateways (use-gateway-boot) seeds the keep-set with those scopes; pruneSecondaryGateways matches registry-scoped entries ONLY on their composite key, while local entries keep matching bare profile names (single-source path unchanged). - gateway-event's 'from the active profile' gates now compare the event's composite scope against the active gateway's connection via the new activeGatewayConnectionId(); untagged local/primary events behave byte-identically. Display-only surfaces that already use roster handles are untouched.
This commit is contained in:
@@ -39,7 +39,13 @@ import {
|
||||
setCurrentCwd,
|
||||
setSessionsLoading
|
||||
} from '@/store/session'
|
||||
import { $attentionSessionIds, $workingSessionIds, resetTileRuntimeBindings } from '@/store/session-states'
|
||||
import {
|
||||
$attentionSessionIds,
|
||||
$workingSessionIds,
|
||||
liveSessionScopes,
|
||||
recordSessionEventScope,
|
||||
resetTileRuntimeBindings
|
||||
} from '@/store/session-states'
|
||||
import { windowProfileOverride } from '@/store/windows'
|
||||
import type { RpcEvent } from '@/types/hermes'
|
||||
|
||||
@@ -395,7 +401,15 @@ export function useGatewayBoot({
|
||||
callbacksRef.current.onGatewayReady(gateway)
|
||||
setPrimaryGateway(gateway, survivor?.profile ?? normalizeProfileKey($activeGatewayProfile.get()))
|
||||
// Secondary (background-profile) sockets funnel into the same handler.
|
||||
configureGatewayRegistry({ onEvent: event => callbacksRef.current.handleGatewayEvent(event) })
|
||||
// Record each event's source scope first: registry-tagged events feed the
|
||||
// (connectionId, profile) keep-set so two sources exposing the same
|
||||
// profile name (every source has a 'default') can't collide.
|
||||
configureGatewayRegistry({
|
||||
onEvent: event => {
|
||||
recordSessionEventScope(event)
|
||||
callbacksRef.current.handleGatewayEvent(event)
|
||||
}
|
||||
})
|
||||
|
||||
const offState = gateway.onState(st => {
|
||||
// Mirror to the composer only while the primary is the active profile —
|
||||
@@ -471,7 +485,11 @@ export function useGatewayBoot({
|
||||
// to idle-reap. The active profile is always spared.
|
||||
const recomputeKeptGateways = () => {
|
||||
const live = new Set([...$workingSessionIds.get(), ...$attentionSessionIds.get()])
|
||||
const keep = new Set<string>()
|
||||
// Registry-scoped (connectionId, profile) scopes with live work. Two
|
||||
// sources can expose the same profile name (every source has a
|
||||
// 'default'), so bare profile names can't represent a non-local
|
||||
// source's liveness without keeping the wrong gateway alive.
|
||||
const keep = liveSessionScopes()
|
||||
|
||||
for (const session of $sessions.get()) {
|
||||
if (live.has(session.id)) {
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
import type { BillingBlock } from '@hermes/shared'
|
||||
import { backendScopeKey } from '@hermes/shared'
|
||||
import type { HermesSkin } from '@hermes/shared/skin'
|
||||
import type { QueryClient } from '@tanstack/react-query'
|
||||
import { type MutableRefObject, useCallback, useEffect, useRef } from 'react'
|
||||
@@ -27,7 +28,7 @@ import { billingCtaLabel, clearBillingBlock, runBillingRecovery, setBillingBlock
|
||||
import { clearClarifyRequest, normalizeChoices, setClarifyRequest, warnDroppedChoices } from '@/store/clarify'
|
||||
import { setSessionCompacting } from '@/store/compaction'
|
||||
import { refreshBackgroundProcesses } from '@/store/composer-status'
|
||||
import { $gateway } from '@/store/gateway'
|
||||
import { $gateway, activeGatewayConnectionId } from '@/store/gateway'
|
||||
import { applyGoalStatusText } from '@/store/goals'
|
||||
import {
|
||||
notifyCronChanged,
|
||||
@@ -335,6 +336,18 @@ export function useGatewayEventHandler(deps: GatewayEventDeps) {
|
||||
(event: RpcEvent) => {
|
||||
const payload = event.payload as GatewayEventPayload | undefined
|
||||
|
||||
// "From the active profile" must mean "from the active SOURCE": every
|
||||
// registered connection exposes a 'default' profile, so a bare profile
|
||||
// comparison attributes gateway B's 'default' events to gateway A's
|
||||
// 'default'. Compare the composite (connectionId, profile) scope with
|
||||
// backendScopeKey — untagged (local/primary) events keep the legacy
|
||||
// bare-profile behavior byte-identical.
|
||||
const fromActiveSource = (): boolean =>
|
||||
(!event.profile ||
|
||||
normalizeProfileKey(event.profile) === normalizeProfileKey($activeGatewayProfile.get())) &&
|
||||
backendScopeKey(event.connectionId ?? null, event.profile ?? null) ===
|
||||
backendScopeKey(activeGatewayConnectionId(), event.profile ?? null)
|
||||
|
||||
const occurredAt =
|
||||
typeof payload?.timestamp === 'number' && Number.isFinite(payload.timestamp)
|
||||
? payload.timestamp
|
||||
@@ -417,11 +430,8 @@ export function useGatewayEventHandler(deps: GatewayEventDeps) {
|
||||
return
|
||||
} else if (event.type === 'skin.changed') {
|
||||
// A runtime skin switch (Hermes activating an authored skin, or `/skin`
|
||||
// on another surface). Only the active profile's change repaints.
|
||||
const fromActiveProfile =
|
||||
!event.profile || normalizeProfileKey(event.profile) === normalizeProfileKey($activeGatewayProfile.get())
|
||||
|
||||
if (fromActiveProfile) {
|
||||
// on another surface). Only the active source+profile's change repaints.
|
||||
if (fromActiveSource()) {
|
||||
ingestBackendSkin(payload as HermesSkin | undefined, { apply: true })
|
||||
}
|
||||
|
||||
@@ -435,12 +445,10 @@ export function useGatewayEventHandler(deps: GatewayEventDeps) {
|
||||
) {
|
||||
// Change-watcher broadcasts (server._broadcast_watched_changes): the
|
||||
// backend's on-disk signature moved. Route to the live-sync ticks the
|
||||
// former pollers now subscribe to. Only the active profile's changes
|
||||
// apply — background profile sockets watch their own homes.
|
||||
const fromActiveChangeProfile =
|
||||
!event.profile || normalizeProfileKey(event.profile) === normalizeProfileKey($activeGatewayProfile.get())
|
||||
|
||||
if (fromActiveChangeProfile) {
|
||||
// former pollers now subscribe to. Only the active source+profile's
|
||||
// changes apply — background profile sockets (and other connections'
|
||||
// gateways) watch their own homes.
|
||||
if (fromActiveSource()) {
|
||||
if (event.type === 'pet.changed') {
|
||||
notifyPetChanged(payload as PetChangeMeta | undefined)
|
||||
} else if (event.type === 'cron.changed') {
|
||||
@@ -504,7 +512,7 @@ export function useGatewayEventHandler(deps: GatewayEventDeps) {
|
||||
isActiveEvent &&
|
||||
typeof payload?.approval_mode === 'string' &&
|
||||
event.profile &&
|
||||
normalizeProfileKey(event.profile) === normalizeProfileKey($activeGatewayProfile.get())
|
||||
fromActiveSource()
|
||||
) {
|
||||
reconcileApprovalModeForProfile(event.profile, payload.approval_mode)
|
||||
}
|
||||
|
||||
111
apps/desktop/src/store/gateway-connection-scope.test.ts
Normal file
111
apps/desktop/src/store/gateway-connection-scope.test.ts
Normal file
@@ -0,0 +1,111 @@
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
|
||||
// Cross-connection event/prune scoping: every registered source exposes a
|
||||
// 'default' profile (the roster force-unshifts it), so any consumer keyed by
|
||||
// the bare profile name conflates two connected gateways' activity. These
|
||||
// tests pin the composite (connectionId, profile) keying:
|
||||
// - pruneSecondaryGateways must NOT keep a registry-scoped socket alive off
|
||||
// another source's same-named profile (and vice versa) — registry entries
|
||||
// match only their composite backendScopeKey scope.
|
||||
// - the session-states scope ledger (recordSessionEventScope /
|
||||
// liveSessionScopes) turns registry-tagged live work into those composite
|
||||
// keep-set entries, and ignores untagged local/primary events.
|
||||
|
||||
const gatewayMocks = vi.hoisted(() => ({
|
||||
closed: [] as string[],
|
||||
setConnection: vi.fn()
|
||||
}))
|
||||
|
||||
vi.mock('@/hermes', () => ({
|
||||
HermesGateway: class {
|
||||
connectionState = 'closed'
|
||||
wsUrl = ''
|
||||
connect = async (wsUrl: string): Promise<void> => {
|
||||
this.wsUrl = wsUrl
|
||||
this.connectionState = 'open'
|
||||
}
|
||||
close = (): void => {
|
||||
gatewayMocks.closed.push(this.wsUrl)
|
||||
this.connectionState = 'closed'
|
||||
}
|
||||
onEvent = vi.fn(() => () => {})
|
||||
onState = vi.fn(() => () => {})
|
||||
}
|
||||
}))
|
||||
vi.mock('@/store/session', () => ({
|
||||
setConnection: gatewayMocks.setConnection,
|
||||
setGatewayState: vi.fn()
|
||||
}))
|
||||
vi.mock('@/store/notify-baseline', () => ({ markNativeNotifyBaseline: vi.fn() }))
|
||||
|
||||
const { closeSecondaryGateways, configureGatewayRegistry, openGatewayForAgent, pruneSecondaryGateways, setPrimaryGateway } =
|
||||
await import('./gateway')
|
||||
|
||||
function installDesktop(): void {
|
||||
;(window as unknown as { hermesDesktop: unknown }).hermesDesktop = {
|
||||
getConnection: vi.fn(async () => ({
|
||||
authMode: 'token',
|
||||
profile: 'default',
|
||||
token: 't',
|
||||
wsUrl: 'wss://local.invalid/api/ws?token=t'
|
||||
})),
|
||||
getConnectionFor: vi.fn(async ({ connectionId, profile }: { connectionId: string; profile: string }) => ({
|
||||
authMode: 'token',
|
||||
connectionId,
|
||||
profile,
|
||||
token: 't',
|
||||
wsUrl: `wss://${connectionId}.invalid/api/ws?profile=${profile}`
|
||||
})),
|
||||
getGatewayWsUrlFor: vi.fn(
|
||||
async ({ connectionId, profile }: { connectionId: string; profile: string }) =>
|
||||
`wss://${connectionId}.invalid/api/ws?profile=${profile}`
|
||||
),
|
||||
touchBackend: vi.fn(async () => undefined)
|
||||
}
|
||||
}
|
||||
|
||||
beforeEach(() => {
|
||||
installDesktop()
|
||||
configureGatewayRegistry({ onEvent: vi.fn() })
|
||||
setPrimaryGateway({ connectionState: 'open' } as never, 'default')
|
||||
gatewayMocks.closed = []
|
||||
})
|
||||
|
||||
afterEach(() => {
|
||||
closeSecondaryGateways()
|
||||
vi.clearAllMocks()
|
||||
delete (window as unknown as { hermesDesktop?: unknown }).hermesDesktop
|
||||
})
|
||||
|
||||
describe('pruneSecondaryGateways with registry-scoped entries', () => {
|
||||
it("does not keep a registry socket alive off another source's same-named profile", async () => {
|
||||
// Gateway B's 'default' — the roster row (connectionId 'homelab', profile
|
||||
// 'default'). The keep-set carries the bare 'default' profile because the
|
||||
// LOCAL source has live work; that must not pin homelab's socket.
|
||||
await openGatewayForAgent('homelab', 'default')
|
||||
|
||||
pruneSecondaryGateways(new Set(['default']))
|
||||
|
||||
expect(gatewayMocks.closed).toEqual(['wss://homelab.invalid/api/ws?profile=default'])
|
||||
})
|
||||
|
||||
it('keeps a registry socket whose composite scope has live work', async () => {
|
||||
await openGatewayForAgent('homelab', 'default')
|
||||
|
||||
pruneSecondaryGateways(new Set(['conn:homelab::default']))
|
||||
|
||||
expect(gatewayMocks.closed).toEqual([])
|
||||
})
|
||||
|
||||
it('still keeps a local (profile-keyed) secondary via its bare profile name', async () => {
|
||||
await openGatewayForAgent(null, 'research')
|
||||
|
||||
pruneSecondaryGateways(new Set(['research']))
|
||||
|
||||
expect(gatewayMocks.closed).toEqual([])
|
||||
|
||||
pruneSecondaryGateways(new Set())
|
||||
|
||||
expect(gatewayMocks.closed).toHaveLength(1)
|
||||
})
|
||||
})
|
||||
@@ -142,6 +142,22 @@ export function activeGateway(): HermesGateway | null {
|
||||
return g.secondaries.get(g.activeKey)?.gateway ?? null
|
||||
}
|
||||
|
||||
/**
|
||||
* The registry connection serving the gateway the user is currently looking
|
||||
* at — null for the local/legacy primary path and for profile-keyed (local)
|
||||
* secondaries. Event consumers pair this with the event's own `connectionId`
|
||||
* tag so "from the active profile" really means "from the active SOURCE":
|
||||
* two connected gateways can both expose a 'default' profile, and a bare
|
||||
* profile comparison attributed gateway B's 'default' activity to gateway A.
|
||||
*/
|
||||
export function activeGatewayConnectionId(): null | string {
|
||||
if (g.activeKey === g.primaryProfile) {
|
||||
return null
|
||||
}
|
||||
|
||||
return g.secondaries.get(g.activeKey)?.connectionId ?? null
|
||||
}
|
||||
|
||||
// Mirror a backend's connection state into the global composer state, but only
|
||||
// when that backend is the one the user is currently looking at. Lets the
|
||||
// composer reflect the active profile's socket without a background reconnect
|
||||
@@ -524,15 +540,17 @@ function restoreActiveToPrimaryIfEvicted(): void {
|
||||
}
|
||||
}
|
||||
|
||||
// Close + evict secondaries whose profile is neither active nor in `keep`
|
||||
// (profiles with a running / needs-input session). Bounds cost to live work.
|
||||
// `keep` carries PROFILE names (session ownership is profile-keyed), so a
|
||||
// registry-scoped entry survives when ITS profile has live work — matching on
|
||||
// the composite key alone would prune every non-local socket the moment the
|
||||
// user looks away.
|
||||
// Close + evict secondaries whose scope is neither active nor in `keep`
|
||||
// (scopes with a running / needs-input session). Bounds cost to live work.
|
||||
// `keep` carries PROFILE names for local/legacy entries and composite
|
||||
// backendScopeKey(connectionId, profile) scopes for registry-sourced live
|
||||
// work. A registry-scoped entry matches ONLY on its composite key: every
|
||||
// source exposes a 'default' profile, so matching a non-local entry on the
|
||||
// bare profile name kept gateway B's 'default' socket alive off gateway A's
|
||||
// 'default' activity (and vice versa) — cross-connection attribution.
|
||||
export function pruneSecondaryGateways(keep: Set<string>): void {
|
||||
for (const [key, entry] of [...g.secondaries]) {
|
||||
if (key === g.activeKey || keep.has(key) || keep.has(entry.profile)) {
|
||||
if (key === g.activeKey || keep.has(key) || (!entry.connectionId && keep.has(entry.profile))) {
|
||||
continue
|
||||
}
|
||||
|
||||
|
||||
74
apps/desktop/src/store/session-states-scopes.test.ts
Normal file
74
apps/desktop/src/store/session-states-scopes.test.ts
Normal file
@@ -0,0 +1,74 @@
|
||||
import { beforeEach, describe, expect, it } from 'vitest'
|
||||
|
||||
import { createClientSessionState } from '@/lib/chat-runtime'
|
||||
import {
|
||||
$sessionStates,
|
||||
clearAllSessionStates,
|
||||
dropSessionState,
|
||||
liveSessionScopes,
|
||||
publishSessionState,
|
||||
recordSessionEventScope
|
||||
} from '@/store/session-states'
|
||||
|
||||
/**
|
||||
* The (connectionId, profile) half of the gateway keep-set. Working/attention
|
||||
* ids are profile-blind, and every registered source exposes a 'default'
|
||||
* profile — so registry-sourced live work must surface as composite
|
||||
* backendScopeKey scopes, while untagged local/primary events contribute
|
||||
* nothing (their liveness keeps flowing through bare profile names).
|
||||
*/
|
||||
|
||||
const state = (patch: Partial<ReturnType<typeof createClientSessionState>> = {}) => ({
|
||||
...createClientSessionState('stored-1'),
|
||||
...patch
|
||||
})
|
||||
|
||||
beforeEach(() => {
|
||||
clearAllSessionStates()
|
||||
$sessionStates.set({})
|
||||
})
|
||||
|
||||
describe('liveSessionScopes', () => {
|
||||
it('maps a registry-tagged busy session to its composite scope', () => {
|
||||
recordSessionEventScope({ connectionId: 'homelab', profile: 'default', session_id: 'rt-1' })
|
||||
publishSessionState('rt-1', state({ busy: true }))
|
||||
|
||||
expect(liveSessionScopes()).toEqual(new Set(['conn:homelab::default']))
|
||||
})
|
||||
|
||||
it('includes needs-input sessions and drops settled ones', () => {
|
||||
recordSessionEventScope({ connectionId: 'homelab', profile: 'default', session_id: 'rt-1' })
|
||||
publishSessionState('rt-1', state({ busy: false, needsInput: true }))
|
||||
|
||||
expect(liveSessionScopes()).toEqual(new Set(['conn:homelab::default']))
|
||||
|
||||
publishSessionState('rt-1', state({ busy: false, needsInput: false }))
|
||||
|
||||
expect(liveSessionScopes()).toEqual(new Set())
|
||||
})
|
||||
|
||||
it('ignores untagged (local/primary) events — no connectionId, no scope', () => {
|
||||
recordSessionEventScope({ profile: 'default', session_id: 'rt-1' })
|
||||
publishSessionState('rt-1', state({ busy: true }))
|
||||
|
||||
expect(liveSessionScopes()).toEqual(new Set())
|
||||
})
|
||||
|
||||
it("keeps two sources' same-named 'default' profiles distinct", () => {
|
||||
recordSessionEventScope({ connectionId: 'homelab', profile: 'default', session_id: 'rt-a' })
|
||||
recordSessionEventScope({ connectionId: 'spark', profile: 'default', session_id: 'rt-b' })
|
||||
publishSessionState('rt-a', state({ busy: true }))
|
||||
publishSessionState('rt-b', state({ busy: true }))
|
||||
|
||||
expect(liveSessionScopes()).toEqual(new Set(['conn:homelab::default', 'conn:spark::default']))
|
||||
})
|
||||
|
||||
it('forgets a dropped runtime session', () => {
|
||||
recordSessionEventScope({ connectionId: 'homelab', profile: 'default', session_id: 'rt-1' })
|
||||
publishSessionState('rt-1', state({ busy: true }))
|
||||
dropSessionState('rt-1')
|
||||
publishSessionState('rt-1', state({ busy: true }))
|
||||
|
||||
expect(liveSessionScopes()).toEqual(new Set())
|
||||
})
|
||||
})
|
||||
@@ -18,6 +18,8 @@
|
||||
|
||||
import { atom, computed } from 'nanostores'
|
||||
|
||||
import { backendScopeKey } from '@hermes/shared'
|
||||
|
||||
import type { ClientSessionState } from '@/app/types'
|
||||
import { findGroup, findGroupOfPane, type LayoutNode } from '@/components/pane-shell/tree/model'
|
||||
import {
|
||||
@@ -55,6 +57,49 @@ import { isSecondaryWindow } from './windows'
|
||||
|
||||
export const $sessionStates = atom<Record<string, ClientSessionState>>({})
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Event-source scopes: which registry connection's socket delivered a runtime
|
||||
// session's events. Working/attention membership alone is profile-blind — two
|
||||
// connected gateways can both expose a 'default' profile, so the gateway
|
||||
// keep-set (pruneSecondaryGateways) must key live work by the composite
|
||||
// (connectionId, profile) scope, not the bare profile name. Recorded at
|
||||
// event fan-in (use-gateway-boot); local/primary events carry no connectionId
|
||||
// and record nothing, so single-source behavior is untouched.
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
const sessionScopeByRuntimeId = new Map<string, string>()
|
||||
|
||||
export function recordSessionEventScope(event: {
|
||||
connectionId?: string
|
||||
profile?: string
|
||||
session_id?: string
|
||||
}): void {
|
||||
if (event.session_id && event.connectionId) {
|
||||
sessionScopeByRuntimeId.set(event.session_id, backendScopeKey(event.connectionId, event.profile))
|
||||
}
|
||||
}
|
||||
|
||||
/** Composite scopes of registry-sourced sessions that are live (busy or
|
||||
* waiting on input) — the (connectionId, profile) half of the gateway
|
||||
* keep-set. Local-source live work keeps flowing through profile names. */
|
||||
export function liveSessionScopes(): Set<string> {
|
||||
const scopes = new Set<string>()
|
||||
|
||||
for (const [runtimeId, state] of Object.entries($sessionStates.get())) {
|
||||
if (!state || (!state.busy && !state.needsInput)) {
|
||||
continue
|
||||
}
|
||||
|
||||
const scope = sessionScopeByRuntimeId.get(runtimeId)
|
||||
|
||||
if (scope) {
|
||||
scopes.add(scope)
|
||||
}
|
||||
}
|
||||
|
||||
return scopes
|
||||
}
|
||||
|
||||
// Stored session ids whose authoritative state is still busy, but whose
|
||||
// runtime has produced no state publish for the watchdog window. Silence is
|
||||
// not completion: long tool calls can legitimately stay quiet, so this is a
|
||||
@@ -302,6 +347,7 @@ export function dropSessionState(runtimeId: string) {
|
||||
// cached runtime is dropped in the meantime.
|
||||
clearWatchdog(runtimeId)
|
||||
clearSessionProviderWait(runtimeId)
|
||||
sessionScopeByRuntimeId.delete(runtimeId)
|
||||
|
||||
const current = $sessionStates.get()
|
||||
setSessionStalled(current[runtimeId]?.storedSessionId, false)
|
||||
@@ -327,6 +373,7 @@ export function clearAllSessionStates() {
|
||||
sessionWatchdogTimers.clear()
|
||||
settledExpiry.clear()
|
||||
clearAllProviderWaits()
|
||||
sessionScopeByRuntimeId.clear()
|
||||
$stalledSessionIds.set([])
|
||||
$sessionStates.set({})
|
||||
}
|
||||
|
||||
@@ -455,6 +455,9 @@ export interface PaginatedSessions {
|
||||
export interface RpcEvent<T = unknown> {
|
||||
payload?: T
|
||||
profile?: string
|
||||
/** Registry connection whose socket delivered the event (renderer-side tag;
|
||||
* absent for the local/legacy primary path). */
|
||||
connectionId?: string
|
||||
session_id?: string
|
||||
type: string
|
||||
}
|
||||
|
||||
@@ -27,6 +27,9 @@ export interface GatewayEvent<P = unknown> {
|
||||
payload?: P
|
||||
/** Renderer-side source tag added by the Desktop gateway registry. */
|
||||
profile?: string
|
||||
/** Registry connection whose socket delivered the event (renderer-side tag;
|
||||
* absent for the local/legacy primary path). */
|
||||
connectionId?: string
|
||||
session_id?: string
|
||||
type: GatewayEventName
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user