From 1619f8f776a75e92b015706dc6f9bf66e53ddf4c Mon Sep 17 00:00:00 2001 From: Marco Beretta <81851188+berry-13@users.noreply.github.com> Date: Fri, 25 Sep 2026 20:59:02 +0200 Subject: [PATCH 1/2] refactor: drive the resumable stream through the ChatTransport Add reconnectToStream to ChatTransport and the sse.js transport: a GET attach with request headers, a token refresh on every 401, a cancel the caller did not issue reported as a dropped connection (status 0), and a closed flag for the foreground check. useResumableSSE now consumes normalized ChatEvents instead of opening sse.js and parsing frames itself; resume, steer and queue bookkeeping are unchanged. --- client/src/hooks/SSE/__tests__/sse.spec.ts | 189 ++++++++- client/src/hooks/SSE/transport/sse.ts | 137 +++++-- client/src/hooks/SSE/useResumableSSE.ts | 373 ++++++++---------- packages/data-provider/src/types/transport.ts | 44 ++- 4 files changed, 496 insertions(+), 247 deletions(-) diff --git a/client/src/hooks/SSE/__tests__/sse.spec.ts b/client/src/hooks/SSE/__tests__/sse.spec.ts index 665545f298b..4f882c436d2 100644 --- a/client/src/hooks/SSE/__tests__/sse.spec.ts +++ b/client/src/hooks/SSE/__tests__/sse.spec.ts @@ -17,13 +17,19 @@ class FakeXHR { withCredentials = false; headers: Record = {}; body: string | null = null; + method = ''; + url = ''; private readonly listeners: Record = {}; addEventListener(type: string, listener: XHRListener) { (this.listeners[type] ??= []).push(listener); } - open() {} + open(method: string, url: string) { + this.method = method; + this.url = url; + } + setRequestHeader(key: string, value: string) { this.headers[key] = value; } @@ -260,3 +266,184 @@ describe('createSSETransport', () => { expect(events).toEqual([]); }); }); + +describe('createSSETransport().reconnectToStream', () => { + const OriginalXHR = global.XMLHttpRequest; + const stream = { + url: '/api/agents/chat/stream/convo-1?resume=true', + headers: { 'X-LibreChat-Generation-Protocol': '2' }, + }; + let xhrs: FakeXHR[]; + let events: ChatEvent[]; + let controller: AbortController; + + const current = () => xhrs[xhrs.length - 1]; + const attach = (token = 'token-1') => + createSSETransport({ token }).reconnectToStream(stream, { + signal: controller.signal, + onEvent: (event) => events.push(event), + }); + + beforeEach(() => { + xhrs = []; + events = []; + controller = new AbortController(); + const FakeConstructor = jest.fn(() => { + const xhr = new FakeXHR(); + xhrs.push(xhr); + return xhr; + }); + global.XMLHttpRequest = Object.assign(FakeConstructor, { + HEADERS_RECEIVED: FakeXHR.HEADERS_RECEIVED, + }) as unknown as typeof XMLHttpRequest; + }); + + afterEach(() => { + global.XMLHttpRequest = OriginalXHR; + jest.restoreAllMocks(); + }); + + it('attaches with a GET carrying the bearer token and the request headers', () => { + attach(); + + expect(current().method).toBe('GET'); + expect(current().url).toBe(stream.url); + expect(current().headers).toEqual({ + Authorization: 'Bearer token-1', + 'X-LibreChat-Generation-Protocol': '2', + }); + }); + + it('normalizes the resume snapshot and the frames behind it', () => { + attach(); + current().receiveHeaders(); + current().write( + message({ sync: true, resumeState: { runSteps: [] }, pendingEvents: [] }) + + message({ event: 'attachment', data: { file_id: 'file-1' } }) + + message({ final: true, reconcile: true, terminalStatus: 'complete' }), + ); + + expect(events).toEqual([ + { type: 'open' }, + { type: 'sync', data: { sync: true, resumeState: { runSteps: [] }, pendingEvents: [] } }, + { type: 'attachment', data: { file_id: 'file-1' } }, + { type: 'final', data: { final: true, reconcile: true, terminalStatus: 'complete' } }, + ]); + }); + + it('emits a server error event with its parsed body and no status', () => { + attach(); + current().write(`event: error\ndata: ${JSON.stringify({ error: 'failed' })}\n\n`); + + expect(events).toEqual([{ type: 'error', status: undefined, data: { error: 'failed' } }]); + }); + + it('emits an HTTP failure with its status and no data for a body that is not JSON', () => { + attach(); + current().status = 404; + current().write('Not Found'); + + expect(events).toEqual([{ type: 'error', status: 404, data: undefined }]); + }); + + it('reports a cancel the caller did not issue as a dropped connection', () => { + const connection = attach(); + current().receiveHeaders(); + + current().abort(); + + expect(events).toEqual([{ type: 'open' }, { type: 'error', status: 0 }]); + expect(connection.closed).toBe(true); + }); + + it('emits abort when the caller closes an open stream, then goes quiet', () => { + const connection = attach(); + current().receiveHeaders(); + const xhr = current(); + + controller.abort(); + xhr.write(message({ final: true })); + + expect(events).toEqual([{ type: 'open' }, { type: 'abort' }]); + expect(connection.closed).toBe(true); + }); + + it('reports closed once the response body ends without a terminal event', () => { + const connection = attach(); + current().receiveHeaders(); + current().write(message({ created: true, message: {} })); + expect(connection.closed).toBe(false); + + current().emit('load'); + + expect(connection.closed).toBe(true); + controller.abort(); + expect(events.map((event) => event.type)).toEqual(['open', 'created']); + }); + + it('refreshes the token on every 401 and reattaches with the request headers', async () => { + const refreshToken = jest + .spyOn(request, 'refreshToken') + .mockResolvedValueOnce({ token: 'token-2' } as never) + .mockResolvedValueOnce({ token: 'token-3' } as never); + const dispatchTokenUpdated = jest + .spyOn(request, 'dispatchTokenUpdatedEvent') + .mockImplementation(() => undefined); + attach(); + + for (const _attempt of [1, 2]) { + current().status = 401; + current().write('Unauthorized'); + await new Promise(process.nextTick); + } + + expect(refreshToken).toHaveBeenCalledTimes(2); + expect(dispatchTokenUpdated).toHaveBeenLastCalledWith('token-3'); + expect(xhrs).toHaveLength(3); + expect(current().method).toBe('GET'); + expect(current().headers).toEqual({ + Authorization: 'Bearer token-3', + 'X-LibreChat-Generation-Protocol': '2', + }); + expect(events).toEqual([]); + }); + + it('reports the 401 when the token refresh fails', async () => { + jest.spyOn(console, 'log').mockImplementation(() => undefined); + jest.spyOn(request, 'refreshToken').mockRejectedValue(new Error('refresh failed')); + attach(); + current().status = 401; + current().write('Unauthorized'); + + await new Promise(process.nextTick); + + expect(xhrs).toHaveLength(1); + expect(events).toEqual([{ type: 'error', status: 401, data: undefined }]); + }); + + it('stays closed when the caller aborts while a 401 refresh is in flight', async () => { + jest.spyOn(request, 'refreshToken').mockResolvedValue({ token: 'token-2' } as never); + const dispatchTokenUpdated = jest + .spyOn(request, 'dispatchTokenUpdatedEvent') + .mockImplementation(() => undefined); + attach(); + current().status = 401; + current().write('Unauthorized'); + + controller.abort(); + await new Promise(process.nextTick); + + expect(xhrs).toHaveLength(1); + expect(dispatchTokenUpdated).not.toHaveBeenCalled(); + expect(events).toEqual([]); + }); + + it('never opens a connection for an already-aborted signal', () => { + controller.abort(); + const connection = attach(); + + expect(xhrs).toHaveLength(0); + expect(connection.closed).toBe(true); + expect(events).toEqual([]); + }); +}); diff --git a/client/src/hooks/SSE/transport/sse.ts b/client/src/hooks/SSE/transport/sse.ts index a0e4c7e1f54..f73b6c0c072 100644 --- a/client/src/hooks/SSE/transport/sse.ts +++ b/client/src/hooks/SSE/transport/sse.ts @@ -1,15 +1,65 @@ import { SSE } from 'sse.js'; import { request } from 'librechat-data-provider'; -import type { ChatFrame, ChatTransport, ChatErrorData } from 'librechat-data-provider'; +import type { + ChatFrame, + ChatEvent, + ChatTransport, + ChatErrorData, + ChatStreamConnection, +} from 'librechat-data-provider'; import { normalizeFrame } from './frames'; type StreamErrorEvent = MessageEvent & { responseCode?: number }; +type EventCallback = (event: ChatEvent) => void; const jsonHeaders = (token?: string) => ({ 'Content-Type': 'application/json', Authorization: `Bearer ${token}`, }); +/** Parses one `message` event and emits it normalized; a frame that is not JSON is skipped. */ +const emitFrame = (onEvent: EventCallback) => (e: MessageEvent) => { + let frame: ChatFrame; + try { + frame = JSON.parse(e.data); + } catch (error) { + console.error('Skipping malformed stream frame:', error); + return; + } + const event = normalizeFrame(frame); + if (event) { + onEvent(event); + } +}; + +/** Resolves to the refreshed token, or `null` when the refresh failed. */ +async function refreshToken(): Promise { + try { + const refreshResponse = await request.refreshToken(); + const refreshedToken = refreshResponse?.token ?? ''; + if (!refreshedToken) { + throw new Error('Token refresh failed.'); + } + return refreshedToken; + } catch (error) { + /* token refresh failed, continue handling the original 401 */ + console.log(error); + return null; + } +} + +/** Failure bodies are often empty or HTML, so a body that is not JSON is `undefined`. */ +function parseErrorBody(body: unknown): ChatErrorData | undefined { + if (typeof body !== 'string' || body === '') { + return undefined; + } + try { + return JSON.parse(body); + } catch { + return undefined; + } +} + /** * POSTs the turn with `sse.js` and normalizes what comes back. Owns the wire: * frame parsing, named `attachment`/`error` events, and one token refresh and @@ -40,19 +90,7 @@ export function createSSETransport({ token }: { token?: string }): ChatTransport } }); - sse.addEventListener('message', (e: MessageEvent) => { - let frame: ChatFrame; - try { - frame = JSON.parse(e.data); - } catch (error) { - console.error('Skipping malformed stream frame:', error); - return; - } - const event = normalizeFrame(frame); - if (event) { - onEvent(event); - } - }); + sse.addEventListener('message', emitFrame(onEvent)); let refreshed = false; /** sse.js marks the connection closed on a 401, but the turn is still @@ -62,27 +100,16 @@ export function createSSETransport({ token }: { token?: string }): ChatTransport if (e.responseCode === 401 && !refreshed) { refreshed = true; refreshing = true; - try { - const refreshResponse = await request.refreshToken(); - refreshing = false; - if (signal.aborted) { - return; - } - const refreshedToken = refreshResponse?.token ?? ''; - if (!refreshedToken) { - throw new Error('Token refresh failed.'); - } + const refreshedToken = await refreshToken(); + refreshing = false; + if (signal.aborted) { + return; + } + if (refreshedToken) { sse.headers = jsonHeaders(refreshedToken); request.dispatchTokenUpdatedEvent(refreshedToken); sse.stream(); return; - } catch (error) { - refreshing = false; - /* token refresh failed, continue handling the original 401 */ - console.log(error); - } - if (signal.aborted) { - return; } } @@ -115,5 +142,53 @@ export function createSSETransport({ token }: { token?: string }): ChatTransport sse.stream(); }, + + reconnectToStream({ url, headers }, { signal, onEvent }): ChatStreamConnection { + if (signal.aborted) { + return { closed: true }; + } + + const sse = new SSE(url, { + headers: { Authorization: `Bearer ${token}`, ...headers }, + method: 'GET', + }); + + sse.addEventListener('open', () => { + onEvent({ type: 'open' }); + }); + + sse.addEventListener('message', emitFrame(onEvent)); + + sse.addEventListener('error', async (e: StreamErrorEvent) => { + if (e.responseCode === 401) { + const refreshedToken = await refreshToken(); + if (signal.aborted) { + return; + } + if (refreshedToken) { + sse.headers = { ...sse.headers, Authorization: `Bearer ${refreshedToken}` }; + request.dispatchTokenUpdatedEvent(refreshedToken); + sse.stream(); + return; + } + } + onEvent({ type: 'error', status: e.responseCode, data: parseErrorBody(e.data) }); + }); + + /** sse.js dispatches `abort` when the XHR is cancelled, by our close or by the user agent. */ + sse.addEventListener('abort', () => { + onEvent(signal.aborted ? { type: 'abort' } : { type: 'error', status: 0 }); + }); + + signal.addEventListener('abort', () => sse.close(), { once: true }); + + sse.stream(); + + return { + get closed() { + return sse.readyState === SSE.CLOSED; + }, + }; + }, }; } diff --git a/client/src/hooks/SSE/useResumableSSE.ts b/client/src/hooks/SSE/useResumableSSE.ts index 2b6aa0f7fa9..2cbdd04793a 100644 --- a/client/src/hooks/SSE/useResumableSSE.ts +++ b/client/src/hooks/SSE/useResumableSSE.ts @@ -1,11 +1,9 @@ import { useEffect, useState, useRef, useCallback } from 'react'; import { v4 } from 'uuid'; -import { SSE } from 'sse.js'; import { useStore } from 'jotai'; import { useQueryClient } from '@tanstack/react-query'; import { useSetRecoilState, useRecoilCallback } from 'recoil'; import { - request, Constants, QueryKeys, ErrorTypes, @@ -27,6 +25,7 @@ import { import type { Agents, TMessage, + ChatEvent, TPayload, TSubmission, TConversation, @@ -36,12 +35,17 @@ import type { TSteerUpdatedEvent, TActivityLabelEvent, TReasoningLabelEvent, + ChatFinalFrame, + TAttachment, + TTokenUsageEvent, + TContextUsageEvent, + ChatStreamConnection, } from 'librechat-data-provider'; import type { ActiveJobsResponse, StreamStatusResponse } from '~/data-provider'; import type { DrainAfterAbort, QueuedMessageOrigin } from '~/store/families'; import type { GenerationProtocolVersion } from '~/data-provider'; import type { EventHandlerParams } from './useEventHandlers'; -import type { TResData } from '~/common'; +import type { TResData, TFinalResData } from '~/common'; import { logger, clearComposerDrafts, @@ -89,10 +93,14 @@ import useEventHandlers, { import { pendingApprovalActionFamily } from '~/components/Chat/approval/state'; import useSteerConvert from '~/hooks/Chat/useSteerConvert'; import { useAuthContext } from '~/hooks/AuthContext'; +import { createSSETransport } from './transport'; import { useFileMapContext } from '~/Providers'; import useUsageHandler from './useUsageHandler'; import store from '~/store'; +/** The step handler predates the wire types and accepts a narrower payload. */ +type StepEvent = Parameters['stepHandler']>[0]; + type ChatHelpers = Pick< EventHandlerParams, 'setMessages' | 'getMessages' | 'setConversation' | 'setIsSubmitting' | 'newConversation' @@ -898,7 +906,7 @@ export default function useResumableSSE( const setShowStopButton = useSetRecoilState(store.showStopButtonByIndex(runIndex)); const setLiveAppliedSteerIds = useSetRecoilState(store.liveAppliedSteerIds); - const sseRef = useRef(null); + const streamRef = useRef(null); /** Removes the foreground re-attach listener owned by the newest * subscription; exactly one is registered at a time. */ const stopForegroundReattachRef = useRef<(() => void) | null>(null); @@ -1462,7 +1470,7 @@ export default function useResumableSSE( * stream the server no longer has — and by the dev-only navigation * simulator below. `finalReceived` covers the frame-carried terminals. */ let subscriptionRetired = false; - const preCreatedStepEvents: Array[0]> = []; + const preCreatedStepEvents: StepEvent[] = []; const replayPreCreatedStepEvents = () => { if (!isCurrentSubscription() || preCreatedStepEvents.length === 0) { return; @@ -1885,34 +1893,17 @@ export default function useResumableSSE( const url = queryString ? `${baseUrl}?${queryString}` : baseUrl; logger.log('ResumableSSE', 'Subscribing to stream:', url, { isResume }); - const sse = new SSE(url, { - headers: { - Authorization: `Bearer ${token}`, - ...generationProtocolHeaders(), - }, - method: 'GET', - }); - sseRef.current = sse; + /** Closing is per connection: the transport tells this hook's own close + * (`abort`) apart from a cancel by the user agent (`error` with status 0). */ + const streamController = new AbortController(); + streamRef.current = streamController; + let connection: ChatStreamConnection | null = null; const isCurrentSubscription = () => lifecycleSignal?.aborted !== true && - sseRef.current === sse && + streamRef.current === streamController && submissionRef.current === currentSubmission; - /** - * Whether THIS connection was closed by this hook. The abort listener - * used to infer that from `reconnectAttemptRef`, but that ref is shared - * across the reconnect ladder and stays raised from the moment a retry is - * scheduled until the replacement connection opens — so a user agent that - * cancelled the replacement before it opened (the ordinary case when the - * retry timer fires while the tab is still backgrounded) read as the - * previous connection's deliberate close, and recovery stopped there. - * Ownership is per connection, so the flag must be too. - */ - let closedByUs = false; - const closeStream = () => { - closedByUs = true; - sse.close(); - }; + const closeStream = () => streamController.abort(); let foregroundStatusCheckInFlight = false; const reattachOnForeground = () => { @@ -1964,7 +1955,7 @@ export default function useResumableSSE( return; } - if (sse.readyState === SSE.CLOSED) { + if (connection?.closed === true) { reattachOnForeground(); return; } @@ -2015,7 +2006,7 @@ export default function useResumableSSE( stopForegroundReattachRef.current = () => document.removeEventListener('visibilitychange', handleForegroundReattach); - sse.addEventListener('open', () => { + const handleOpen = () => { if (!isCurrentSubscription()) { return; } @@ -2025,19 +2016,16 @@ export default function useResumableSSE( setIsSubmitting(true); setShowStopButton(generationCreatedAt != null); reconnectAttemptRef.current = 0; - }); + }; - sse.addEventListener('message', async (e: MessageEvent) => { + const handleFrame = async (event: ChatEvent) => { try { - if (!isCurrentSubscription()) { - return; - } - const data = JSON.parse(e.data); - if (finalReceived) { + if (!isCurrentSubscription() || finalReceived) { return; } - if (data.final === true && data.reconcile === true) { + if (event.type === 'final' && event.data.reconcile === true) { + const { data } = event; if ( generationProtocolVersion !== GENERATION_PROTOCOL_VERSION || !supportsGenerationProtocolV2(data) @@ -2054,7 +2042,8 @@ export default function useResumableSSE( return; } - if (data.final != null) { + if (event.type === 'final') { + const { data } = event; finalReceived = true; const finalConvoId = data.conversation?.conversationId ?? @@ -2119,7 +2108,7 @@ export default function useResumableSSE( ); let finalHandled = false; try { - finalHandler(data, currentSubmission as EventSubmission); + finalHandler(data as TFinalResData, currentSubmission as EventSubmission); finalHandled = true; finalizeUsage(data, { ...currentSubmission, userMessage }); } catch (error) { @@ -2169,7 +2158,8 @@ export default function useResumableSSE( return; } - if (data.created != null) { + if (event.type === 'created') { + const { data } = event; logger.log('ResumableSSE', 'Received CREATED event', { messageId: data.message?.messageId, conversationId: data.message?.conversationId, @@ -2202,60 +2192,57 @@ export default function useResumableSSE( return; } - if (data.event === 'attachment' && data.data) { + if (event.type === 'attachment' && event.data) { attachmentHandler({ - data: data.data, + data: event.data as TAttachment, submission: currentSubmission as EventSubmission, }); return; } - if (data.event === 'title') { - titleHandler(data); + if (event.type === 'title') { + titleHandler(event.data); return; } - if (data.event === UsageEvents.ON_CONTEXT_USAGE) { - contextHandler(data.data, { ...currentSubmission, userMessage }); + if (event.type === 'context_usage') { + contextHandler(event.data, { ...currentSubmission, userMessage }); return; } - if (data.event === UsageEvents.ON_TOKEN_USAGE) { - usageHandler(data.data, { ...currentSubmission, userMessage }); + if (event.type === 'token_usage') { + usageHandler(event.data, { ...currentSubmission, userMessage }); return; } - if (data.event === ApprovalEvents.ON_PENDING_ACTION) { - applyPendingActionToMessages(data.data as Agents.PendingAction); + if (event.type === 'pending_action') { + applyPendingActionToMessages(event.data); setIsSubmitting(true); return; } - if (data.event === SteerEvents.ON_STEER_APPLIED) { - applySteerToMessages(data.data as TSteerAppliedEvent); - return; - } - - if (data.event === SteerEvents.ON_STEER_UPDATED) { - updateSteerChips(data.data as TSteerUpdatedEvent); + if (event.type === 'steer_applied') { + applySteerToMessages(event.data); return; } - if (data.event === ActivityLabelEvents.ON_ACTIVITY_LABEL) { - applyActivityLabelToMessages(data.data as TActivityLabelEvent); + if (event.type === 'steer_updated') { + updateSteerChips(event.data); return; } - if (data.event === ReasoningLabelEvents.ON_REASONING_LABEL) { - applyReasoningLabelToMessages(data.data as TReasoningLabelEvent); + if (event.type === 'activity_label') { + applyActivityLabelToMessages(event.data); return; } - if (data.event === ReasoningLabelEvents.ON_REASONING_LABEL_ATTEMPT) { + if (event.type === 'reasoning_label') { + applyReasoningLabelToMessages(event.data); return; } - if (data.event != null) { + if (event.type === 'step') { + const data = event.data as StepEvent; if ( data.event === StepEvents.ON_MESSAGE_DELTA || data.event === StepEvents.ON_REASONING_DELTA @@ -2275,7 +2262,8 @@ export default function useResumableSSE( return; } - if (data.sync != null) { + if (event.type === 'sync') { + const { data } = event; logger.log('ResumableSSE', 'SYNC received', { runSteps: data.resumeState?.runSteps?.length ?? 0, pendingEvents: data.pendingEvents?.length ?? 0, @@ -2498,12 +2486,10 @@ export default function useResumableSSE( * normally on the restored stream and updates its row in place. */ prunePtcTraces(); - if (data.resumeState?.replayEvents?.length > 0) { - logger.log( - 'ResumableSSE', - `Replaying ${data.resumeState.replayEvents.length} resume events`, - ); - for (const replayEvent of data.resumeState.replayEvents) { + const replayEvents = data.resumeState?.replayEvents ?? []; + if (replayEvents.length > 0) { + logger.log('ResumableSSE', `Replaying ${replayEvents.length} resume events`); + for (const replayEvent of replayEvents) { const replayStepId = getStepEventId(replayEvent); if ( replayStepId != null && @@ -2513,9 +2499,9 @@ export default function useResumableSSE( continue; } if (replayEvent.event === UsageEvents.ON_CONTEXT_USAGE) { - contextHandler(replayEvent.data, resumeSubmission); + contextHandler(replayEvent.data as TContextUsageEvent, resumeSubmission); } else if (replayEvent.event === UsageEvents.ON_TOKEN_USAGE) { - usageHandler(replayEvent.data, resumeSubmission); + usageHandler(replayEvent.data as TTokenUsageEvent, resumeSubmission); } else if (replayEvent.event === ApprovalEvents.ON_PENDING_ACTION) { // A pause that landed after the resume snapshot must still render its // controls (mirror the live handler), not fall through to stepHandler. @@ -2535,16 +2521,29 @@ export default function useResumableSSE( replayEvent.event === StepEvents.ON_MESSAGE_DELTA || replayEvent.event === StepEvents.ON_REASONING_DELTA ) { - tapStream(replayEvent.data, resumeSubmission); + tapStream(replayEvent.data as Agents.MessageDeltaEvent, resumeSubmission); } - stepHandler(replayEvent, resumeSubmission); + stepHandler(replayEvent as StepEvent, resumeSubmission); } } } - if (data.pendingEvents?.length > 0) { - logger.log('ResumableSSE', `Replaying ${data.pendingEvents.length} pending events`); - for (const pendingEvent of data.pendingEvents) { + const pendingEvents = data.pendingEvents ?? []; + if (pendingEvents.length > 0) { + logger.log('ResumableSSE', `Replaying ${pendingEvents.length} pending events`); + for (const pendingEvent of pendingEvents) { + if (!('event' in pendingEvent)) { + if ('type' in pendingEvent) { + /** Gap output streamed past the resume snapshot must reach the + * live estimate too, not just the message UI */ + tapContent( + 'text' in pendingEvent ? pendingEvent.text : undefined, + resumeSubmission, + ); + contentHandler({ data: pendingEvent, submission: resumeSubmission }); + } + continue; + } if (pendingEvent.event === 'title') { titleHandler(pendingEvent); } else if (pendingEvent.event === UsageEvents.ON_CONTEXT_USAGE) { @@ -2574,12 +2573,7 @@ export default function useResumableSSE( ) { tapStream(pendingEvent.data, resumeSubmission); } - stepHandler(pendingEvent, resumeSubmission); - } else if (pendingEvent.type != null) { - /** Gap output streamed past the resume snapshot must reach the - * live estimate too, not just the message UI */ - tapContent(pendingEvent.text, resumeSubmission); - contentHandler({ data: pendingEvent, submission: resumeSubmission }); + stepHandler(pendingEvent as StepEvent, resumeSubmission); } } } @@ -2589,8 +2583,10 @@ export default function useResumableSSE( return; } - if (data.type != null) { - const { text, index } = data; + if (event.type === 'content') { + const { data } = event; + const { index } = data; + const text = 'text' in data ? data.text : undefined; if (text != null && index !== textIndex) { textIndex = index; } @@ -2599,13 +2595,14 @@ export default function useResumableSSE( return; } - if (data.message != null) { + if (event.type === 'text') { + const { data } = event; const text = data.text ?? data.response; const initialResponse = { ...(currentSubmission.initialResponse as TMessage), parentMessageId: data.parentMessageId, messageId: data.messageId, - }; + } as TMessage; /** Legacy non-agent streams send cumulative text here — feed the * live estimate like the content path above */ const textSubmission = { ...currentSubmission, userMessage, initialResponse }; @@ -2617,7 +2614,7 @@ export default function useResumableSSE( } catch (error) { logger.error('ResumableSSE', 'Error processing message:', error); } - }); + }; async function handoffToReplacement( conversationId: string, @@ -2812,12 +2809,12 @@ export default function useResumableSSE( return true; } - const reconcileGenerationLifecycle = async (event: { - reconcileReason?: string; - terminalStatus?: 'complete' | 'error' | 'aborted'; - generationCreatedAt?: number; - conversation?: { conversationId?: string }; - }): Promise => { + const reconcileGenerationLifecycle = async ( + event: Pick< + ChatFinalFrame, + 'reconcileReason' | 'terminalStatus' | 'generationCreatedAt' | 'conversation' + >, + ): Promise => { if (!isCurrentSubscription()) { return; } @@ -3074,21 +3071,24 @@ export default function useResumableSSE( /** * Error event handler - handles BOTH: - * 1. HTTP-level errors (responseCode present) - 404, 401, network failures + * 1. HTTP-level errors (status present) - 404, 409, network failures (0) * 2. Server-sent error events (event: error with data) - known errors like ViolationTypes/ErrorTypes * - * Order matters: check responseCode first since HTTP errors may also include data + * Order matters: check the status first since HTTP errors may also include data. + * The transport has already retried a 401 with a refreshed token. */ - const handleTransportFailure = async (e: MessageEvent) => { + const handleTransportFailure = async ({ + status: responseCode, + data, + }: Pick, 'status' | 'data'>) => { if (!isCurrentSubscription()) { return; } - const responseCode = (e as MessageEvent & { responseCode?: number }).responseCode; if (finalReceived) { logger.log('ResumableSSE', 'Ignoring error after FINAL event', { responseCode, - hasData: !!e.data, + hasData: data != null, }); return; } @@ -3353,41 +3353,14 @@ export default function useResumableSSE( return; } - // Check for 401 and try to refresh token (same pattern as useSSE) - if (responseCode === 401) { - try { - const refreshResponse = await request.refreshToken(); - if (!isCurrentSubscription()) { - return; - } - const newToken = refreshResponse?.token ?? ''; - if (!newToken) { - throw new Error('Token refresh failed.'); - } - sse.headers = { - ...sse.headers, - ...generationProtocolHeaders(), - Authorization: `Bearer ${newToken}`, - }; - request.dispatchTokenUpdatedEvent(newToken); - sse.stream(); - return; - } catch (error) { - if (!isCurrentSubscription()) { - return; - } - logger.log('ResumableSSE', 'Token refresh failed:', error); - } - } - /** * Server-sent error event (event: error with data) - no responseCode. * These are known errors (ErrorTypes, ViolationTypes) that should be displayed to user. - * Only check e.data if there's no HTTP responseCode, since HTTP errors may also have body data. + * Only check the data if there's no HTTP status, since HTTP errors may also have body data. * Note: responseCode === 0 means transport failure (connection dropped) - treat as network error, * not a server-sent error payload. Use `== null` to only match undefined/null (no HTTP status). */ - if (responseCode == null && e.data) { + if (responseCode == null && data != null) { finalReceived = true; const recoveryConvoId = currentSubmission.conversation?.conversationId ?? currentStreamId; if ( @@ -3396,7 +3369,7 @@ export default function useResumableSSE( ) { return; } - logger.log('ResumableSSE', 'Server-sent error event received:', e.data); + logger.log('ResumableSSE', 'Server-sent error event received:', data); cancelSteerRetryFrames(); closeStream(); /** FLUSH (not cancel): the error card below is built from the cache @@ -3414,42 +3387,34 @@ export default function useResumableSSE( removeConvoFromAllQueries(queryClient, currentStreamId); } - let errorSupportsV2 = false; - try { - const errorData = JSON.parse(e.data); - errorSupportsV2 = supportsGenerationProtocolV2(errorData); - const errorString = errorData.error ?? errorData.message ?? JSON.stringify(errorData); + const errorSupportsV2 = supportsGenerationProtocolV2(data); + const errorString = + typeof data === 'string' + ? JSON.stringify(data) + : (data.error ?? data.message ?? JSON.stringify(data)); - // Check if it's a known error type (ViolationTypes or ErrorTypes) - let isKnownError = false; - try { - const parsed = - typeof errorString === 'string' ? JSON.parse(errorString) : errorString; - const errorType = parsed?.type ?? parsed?.code; - if (errorType) { - const violationValues = Object.values(ViolationTypes) as string[]; - const errorTypeValues = Object.values(ErrorTypes) as string[]; - isKnownError = - violationValues.includes(errorType) || errorTypeValues.includes(errorType); - } - } catch { - // Not JSON or parsing failed - treat as generic error + // Check if it's a known error type (ViolationTypes or ErrorTypes) + let isKnownError = false; + try { + const parsed = typeof errorString === 'string' ? JSON.parse(errorString) : errorString; + const errorType = parsed?.type ?? parsed?.code; + if (errorType) { + const violationValues = Object.values(ViolationTypes) as string[]; + const errorTypeValues = Object.values(ErrorTypes) as string[]; + isKnownError = + violationValues.includes(errorType) || errorTypeValues.includes(errorType); } + } catch { + // Not JSON or parsing failed - treat as generic error + } - logger.log('ResumableSSE', 'Error type check:', { isKnownError, errorString }); + logger.log('ResumableSSE', 'Error type check:', { isKnownError, errorString }); - // Display the error to user via errorHandler - errorHandler({ - data: { text: errorString } as unknown as Parameters[0]['data'], - submission: currentSubmission as EventSubmission, - }); - } catch (parseError) { - logger.error('ResumableSSE', 'Failed to parse server error:', parseError); - errorHandler({ - data: { text: e.data } as unknown as Parameters[0]['data'], - submission: currentSubmission as EventSubmission, - }); - } + // Display the error to user via errorHandler + errorHandler({ + data: { text: errorString } as unknown as Parameters[0]['data'], + submission: currentSubmission as EventSubmission, + }); setIsSubmitting(false); setShowStopButton(false); @@ -3515,7 +3480,7 @@ export default function useResumableSSE( // Network failure or unknown HTTP error - attempt reconnection with backoff logger.log('ResumableSSE', 'Stream error (network failure) - will attempt reconnect', { responseCode, - hasData: !!e.data, + hasData: data != null, }); if (reconnectAttemptRef.current < MAX_RETRIES) { @@ -3777,44 +3742,18 @@ export default function useResumableSSE( } }; - sse.addEventListener('error', handleTransportFailure); - /** - * Abort event - fired when the underlying XHR is cancelled, either by one - * of this hook's own closes or by the user agent. + * The transport emits `abort` only for this hook's own closes. A cancel + * the user agent issued (a backgrounded or frozen mobile tab, while the + * generation keeps running server-side) arrives as an `error` with status + * 0 instead, so it climbs the same reconnect ladder as any dropped + * connection rather than stranding the pane on its partial content. */ - sse.addEventListener('abort', () => { + const handleClose = () => { if (!isCurrentSubscription()) { return; } - /** - * A cancellation this hook did not issue came from the user agent, - * which cancels in-flight requests when a mobile browser is - * backgrounded or the page is frozen — and the generation it was - * carrying is still running server-side. - * - * Treating that as a deliberate close is what strands the response: - * the pane goes idle holding whatever partial content arrived before - * the switch, looking finished, and nothing re-reads the conversation - * until a reload or a navigation remounts the messages query. It is a - * dropped connection by every meaningful measure, so hand it to the - * transport-failure path verbatim rather than re-deriving a ladder - * beside it: that one already climbs its backoff, adjudicates the - * retry ceiling against durable status, and terminalizes into the - * refetch when the job turns out to have finished meanwhile. - */ - if (!closedByUs) { - logger.log( - 'ResumableSSE', - 'Stream aborted by the user agent - recovering as transport failure', - ); - void handleTransportFailure({ - responseCode: 0, - } as MessageEvent & { responseCode?: number }); - return; - } - if (replacementHandoffRef.current) { logger.log('ResumableSSE', 'Stream closed for generation handoff - preserving state'); return; @@ -3842,25 +3781,41 @@ export default function useResumableSSE( * merge into the next response in this conversation. On a resume the * collected usage is re-folded via backfillUsage, so nothing is lost. */ resetLive({ ...currentSubmission, userMessage }); - }); + }; - // Start the SSE connection - sse.stream(); + connection = createSSETransport({ token }).reconnectToStream( + { url, headers: generationProtocolHeaders() }, + { + signal: streamController.signal, + onEvent: (event) => { + switch (event.type) { + case 'open': + handleOpen(); + return; + case 'error': + void handleTransportFailure(event); + return; + case 'abort': + handleClose(); + return; + default: + void handleFrame(event); + } + }, + }, + ); // Debug hooks for testing reconnection vs clean close behavior (dev only) if (import.meta.env.DEV) { const debugWindow = window as Window & { - __sse?: SSE; __killNetwork?: () => void; __closeClean?: () => void; }; - debugWindow.__sse = sse; /** Simulate network drop - triggers error event → reconnection */ debugWindow.__killNetwork = () => { logger.log('Debug', 'Simulating network drop...'); - // @ts-ignore - sse.js types are incorrect, dispatchEvent actually takes Event - sse.dispatchEvent(new Event('error')); + void handleTransportFailure({}); }; /** Simulate clean close (navigation away) - triggers abort event → no reconnection */ @@ -4150,10 +4105,8 @@ export default function useResumableSSE( stopForegroundReattachRef.current?.(); stopForegroundReattachRef.current = null; // Close SSE but do NOT dispatch cancel - navigation should not abort - if (sseRef.current) { - sseRef.current.close(); - sseRef.current = null; - } + streamRef.current?.abort(); + streamRef.current = null; setStreamId(null); reconnectAttemptRef.current = 0; submissionRef.current = null; @@ -4694,10 +4647,8 @@ export default function useResumableSSE( stopForegroundReattachRef.current = null; // Reset reconnect counter before closing (so abort handler doesn't think we're reconnecting) reconnectAttemptRef.current = 0; - if (sseRef.current) { - sseRef.current.close(); - sseRef.current = null; - } + streamRef.current?.abort(); + streamRef.current = null; // Clear handler maps to prevent memory leaks and stale state clearStepMaps(); // Reset UI state on ordinary cleanup. A generation handoff already diff --git a/packages/data-provider/src/types/transport.ts b/packages/data-provider/src/types/transport.ts index 369842cddfe..02eeb93a9d6 100644 --- a/packages/data-provider/src/types/transport.ts +++ b/packages/data-provider/src/types/transport.ts @@ -224,8 +224,13 @@ export type ChatEvent = | { type: 'content'; data: ChatContentFrame } /** AI SDK: `text-delta`, except the text is cumulative rather than a delta. */ | { type: 'text'; data: ChatTextFrame } - /** AI SDK: `error`. `data` is `undefined` when the error body was not JSON. */ - | { type: 'error'; data?: ChatErrorData | null } + /** + * AI SDK: `error`. `data` is `undefined` when the error body was not JSON. + * `status` is set by `reconnectToStream` only: the HTTP status of a failed + * connection, `0` when it dropped (including a cancel the caller did not + * issue), and absent for an error event the server wrote into the stream. + */ + | { type: 'error'; data?: ChatErrorData | null; status?: number } /** AI SDK: `abort`. The caller closed a stream that was still open. */ | { type: 'abort' }; @@ -247,15 +252,46 @@ export type ChatTransportOptions = { onEvent: (event: ChatEvent) => void; }; +/** Attaches to a generation already running on the server. */ +export type ChatStreamRequest = { + /** The stream route, with the resume cursor and generation fence in its query. */ + url: string; + /** Sent beside the bearer token and kept across a token refresh. */ + headers?: Record; +}; + +/** The handle `reconnectToStream` returns for one attachment. */ +export interface ChatStreamConnection { + /** + * Whether the connection has closed. A response body that simply ends + * dispatches no event, so this is the only way to see that it did. + */ + readonly closed: boolean; +} + /** * Carries one turn from request to terminal event. Implementations own the * wire (connection, framing, auth refresh); callers only see {@link ChatEvent}s. * * AI SDK: `ChatTransport`. `send` corresponds to `sendMessages`, returning * through a callback rather than a `ReadableStream` so handlers keep running - * synchronously inside the frame that produced them. `reconnectToStream` has - * no counterpart yet; resume still lives in `useResumableSSE`. + * synchronously inside the frame that produced them. */ export interface ChatTransport { send(request: TRequest, options: ChatTransportOptions): void; + /** + * Attaches to a running generation and reports it through the same events + * as `send`. Aborting the signal while the stream is open emits + * `{ type: 'abort' }`; a cancel the caller did not issue (a backgrounded or + * frozen tab) is a dropped connection and emits `{ type: 'error', status: 0 }`. + * Each 401 refreshes the token and reattaches on the same handle; a refresh + * that fails is reported as the 401. + * + * AI SDK: `reconnectToStream`, which resolves to a stream (or `null` when + * nothing is running); here the caller learns that from a 404 `error`. + */ + reconnectToStream( + request: ChatStreamRequest, + options: ChatTransportOptions, + ): ChatStreamConnection; } From 82901a1408b99fe37a40702c7bfcad9322b8713b Mon Sep 17 00:00:00 2001 From: Marco Beretta <81851188+berry-13@users.noreply.github.com> Date: Fri, 25 Sep 2026 21:07:31 +0200 Subject: [PATCH 2/2] fix: bound the resume 401 refresh and keep plain-text server errors Refresh the token once per attachment and report a second 401, so the hook's reconnect budget bounds a stream that keeps answering 401. A server-written error event that is not JSON keeps its raw text instead of reading as a dropped connection. --- client/src/hooks/SSE/__tests__/sse.spec.ts | 47 +++++++++++++------ client/src/hooks/SSE/transport/sse.ts | 19 ++++++-- client/src/hooks/SSE/useResumableSSE.ts | 4 +- packages/data-provider/src/types/transport.ts | 6 ++- 4 files changed, 52 insertions(+), 24 deletions(-) diff --git a/client/src/hooks/SSE/__tests__/sse.spec.ts b/client/src/hooks/SSE/__tests__/sse.spec.ts index 4f882c436d2..77f0f1ea629 100644 --- a/client/src/hooks/SSE/__tests__/sse.spec.ts +++ b/client/src/hooks/SSE/__tests__/sse.spec.ts @@ -338,6 +338,13 @@ describe('createSSETransport().reconnectToStream', () => { expect(events).toEqual([{ type: 'error', status: undefined, data: { error: 'failed' } }]); }); + it('keeps the raw text of a server error event that is not JSON', () => { + attach(); + current().write('event: error\ndata: service unavailable\n\n'); + + expect(events).toEqual([{ type: 'error', status: undefined, data: 'service unavailable' }]); + }); + it('emits an HTTP failure with its status and no data for a body that is not JSON', () => { attach(); current().status = 404; @@ -381,33 +388,45 @@ describe('createSSETransport().reconnectToStream', () => { expect(events.map((event) => event.type)).toEqual(['open', 'created']); }); - it('refreshes the token on every 401 and reattaches with the request headers', async () => { - const refreshToken = jest - .spyOn(request, 'refreshToken') - .mockResolvedValueOnce({ token: 'token-2' } as never) - .mockResolvedValueOnce({ token: 'token-3' } as never); + it('refreshes the token on a 401 and reattaches with the request headers', async () => { + jest.spyOn(request, 'refreshToken').mockResolvedValue({ token: 'token-2' } as never); const dispatchTokenUpdated = jest .spyOn(request, 'dispatchTokenUpdatedEvent') .mockImplementation(() => undefined); attach(); + current().status = 401; + current().write('Unauthorized'); - for (const _attempt of [1, 2]) { - current().status = 401; - current().write('Unauthorized'); - await new Promise(process.nextTick); - } + await new Promise(process.nextTick); - expect(refreshToken).toHaveBeenCalledTimes(2); - expect(dispatchTokenUpdated).toHaveBeenLastCalledWith('token-3'); - expect(xhrs).toHaveLength(3); + expect(dispatchTokenUpdated).toHaveBeenCalledWith('token-2'); + expect(xhrs).toHaveLength(2); expect(current().method).toBe('GET'); expect(current().headers).toEqual({ - Authorization: 'Bearer token-3', + Authorization: 'Bearer token-2', 'X-LibreChat-Generation-Protocol': '2', }); expect(events).toEqual([]); }); + it('reports a second 401 instead of refreshing again', async () => { + const refreshToken = jest + .spyOn(request, 'refreshToken') + .mockResolvedValue({ token: 'token-2' } as never); + jest.spyOn(request, 'dispatchTokenUpdatedEvent').mockImplementation(() => undefined); + attach(); + + for (const _attempt of [1, 2]) { + current().status = 401; + current().write('Unauthorized'); + await new Promise(process.nextTick); + } + + expect(refreshToken).toHaveBeenCalledTimes(1); + expect(xhrs).toHaveLength(2); + expect(events).toEqual([{ type: 'error', status: 401, data: undefined }]); + }); + it('reports the 401 when the token refresh fails', async () => { jest.spyOn(console, 'log').mockImplementation(() => undefined); jest.spyOn(request, 'refreshToken').mockRejectedValue(new Error('refresh failed')); diff --git a/client/src/hooks/SSE/transport/sse.ts b/client/src/hooks/SSE/transport/sse.ts index f73b6c0c072..5970a15ec53 100644 --- a/client/src/hooks/SSE/transport/sse.ts +++ b/client/src/hooks/SSE/transport/sse.ts @@ -48,15 +48,18 @@ async function refreshToken(): Promise { } } -/** Failure bodies are often empty or HTML, so a body that is not JSON is `undefined`. */ -function parseErrorBody(body: unknown): ChatErrorData | undefined { +/** + * An HTTP failure body is often empty or HTML, so one that is not JSON is + * `undefined`. An error event the server wrote keeps its raw text instead. + */ +function parseErrorBody(body: unknown, status?: number): ChatErrorData | undefined { if (typeof body !== 'string' || body === '') { return undefined; } try { return JSON.parse(body); } catch { - return undefined; + return status == null ? body : undefined; } } @@ -159,8 +162,10 @@ export function createSSETransport({ token }: { token?: string }): ChatTransport sse.addEventListener('message', emitFrame(onEvent)); + let refreshed = false; sse.addEventListener('error', async (e: StreamErrorEvent) => { - if (e.responseCode === 401) { + if (e.responseCode === 401 && !refreshed) { + refreshed = true; const refreshedToken = await refreshToken(); if (signal.aborted) { return; @@ -172,7 +177,11 @@ export function createSSETransport({ token }: { token?: string }): ChatTransport return; } } - onEvent({ type: 'error', status: e.responseCode, data: parseErrorBody(e.data) }); + onEvent({ + type: 'error', + status: e.responseCode, + data: parseErrorBody(e.data, e.responseCode), + }); }); /** sse.js dispatches `abort` when the XHR is cancelled, by our close or by the user agent. */ diff --git a/client/src/hooks/SSE/useResumableSSE.ts b/client/src/hooks/SSE/useResumableSSE.ts index 2cbdd04793a..c7cdf01e636 100644 --- a/client/src/hooks/SSE/useResumableSSE.ts +++ b/client/src/hooks/SSE/useResumableSSE.ts @@ -3389,9 +3389,7 @@ export default function useResumableSSE( const errorSupportsV2 = supportsGenerationProtocolV2(data); const errorString = - typeof data === 'string' - ? JSON.stringify(data) - : (data.error ?? data.message ?? JSON.stringify(data)); + typeof data === 'string' ? data : (data.error ?? data.message ?? JSON.stringify(data)); // Check if it's a known error type (ViolationTypes or ErrorTypes) let isKnownError = false; diff --git a/packages/data-provider/src/types/transport.ts b/packages/data-provider/src/types/transport.ts index 02eeb93a9d6..af3751edfab 100644 --- a/packages/data-provider/src/types/transport.ts +++ b/packages/data-provider/src/types/transport.ts @@ -284,8 +284,10 @@ export interface ChatTransport { * as `send`. Aborting the signal while the stream is open emits * `{ type: 'abort' }`; a cancel the caller did not issue (a backgrounded or * frozen tab) is a dropped connection and emits `{ type: 'error', status: 0 }`. - * Each 401 refreshes the token and reattaches on the same handle; a refresh - * that fails is reported as the 401. + * The first 401 refreshes the token and reattaches on the same handle; a + * failed refresh or a second 401 is reported as the 401, so the caller's own + * retry budget bounds it. A server-written error that is not JSON arrives as + * its raw text. * * AI SDK: `reconnectToStream`, which resolves to a stream (or `null` when * nothing is running); here the caller learns that from a 404 `error`.