diff --git a/frontend-modern/src/api/__tests__/aiChat.test.ts b/frontend-modern/src/api/__tests__/aiChat.test.ts index 3e5b14194..01381f393 100644 --- a/frontend-modern/src/api/__tests__/aiChat.test.ts +++ b/frontend-modern/src/api/__tests__/aiChat.test.ts @@ -1146,6 +1146,35 @@ describe('AIChatAPI', () => { clearTimeoutSpy.mockRestore(); }); + it('rejects stalled chat stream reads instead of emitting synthetic completion', async () => { + vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout'] }); + const read = vi.fn( + () => new Promise>(() => undefined), + ); + const releaseLock = vi.fn(); + const onEvent = vi.fn(); + + apiFetchMock.mockResolvedValueOnce({ + ok: true, + body: { + getReader: () => ({ read, releaseLock }), + }, + } as unknown as Response); + + const streamPromise = AIChatAPI.chat('hello', undefined, undefined, onEvent); + const expectedRejection = expect(streamPromise).rejects.toThrow( + 'Pulse Assistant stream timed out waiting for provider data.', + ); + await flushMicrotasks(); + await vi.advanceTimersByTimeAsync(300000); + + await expectedRejection; + expect(read).toHaveBeenCalledTimes(1); + expect(releaseLock).toHaveBeenCalledTimes(1); + expect(onEvent).not.toHaveBeenCalledWith({ type: 'done' }); + expect(logger.warn).toHaveBeenCalledWith('[AI Chat] Stream timeout'); + }); + it('ignores invalid chat stream events through the shared JSON-text helper', async () => { const encoder = new TextEncoder(); const read = vi diff --git a/frontend-modern/src/api/__tests__/streaming.test.ts b/frontend-modern/src/api/__tests__/streaming.test.ts index 0039441a9..2ec1edffe 100644 --- a/frontend-modern/src/api/__tests__/streaming.test.ts +++ b/frontend-modern/src/api/__tests__/streaming.test.ts @@ -235,4 +235,39 @@ describe('consumeJSONEventStream', () => { await streamPromise; expect(events).toEqual(['tool_start', 'tool_progress', 'done']); }); + + it('surfaces stalled read timeouts without reporting normal completion', async () => { + vi.useFakeTimers(); + const read = vi.fn( + () => new Promise>(() => undefined), + ); + const releaseLock = vi.fn(); + const onEvent = vi.fn(); + const onTimeout = vi.fn(); + const onComplete = vi.fn(); + + const streamPromise = consumeJSONEventStream<{ type: string }>( + { + body: { + getReader: () => ({ read, releaseLock }), + }, + } as unknown as Response, + { + onEvent, + onTimeout, + onComplete, + timeoutMs: 1000, + }, + ); + + await flushMicrotasks(); + await vi.advanceTimersByTimeAsync(1000); + await streamPromise; + + expect(read).toHaveBeenCalledTimes(1); + expect(onTimeout).toHaveBeenCalledTimes(1); + expect(onComplete).not.toHaveBeenCalled(); + expect(onEvent).not.toHaveBeenCalled(); + expect(releaseLock).toHaveBeenCalledTimes(1); + }); }); diff --git a/frontend-modern/src/api/aiChat.ts b/frontend-modern/src/api/aiChat.ts index 915bb1a91..ea5ae2a4f 100644 --- a/frontend-modern/src/api/aiChat.ts +++ b/frontend-modern/src/api/aiChat.ts @@ -408,6 +408,7 @@ export class AIChatAPI { }, onTimeout: () => { logger.warn('[AI Chat] Stream timeout'); + throw new Error('Pulse Assistant stream timed out waiting for provider data.'); }, onComplete: () => { onEvent({ type: 'done' }); diff --git a/frontend-modern/src/api/streaming.ts b/frontend-modern/src/api/streaming.ts index 2effd8214..13b73edb5 100644 --- a/frontend-modern/src/api/streaming.ts +++ b/frontend-modern/src/api/streaming.ts @@ -65,6 +65,12 @@ export async function consumeJSONEventStream( let buffer = ''; const timeoutMs = options.timeoutMs ?? 300000; let lastEventTime = Date.now(); + let timedOut = false; + + const markTimedOut = () => { + timedOut = true; + options.onTimeout?.(); + }; const readWithTimeout = async (): Promise> => { let timeoutId: ReturnType | undefined; @@ -122,7 +128,7 @@ export async function consumeJSONEventStream( try { for (;;) { if (Date.now() - lastEventTime > timeoutMs) { - options.onTimeout?.(); + markTimedOut(); break; } @@ -131,6 +137,7 @@ export async function consumeJSONEventStream( result = await readWithTimeout(); } catch (error) { if ((error as Error).message === 'Read timeout') { + markTimedOut(); break; } throw error; @@ -148,6 +155,10 @@ export async function consumeJSONEventStream( } } + if (timedOut) { + return; + } + const trailing = buffer.trim(); if (trailing.startsWith('data: ')) { const event = parseJSONTextSafe(trailing.slice(6));