diff --git a/docs/design/2026-08-10-transactional-webui-session-switching.md b/docs/design/2026-08-10-transactional-webui-session-switching.md index 06ab2e95949..75b2a68e166 100644 --- a/docs/design/2026-08-10-transactional-webui-session-switching.md +++ b/docs/design/2026-08-10-transactional-webui-session-switching.md @@ -12,9 +12,9 @@ Modern transactional behavior requires a successful capability snapshot that adv ## Coordinator -Each provider owns one raw restore slot and one desired intent. Equivalent requests coalesce. A newer target rejects the prior public intent and replaces the queued intent, while an already-running SDK request continues to settlement because it is not cancellable. Its result is adopted only when it still matches the latest target; otherwise its attachment is detached once on a best-effort basis. The queued deadline begins when the caller requests the switch, so an expired target never starts a restore. +Each provider owns one raw restore slot and one desired intent. Restore equivalence includes the normalized session and workspace plus the effective replay shape: `resume/none`, `load/all`, or `load/recent(N)`. The provider snapshots the effective page when admitting the intent, after applying daemon pagination capability, and uses that snapshot for the initial request, queued execution, and retries. Only exact shapes coalesce; a non-equivalent request rejects the prior public intent, permanently marks any different raw result as superseded, and replaces the queued intent, while the already-running SDK request continues to settlement because it is not cancellable. The superseded result is never adopted even if a later intent returns to its shape; its attachment is detached once on a best-effort basis. A timed-out raw request that has not been superseded by a different shape may still satisfy an exact-shape retry. The queued deadline begins when the caller requests the switch, so an expired target never starts a restore. -Commit is guarded by the desired intent, absolute deadline, provider environment, local lifecycle, source logical identity, and restored target identity. Timeout, SDK failure, supersede, staging failure, and commit are explicit competing terminal states rather than an implicit `Promise.race`. +Commit is guarded by the desired intent, absolute deadline, provider environment, local lifecycle, source logical identity, and restored target identity. A same-shape retry may adopt a late raw result only when an ordinary timeout left the lifecycle unchanged; an explicit lifecycle cancellation fences that result even if a later intent returns to the same shape. Timeout, SDK failure, supersede, staging failure, and commit are explicit competing terminal states rather than an implicit `Promise.race`. ## Staging and commit diff --git a/integration-tests/cli/qwen-serve-webui-session-switching.test.ts b/integration-tests/cli/qwen-serve-webui-session-switching.test.ts index 1e2c8d7d5fb..eca89c76c62 100644 --- a/integration-tests/cli/qwen-serve-webui-session-switching.test.ts +++ b/integration-tests/cli/qwen-serve-webui-session-switching.test.ts @@ -255,6 +255,125 @@ describe('qwen serve WebUI transactional session switching', () => { } }, 30_000); + it('serializes non-equivalent load and resume requests for one target', async () => { + const originalFetch = globalThis.fetch; + const state = await setup(); + const loadResponseReady = deferred(); + const releaseLoadResponse = deferred(); + let loadRequests = 0; + let resumeRequests = 0; + let targetDetachRequests = 0; + let staleClientId: string | undefined; + const detachedClientIds: Array = []; + let loadOutcome: Promise | undefined; + let resumeOutcome: Promise | undefined; + try { + globalThis.fetch = async (input, init) => { + const request = + input instanceof Request ? input : new Request(input, init); + const pathname = new URL(request.url).pathname; + const targetPath = `/session/${encodeURIComponent(state.target.sessionId)}`; + if ( + request.method === 'POST' && + pathname.endsWith(`${targetPath}/resume`) + ) { + resumeRequests += 1; + } + if ( + request.method === 'POST' && + pathname.endsWith(`${targetPath}/detach`) + ) { + targetDetachRequests += 1; + detachedClientIds.push(request.headers.get('X-Qwen-Client-Id')); + } + const response = await originalFetch(request); + if ( + request.method === 'POST' && + pathname.endsWith(`${targetPath}/load`) + ) { + loadRequests += 1; + const payload = (await response.clone().json()) as { + clientId?: unknown; + }; + staleClientId = + typeof payload.clientId === 'string' ? payload.clientId : undefined; + loadResponseReady.resolve(); + await releaseLoadResponse.promise; + } + return response; + }; + act(() => { + loadOutcome = state + .getActions() + .loadSession(state.target.sessionId, { + workspaceCwd: state.workspace, + }) + .then( + () => undefined, + (error: unknown) => error, + ); + }); + await loadResponseReady.promise; + + act(() => { + resumeOutcome = state + .getActions() + .resumeSession(state.target.sessionId, { + workspaceCwd: state.workspace, + }); + }); + expect(await loadOutcome).toMatchObject({ name: 'AbortError' }); + expect(loadRequests).toBe(1); + expect(resumeRequests).toBe(0); + expect(state.getConnection()).toMatchObject({ + status: 'connected', + sessionId: state.source.sessionId, + sessionTransition: { phase: 'queued', operation: 'resume' }, + }); + + await activeDaemon!.client.prompt(state.source.sessionId, { + prompt: [{ type: 'text', text: 'source remains live while queued' }], + }); + await waitFor( + () => + JSON.stringify(state.getBlocks()).includes( + 'source remains live while queued', + ), + 'source event while resume is queued', + ); + + await act(async () => { + releaseLoadResponse.resolve(); + await resumeOutcome; + }); + expect(resumeRequests).toBe(1); + await waitFor( + () => targetDetachRequests === 1, + 'stale target attachment cleanup', + ); + expect(staleClientId).toBeTruthy(); + expect(detachedClientIds).toEqual([staleClientId]); + expect(state.getConnection()).toMatchObject({ + status: 'connected', + sessionId: state.target.sessionId, + }); + expect(state.getConnection()?.clientId).toBeTruthy(); + expect(state.getConnection()?.clientId).not.toBe(staleClientId); + } finally { + globalThis.fetch = originalFetch; + releaseLoadResponse.resolve(); + await loadOutcome?.catch(() => undefined); + await resumeOutcome?.catch(() => undefined); + if (root) { + await act(async () => root?.unmount()); + root = undefined; + } + await activeDaemon?.dispose(); + activeDaemon = undefined; + fs.rmSync(state.workspace, { recursive: true, force: true }); + } + }, 30_000); + it('preserves the source after a structured target timeout', async () => { const originalFetch = globalThis.fetch; const state = await setup(); diff --git a/packages/acp-bridge/src/bridge.test.ts b/packages/acp-bridge/src/bridge.test.ts index d7801644412..2bca67d1958 100644 --- a/packages/acp-bridge/src/bridge.test.ts +++ b/packages/acp-bridge/src/bridge.test.ts @@ -5294,6 +5294,204 @@ describe('createAcpSessionBridge', () => { await bridge.shutdown(); }); + it('coalesces response restores with the same history page', async () => { + const load = deferred(); + const handle = makeChannel({ loadSessionImpl: () => load.promise }); + const bridge = makeBridge({ channelFactory: async () => handle.channel }); + + const first = bridge.loadSession({ + sessionId: 'coalesce-page', + workspaceCwd: WS_A, + historyReplay: 'response', + historyPageSize: 100, + }); + for (let i = 0; i < 50 && handle.agent.loadSessionCalls.length !== 1; i++) { + await new Promise((resolve) => setTimeout(resolve, 10)); + } + const second = bridge.loadSession({ + sessionId: 'coalesce-page', + workspaceCwd: WS_A, + historyReplay: 'response', + historyPageSize: 100, + }); + + load.resolve({ _meta: { tag: 'same-page' } }); + const [owner, waiter] = await Promise.all([first, second]); + + expect(handle.agent.loadSessionCalls).toHaveLength(1); + expect(handle.agent.loadSessionCalls[0]?._meta).toMatchObject({ + 'qwen.session.loadReplayMode': 'bulk', + 'qwen.session.loadReplayPageSize': 100, + }); + expect(owner).toMatchObject({ + attached: false, + state: { _meta: { tag: 'same-page' } }, + }); + expect(waiter).toMatchObject({ + attached: true, + state: { _meta: { tag: 'same-page' } }, + }); + expect(waiter.clientId).not.toBe(owner.clientId); + + await bridge.shutdown(); + }); + + it.each([0, 501, 1.5])( + 'rejects invalid response history page %s before restore', + async (historyPageSize) => { + const handle = makeChannel(); + const bridge = makeBridge({ + channelFactory: async () => handle.channel, + }); + + await expect( + bridge.loadSession({ + sessionId: 'invalid-page', + workspaceCwd: WS_A, + historyReplay: 'response', + historyPageSize, + }), + ).rejects.toThrow('Invalid historyPageSize'); + expect(handle.agent.loadSessionCalls).toHaveLength(0); + await bridge.shutdown(); + }, + ); + + it.each([ + ['a different explicit page', 500], + ['an omitted page', undefined], + ] as const)( + 'rejects response restore coalescing with %s', + async (_label, historyPageSize) => { + const load = deferred(); + const handle = makeChannel({ loadSessionImpl: () => load.promise }); + const bridge = makeBridge({ + channelFactory: async () => handle.channel, + }); + + const first = bridge.loadSession({ + sessionId: 'mismatched-page', + workspaceCwd: WS_A, + historyReplay: 'response', + historyPageSize: 100, + }); + for ( + let i = 0; + i < 50 && handle.agent.loadSessionCalls.length !== 1; + i++ + ) { + await new Promise((resolve) => setTimeout(resolve, 10)); + } + + await expect( + bridge.loadSession({ + sessionId: 'mismatched-page', + workspaceCwd: WS_A, + historyReplay: 'response', + ...(historyPageSize !== undefined ? { historyPageSize } : {}), + }), + ).rejects.toBeInstanceOf(RestoreInProgressError); + + load.resolve({}); + const restored = await first; + expect(handle.agent.loadSessionCalls).toHaveLength(1); + await bridge.killSession(restored.sessionId, { + requireZeroAttaches: true, + }); + expect(bridge.sessionCount).toBe(0); + await bridge.shutdown(); + }, + ); + + it('ignores history pages when coalescing streamed loads', async () => { + const load = deferred(); + const handle = makeChannel({ loadSessionImpl: () => load.promise }); + const bridge = makeBridge({ channelFactory: async () => handle.channel }); + + const first = bridge.loadSession({ + sessionId: 'stream-pages-ignored', + workspaceCwd: WS_A, + historyReplay: 'stream', + historyPageSize: 100, + }); + for (let i = 0; i < 50 && handle.agent.loadSessionCalls.length !== 1; i++) { + await new Promise((resolve) => setTimeout(resolve, 10)); + } + const second = bridge.loadSession({ + sessionId: 'stream-pages-ignored', + workspaceCwd: WS_A, + historyReplay: 'stream', + historyPageSize: 500, + }); + + load.resolve({}); + await Promise.all([first, second]); + expect(handle.agent.loadSessionCalls).toHaveLength(1); + expect(handle.agent.loadSessionCalls[0]?._meta).not.toHaveProperty( + 'qwen.session.loadReplayPageSize', + ); + await bridge.shutdown(); + }); + + it('ignores history pages when coalescing resumes', async () => { + const resume = deferred(); + const handle = makeChannel({ + resumeSessionImpl: () => resume.promise, + }); + const bridge = makeBridge({ channelFactory: async () => handle.channel }); + + const first = bridge.resumeSession({ + sessionId: 'resume-pages-ignored', + workspaceCwd: WS_A, + historyPageSize: 100, + }); + for ( + let i = 0; + i < 50 && handle.agent.resumeSessionCalls.length !== 1; + i++ + ) { + await new Promise((resolve) => setTimeout(resolve, 10)); + } + const second = bridge.resumeSession({ + sessionId: 'resume-pages-ignored', + workspaceCwd: WS_A, + historyPageSize: 500, + }); + + resume.resolve({}); + await Promise.all([first, second]); + expect(handle.agent.resumeSessionCalls).toHaveLength(1); + await bridge.shutdown(); + }); + + it('rejects restore coalescing with different inherited-history policies', async () => { + const load = deferred(); + const handle = makeChannel({ loadSessionImpl: () => load.promise }); + const bridge = makeBridge({ channelFactory: async () => handle.channel }); + + const first = bridge.loadSession({ + sessionId: 'coalesce-inherited-policy', + workspaceCwd: WS_A, + hideInheritedHistory: true, + }); + for (let i = 0; i < 50 && handle.agent.loadSessionCalls.length !== 1; i++) { + await new Promise((resolve) => setTimeout(resolve, 10)); + } + + await expect( + bridge.loadSession({ + sessionId: 'coalesce-inherited-policy', + workspaceCwd: WS_A, + hideInheritedHistory: false, + }), + ).rejects.toBeInstanceOf(RestoreInProgressError); + + load.resolve({}); + await first; + expect(handle.agent.loadSessionCalls).toHaveLength(1); + await bridge.shutdown(); + }); + it('rejects coalescing load requests with incompatible replay modes', async () => { let releaseLoad: ((value: LoadSessionResponse) => void) | undefined; const factory: ChannelFactory = async () => diff --git a/packages/acp-bridge/src/bridge.ts b/packages/acp-bridge/src/bridge.ts index 2b500f6579d..c097f65fb64 100644 --- a/packages/acp-bridge/src/bridge.ts +++ b/packages/acp-bridge/src/bridge.ts @@ -29,6 +29,7 @@ import { PRIVATE_ACP_CAPABILITY_ENV, PRIVATE_PARENT_CAPABILITY_META_KEY, SESSION_ARTIFACT_PERSISTENCE_VERSION, + SESSION_TRANSCRIPT_MAX_LIMIT, TrustGateError, normalizeSnapshotPayload, ShellExecutionService, @@ -2589,6 +2590,7 @@ export function createAcpSessionBridge(opts: BridgeOptions): AcpSessionBridge { interface InFlightRestore { action: 'load' | 'resume'; historyReplay: 'stream' | 'response'; + historyPageSize?: number; hideInheritedHistory: boolean; publicPromise: Promise; settlementPromise: Promise; @@ -5087,6 +5089,20 @@ export function createAcpSessionBridge(opts: BridgeOptions): AcpSessionBridge { } const historyReplay = action === 'load' ? (req.historyReplay ?? 'stream') : 'stream'; + const historyPageSize = + action === 'load' && historyReplay === 'response' + ? req.historyPageSize + : undefined; + if ( + historyPageSize !== undefined && + (!Number.isSafeInteger(historyPageSize) || + historyPageSize < 1 || + historyPageSize > SESSION_TRANSCRIPT_MAX_LIMIT) + ) { + throw new Error( + `Invalid historyPageSize; expected 1..${SESSION_TRANSCRIPT_MAX_LIMIT}`, + ); + } const hideInheritedHistory = action === 'load' && req.hideInheritedHistory === true; @@ -5103,8 +5119,8 @@ export function createAcpSessionBridge(opts: BridgeOptions): AcpSessionBridge { ); } const replayFields = - action === 'load' && req.historyPageSize !== undefined - ? await refreshedReplayFieldsFor(existing, req.historyPageSize) + historyPageSize !== undefined + ? await refreshedReplayFieldsFor(existing, historyPageSize) : replayFieldsFor(existing, action); // Backfill a pagination anchor when the snapshot's truncation // marker carries no recordId (live session, in-flight turn capped @@ -5156,16 +5172,15 @@ export function createAcpSessionBridge(opts: BridgeOptions): AcpSessionBridge { const inFlight = inFlightRestores.get(req.sessionId); if (inFlight) { - // Cross-action races BOTH ways must reject. A `resume` arriving - // while a `load` is in flight cannot quietly coalesce: load - // returns compacted replay + watermark while resume returns only - // a watermark — mixing the two on a shared EventBus would give - // the resume client unexpected replay data or the load client a - // missing snapshot. Same-action coalescing is unaffected. + // Cold restores only coalesce when their effective request shapes + // match. Sharing across actions, replay transports, response pages, or + // inherited-history policies can return replay selected for another + // caller. Same-shape coalescing is unaffected. if ( inFlight.lifecycle.phase === 'abandoned' || action !== inFlight.action || historyReplay !== inFlight.historyReplay || + historyPageSize !== inFlight.historyPageSize || hideInheritedHistory !== inFlight.hideInheritedHistory ) { // An abandoned restore is fenced until the real ACP request and its @@ -5482,10 +5497,9 @@ export function createAcpSessionBridge(opts: BridgeOptions): AcpSessionBridge { ...(historyReplay === 'response' ? { [LOAD_REPLAY_MODE_META_KEY]: LOAD_REPLAY_BULK_MODE, - ...(req.historyPageSize !== undefined + ...(historyPageSize !== undefined ? { - [LOAD_REPLAY_PAGE_SIZE_META_KEY]: - req.historyPageSize, + [LOAD_REPLAY_PAGE_SIZE_META_KEY]: historyPageSize, } : {}), } @@ -5713,7 +5727,7 @@ export function createAcpSessionBridge(opts: BridgeOptions): AcpSessionBridge { restoredArtifactSnapshot === undefined || artifactRestoreFailed, }); if ( - req.historyPageSize !== undefined && + historyPageSize !== undefined && entry.events .snapshotReplay() ?.compactedTurns.some((event) => event.type === 'history_truncated') @@ -5829,6 +5843,7 @@ export function createAcpSessionBridge(opts: BridgeOptions): AcpSessionBridge { inFlightRestores.set(req.sessionId, { action, historyReplay, + ...(historyPageSize !== undefined ? { historyPageSize } : {}), hideInheritedHistory, publicPromise: promise, settlementPromise, diff --git a/packages/acp-bridge/src/bridgeTypes.ts b/packages/acp-bridge/src/bridgeTypes.ts index f436fa6af34..43b93b6abce 100644 --- a/packages/acp-bridge/src/bridgeTypes.ts +++ b/packages/acp-bridge/src/bridgeTypes.ts @@ -157,7 +157,7 @@ export interface BridgeRestoreSessionRequest { workspaceCwd: string; /** Optional echo of a daemon-issued client id for this session. */ clientId?: string; - /** Internal replay transport for `session/load`; defaults to bulk response. */ + /** Internal replay transport for `session/load`; defaults to stream. */ historyReplay?: 'stream' | 'response'; /** Optional newest persisted-record page requested for response replay. */ historyPageSize?: number; diff --git a/packages/webui/src/daemon/session/DaemonSessionProvider.test.tsx b/packages/webui/src/daemon/session/DaemonSessionProvider.test.tsx index e6720d4715e..e9f544a52b7 100644 --- a/packages/webui/src/daemon/session/DaemonSessionProvider.test.tsx +++ b/packages/webui/src/daemon/session/DaemonSessionProvider.test.tsx @@ -252,6 +252,14 @@ const sdkMocks = vi.hoisted(() => { _clientId?: string, ): Promise => takeSession(client), ); + static resume = vi.fn( + async ( + client: unknown, + _sessionId: string, + _opts?: unknown, + _clientId?: string, + ): Promise => takeSession(client), + ); } return { @@ -390,6 +398,11 @@ const sdkMocks = vi.hoisted(() => { async (client: unknown, _sessionId: string): Promise => takeSession(client), ); + MockDaemonSessionClient.resume.mockReset(); + MockDaemonSessionClient.resume.mockImplementation( + async (client: unknown, _sessionId: string): Promise => + takeSession(client), + ); }, }; }); @@ -7380,7 +7393,7 @@ describe('DaemonSessionProvider', () => { it('retries a session switch while the target session is closing', async () => { sdkMocks.capabilities.mockResolvedValue({ workspaceCwd: '/mock-workspace', - features: ['client_identity'], + features: ['client_identity', 'session_transcript_pagination'], }); const firstSession = createMockSession({ sessionId: 'session-a' }); const secondSession = createMockSession({ sessionId: 'session-b' }); @@ -7399,6 +7412,7 @@ describe('DaemonSessionProvider', () => { await renderWithProvider(, { autoConnect: true, sessionId: 'session-a', + historyPageSize: 100, reconnectDelayMs: 10, maxReconnectDelayMs: 100, }); @@ -7440,6 +7454,20 @@ describe('DaemonSessionProvider', () => { targetSessionId: 'session-b', }, }); + act(() => { + root?.render( + + + , + ); + }); await act(async () => { await vi.advanceTimersByTimeAsync(10); @@ -7452,6 +7480,9 @@ describe('DaemonSessionProvider', () => { await flushPromises(); }); expect(sdkMocks.MockDaemonSessionClient.load).toHaveBeenCalledTimes(3); + for (const call of sdkMocks.MockDaemonSessionClient.load.mock.calls) { + expect(call[2]).toMatchObject({ historyPageSize: 100 }); + } await act(async () => { await expect(switched).resolves.toBeUndefined(); @@ -10065,7 +10096,7 @@ describe('DaemonSessionProvider', () => { ); }); - it('coalesces load and resume for the same transactional target', async () => { + it('restores a standalone transactional resume through session/resume', async () => { vi.stubGlobal( 'fetch', vi.fn(async () => new Response(null, { status: 204 })), @@ -10077,16 +10108,60 @@ describe('DaemonSessionProvider', () => { sdkMocks.sessions.push( createMockSession({ sessionId: 'session-a', clientId: 'client-a' }), ); - const target = createDeferred(); let actions: DaemonSessionActions | undefined; + let connection: DaemonConnectionState | undefined; function Harness() { actions = useDaemonActions(); + connection = useDaemonConnection(); return null; } await renderWithProvider(, { autoConnect: true }); sdkMocks.MockDaemonSessionClient.load.mockClear(); + sdkMocks.MockDaemonSessionClient.resume.mockResolvedValueOnce( + createMockSession({ sessionId: 'session-b', clientId: 'client-b' }), + ); + + await act(async () => { + await requireActions(actions).resumeSession('session-b'); + await flushPromises(); + }); + + expect(sdkMocks.MockDaemonSessionClient.resume).toHaveBeenCalledOnce(); + expect(sdkMocks.MockDaemonSessionClient.load).not.toHaveBeenCalled(); + expect(connection).toMatchObject({ + status: 'connected', + sessionId: 'session-b', + clientId: 'client-b', + }); + }); + + it('coalesces transactional loads only when their replay shapes match', async () => { + vi.stubGlobal( + 'fetch', + vi.fn(async () => new Response(null, { status: 204 })), + ); + sdkMocks.capabilities.mockResolvedValue({ + workspaceCwd: '/mock-workspace', + features: ['client_identity', 'session_transcript_pagination'], + }); + sdkMocks.sessions.push( + createMockSession({ sessionId: 'session-a', clientId: 'client-a' }), + ); + const target = createDeferred(); + let actions: DaemonSessionActions | undefined; + + function Harness() { + actions = useDaemonActions(); + return null; + } + + await renderWithProvider(, { + autoConnect: true, + historyPageSize: 100, + }); + sdkMocks.MockDaemonSessionClient.load.mockClear(); sdkMocks.MockDaemonSessionClient.load.mockImplementation( async () => target.promise, ); @@ -10094,11 +10169,454 @@ describe('DaemonSessionProvider', () => { let second!: Promise; act(() => { first = requireActions(actions).loadSession('session-b'); - second = requireActions(actions).resumeSession('session-b'); + second = requireActions(actions).loadSession('session-b'); + }); + await act(async () => flushPromises()); + expect(sdkMocks.MockDaemonSessionClient.load).toHaveBeenCalledOnce(); + expect(sdkMocks.MockDaemonSessionClient.load).toHaveBeenCalledWith( + expect.anything(), + 'session-b', + expect.objectContaining({ historyPageSize: 100 }), + expect.any(String), + ); + + await act(async () => { + target.resolve( + createMockSession({ sessionId: 'session-b', clientId: 'client-b' }), + ); + await Promise.all([first, second]); + await flushPromises(); + }); + }); + + it('serializes load and resume requests for the same target', async () => { + const detachFetch = vi.fn( + async (_input: RequestInfo | URL, _init?: RequestInit) => + new Response(null, { status: 204 }), + ); + vi.stubGlobal('fetch', detachFetch); + sdkMocks.capabilities.mockResolvedValue({ + workspaceCwd: '/mock-workspace', + features: ['client_identity'], + }); + sdkMocks.sessions.push( + createMockSession({ sessionId: 'session-a', clientId: 'client-a' }), + ); + const loadTarget = createDeferred(); + const resumeTarget = createDeferred(); + let actions: DaemonSessionActions | undefined; + let connection: DaemonConnectionState | undefined; + + function Harness() { + actions = useDaemonActions(); + connection = useDaemonConnection(); + return null; + } + + await renderWithProvider(, { autoConnect: true }); + sdkMocks.MockDaemonSessionClient.load.mockClear(); + sdkMocks.MockDaemonSessionClient.load.mockImplementationOnce( + async () => loadTarget.promise, + ); + sdkMocks.MockDaemonSessionClient.resume.mockImplementationOnce( + async () => resumeTarget.promise, + ); + let loadOutcome!: Promise; + let resume!: Promise; + act(() => { + loadOutcome = requireActions(actions) + .loadSession('session-b') + .catch((error: unknown) => error); + resume = requireActions(actions).resumeSession('session-b'); + }); + await act(async () => flushPromises()); + + await expect(loadOutcome).resolves.toMatchObject({ name: 'AbortError' }); + expect(sdkMocks.MockDaemonSessionClient.load).toHaveBeenCalledOnce(); + expect(sdkMocks.MockDaemonSessionClient.resume).not.toHaveBeenCalled(); + expect(connection).toMatchObject({ + status: 'connected', + sessionId: 'session-a', + sessionTransition: { phase: 'queued', operation: 'resume' }, + }); + + await act(async () => { + loadTarget.resolve( + createMockSession({ + sessionId: 'session-b', + clientId: 'client-b-stale', + }), + ); + await flushPromises(); + }); + expect(sdkMocks.MockDaemonSessionClient.resume).toHaveBeenCalledOnce(); + const staleDetaches = detachFetch.mock.calls.filter( + ([, init]) => + new Headers(init?.headers).get('X-Qwen-Client-Id') === 'client-b-stale', + ); + expect(staleDetaches).toHaveLength(1); + + await act(async () => { + resumeTarget.resolve( + createMockSession({ sessionId: 'session-b', clientId: 'client-b' }), + ); + await resume; + await flushPromises(); + }); + expect(connection).toMatchObject({ + status: 'connected', + sessionId: 'session-b', + clientId: 'client-b', + }); + }); + + it('snapshots the effective page for a queued transactional load', async () => { + vi.stubGlobal( + 'fetch', + vi.fn(async () => new Response(null, { status: 204 })), + ); + sdkMocks.capabilities.mockResolvedValue({ + workspaceCwd: '/mock-workspace', + features: ['client_identity', 'session_transcript_pagination'], + }); + sdkMocks.sessions.push( + createMockSession({ sessionId: 'session-a', clientId: 'client-a' }), + ); + const firstTarget = createDeferred(); + const secondTarget = createDeferred(); + let actions: DaemonSessionActions | undefined; + + function Harness() { + actions = useDaemonActions(); + return null; + } + + const renderPageSize = (historyPageSize: number) => + root?.render( + + + , + ); + await renderWithProvider(, { + autoConnect: true, + historyPageSize: 100, + }); + sdkMocks.MockDaemonSessionClient.load.mockClear(); + sdkMocks.MockDaemonSessionClient.load + .mockImplementationOnce(async () => firstTarget.promise) + .mockImplementationOnce(async () => secondTarget.promise); + let firstOutcome!: Promise; + let second!: Promise; + act(() => { + firstOutcome = requireActions(actions) + .loadSession('session-b') + .catch((error: unknown) => error); + renderPageSize(500); + }); + await act(async () => flushPromises()); + act(() => { + second = requireActions(actions).loadSession('session-b'); + renderPageSize(25); + }); + await act(async () => flushPromises()); + + await expect(firstOutcome).resolves.toMatchObject({ name: 'AbortError' }); + expect(sdkMocks.MockDaemonSessionClient.load).toHaveBeenCalledOnce(); + + await act(async () => { + firstTarget.resolve( + createMockSession({ sessionId: 'session-b', clientId: 'stale-client' }), + ); + await flushPromises(); + }); + expect(sdkMocks.MockDaemonSessionClient.load).toHaveBeenCalledTimes(2); + expect( + sdkMocks.MockDaemonSessionClient.load.mock.calls[1]?.[2], + ).toMatchObject({ historyPageSize: 500 }); + + await act(async () => { + secondTarget.resolve( + createMockSession({ sessionId: 'session-b', clientId: 'client-b' }), + ); + await second; + await flushPromises(); + }); + }); + + it('does not reuse a superseded result when the latest shape matches it', async () => { + const detachFetch = vi.fn( + async (_input: RequestInfo | URL, _init?: RequestInit) => + new Response(null, { status: 204 }), + ); + vi.stubGlobal('fetch', detachFetch); + sdkMocks.capabilities.mockResolvedValue({ + workspaceCwd: '/mock-workspace', + features: ['client_identity', 'session_transcript_pagination'], + }); + sdkMocks.sessions.push( + createMockSession({ sessionId: 'session-a', clientId: 'client-a' }), + ); + const staleTarget = createDeferred(); + const latestTarget = createDeferred(); + let actions: DaemonSessionActions | undefined; + let connection: DaemonConnectionState | undefined; + + function Harness() { + actions = useDaemonActions(); + connection = useDaemonConnection(); + return null; + } + + const renderPageSize = (historyPageSize: number) => + root?.render( + + + , + ); + await renderWithProvider(, { + autoConnect: true, + historyPageSize: 100, + }); + sdkMocks.MockDaemonSessionClient.load.mockClear(); + sdkMocks.MockDaemonSessionClient.load + .mockImplementationOnce(async () => staleTarget.promise) + .mockImplementationOnce(async () => latestTarget.promise); + + let firstOutcome!: Promise; + let middleOutcome!: Promise; + let latest!: Promise; + act(() => { + firstOutcome = requireActions(actions) + .loadSession('session-b') + .catch((error: unknown) => error); + renderPageSize(500); }); await act(async () => flushPromises()); + act(() => { + middleOutcome = requireActions(actions) + .loadSession('session-b') + .catch((error: unknown) => error); + renderPageSize(100); + }); + await act(async () => flushPromises()); + act(() => { + latest = requireActions(actions).loadSession('session-b'); + }); + await act(async () => flushPromises()); + + await expect(firstOutcome).resolves.toMatchObject({ name: 'AbortError' }); + await expect(middleOutcome).resolves.toMatchObject({ name: 'AbortError' }); expect(sdkMocks.MockDaemonSessionClient.load).toHaveBeenCalledOnce(); + await act(async () => { + staleTarget.resolve( + createMockSession({ + sessionId: 'session-b', + clientId: 'stale-client', + }), + ); + await flushPromises(); + }); + expect(connection).toMatchObject({ + status: 'connected', + sessionId: 'session-a', + }); + expect(sdkMocks.MockDaemonSessionClient.load).toHaveBeenCalledTimes(2); + expect( + sdkMocks.MockDaemonSessionClient.load.mock.calls[1]?.[2], + ).toMatchObject({ historyPageSize: 100 }); + const staleDetaches = detachFetch.mock.calls.filter( + ([, init]) => + new Headers(init?.headers).get('X-Qwen-Client-Id') === 'stale-client', + ); + expect(staleDetaches).toHaveLength(1); + + await act(async () => { + latestTarget.resolve( + createMockSession({ sessionId: 'session-b', clientId: 'client-b' }), + ); + await latest; + await flushPromises(); + }); + expect(connection).toMatchObject({ + status: 'connected', + sessionId: 'session-b', + clientId: 'client-b', + }); + }); + + it('does not reuse a raw result across a lifecycle cancellation', async () => { + vi.stubGlobal( + 'fetch', + vi.fn(async () => new Response(null, { status: 204 })), + ); + sdkMocks.capabilities.mockResolvedValue({ + workspaceCwd: '/mock-workspace', + features: ['client_identity'], + }); + sdkMocks.sessions.push( + createMockSession({ sessionId: 'session-a', clientId: 'client-a' }), + ); + const staleTarget = createDeferred(); + const latestTarget = createDeferred(); + const reloadedSource = createMockSession({ + sessionId: 'session-a', + clientId: 'client-a-reloaded', + }); + let targetLoadCount = 0; + let actions: DaemonSessionActions | undefined; + let connection: DaemonConnectionState | undefined; + + function Harness() { + actions = useDaemonActions(); + connection = useDaemonConnection(); + return null; + } + + await renderWithProvider(, { + autoConnect: true, + sessionId: 'session-a', + maxQueued: 1024, + }); + sdkMocks.MockDaemonSessionClient.load.mockClear(); + sdkMocks.MockDaemonSessionClient.load.mockImplementation( + async (_client: unknown, sessionId: string) => { + if (sessionId === 'session-a') return reloadedSource; + targetLoadCount += 1; + return targetLoadCount === 1 + ? staleTarget.promise + : latestTarget.promise; + }, + ); + + let firstOutcome!: Promise; + act(() => { + firstOutcome = requireActions(actions) + .loadSession('session-b') + .catch((error: unknown) => error); + }); + await act(async () => flushPromises()); + expect(targetLoadCount).toBe(1); + + act(() => { + root?.render( + + + , + ); + }); + await act(async () => flushPromises()); + await expect(firstOutcome).resolves.toMatchObject({ name: 'AbortError' }); + expect(connection).toMatchObject({ + status: 'connected', + sessionId: 'session-a', + clientId: 'client-a-reloaded', + }); + + let retry!: Promise; + act(() => { + retry = requireActions(actions).loadSession('session-b'); + }); + await act(async () => flushPromises()); + expect(targetLoadCount).toBe(1); + + await act(async () => { + staleTarget.resolve( + createMockSession({ + sessionId: 'session-b', + clientId: 'client-b-stale', + }), + ); + await flushPromises(); + }); + expect(targetLoadCount).toBe(2); + expect(connection).toMatchObject({ + status: 'connected', + sessionId: 'session-a', + }); + + await act(async () => { + latestTarget.resolve( + createMockSession({ sessionId: 'session-b', clientId: 'client-b' }), + ); + await retry; + await flushPromises(); + }); + expect(connection).toMatchObject({ + status: 'connected', + sessionId: 'session-b', + clientId: 'client-b', + }); + }); + + it('normalizes configured pages to load/all without pagination support', async () => { + vi.stubGlobal( + 'fetch', + vi.fn(async () => new Response(null, { status: 204 })), + ); + sdkMocks.capabilities.mockResolvedValue({ + workspaceCwd: '/mock-workspace', + features: ['client_identity'], + }); + sdkMocks.sessions.push( + createMockSession({ sessionId: 'session-a', clientId: 'client-a' }), + ); + const target = createDeferred(); + let actions: DaemonSessionActions | undefined; + + function Harness() { + actions = useDaemonActions(); + return null; + } + + await renderWithProvider(, { + autoConnect: true, + historyPageSize: 100, + }); + sdkMocks.MockDaemonSessionClient.load.mockClear(); + sdkMocks.MockDaemonSessionClient.load.mockImplementationOnce( + async () => target.promise, + ); + let first!: Promise; + let second!: Promise; + act(() => { + first = requireActions(actions).loadSession('session-b'); + root?.render( + + + , + ); + }); + await act(async () => flushPromises()); + act(() => { + second = requireActions(actions).loadSession('session-b'); + }); + await act(async () => flushPromises()); + + expect(sdkMocks.MockDaemonSessionClient.load).toHaveBeenCalledOnce(); + expect( + sdkMocks.MockDaemonSessionClient.load.mock.calls[0]?.[2], + ).not.toHaveProperty('historyPageSize'); await act(async () => { target.resolve( createMockSession({ sessionId: 'session-b', clientId: 'client-b' }), diff --git a/packages/webui/src/daemon/session/DaemonSessionProvider.tsx b/packages/webui/src/daemon/session/DaemonSessionProvider.tsx index 8ac734cee95..d3361170ff6 100644 --- a/packages/webui/src/daemon/session/DaemonSessionProvider.tsx +++ b/packages/webui/src/daemon/session/DaemonSessionProvider.tsx @@ -193,6 +193,8 @@ interface CrossSessionTarget { } interface CrossSessionIntent extends CrossSessionTarget { key: string; + effectiveHistoryPageSize?: number; + resultSuperseded?: true; source: DaemonSessionClient; baseUrl: string; token?: string; @@ -213,8 +215,16 @@ const STAGING_BATCH_SIZE = 512; function crossSessionKey( sessionId: string, workspaceCwd: string | undefined, + mode: CrossSessionTarget['mode'], + historyPageSize: number | undefined, ): string { - return `${sessionId}\0${normalizeWorkspaceIdentity(workspaceCwd)}`; + const replayShape = + mode === 'resume' + ? 'resume:none' + : historyPageSize === undefined + ? 'load:all' + : `load:recent:${historyPageSize}`; + return `${sessionId}\0${normalizeWorkspaceIdentity(workspaceCwd)}\0${replayShape}`; } function transitionState( target: CrossSessionTarget, @@ -3433,10 +3443,8 @@ export function DaemonSessionProvider(props: DaemonSessionProviderProps) { { workspaceCwd: intent.workspaceCwd, timeoutMs, - ...(intent.mode === 'load' && - historyPageSizeRef.current !== undefined && - capabilities.features.includes(SESSION_TRANSCRIPT_PAGINATION_FEATURE) - ? { historyPageSize: historyPageSizeRef.current } + ...(intent.effectiveHistoryPageSize !== undefined + ? { historyPageSize: intent.effectiveHistoryPageSize } : {}), }, requestClientId, @@ -3459,8 +3467,10 @@ export function DaemonSessionProvider(props: DaemonSessionProviderProps) { return; } if ( + intent.resultSuperseded === true || latest?.key !== intent.key || - latest?.environmentGeneration !== intent.environmentGeneration + latest.lifecycle !== intent.lifecycle || + latest.environmentGeneration !== intent.environmentGeneration ) { retireAttachment(candidate, intent); return; @@ -3605,7 +3615,20 @@ export function DaemonSessionProvider(props: DaemonSessionProviderProps) { ), ); } - const key = crossSessionKey(request.sessionId, request.workspaceCwd); + const effectiveHistoryPageSize = + request.mode === 'load' && + historyPageSizeRef.current !== undefined && + capabilities.features.includes(SESSION_TRANSCRIPT_PAGINATION_FEATURE) + ? historyPageSizeRef.current + : undefined; + const key = crossSessionKey( + request.sessionId, + request.workspaceCwd, + request.mode, + effectiveHistoryPageSize, + ); + const raw = rawTransitionRef.current; + if (raw && raw.key !== key) raw.resultSuperseded = true; const current = desiredTransitionRef.current; if (current?.key === key) return current.promise; if (current) { @@ -3630,6 +3653,9 @@ export function DaemonSessionProvider(props: DaemonSessionProviderProps) { : Date.now() + timeouts.watchdogTimeoutMs; const intent: CrossSessionIntent = { key, + ...(effectiveHistoryPageSize !== undefined + ? { effectiveHistoryPageSize } + : {}), ...request, source, baseUrl: resolvedBaseUrl,