diff --git a/.changeset/calm-streams-close.md b/.changeset/calm-streams-close.md new file mode 100644 index 0000000000..0c673cc183 --- /dev/null +++ b/.changeset/calm-streams-close.md @@ -0,0 +1,5 @@ +--- +'@tanstack/router-core': patch +--- + +Close SSR stream transforms gracefully when their lifetime watchdog expires diff --git a/packages/router-core/src/ssr/transformStreamWithRouter.ts b/packages/router-core/src/ssr/transformStreamWithRouter.ts index 4393f955b3..d78052cc90 100644 --- a/packages/router-core/src/ssr/transformStreamWithRouter.ts +++ b/packages/router-core/src/ssr/transformStreamWithRouter.ts @@ -311,7 +311,7 @@ function makeFastPathStream( console.warn( `SSR stream transform exceeded maximum lifetime (${lifetimeMs}ms), forcing cleanup`, ) - safeError(err) + safeClose() cleanup(err) } }, lifetimeMs) @@ -735,7 +735,7 @@ function makeMainStream( console.warn( `SSR stream transform exceeded maximum lifetime (${lifetimeMs}ms), forcing cleanup`, ) - safeError(err) + safeClose() cleanup(err) } }, lifetimeMs) diff --git a/packages/router-core/tests/transformStreamWithRouter.test.ts b/packages/router-core/tests/transformStreamWithRouter.test.ts index 13b0072e24..7f4263c823 100644 --- a/packages/router-core/tests/transformStreamWithRouter.test.ts +++ b/packages/router-core/tests/transformStreamWithRouter.test.ts @@ -639,37 +639,80 @@ describe('transformStreamWithRouter — cleanup side-effects', () => { } }) - test('lifetime timeout cancels upstream and runs cleanup once', async () => { - vi.useFakeTimers() - try { - const { router, cleanupCalls } = makeRouter({ - isSerializationFinished: () => true, - takeBufferedHtml: () => undefined, - }) - const upstream = makeManualUpstream() - - const out = transformStreamWithRouter( - router as any, - upstream.stream as any, - { lifetimeMs: 10 }, - ) - - // Do NOT consume. Advance fake time past lifetimeMs deterministically. - await vi.advanceTimersByTimeAsync(15) - - expect(upstream.cancelled.value).toBe(true) - expect(cleanupCalls.count).toBe(1) - - // Drain (read errors silently) so vitest doesn't see an unhandled error. - const reader = ( - out as any - ).getReader() as ReadableStreamDefaultReader - reader.read().catch(() => {}) - reader.releaseLock() - } finally { - vi.useRealTimers() - } - }) + const lifetimeTimeoutCases = [ + { + name: 'fast path with an active reader', + reserveStreamFastPath: true, + activeReader: true, + }, + { + name: 'fast path without an active reader', + reserveStreamFastPath: true, + activeReader: false, + }, + { + name: 'main path with an active reader', + reserveStreamFastPath: false, + activeReader: true, + }, + { + name: 'main path without an active reader', + reserveStreamFastPath: false, + activeReader: false, + }, + ] as const + + test.each(lifetimeTimeoutCases)( + 'lifetime timeout closes $name and cleans up', + async ({ reserveStreamFastPath, activeReader }) => { + vi.useFakeTimers() + const warnSpy = vi.spyOn(console, 'warn').mockImplementation(() => {}) + try { + const { router, cleanupCalls } = makeRouter({ + isSerializationFinished: () => true, + reserveStreamFastPath: () => reserveStreamFastPath, + takeBufferedHtml: () => undefined, + }) + const upstream = makeManualUpstream() + + const out = transformStreamWithRouter( + router as any, + upstream.stream as any, + { lifetimeMs: 10 }, + ) + if (activeReader) { + const reader = out.getReader() + const pendingRead = reader.read() + + await vi.advanceTimersByTimeAsync(15) + + await expect(pendingRead).resolves.toEqual({ + done: true, + value: undefined, + }) + reader.releaseLock() + } else { + await vi.advanceTimersByTimeAsync(15) + + const reader = out.getReader() + await expect(reader.read()).resolves.toEqual({ + done: true, + value: undefined, + }) + reader.releaseLock() + } + expect(upstream.cancelled.value).toBe(true) + expect(upstream.cancelled.reason).toEqual( + new Error('Stream lifetime exceeded'), + ) + expect(cleanupCalls.count).toBe(1) + expect(warnSpy).toHaveBeenCalledOnce() + } finally { + warnSpy.mockRestore() + vi.useRealTimers() + } + }, + ) test('upstream cancel() that rejects does not produce an unhandled rejection', async () => { // Some upstream sources may reject when cancel() is called (e.g. their