From 6fdf873464ea25afecfae8ae534f337e774dc9e7 Mon Sep 17 00:00:00 2001 From: Deus Date: Tue, 25 Aug 2026 18:02:26 +0200 Subject: [PATCH] fix(desktop): serialize profile switches and drain edit redials --- .../gateway-connection-lifecycle.test.ts | 39 ++++++++++++ apps/desktop/src/store/gateway.ts | 62 +++++++++++++++++-- .../store/profile-agent-activation.test.ts | 55 ++++++++++++++++ apps/desktop/src/store/profile.ts | 14 +++-- 4 files changed, 158 insertions(+), 12 deletions(-) diff --git a/apps/desktop/src/store/gateway-connection-lifecycle.test.ts b/apps/desktop/src/store/gateway-connection-lifecycle.test.ts index 2ef43a617a..edd0c11215 100644 --- a/apps/desktop/src/store/gateway-connection-lifecycle.test.ts +++ b/apps/desktop/src/store/gateway-connection-lifecycle.test.ts @@ -64,6 +64,7 @@ const { openGatewayForProfile, pruneSecondaryGateways, reconnectSecondaryGateways, + retainGatewayForAgent, retireLocalProfileGateways, setPrimaryGateway } = await import('./gateway') @@ -166,6 +167,44 @@ describe('disposeSecondariesForConnection', () => { expect(getConnectionFor).toHaveBeenCalledTimes(2) }) + it('defers edit redials until request and foreground owners release the old sockets', async () => { + const foregroundScopes = new Set() + + const getConnectionFor = vi.fn(async ({ connectionId, profile }: { connectionId: string; profile: string }) => + descriptorFor(connectionId, profile) + ) + + configureGatewayRegistry({ foregroundScopes: () => foregroundScopes, onEvent: vi.fn() } as never) + installDesktop({ getConnectionFor }) + + await ensureGatewayForAgent('homelab', 'default') + await ensureGatewayForAgent('office', 'default') + const release = await retainGatewayForAgent('homelab', 'default') + await openGatewayForAgent('homelab', 'work') + foregroundScopes.add('conn:homelab::work') + + const retainedSocket = gatewayMocks.instances[0] + const foregroundSocket = gatewayMocks.instances[2] + + disposeSecondariesForConnection('homelab', { redial: true }) + + // An edit may need a new endpoint, but it cannot sever an in-flight turn + // or a mounted runtime's owner socket. No replacement is dialed yet. + expect(retainedSocket.close).not.toHaveBeenCalled() + expect(foregroundSocket.close).not.toHaveBeenCalled() + expect(gatewayMocks.connect).toHaveBeenCalledTimes(3) + + release() + await vi.waitFor(() => expect(gatewayMocks.connect).toHaveBeenCalledTimes(4)) + expect(retainedSocket.close).toHaveBeenCalledOnce() + expect(foregroundSocket.close).not.toHaveBeenCalled() + + foregroundScopes.clear() + pruneSecondaryGateways(new Set(['conn:homelab::default'])) + await vi.waitFor(() => expect(gatewayMocks.connect).toHaveBeenCalledTimes(5)) + expect(foregroundSocket.close).toHaveBeenCalledOnce() + }) + it('is a no-op for blank or unknown connection ids', async () => { installDesktop({ getConnectionFor: vi.fn(async ({ connectionId, profile }: { connectionId: string; profile: string }) => diff --git a/apps/desktop/src/store/gateway.ts b/apps/desktop/src/store/gateway.ts index b6d12e3ecd..f51c6ab99c 100644 --- a/apps/desktop/src/store/gateway.ts +++ b/apps/desktop/src/store/gateway.ts @@ -69,6 +69,8 @@ interface Secondary { reconnectTimer: ReturnType | null reconnectAttempt: number reconnecting: boolean + /** A material connection edit is waiting for live owners to drain. */ + pendingConnectionRedial: boolean /** * True when a foreground/prewarmed consumer owns this entry beyond one RPC. * Guards ONLY the dispose-at-refcount-0 paths (request/relay leases), never @@ -611,6 +613,7 @@ function createSecondary(profile: string, connectionId: null | string = null): S reconnectTimer: null, reconnectAttempt: 0, reconnecting: false, + pendingConnectionRedial: false, retained: false, relayRetainCount: 0, wantOpen: true, @@ -840,6 +843,7 @@ export async function requestGatewayForAgent( entry.activeRequests = Math.max(0, entry.activeRequests - 1) if ( + !drainPendingConnectionRedial(entry) && entry.activeRequests === 0 && !entry.retained && !relayRetained(entry) && @@ -887,6 +891,36 @@ function relayRetained(entry: Secondary): boolean { return Number.isFinite(entry.relayRetainCount) && entry.relayRetainCount > 0 } +/** + * Finish a material-edit redial once no request, relay, or foreground surface + * still owns the old socket. Removal deliberately bypasses this drain: a + * deleted source can never become valid again and must fail-stop immediately. + */ +function drainPendingConnectionRedial(entry: Secondary): boolean { + if ( + entry.pendingConnectionRedial !== true || + entry.activeRequests > 0 || + relayRetained(entry) || + foregroundPinned(entry) || + g.secondaries.get(entry.scope) !== entry + ) { + return false + } + + entry.pendingConnectionRedial = false + const wasActive = g.activeKey === entry.scope + disposeSecondary(entry) + g.secondaries.delete(entry.scope) + + const reopen = wasActive + ? ensureGatewayForAgent(entry.connectionId, entry.profile) + : openGatewayForAgent(entry.connectionId, entry.profile) + + void reopen.catch(() => undefined) + + return true +} + /** * Pin the pooled socket for one relay route open across drain ticks. Returns * a once-only release. Local routes (null/empty or explicit `local` source) @@ -926,6 +960,7 @@ export function retainGatewayForRelay(connectionId: null | string, profile: stri entry.relayRetainCount = Math.max(0, (entry.relayRetainCount || 0) - 1) if ( + !drainPendingConnectionRedial(entry) && entry.relayRetainCount === 0 && entry.activeRequests === 0 && !entry.retained && @@ -992,6 +1027,10 @@ export async function retainGatewayForAgent(connectionId: null | string, profile released = true entry.activeRequests = Math.max(0, entry.activeRequests - 1) + if (drainPendingConnectionRedial(entry)) { + return + } + if ( entry.activeRequests === 0 && !entry.retained && @@ -1471,6 +1510,10 @@ export function pruneSecondaryGateways(keep: Set): void { const now = Date.now() for (const [key, entry] of [...g.secondaries]) { + if (drainPendingConnectionRedial(entry)) { + continue + } + if ( key === g.activeKey || keep.has(key) || @@ -1604,12 +1647,13 @@ export function retireLocalProfileGateways(profile: string): void { } } -// Registry lifecycle: a connection was removed or materially edited. Dispose -// every secondary scoped to it (a removed remote/cloud source has no local -// process to die, so without this its WebSocket stays open streaming ghost -// events). With `redial` (the edit case) each disposed profile is re-dialed -// through the normal open path so the fresh socket targets the NEW endpoint; -// the active scope re-activates so the foreground keeps painting. +// Registry lifecycle: a connection was removed or materially edited. Removal +// disposes every scoped secondary immediately (a removed remote/cloud source +// has no local process to die, so otherwise its WebSocket streams ghost +// events). A material edit redials each profile through the normal open path so +// fresh sockets target the NEW endpoint, but request/relay leases and mounted +// foreground runtimes keep their old socket until they drain; the active scope +// re-activates when its replacement is safe to publish. export function disposeSecondariesForConnection(connectionId: string, opts: { redial?: boolean } = {}): void { const id = String(connectionId || '').trim() let activeInvalidated = false @@ -1626,6 +1670,12 @@ export function disposeSecondariesForConnection(connectionId: string, opts: { re const wasActive = key === g.activeKey activeInvalidated ||= wasActive + if (opts.redial && (entry.activeRequests > 0 || relayRetained(entry) || foregroundPinned(entry))) { + entry.pendingConnectionRedial = true + + continue + } + disposeSecondary(entry) g.secondaries.delete(key) diff --git a/apps/desktop/src/store/profile-agent-activation.test.ts b/apps/desktop/src/store/profile-agent-activation.test.ts index 74408a2619..c8838b39ce 100644 --- a/apps/desktop/src/store/profile-agent-activation.test.ts +++ b/apps/desktop/src/store/profile-agent-activation.test.ts @@ -162,6 +162,61 @@ describe('ensureGatewayAgent shares the gatewaySwitch mutex with profile switche expect($connection.get()?.profile).toBe('research') }) + it('serializes every profile waiter after three overlapping activations wake together', async () => { + const firstGate = deferred() + const secondGate = deferred() + const thirdGate = deferred() + const order: string[] = [] + + ensureGatewayForProfile.mockImplementation(async (profile: string) => { + order.push(`start:${profile}`) + + if (profile === 'worker') { + await firstGate.promise + } else if (profile === 'research') { + await secondGate.promise + } else if (profile === 'writer') { + await thirdGate.promise + } + + order.push(`finish:${profile}`) + }) + getConnection.mockImplementation(async profile => localConn({ profile: profile || 'default' })) + + const first = ensureGatewayProfile('worker') + await Promise.resolve() + const second = ensureGatewayProfile('research') + const third = ensureGatewayProfile('writer') + await Promise.resolve() + + expect(order).toEqual(['start:worker']) + + firstGate.resolve() + await first + await vi.waitFor(() => expect(order).toContain('start:research')) + + // The third activation was requested last, but it must not enter while the + // second waiter owns the switch mutex after both woke behind `worker`. + expect(order).not.toContain('start:writer') + + secondGate.resolve() + await second + await vi.waitFor(() => expect(order).toContain('start:writer')) + thirdGate.resolve() + await third + + expect(order).toEqual([ + 'start:worker', + 'finish:worker', + 'start:research', + 'finish:research', + 'start:writer', + 'finish:writer' + ]) + expect($activeGatewayProfile.get()).toBe('writer') + expect($connection.get()?.profile).toBe('writer') + }) + it('serializes a profile switch behind an in-flight agent activation', async () => { const agentGate = deferred() const order: string[] = [] diff --git a/apps/desktop/src/store/profile.ts b/apps/desktop/src/store/profile.ts index e0e72cbc91..0ca214a927 100644 --- a/apps/desktop/src/store/profile.ts +++ b/apps/desktop/src/store/profile.ts @@ -470,14 +470,16 @@ export async function ensureGatewayProfile(profile: string | null | undefined): return } - // Serialize concurrent activations so two rapid session switches don't race - // the active pointer. - if (gatewaySwitch) { + // Serialize concurrent activations so rapid session switches cannot race the + // active pointer. Re-acquire after every wake: multiple waiters can observe + // the same settled switch, and the first one starts the next switch before + // the others resume. + while (gatewaySwitch) { await gatewaySwitch.catch(() => undefined) + } - if (normalizeProfileKey($activeGatewayProfile.get()) === target && $gateway.get()) { - return - } + if (normalizeProfileKey($activeGatewayProfile.get()) === target && $gateway.get()) { + return } $gatewaySwapTarget.set(target)