diff --git a/packages/router-ssr-query-core/src/index.ts b/packages/router-ssr-query-core/src/index.ts index 685653dea0c..49bdd5edf05 100644 --- a/packages/router-ssr-query-core/src/index.ts +++ b/packages/router-ssr-query-core/src/index.ts @@ -1,4 +1,5 @@ import { + notifyManager, dehydrate as queryDehydrate, hydrate as queryHydrate, } from '@tanstack/query-core' @@ -48,6 +49,7 @@ export function setupCoreRouterSsrQueryIntegration({ let unsubscribe: (() => void) | undefined = undefined let cleanupRegistered = false let tornDown = false + let pendingQueryHashes: Set | undefined const teardown = () => { if (tornDown) return @@ -78,6 +80,8 @@ export function setupCoreRouterSsrQueryIntegration({ // ignore } sentQueries.clear() + pendingQueryHashes?.clear() + pendingQueryHashes = undefined } // Register teardown as soon as SSR attaches. attachRouterServerSsrUtils() @@ -100,6 +104,7 @@ export function setupCoreRouterSsrQueryIntegration({ router.options.dehydrate = async (): Promise => { router.serverSsr!.onRenderFinished(() => { + flushPendingQueries() if (!queryStream.isClosed()) queryStream.close() unsubscribe?.() unsubscribe = undefined @@ -135,6 +140,43 @@ export function setupCoreRouterSsrQueryIntegration({ }, }) + const flushPendingQueries = () => { + const queryHashes = pendingQueryHashes + pendingQueryHashes = undefined + if ( + tornDown || + queryStream.isClosed() || + queryHashes === undefined || + queryHashes.size === 0 + ) { + return + } + + const dehydratedQuery = queryDehydrate(queryClient, { + ...dehydrateOptions, + shouldDehydrateQuery: (query) => { + if (!queryHashes.has(query.queryHash)) { + return false + } + + return ( + (ogClientOptions.dehydrate?.shouldDehydrateQuery?.(query) ?? + true) && + (dehydrateOptions?.shouldDehydrateQuery?.(query) ?? true) + ) + }, + }) + + if (dehydratedQuery.queries.length === 0) { + return + } + + dehydratedQuery.queries.forEach((query) => { + sentQueries.add(query.queryHash) + }) + queryStream.enqueue(dehydratedQuery) + } + unsubscribe = queryClient.getQueryCache().subscribe((event) => { // before rendering starts, we do not stream individual queries // instead we dehydrate the entire query client in router's dehydrate() @@ -155,27 +197,19 @@ export function setupCoreRouterSsrQueryIntegration({ ) return } - const dehydratedQuery = queryDehydrate(queryClient, { - ...dehydrateOptions, - shouldDehydrateQuery: (query) => { - if (query.queryHash !== event.query.queryHash) { - return false + if (!pendingQueryHashes) { + pendingQueryHashes = new Set() + // QueryCache listeners run inside notifyManager.batch. Scheduling the + // flush collects queries completed during the same scheduler turn. + notifyManager.schedule(() => { + try { + flushPendingQueries() + } catch (err) { + queryStream.error(err) } - - return ( - (ogClientOptions.dehydrate?.shouldDehydrateQuery?.(query) ?? - true) && - (dehydrateOptions?.shouldDehydrateQuery?.(query) ?? true) - ) - }, - }) - - if (dehydratedQuery.queries.length === 0) { - return + }) } - - sentQueries.add(event.query.queryHash) - queryStream.enqueue(dehydratedQuery) + pendingQueryHashes.add(event.query.queryHash) }) // on the client } else { diff --git a/packages/router-ssr-query-core/tests/index.test.ts b/packages/router-ssr-query-core/tests/index.test.ts index a167dbd03e1..075debf94f4 100644 --- a/packages/router-ssr-query-core/tests/index.test.ts +++ b/packages/router-ssr-query-core/tests/index.test.ts @@ -1,25 +1,40 @@ import { QueryClient } from '@tanstack/query-core' +import { + BaseRootRoute, + RouterCore, + createNonReactiveMutableStore, + createNonReactiveReadonlyStore, +} from '@tanstack/router-core' +import { attachRouterServerSsrUtils } from '@tanstack/router-core/ssr/server' import { afterEach, describe, expect, it, vi } from 'vitest' import { setupCoreRouterSsrQueryIntegration } from '../src' +import type { GetStoreConfig } from '@tanstack/router-core' -type TestRouter = { - isServer: boolean - options: { - dehydrate?: () => unknown | Promise - hydrate?: (dehydrated: any) => unknown | Promise - } - serverSsr?: { - isDehydrated: () => boolean - onRenderFinished: (listener: () => void) => void - onCleanup: (listener: () => void) => void - } - serverSsrLifecycle?: { - onServerSsrAttach: Array< - (serverSsr: NonNullable) => void - > +const getStoreConfig: GetStoreConfig = () => ({ + createMutableStore: createNonReactiveMutableStore, + createReadonlyStore: createNonReactiveReadonlyStore, + batch: (fn) => fn(), +}) + +function createTestRouter(isServer: boolean) { + const router = new RouterCore( + { + routeTree: new BaseRootRoute({}), + isServer, + }, + getStoreConfig, + ) + + // RouterCore exposes test routers globally in jsdom, defeating the GC tests. + if (Reflect.get(globalThis, '__TSR_ROUTER__') === router) { + Reflect.deleteProperty(globalThis, '__TSR_ROUTER__') } + + return router } +type TestRouter = ReturnType + type ServerRouterFixture = { router: TestRouter finishRender: () => void @@ -30,46 +45,48 @@ type ServerRouterFixture = { } function createServerRouter(): ServerRouterFixture { - const renderFinishedListeners = new Array<() => void>() - const cleanupListeners = new Array<() => void>() + const router = createTestRouter(true) + let serverSsr: NonNullable | undefined let dehydrated = false - const serverSsr = { - isDehydrated: () => dehydrated, - onRenderFinished: (listener: () => void) => { - renderFinishedListeners.push(listener) - }, - onCleanup: (listener: () => void) => { - cleanupListeners.push(listener) - }, + let cleanupListenerCount = 0 + + router.serverSsrLifecycle = { + onServerSsrAttach: [ + (attachedServerSsr) => { + serverSsr = attachedServerSsr + attachedServerSsr.isDehydrated = () => dehydrated + + const onCleanup = attachedServerSsr.onCleanup + attachedServerSsr.onCleanup = (listener) => { + cleanupListenerCount++ + onCleanup(() => { + try { + listener() + } finally { + cleanupListenerCount-- + } + }) + } + }, + ], } - const result: ServerRouterFixture = { - router: { - isServer: true, - options: {}, - serverSsr, - }, + return { + router, finishRender: () => { - renderFinishedListeners.splice(0).forEach((listener) => listener()) + serverSsr?.setRenderFinished() }, triggerCleanup: () => { - cleanupListeners.splice(0).forEach((listener) => listener()) + serverSsr?.cleanup() }, attachServerSsr: () => { - result.router.serverSsr = serverSsr - result.router.serverSsrLifecycle?.onServerSsrAttach.forEach( - (listener) => { - listener(serverSsr) - }, - ) + attachRouterServerSsrUtils({ router, manifest: undefined }) }, setDehydrated: (value: boolean) => { dehydrated = value }, - cleanupListenerCount: () => cleanupListeners.length, + cleanupListenerCount: () => cleanupListenerCount, } - - return result } async function readStream(stream: ReadableStream): Promise> { @@ -149,9 +166,8 @@ describe('setupCoreRouterSsrQueryIntegration', () => { const { router, finishRender, attachServerSsr, setDehydrated } = createServerRouter() - router.serverSsr = undefined setupCoreRouterSsrQueryIntegration({ - router: router as any, + router, queryClient, dehydrateOptions: { serializeData: (data) => `${data}-serialized`, @@ -210,13 +226,10 @@ describe('setupCoreRouterSsrQueryIntegration', () => { it('uses custom hydrate options for the initial payload and streamed queries', async () => { const queryClient = track(new QueryClient()) - const router: TestRouter = { - isServer: false, - options: {}, - } + const router = createTestRouter(false) setupCoreRouterSsrQueryIntegration({ - router: router as any, + router, queryClient, hydrateOptions: { defaultOptions: { @@ -261,6 +274,101 @@ describe('setupCoreRouterSsrQueryIntegration', () => { expect(queryClient.getQueryData(['initial'])).toBe('initial-hydrated') expect(queryClient.getQueryData(['streamed'])).toBe('stream-hydrated') }) + + it('batches queries by scheduler turn and only streams them once', async () => { + const queryClient = track(new QueryClient()) + const { router, finishRender, attachServerSsr, setDehydrated } = + createServerRouter() + + setupCoreRouterSsrQueryIntegration({ + router, + queryClient, + }) + attachServerSsr() + + const dehydrated = (await router.options.dehydrate?.()) as { + queryStream: ReadableStream<{ + queries: Array<{ queryKey: Array }> + }> + } + const streamedQueriesPromise = readStream(dehydrated.queryStream) + const firstDeferred = createDeferred() + const secondDeferred = createDeferred() + const laterDeferred = createDeferred() + + setDehydrated(true) + const firstPromise = queryClient.fetchQuery({ + queryKey: ['first'], + queryFn: () => firstDeferred.promise, + }) + const secondPromise = queryClient.fetchQuery({ + queryKey: ['second'], + queryFn: () => secondDeferred.promise, + }) + + firstDeferred.resolve('first-data') + secondDeferred.resolve('second-data') + await Promise.all([firstPromise, secondPromise]) + await new Promise((resolve) => setTimeout(resolve, 0)) + + await queryClient.fetchQuery({ + queryKey: ['first'], + queryFn: () => 'first-data-updated', + }) + expect(queryClient.getQueryData(['first'])).toBe('first-data-updated') + + const laterPromise = queryClient.fetchQuery({ + queryKey: ['later'], + queryFn: () => laterDeferred.promise, + }) + laterDeferred.resolve('later-data') + await laterPromise + finishRender() + + const streamedQueries = await streamedQueriesPromise + + expect(streamedQueries).toHaveLength(2) + expect(streamedQueries[0]?.queries.map((query) => query.queryKey)).toEqual([ + ['first'], + ['second'], + ]) + expect(streamedQueries[1]?.queries.map((query) => query.queryKey)).toEqual([ + ['later'], + ]) + }) + + it('does not stream pending queries after request cleanup', async () => { + const queryClient = track(new QueryClient()) + const { router, triggerCleanup, attachServerSsr, setDehydrated } = + createServerRouter() + + setupCoreRouterSsrQueryIntegration({ + router, + queryClient, + }) + attachServerSsr() + + const dehydrated = (await router.options.dehydrate?.()) as { + queryStream: ReadableStream<{ + queries: Array<{ queryKey: Array }> + }> + } + const streamedQueriesPromise = readStream(dehydrated.queryStream) + const deferred = createDeferred() + + setDehydrated(true) + const queryPromise = queryClient.fetchQuery({ + queryKey: ['pending'], + queryFn: () => deferred.promise, + }) + + deferred.resolve('pending-data') + await queryPromise + triggerCleanup() + + await new Promise((resolve) => setTimeout(resolve, 0)) + expect(await streamedQueriesPromise).toEqual([]) + }) }) // GC reclamation tests are non-deterministic by nature (V8 makes no @@ -293,9 +401,8 @@ describe.runIf(gcTestsEnabled)('SSR memory: GC reclamation', () => { let serverRouter: ReturnType | null = createServerRouter() - serverRouter.router.serverSsr = undefined setupCoreRouterSsrQueryIntegration({ - router: serverRouter.router as any, + router: serverRouter.router, queryClient, }) serverRouter.attachServerSsr() @@ -334,9 +441,8 @@ describe.runIf(gcTestsEnabled)('SSR memory: GC reclamation', () => { let serverRouter: ReturnType | null = createServerRouter() - serverRouter.router.serverSsr = undefined setupCoreRouterSsrQueryIntegration({ - router: serverRouter.router as any, + router: serverRouter.router, queryClient, }) serverRouter.attachServerSsr() @@ -376,9 +482,8 @@ describe('SSR cleanup: deterministic behavior', () => { const { router, triggerCleanup, attachServerSsr, setDehydrated } = createServerRouter() - router.serverSsr = undefined setupCoreRouterSsrQueryIntegration({ - router: router as any, + router, queryClient, }) attachServerSsr() @@ -404,9 +509,8 @@ describe('SSR cleanup: deterministic behavior', () => { const { router, triggerCleanup, attachServerSsr, setDehydrated } = createServerRouter() - router.serverSsr = undefined setupCoreRouterSsrQueryIntegration({ - router: router as any, + router, queryClient, }) attachServerSsr() @@ -450,9 +554,8 @@ describe('SSR cleanup: deterministic behavior', () => { cleanupListenerCount, } = createServerRouter() - router.serverSsr = undefined setupCoreRouterSsrQueryIntegration({ - router: router as any, + router, queryClient, }) attachServerSsr() @@ -490,10 +593,8 @@ describe('SSR cleanup: deterministic behavior', () => { cleanupListenerCount, } = createServerRouter() // Detach to simulate pre-attach state. - router.serverSsr = undefined - setupCoreRouterSsrQueryIntegration({ - router: router as any, + router, queryClient, })