diff --git a/apps/desktop/src/app/gateway/hooks/use-gateway-boot.ts b/apps/desktop/src/app/gateway/hooks/use-gateway-boot.ts index dc988eb089..a2dd1565dc 100644 --- a/apps/desktop/src/app/gateway/hooks/use-gateway-boot.ts +++ b/apps/desktop/src/app/gateway/hooks/use-gateway-boot.ts @@ -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() + // 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)) { diff --git a/apps/desktop/src/app/session/hooks/use-message-stream/gateway-event.ts b/apps/desktop/src/app/session/hooks/use-message-stream/gateway-event.ts index 2290e31504..f8344317bf 100644 --- a/apps/desktop/src/app/session/hooks/use-message-stream/gateway-event.ts +++ b/apps/desktop/src/app/session/hooks/use-message-stream/gateway-event.ts @@ -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) } diff --git a/apps/desktop/src/store/gateway-connection-scope.test.ts b/apps/desktop/src/store/gateway-connection-scope.test.ts new file mode 100644 index 0000000000..b930b5b4cf --- /dev/null +++ b/apps/desktop/src/store/gateway-connection-scope.test.ts @@ -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 => { + 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) + }) +}) diff --git a/apps/desktop/src/store/gateway.ts b/apps/desktop/src/store/gateway.ts index 1c2173b56c..9ad578592d 100644 --- a/apps/desktop/src/store/gateway.ts +++ b/apps/desktop/src/store/gateway.ts @@ -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): 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 } diff --git a/apps/desktop/src/store/session-states-scopes.test.ts b/apps/desktop/src/store/session-states-scopes.test.ts new file mode 100644 index 0000000000..b5b96474c7 --- /dev/null +++ b/apps/desktop/src/store/session-states-scopes.test.ts @@ -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> = {}) => ({ + ...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()) + }) +}) diff --git a/apps/desktop/src/store/session-states.ts b/apps/desktop/src/store/session-states.ts index 6407678bf9..3bd5fc4e13 100644 --- a/apps/desktop/src/store/session-states.ts +++ b/apps/desktop/src/store/session-states.ts @@ -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>({}) +// --------------------------------------------------------------------------- +// 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() + +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 { + const scopes = new Set() + + 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({}) } diff --git a/apps/desktop/src/types/hermes.ts b/apps/desktop/src/types/hermes.ts index 5535265d2c..c180cb08a5 100644 --- a/apps/desktop/src/types/hermes.ts +++ b/apps/desktop/src/types/hermes.ts @@ -455,6 +455,9 @@ export interface PaginatedSessions { export interface RpcEvent { 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 } diff --git a/apps/shared/src/json-rpc-gateway.ts b/apps/shared/src/json-rpc-gateway.ts index 6948db62e7..8a2713a3ee 100644 --- a/apps/shared/src/json-rpc-gateway.ts +++ b/apps/shared/src/json-rpc-gateway.ts @@ -27,6 +27,9 @@ export interface GatewayEvent

{ 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 }