fix(desktop): serialize profile switches and drain edit redials
This commit is contained in:
@@ -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<string>()
|
||||
|
||||
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 }) =>
|
||||
|
||||
@@ -69,6 +69,8 @@ interface Secondary {
|
||||
reconnectTimer: ReturnType<typeof setTimeout> | 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<T>(
|
||||
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<string>): 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)
|
||||
|
||||
|
||||
@@ -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[] = []
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user