Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
72 changes: 53 additions & 19 deletions packages/router-ssr-query-core/src/index.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import {
notifyManager,
dehydrate as queryDehydrate,
hydrate as queryHydrate,
} from '@tanstack/query-core'
Expand Down Expand Up @@ -48,6 +49,7 @@ export function setupCoreRouterSsrQueryIntegration<TRouter extends AnyRouter>({
let unsubscribe: (() => void) | undefined = undefined
let cleanupRegistered = false
let tornDown = false
let pendingQueryHashes: Set<string> | undefined

const teardown = () => {
if (tornDown) return
Expand Down Expand Up @@ -78,6 +80,8 @@ export function setupCoreRouterSsrQueryIntegration<TRouter extends AnyRouter>({
// ignore
}
sentQueries.clear()
pendingQueryHashes?.clear()
pendingQueryHashes = undefined
}

// Register teardown as soon as SSR attaches. attachRouterServerSsrUtils()
Expand All @@ -100,6 +104,7 @@ export function setupCoreRouterSsrQueryIntegration<TRouter extends AnyRouter>({
router.options.dehydrate =
async (): Promise<DehydratedRouterQueryState> => {
router.serverSsr!.onRenderFinished(() => {
flushPendingQueries()
if (!queryStream.isClosed()) queryStream.close()
unsubscribe?.()
unsubscribe = undefined
Expand Down Expand Up @@ -135,6 +140,43 @@ export function setupCoreRouterSsrQueryIntegration<TRouter extends AnyRouter>({
},
})

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()
Expand All @@ -155,27 +197,19 @@ export function setupCoreRouterSsrQueryIntegration<TRouter extends AnyRouter>({
)
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 {
Expand Down
Loading