diff --git a/packages/webui/src/daemon/session/DaemonSessionProvider.test.tsx b/packages/webui/src/daemon/session/DaemonSessionProvider.test.tsx index 4448aa016f1..35619f8a5cd 100644 --- a/packages/webui/src/daemon/session/DaemonSessionProvider.test.tsx +++ b/packages/webui/src/daemon/session/DaemonSessionProvider.test.tsx @@ -6934,6 +6934,173 @@ describe('DaemonSessionProvider', () => { expect(blocks).toEqual([]); }); + it('ignores streamed events from a replaced same-id attachment', async () => { + vi.stubGlobal( + 'fetch', + vi.fn(async () => new Response(null, { status: 204 })), + ); + const streamStarted = createDeferred(); + const releaseOldEvent = createDeferred(); + const source = createMockSession({ + sessionId: 'session-a', + clientId: 'client-a', + events: async function* staleEvents() { + streamStarted.resolve(); + await releaseOldEvent.promise; + yield { + id: 1, + v: 1, + type: 'session_update', + sessionId: 'session-a', + data: { + update: { + sessionUpdate: 'agent_message_chunk', + content: { type: 'text', text: 'stale output' }, + }, + }, + }; + }, + }); + const target = createMockSession({ + sessionId: 'session-a', + clientId: 'client-b', + replaySnapshot: createTextReplaySnapshot('replacement transcript'), + events: createIdleEvents(), + }); + sdkMocks.sessions.push(source, target); + let blocks: readonly DaemonTranscriptBlock[] = []; + let connection: DaemonConnectionState | undefined; + + function Harness() { + blocks = useDaemonTranscriptBlocks(); + connection = useDaemonConnection(); + return null; + } + + await renderWithProvider(, { + autoConnect: true, + sessionId: 'session-a', + clientId: 'client-a', + }); + await act(async () => { + await streamStarted.promise; + await flushPromises(); + }); + + act(() => { + root?.render( + + + , + ); + }); + await act(async () => { + await flushPromises(); + }); + expect(connection).toMatchObject({ + sessionId: 'session-a', + clientId: 'client-b', + }); + expect(blocks).toMatchObject([ + { kind: 'assistant', text: 'replacement transcript' }, + ]); + + releaseOldEvent.resolve(); + await act(async () => { + await flushPromises(); + await flushTranscriptDispatch(); + }); + + expect(JSON.stringify(blocks)).not.toContain('stale output'); + expect(blocks).toMatchObject([ + { kind: 'assistant', text: 'replacement transcript' }, + ]); + }); + + it('ignores session_closed from a replaced same-id attachment', async () => { + vi.stubGlobal( + 'fetch', + vi.fn(async () => new Response(null, { status: 204 })), + ); + const streamStarted = createDeferred(); + const releaseOldEvent = createDeferred(); + const source = createMockSession({ + sessionId: 'session-a', + clientId: 'client-a', + events: async function* staleEvents() { + streamStarted.resolve(); + await releaseOldEvent.promise; + yield { + id: 1, + v: 1, + type: 'session_closed', + sessionId: 'session-a', + data: { reason: 'client_close' }, + }; + }, + }); + const target = createMockSession({ + sessionId: 'session-a', + clientId: 'client-b', + replaySnapshot: createTextReplaySnapshot('replacement transcript'), + events: createIdleEvents(), + }); + sdkMocks.sessions.push(source, target); + let blocks: readonly DaemonTranscriptBlock[] = []; + let connection: DaemonConnectionState | undefined; + + function Harness() { + blocks = useDaemonTranscriptBlocks(); + connection = useDaemonConnection(); + return null; + } + + await renderWithProvider(, { + autoConnect: true, + sessionId: 'session-a', + clientId: 'client-a', + }); + await act(async () => { + await streamStarted.promise; + await flushPromises(); + }); + + act(() => { + root?.render( + + + , + ); + }); + await act(async () => { + await flushPromises(); + }); + + releaseOldEvent.resolve(); + await act(async () => { + await flushPromises(); + }); + + expect(connection).toMatchObject({ + status: 'connected', + sessionId: 'session-a', + clientId: 'client-b', + }); + expect(blocks).toMatchObject([ + { kind: 'assistant', text: 'replacement transcript' }, + ]); + }); + it('clears connection state before detach resolves', async () => { const detached = createDeferred(); const firstSession = createMockSession({ @@ -7305,6 +7472,90 @@ describe('DaemonSessionProvider', () => { }); }); + it('ignores a late heartbeat failure from a replaced same-id attachment', async () => { + vi.stubGlobal( + 'fetch', + vi.fn(async () => new Response(null, { status: 204 })), + ); + sdkMocks.capabilities.mockResolvedValue({ + v: 1, + mode: 'http-bridge', + features: ['client_heartbeat'], + modelServices: [], + workspaceCwd: '/mock-workspace', + }); + const releaseHeartbeat = createDeferred(); + const sourceHeartbeat = vi.fn(async () => { + await releaseHeartbeat.promise; + throw Object.assign(new Error('source session gone'), { status: 410 }); + }); + sdkMocks.sessions.push( + createMockSession({ + sessionId: 'session-a', + clientId: 'client-a', + heartbeat: sourceHeartbeat, + events: createIdleEvents(), + }), + createMockSession({ + sessionId: 'session-a', + clientId: 'client-b', + events: createIdleEvents(), + }), + ); + let connection: DaemonConnectionState | undefined; + + function Harness() { + connection = useDaemonConnection(); + return null; + } + + await renderWithProvider(, { + autoConnect: true, + sessionId: 'session-a', + clientId: 'client-a', + heartbeatIntervalMs: 20, + heartbeatFailureThreshold: 1, + }); + await vi.waitFor(() => expect(sourceHeartbeat).toHaveBeenCalled()); + + act(() => { + root?.render( + + + , + ); + }); + await act(async () => { + await flushPromises(); + }); + expect(connection).toMatchObject({ + status: 'connected', + sessionId: 'session-a', + clientId: 'client-b', + }); + + releaseHeartbeat.resolve(); + await act(async () => { + await wait(30); + await flushPromises(); + }); + + expect(connection).toMatchObject({ + status: 'connected', + sessionId: 'session-a', + clientId: 'client-b', + missingSession: false, + }); + expect(connection?.error).toBeUndefined(); + }); + it('clears stale sessions on terminal HTTP heartbeat errors', async () => { sdkMocks.capabilities.mockResolvedValue({ v: 1, @@ -7694,6 +7945,10 @@ describe('DaemonSessionProvider', () => { }); it('keeps the current transcript when a same-session reload is aborted', async () => { + vi.stubGlobal( + 'fetch', + vi.fn(async () => new Response(null, { status: 204 })), + ); const replacement = createDeferred(); const currentSession = createMockSession({ sessionId: 'session-a', @@ -7717,27 +7972,164 @@ describe('DaemonSessionProvider', () => { async () => replacement.promise, ); const controller = new AbortController(); - const reload = requireActions(actions).reloadSession(controller.signal); + let reloadOutcome: Promise | undefined; + act(() => { + reloadOutcome = requireActions(actions) + .reloadSession(controller.signal) + .then( + () => undefined, + (error: unknown) => error, + ); + }); await act(async () => { await flushPromises(); }); - controller.abort(); const refreshedSession = createMockSession({ sessionId: 'session-a', replaySnapshot: createTextReplaySnapshot('replacement transcript'), }); - replacement.resolve(refreshedSession); + let outcome: unknown; await act(async () => { - await expect(reload).rejects.toMatchObject({ name: 'AbortError' }); + controller.abort(); + replacement.resolve(refreshedSession); + outcome = await reloadOutcome; await flushPromises(); }); + expect(outcome).toMatchObject({ name: 'AbortError' }); expect(blocks).toMatchObject([ { kind: 'assistant', text: 'current transcript' }, ]); expect(currentSession.detach).not.toHaveBeenCalled(); expect(refreshedSession.detach).toHaveBeenCalledOnce(); + + await act(async () => { + root?.unmount(); + root = null; + await flushPromises(); + }); + }); + + it('detaches the source attachment once after a same-session reload', async () => { + const detachFetch = vi.fn( + async (_input: RequestInfo | URL, _init?: RequestInit) => + new Response(null, { status: 204 }), + ); + vi.stubGlobal('fetch', detachFetch); + const currentSession = createMockSession({ + sessionId: 'session-a', + clientId: 'client-a', + events: createIdleEvents(), + }); + const replacementSession = createMockSession({ + sessionId: 'session-a', + clientId: 'client-b', + events: createIdleEvents(), + }); + sdkMocks.sessions.push(currentSession); + let actions: DaemonSessionActions | undefined; + + function Harness() { + actions = useDaemonActions(); + return null; + } + + await renderWithProvider(, { autoConnect: true }); + await act(async () => { + await flushPromises(); + }); + sdkMocks.MockDaemonSessionClient.load.mockResolvedValueOnce( + replacementSession, + ); + + let reload: Promise | undefined; + act(() => { + reload = requireActions(actions).reloadSession( + new AbortController().signal, + ); + }); + await act(async () => { + await flushPromises(); + }); + await expect(reload).resolves.toBeUndefined(); + + await act(async () => { + root?.unmount(); + root = null; + await flushPromises(); + }); + + expect(currentSession.detach).toHaveBeenCalledOnce(); + expect(detachFetch).toHaveBeenCalledOnce(); + expect( + new Headers(detachFetch.mock.calls[0][1]?.headers).get( + 'X-Qwen-Client-Id', + ), + ).toBe('client-b'); + }); + + it('retires both attachments when unmounted during a same-session reload', async () => { + const detachFetch = vi.fn( + async (_input: RequestInfo | URL, _init?: RequestInit) => + new Response(null, { status: 204 }), + ); + vi.stubGlobal('fetch', detachFetch); + const replacement = createDeferred(); + const source = createMockSession({ + sessionId: 'session-a', + clientId: 'client-a', + events: createIdleEvents(), + }); + const target = createMockSession({ + sessionId: 'session-a', + clientId: 'client-b', + events: createIdleEvents(), + }); + sdkMocks.sessions.push(source); + let actions: DaemonSessionActions | undefined; + + function Harness() { + actions = useDaemonActions(); + return null; + } + + await renderWithProvider(, { autoConnect: true }); + sdkMocks.MockDaemonSessionClient.load.mockClear(); + sdkMocks.MockDaemonSessionClient.load.mockImplementationOnce( + async () => replacement.promise, + ); + let reloadOutcome: Promise | undefined; + act(() => { + reloadOutcome = requireActions(actions) + .reloadSession(new AbortController().signal) + .then( + () => undefined, + (error: unknown) => error, + ); + }); + await act(async () => { + await flushPromises(); + }); + expect(sdkMocks.MockDaemonSessionClient.load).toHaveBeenCalledOnce(); + + await act(async () => { + root?.unmount(); + root = null; + }); + await expect(reloadOutcome).resolves.toMatchObject({ name: 'AbortError' }); + + replacement.resolve(target); + await act(async () => { + await flushPromises(); + }); + + const detachedClientIds = detachFetch.mock.calls.map(([, init]) => + new Headers(init?.headers).get('X-Qwen-Client-Id'), + ); + expect(detachedClientIds).toEqual( + expect.arrayContaining(['client-a', 'client-b']), + ); }); it('clears transcript immediately for default session switches', async () => { @@ -7926,7 +8318,11 @@ describe('DaemonSessionProvider', () => { }); }); - it('ignores stale metadata after switching to another session', async () => { + it('ignores stale metadata from a replaced same-id attachment', async () => { + vi.stubGlobal( + 'fetch', + vi.fn(async () => new Response(null, { status: 204 })), + ); const providersA = createDeferred(); const commandsA = createDeferred>>(); @@ -7936,13 +8332,15 @@ describe('DaemonSessionProvider', () => { sdkMocks.sessions.push( createMockSession({ sessionId: 'session-a', + clientId: 'client-a', replaySnapshot: createTextReplaySnapshot('session a transcript'), supportedCommands: vi.fn(() => commandsA.promise), context: vi.fn(() => contextA.promise), }), createMockSession({ - sessionId: 'session-b', - replaySnapshot: createTextReplaySnapshot('session b transcript'), + sessionId: 'session-a', + clientId: 'client-b', + replaySnapshot: createTextReplaySnapshot('replacement transcript'), }), ); let blocks: readonly DaemonTranscriptBlock[] = []; @@ -7957,6 +8355,7 @@ describe('DaemonSessionProvider', () => { await renderWithProvider(, { autoConnect: true, sessionId: 'session-a', + clientId: 'client-a', }); await act(async () => { await flushPromises(); @@ -7971,7 +8370,8 @@ describe('DaemonSessionProvider', () => { , @@ -7981,11 +8381,12 @@ describe('DaemonSessionProvider', () => { await flushPromises(); }); expect(connection).toMatchObject({ - sessionId: 'session-b', + sessionId: 'session-a', + clientId: 'client-b', loadingTranscript: undefined, }); expect(blocks).toMatchObject([ - { kind: 'assistant', text: 'session b transcript' }, + { kind: 'assistant', text: 'replacement transcript' }, ]); providersA.resolve({ @@ -8010,10 +8411,13 @@ describe('DaemonSessionProvider', () => { await flushPromises(); }); - expect(connection).toMatchObject({ sessionId: 'session-b' }); + expect(connection).toMatchObject({ + sessionId: 'session-a', + clientId: 'client-b', + }); expect(connection?.currentModel).not.toBe('stale-model'); expect(blocks).toMatchObject([ - { kind: 'assistant', text: 'session b transcript' }, + { kind: 'assistant', text: 'replacement transcript' }, ]); }); diff --git a/packages/webui/src/daemon/session/DaemonSessionProvider.tsx b/packages/webui/src/daemon/session/DaemonSessionProvider.tsx index 1cb47e3297e..80e6778127e 100644 --- a/packages/webui/src/daemon/session/DaemonSessionProvider.tsx +++ b/packages/webui/src/daemon/session/DaemonSessionProvider.tsx @@ -444,7 +444,7 @@ const TERMINAL_SESSION_HTTP_STATUSES = new Set([ ]); interface HeartbeatFailureState { - sessionId?: string; + session?: DaemonSessionClient; consecutiveFailures: number; lastHttpError?: { status: number; message: string }; } @@ -578,9 +578,9 @@ export function DaemonSessionProvider(props: DaemonSessionProviderProps) { consecutiveFailures: 0, }); const manualSessionClearRef = useRef(false); - const skipNextCleanupDetachSessionIdRef = useRef( - undefined, - ); + const skipNextCleanupDetachSessionRef = useRef< + DaemonSessionClient | undefined + >(undefined); const settledRestoredActivePromptSessionsRef = useRef< WeakSet >(new WeakSet()); @@ -687,7 +687,8 @@ export function DaemonSessionProvider(props: DaemonSessionProviderProps) { restoreMode === 'load' && restoreSessionId !== undefined && restoreSessionId === sessionRef.current?.sessionId && - restoreSessionId === skipNextCleanupDetachSessionIdRef.current; + sessionRef.current === skipNextCleanupDetachSessionRef.current; + let runnerSession = sessionRef.current; // ── Batched transcript dispatch ──────────────────────────────── // The live SSE loop dispatches transcript events through this batcher @@ -706,9 +707,17 @@ export function DaemonSessionProvider(props: DaemonSessionProviderProps) { // streaming stays at ~one dispatch per network chunk. let pendingTranscriptEvents: DaemonUiEvent[] = []; let transcriptFlushTimer: ReturnType | undefined; - const runTranscriptFlush = () => { + const runTranscriptFlush = (force = false) => { transcriptFlushTimer = undefined; if (pendingTranscriptEvents.length === 0) return; + if ( + !force && + runnerSession !== undefined && + sessionRef.current !== runnerSession + ) { + pendingTranscriptEvents = []; + return; + } const batch = pendingTranscriptEvents; pendingTranscriptEvents = []; // Swallow a reducer throw (log it) so it cannot escape as an uncaught @@ -740,7 +749,7 @@ export function DaemonSessionProvider(props: DaemonSessionProviderProps) { // buffered content stays correctly ordered ahead of it. const flushTranscriptSync = () => { cancelTranscriptFlush(); - runTranscriptFlush(); + runTranscriptFlush(true); }; const dispatchTranscriptNow = (events: DaemonUiEvent | DaemonUiEvent[]) => { flushTranscriptSync(); @@ -1168,11 +1177,8 @@ export function DaemonSessionProvider(props: DaemonSessionProviderProps) { new DOMException('Session load cancelled', 'AbortError'), ); } - if ( - skipNextCleanupDetachSessionIdRef.current === - nextSession.sessionId - ) { - skipNextCleanupDetachSessionIdRef.current = undefined; + if (skipNextCleanupDetachSessionRef.current === previousSession) { + skipNextCleanupDetachSessionRef.current = undefined; } loadingRequestedSession = false; if (previousSession?.sessionId === nextSession.sessionId) { @@ -1225,10 +1231,9 @@ export function DaemonSessionProvider(props: DaemonSessionProviderProps) { ); } if ( - skipNextCleanupDetachSessionIdRef.current === - nextSession.sessionId + skipNextCleanupDetachSessionRef.current === previousSession ) { - skipNextCleanupDetachSessionIdRef.current = undefined; + skipNextCleanupDetachSessionRef.current = undefined; } if (previousSession?.sessionId === nextSession.sessionId) { session = previousSession; @@ -1297,6 +1302,7 @@ export function DaemonSessionProvider(props: DaemonSessionProviderProps) { } const activeSession = session; + runnerSession = activeSession; // Prompt activity is session state returned by /load. Surface it // immediately so a refreshed page shows the running state without // waiting for auxiliary data such as providers, commands, or context. @@ -1692,11 +1698,8 @@ export function DaemonSessionProvider(props: DaemonSessionProviderProps) { if (pendingLoadToResolve.timeout !== undefined) { clearTimeout(pendingLoadToResolve.timeout); } - if ( - skipNextCleanupDetachSessionIdRef.current === - activeSession.sessionId - ) { - skipNextCleanupDetachSessionIdRef.current = undefined; + if (skipNextCleanupDetachSessionRef.current === activeSession) { + skipNextCleanupDetachSessionRef.current = undefined; } pendingLoadToResolve.resolve(); } @@ -1726,6 +1729,13 @@ export function DaemonSessionProvider(props: DaemonSessionProviderProps) { : activeSession.context(), gitPromise, ]); + if ( + disposed || + abort.signal.aborted || + sessionRef.current !== activeSession + ) { + return; + } const providers = providerResult?.status === 'fulfilled' ? providerResult.value @@ -1778,7 +1788,12 @@ export function DaemonSessionProvider(props: DaemonSessionProviderProps) { getCurrentMode(context) ?? providerModelStatus.currentMode; setConnection((current) => { - if (current.sessionId !== activeSession.sessionId) return current; + if ( + sessionRef.current !== activeSession || + current.sessionId !== activeSession.sessionId + ) { + return current; + } return { ...current, status: 'connected', @@ -1914,7 +1929,7 @@ export function DaemonSessionProvider(props: DaemonSessionProviderProps) { maxQueued, ...(sseConnectReason ? { sseConnectReason } : {}), })) { - if (sessionRef.current?.sessionId !== activeSession.sessionId) { + if (sessionRef.current !== activeSession) { break; } if (!sawEvent) { @@ -1943,9 +1958,11 @@ export function DaemonSessionProvider(props: DaemonSessionProviderProps) { // Keep the sidechannel for queue dedupe, but still normalize the // event below so chat UIs can render the inserted-message status. publishSidechannelMidTurnInjected(midTurnInjected); + if (sessionRef.current !== activeSession) break; } if (isPendingPromptEvent(event)) { publishPendingPromptEvent(event); + if (sessionRef.current !== activeSession) break; if (event.type === 'pending_prompt_started') { clearPassiveAssistantDoneTimer(passiveAssistantDoneTimerRef); setPromptStatus('waiting'); @@ -2219,6 +2236,7 @@ export function DaemonSessionProvider(props: DaemonSessionProviderProps) { break; } } catch (error) { + if (sessionRef.current !== activeSession) break; const message = error instanceof Error ? error.message : String(error); addNotice({ @@ -2236,6 +2254,15 @@ export function DaemonSessionProvider(props: DaemonSessionProviderProps) { ); } } + if ( + sessionRef.current !== activeSession && + !resyncRequested && + !userDeletedSession + ) { + clearPendingTranscriptEvents(); + clearEventStream(); + return; + } // The stream ended or broke: apply any buffered transcript events now // so post-loop handling (and consumers reading the snapshot) see a // complete transcript without waiting for the scheduled flush. @@ -2281,7 +2308,7 @@ export function DaemonSessionProvider(props: DaemonSessionProviderProps) { nextSseConnectReason = 'stream_end'; // Keep the session handle after a normal SSE close so the next // subscription can resume from DaemonSessionClient.lastEventId. - if (sessionRef.current?.sessionId === activeSession.sessionId) { + if (sessionRef.current === activeSession) { console.debug('[DaemonSessionProvider] SSE stream ended'); if (!hasSessionActivePrompt()) { // A transport close is only a safe "done" signal for passive @@ -2310,6 +2337,10 @@ export function DaemonSessionProvider(props: DaemonSessionProviderProps) { } catch (error) { const restartRequested = eventStream?.restartRequested === true; clearEventStream(); + if (session && sessionRef.current !== session) { + clearPendingTranscriptEvents(); + return; + } if (restartRequested && !disposed && !abort.signal.aborted) { flushTranscriptSync(); nextSseConnectReason = 'prompt_restart'; @@ -2351,10 +2382,10 @@ export function DaemonSessionProvider(props: DaemonSessionProviderProps) { pendingLoad.sessionId === reconnectSessionId) ) { if ( - skipNextCleanupDetachSessionIdRef.current === + skipNextCleanupDetachSessionRef.current?.sessionId === pendingLoad.sessionId ) { - skipNextCleanupDetachSessionIdRef.current = undefined; + skipNextCleanupDetachSessionRef.current = undefined; } pendingSessionLoadRef.current = undefined; if (pendingLoad.timeout !== undefined) { @@ -2499,26 +2530,34 @@ export function DaemonSessionProvider(props: DaemonSessionProviderProps) { void run(); return () => { - const session = sessionRef.current; + const session = runnerSession; disposed = true; abort.abort(); - // Apply buffered transcript events before tearing down (do NOT drop them). - // The SSE client advances lastSeenEventId as each event is *yielded* — - // before our batched dispatch runs — so dropping here would make a - // same-session incremental resume (the keepSessionForNextEffect path) - // skip them permanently. Flushing is free and safe: the store notifies - // via queueMicrotask, so there is no synchronous setState; on unmount it - // is an orphaned dispatch and on a session switch the next run resets the - // store anyway. - flushTranscriptSync(); - hasCurrentSessionActivePromptRef.current = () => false; - setPromptStatus('idle'); - clearPassiveAssistantDoneTimer(passiveAssistantDoneTimerRef); + const ownsCurrentSession = + session !== undefined && sessionRef.current === session; + const ownsEmptyState = + session === undefined && sessionRef.current === undefined; const keepSessionForNextEffect = - session?.sessionId === skipNextCleanupDetachSessionIdRef.current; + ownsCurrentSession && + session === skipNextCleanupDetachSessionRef.current; const isUnmounting = !mountedRef.current; + if (ownsCurrentSession || ownsEmptyState) { + // A same-attachment effect restart must flush events already yielded by + // the SSE client, because its resume cursor has advanced past them. + flushTranscriptSync(); + } else { + // A replacement attachment owns the store now. Never let the old + // runner's pending macrotask append its buffered events to that owner. + clearPendingTranscriptEvents(); + } + if (ownsCurrentSession && (!keepSessionForNextEffect || isUnmounting)) { + hasCurrentSessionActivePromptRef.current = () => false; + setPromptStatus('idle'); + clearPassiveAssistantDoneTimer(passiveAssistantDoneTimerRef); + } if ( pendingSessionLoadRef.current && + (ownsCurrentSession || ownsEmptyState) && (!keepSessionForNextEffect || isUnmounting) ) { if (pendingSessionLoadRef.current.timeout !== undefined) { @@ -2529,7 +2568,11 @@ export function DaemonSessionProvider(props: DaemonSessionProviderProps) { ); pendingSessionLoadRef.current = undefined; } - if ((!keepSessionForNextEffect || isUnmounting) && session?.clientId) { + if ( + ownsCurrentSession && + (!keepSessionForNextEffect || isUnmounting) && + session.clientId + ) { void detachDaemonClient({ baseUrl: resolvedBaseUrl!, token: resolvedToken, @@ -2539,7 +2582,7 @@ export function DaemonSessionProvider(props: DaemonSessionProviderProps) { console.warn('[DaemonSessionProvider] detach failed:', err), ); } - if (!keepSessionForNextEffect || isUnmounting) { + if (ownsCurrentSession && (!keepSessionForNextEffect || isUnmounting)) { sessionRef.current = undefined; } }; @@ -2577,21 +2620,27 @@ export function DaemonSessionProvider(props: DaemonSessionProviderProps) { ) { return undefined; } - if (heartbeatFailureStateRef.current.sessionId !== connection.sessionId) { - heartbeatFailureStateRef.current = { - sessionId: connection.sessionId, - consecutiveFailures: 0, - }; - } - const heartbeatFailureState = heartbeatFailureStateRef.current; let disposed = false; const timer = setInterval(() => { const session = sessionRef.current; if (!session) return; + if (heartbeatFailureStateRef.current.session !== session) { + heartbeatFailureStateRef.current = { + session, + consecutiveFailures: 0, + }; + } + const heartbeatFailureState = heartbeatFailureStateRef.current; session .heartbeat() .then(() => { - if (disposed) return; + if ( + disposed || + sessionRef.current !== session || + heartbeatFailureStateRef.current !== heartbeatFailureState + ) { + return; + } if ( heartbeatFailureState.consecutiveFailures >= heartbeatFailureThreshold @@ -2611,7 +2660,13 @@ export function DaemonSessionProvider(props: DaemonSessionProviderProps) { heartbeatFailureState.lastHttpError = undefined; }) .catch((error: unknown) => { - if (disposed) return; + if ( + disposed || + sessionRef.current !== session || + heartbeatFailureStateRef.current !== heartbeatFailureState + ) { + return; + } heartbeatFailureState.consecutiveFailures += 1; const message = error instanceof Error ? error.message : 'Session heartbeat failed'; @@ -2660,7 +2715,7 @@ export function DaemonSessionProvider(props: DaemonSessionProviderProps) { activePromptsRef.current.delete(deadSessionId); clearPassiveAssistantDoneTimer(passiveAssistantDoneTimerRef); setPromptStatus('idle'); - if (sessionRef.current?.sessionId === deadSessionId) { + if (sessionRef.current === session) { if (missingSession) { manualSessionClearRef.current = true; } @@ -2712,7 +2767,7 @@ export function DaemonSessionProvider(props: DaemonSessionProviderProps) { pendingSessionLoadIdRef, heartbeatSupportedRef, manualSessionClearRef, - skipNextCleanupDetachSessionIdRef, + skipNextCleanupDetachSessionRef, passiveAssistantDoneTimerRef, hasSessionActivePrompt: () => hasCurrentSessionActivePromptRef.current(), diff --git a/packages/webui/src/daemon/session/actions.test.ts b/packages/webui/src/daemon/session/actions.test.ts index 6f4e7f59e2a..f72ffeca6ea 100644 --- a/packages/webui/src/daemon/session/actions.test.ts +++ b/packages/webui/src/daemon/session/actions.test.ts @@ -852,6 +852,147 @@ describe('createDaemonSessionActions', () => { await expect(prompt).resolves.toEqual({ stopReason: 'cancelled' }); expect(restartEventStream).not.toHaveBeenCalled(); }); + + it('does not apply a late model update to a replacement attachment', async () => { + const source = createMockSession('session-a', 'client-a'); + const target = createMockSession('session-a', 'client-b'); + const result = { applied: true }; + const deferred = createDeferred(); + source.setModel.mockReturnValueOnce(deferred.promise); + const { actions, getConnection, replaceConnection, sessionRef } = + createActionsHarness({ + connection: { + status: 'connected', + sessionId: source.sessionId, + clientId: source.clientId, + currentModel: 'source-model', + }, + session: source, + }); + + const pending = actions.setModel('late-source-model'); + sessionRef.current = target as unknown as DaemonSessionClient; + replaceConnection({ + status: 'connected', + sessionId: target.sessionId, + clientId: target.clientId, + currentModel: 'target-model', + }); + deferred.resolve(result); + + await expect(pending).resolves.toBe(result); + expect(getConnection()).toMatchObject({ + clientId: 'client-b', + currentModel: 'target-model', + }); + }); + + it('does not apply a late approval mode to a replacement attachment', async () => { + const source = createMockSession('session-a', 'client-a'); + const target = createMockSession('session-a', 'client-b'); + const result = { + sessionId: 'session-a', + mode: 'yolo', + previous: 'default', + persisted: false, + }; + const deferred = createDeferred(); + source.client.setSessionApprovalMode.mockReturnValueOnce(deferred.promise); + const { actions, getConnection, replaceConnection, sessionRef } = + createActionsHarness({ + connection: { + status: 'connected', + sessionId: source.sessionId, + clientId: source.clientId, + currentMode: 'default', + }, + session: source, + }); + + const pending = actions.setApprovalMode('yolo'); + sessionRef.current = target as unknown as DaemonSessionClient; + replaceConnection({ + status: 'connected', + sessionId: target.sessionId, + clientId: target.clientId, + currentMode: 'plan', + }); + deferred.resolve(result); + + await expect(pending).resolves.toBe(result); + expect(getConnection()).toMatchObject({ + clientId: 'client-b', + currentMode: 'plan', + }); + }); + + it('does not apply late commands to a replacement attachment', async () => { + const source = createMockSession('session-a', 'client-a'); + const target = createMockSession('session-a', 'client-b'); + const status = supportedCommandsStatus('session-a', 'source-command'); + const deferred = createDeferred(); + source.supportedCommands.mockReturnValueOnce(deferred.promise); + const targetStatus = supportedCommandsStatus('session-a', 'target-command'); + const { actions, getConnection, replaceConnection, sessionRef } = + createActionsHarness({ session: source }); + + const pending = actions.refreshCommands(); + sessionRef.current = target as unknown as DaemonSessionClient; + replaceConnection({ + status: 'connected', + sessionId: target.sessionId, + clientId: target.clientId, + commands: [commandInfo('target-command')], + skills: ['target-skill'], + supportedCommands: targetStatus, + }); + deferred.resolve(status); + + await expect(pending).resolves.toBeUndefined(); + expect(getConnection()).toMatchObject({ + clientId: 'client-b', + commands: [commandInfo('target-command')], + skills: ['target-skill'], + supportedCommands: targetStatus, + }); + }); + + it('does not apply late context to a replacement attachment', async () => { + const source = createMockSession('session-a', 'client-a'); + const target = createMockSession('session-a', 'client-b'); + const context = { + ...contextStatus('session-a'), + state: { + models: { currentModelId: 'source-model' }, + modes: { currentModeId: 'source-mode' }, + }, + }; + const deferred = createDeferred(); + source.context.mockReturnValueOnce(deferred.promise); + const targetContext = contextStatus('session-a'); + const { actions, getConnection, replaceConnection, sessionRef } = + createActionsHarness({ session: source }); + + const pending = actions.getContext(); + sessionRef.current = target as unknown as DaemonSessionClient; + replaceConnection({ + status: 'connected', + sessionId: target.sessionId, + clientId: target.clientId, + context: targetContext, + currentModel: 'target-model', + currentMode: 'target-mode', + }); + deferred.resolve(context); + + await expect(pending).resolves.toBe(context); + expect(getConnection()).toMatchObject({ + clientId: 'client-b', + context: targetContext, + currentModel: 'target-model', + currentMode: 'target-mode', + }); + }); }); function createActionsHarness( @@ -873,6 +1014,9 @@ function createActionsHarness( status: 'connected', workspaceCwd: '/workspace', }; + const replaceConnection = (next: DaemonConnectionState) => { + connection = next; + }; const sessionRef = { current: opts.session as unknown as DaemonSessionClient | undefined, }; @@ -898,7 +1042,7 @@ function createActionsHarness( pendingSessionLoadIdRef: { current: 0 }, heartbeatSupportedRef: { current: false }, manualSessionClearRef: opts.manualSessionClearRef ?? { current: false }, - skipNextCleanupDetachSessionIdRef: { current: undefined }, + skipNextCleanupDetachSessionRef: { current: undefined }, passiveAssistantDoneTimerRef: { current: undefined }, getCreateSessionRequest: () => ({ workspaceCwd: '/workspace' }), createDetachedSession: (opts.createDetachedSession ?? @@ -930,25 +1074,37 @@ function createActionsHarness( activePromptsRef, getConnection: () => connection, pendingSessionLoadRef, + replaceConnection, sessionRef, store, }; } -function createMockSession(sessionId: string) { +function createMockSession( + sessionId: string, + clientId = `client-${sessionId}`, +) { return { sessionId, workspaceCwd: '/workspace', - clientId: `client-${sessionId}`, + clientId, client: { createOrAttachSession: vi.fn(), - setSessionApprovalMode: vi.fn(), + setSessionApprovalMode: vi.fn(async () => ({ + sessionId, + mode: 'default', + previous: 'default', + persisted: false, + })), listWorkspaceSessions: vi.fn(), closeSession: vi.fn(), }, cancel: vi.fn(async () => undefined), + context: vi.fn(async () => contextStatus(sessionId)), detach: vi.fn(async () => undefined), + setModel: vi.fn(async () => ({})), submitPrompt: vi.fn(async () => ({ promptId: 'prompt-1' })), + supportedCommands: vi.fn(async () => supportedCommandsStatus(sessionId)), tasks: vi.fn(async () => ({ v: 1 as const, sessionId, tasks: [] })), }; } diff --git a/packages/webui/src/daemon/session/actions.ts b/packages/webui/src/daemon/session/actions.ts index 1ded9b4bfc2..ba7bc45dd26 100644 --- a/packages/webui/src/daemon/session/actions.ts +++ b/packages/webui/src/daemon/session/actions.ts @@ -92,7 +92,7 @@ export interface CreateDaemonSessionActionsArgs { pendingSessionLoadIdRef: RefBox; heartbeatSupportedRef: RefBox; manualSessionClearRef: RefBox; - skipNextCleanupDetachSessionIdRef: RefBox; + skipNextCleanupDetachSessionRef: RefBox; passiveAssistantDoneTimerRef: TimerRef; getCreateSessionRequest: () => CreateSessionRequest; createDetachedSession: ( @@ -163,7 +163,7 @@ export function createDaemonSessionActions({ pendingSessionLoadIdRef, heartbeatSupportedRef, manualSessionClearRef, - skipNextCleanupDetachSessionIdRef, + skipNextCleanupDetachSessionRef, passiveAssistantDoneTimerRef, getCreateSessionRequest, createDetachedSession, @@ -196,10 +196,10 @@ export function createDaemonSessionActions({ clearPassiveAssistantDoneTimer(passiveAssistantDoneTimerRef); if (pendingSessionLoadRef.current) { if ( - skipNextCleanupDetachSessionIdRef.current === + skipNextCleanupDetachSessionRef.current?.sessionId === pendingSessionLoadRef.current.sessionId ) { - skipNextCleanupDetachSessionIdRef.current = undefined; + skipNextCleanupDetachSessionRef.current = undefined; } clearPendingLoadTimeout(pendingSessionLoadRef.current); pendingSessionLoadRef.current.reject( @@ -315,8 +315,14 @@ export function createDaemonSessionActions({ ); }); if (reloadingCurrentSession) { - skipNextCleanupDetachSessionIdRef.current = sessionId; - void loadPromise.then(detachCurrentSession, () => undefined); + skipNextCleanupDetachSessionRef.current = currentSession; + void loadPromise + .then(detachCurrentSession, () => undefined) + .finally(() => { + if (skipNextCleanupDetachSessionRef.current === currentSession) { + skipNextCleanupDetachSessionRef.current = undefined; + } + }); } else { void detachCurrentSession(); } @@ -575,7 +581,9 @@ export function createDaemonSessionActions({ session.setModel(modelId), 'Set model timed out', ); - setConnection((current) => ({ ...current, currentModel: modelId })); + if (sessionRef.current === session) { + setConnection((current) => ({ ...current, currentModel: modelId })); + } return result; } catch (error) { throw dispatchActionError( @@ -602,10 +610,12 @@ export function createDaemonSessionActions({ }), 'Set approval mode timed out', ); - setConnection((current) => ({ - ...current, - currentMode: result.mode || mode, - })); + if (sessionRef.current === session) { + setConnection((current) => ({ + ...current, + currentMode: result.mode || mode, + })); + } return result; } catch (error) { throw dispatchActionError( @@ -786,7 +796,7 @@ export function createDaemonSessionActions({ } persistStableClientId(nextSession.clientId, nextSession.sessionId); sessionRef.current = nextSession; - skipNextCleanupDetachSessionIdRef.current = nextSession.sessionId; + skipNextCleanupDetachSessionRef.current = nextSession; setConnection((current) => ({ ...current, status: 'connected', @@ -904,13 +914,15 @@ export function createDaemonSessionActions({ session.supportedCommands(), 'Refresh commands timed out', ); - const { commands, skills } = mapSupportedCommands(status); - setConnection((current) => ({ - ...current, - commands, - skills, - supportedCommands: status, - })); + if (sessionRef.current === session) { + const { commands, skills } = mapSupportedCommands(status); + setConnection((current) => ({ + ...current, + commands, + skills, + supportedCommands: status, + })); + } } catch (error) { throw dispatchActionError( addNotice, @@ -933,14 +945,16 @@ export function createDaemonSessionActions({ session.context(), 'Load context timed out', ); - setConnection((current) => ({ - ...current, - context, - currentMode: - getModeFromSessionContext(context) ?? current.currentMode, - currentModel: - getModelFromSessionContext(context) ?? current.currentModel, - })); + if (sessionRef.current === session) { + setConnection((current) => ({ + ...current, + context, + currentMode: + getModeFromSessionContext(context) ?? current.currentMode, + currentModel: + getModelFromSessionContext(context) ?? current.currentModel, + })); + } return context; } catch (error) { throw dispatchActionError(