From d14ccb5f836bc1ce18c85b8798b9a90ede5186aa Mon Sep 17 00:00:00 2001 From: Kyle Mathews Date: Mon, 13 Jul 2026 14:01:11 -0600 Subject: [PATCH 1/4] test: characterize query cancellation cleanup --- packages/query-db-collection/src/query.ts | 41 +++- .../query-db-collection/tests/query.test.ts | 179 ++++++++++++++++++ 2 files changed, 216 insertions(+), 4 deletions(-) diff --git a/packages/query-db-collection/src/query.ts b/packages/query-db-collection/src/query.ts index 7b4578d9e..18b59d63a 100644 --- a/packages/query-db-collection/src/query.ts +++ b/packages/query-db-collection/src/query.ts @@ -745,6 +745,7 @@ export function queryCollectionOptions( // queryKey → QueryObserver's unsubscribe function const unsubscribes = new Map void>() + const pendingReadyUnsubscribes = new Map void>>() // queryKey → reference count (how many loadSubset calls are active) // Reference counting for QueryObserver lifecycle management @@ -1174,6 +1175,27 @@ export function queryCollectionOptions( } }) + const trackPendingReadyUnsubscribe = ( + hashedQueryKey: string, + unsubscribe: () => void, + ) => { + const pending = pendingReadyUnsubscribes.get(hashedQueryKey) ?? new Set() + pending.add(unsubscribe) + pendingReadyUnsubscribes.set(hashedQueryKey, pending) + } + + const releasePendingReadyUnsubscribe = ( + hashedQueryKey: string, + unsubscribe: () => void, + ) => { + unsubscribe() + const pending = pendingReadyUnsubscribes.get(hashedQueryKey) + pending?.delete(unsubscribe) + if (pending?.size === 0) { + pendingReadyUnsubscribes.delete(hashedQueryKey) + } + } + const createQueryFromOpts = ( opts: LoadSubsetOptions = {}, queryFunction: typeof queryFn = queryFn, @@ -1239,14 +1261,15 @@ export function queryCollectionOptions( // Use a microtask in case `subscribe` is called synchronously, before `unsubscribe` is initialized queueMicrotask(() => { if (result.isSuccess) { - unsubscribe() + releasePendingReadyUnsubscribe(hashedQueryKey, unsubscribe) resolve() } else if (result.isError) { - unsubscribe() + releasePendingReadyUnsubscribe(hashedQueryKey, unsubscribe) reject(result.error) } }) }) + trackPendingReadyUnsubscribe(hashedQueryKey, unsubscribe) }) } } @@ -1322,14 +1345,15 @@ export function queryCollectionOptions( // Use a microtask in case `subscribe` is called synchronously, before `unsubscribe` is initialized queueMicrotask(() => { if (result.isSuccess) { - unsubscribe() + releasePendingReadyUnsubscribe(hashedQueryKey, unsubscribe) resolve() } else if (result.isError) { - unsubscribe() + releasePendingReadyUnsubscribe(hashedQueryKey, unsubscribe) reject(result.error) } }) }) + trackPendingReadyUnsubscribe(hashedQueryKey, unsubscribe) }) // If sync has started or there are subscribers to the collection, subscribe to the query straight away @@ -1590,9 +1614,17 @@ export function queryCollectionOptions( * Perform row-level cleanup and remove all tracking for a query. * Callers are responsible for ensuring the query is safe to cleanup. */ + const unsubscribePendingReadyListeners = (hashedQueryKey: string) => { + pendingReadyUnsubscribes.get(hashedQueryKey)?.forEach((unsubscribe) => { + unsubscribe() + }) + pendingReadyUnsubscribes.delete(hashedQueryKey) + } + const cleanupQueryInternal = (hashedQueryKey: string) => { unsubscribes.get(hashedQueryKey)?.() unsubscribes.delete(hashedQueryKey) + unsubscribePendingReadyListeners(hashedQueryKey) cancelPersistedRetentionExpiry(hashedQueryKey) retainedQueriesPendingRevalidation.delete(hashedQueryKey) @@ -1655,6 +1687,7 @@ export function queryCollectionOptions( // Drop our subscription so hasListeners reflects only active consumers unsubscribes.get(hashedQueryKey)?.() unsubscribes.delete(hashedQueryKey) + unsubscribePendingReadyListeners(hashedQueryKey) } const hasListeners = observer?.hasListeners() ?? false diff --git a/packages/query-db-collection/tests/query.test.ts b/packages/query-db-collection/tests/query.test.ts index c711d661f..e20c9cd98 100644 --- a/packages/query-db-collection/tests/query.test.ts +++ b/packages/query-db-collection/tests/query.test.ts @@ -50,6 +50,16 @@ const getKey = (item: TestItem) => item.id // Helper to advance timers and allow microtasks to flush const flushPromises = () => new Promise((resolve) => setTimeout(resolve, 0)) +function createDeferred() { + let resolve!: (value: T) => void + let reject!: (reason?: unknown) => void + const promise = new Promise((resolvePromise, rejectPromise) => { + resolve = resolvePromise + reject = rejectPromise + }) + return { promise, resolve, reject } +} + function createInMemorySyncMetadataApi< TKey extends string | number = string | number, TItem extends object = Record, @@ -2372,6 +2382,175 @@ describe(`QueryCollection`, () => { expect(collection.size).toBe(0) }) + describe(`query cancellation and subset cleanup lifecycle`, () => { + const createSubset = (collection: Collection) => + createLiveQueryCollection({ + query: (q) => + q + .from({ item: collection }) + .where(({ item }) => eq(item.id, `1`)) + .select(({ item }) => ({ id: item.id, name: item.name })), + }) + + it(`forwards Query Core's signal and aborts an in-flight eager query on collection cleanup`, async () => { + const deferred = createDeferred>() + let signal: AbortSignal | undefined + const collection = createCollection( + queryCollectionOptions({ + id: `signal-forwarding-cleanup-test`, + queryClient, + queryKey: [`signal-forwarding-cleanup-test`], + queryFn: (context) => { + signal = context.signal + // Reading the signal makes Query Core treat the request as cancellable. + void context.signal.aborted + return deferred.promise + }, + getKey, + startSync: true, + }), + ) + + await vi.waitFor(() => expect(signal).toBeDefined()) + expect(signal?.aborted).toBe(false) + + await collection.cleanup() + + expect(signal?.aborted).toBe(true) + expect(collection.size).toBe(0) + // Query Core may retain the cancelled cache entry, but cleanup releases every observer. + expect( + queryClient + .getQueryCache() + .find({ queryKey: [`signal-forwarding-cleanup-test`] }) + ?.getObserversCount() ?? 0, + ).toBe(0) + }) + + it(`does not rematerialize a late fulfilled result after an in-flight subset is unloaded`, async () => { + const deferred = createDeferred>() + const collection = createCollection( + queryCollectionOptions({ + id: `late-subset-result-test`, + queryClient, + queryKey: [`late-subset-result-test`], + queryFn: () => deferred.promise, + getKey, + syncMode: `on-demand`, + }), + ) + const liveQuery = createSubset(collection) + void liveQuery.preload().catch(() => undefined) + + await vi.waitFor(() => expect(queryClient.isFetching()).toBe(1)) + await liveQuery.cleanup() + deferred.resolve([{ id: `1`, name: `Late item` }]) + await flushPromises() + + expect(collection.size).toBe(0) + expect(queryClient.isFetching()).toBe(0) + expect( + queryClient + .getQueryCache() + .findAll({ + queryKey: [`late-subset-result-test`], + })[0] + ?.getObserversCount() ?? 0, + ).toBe(0) + await collection.cleanup() + }) + + it(`deterministically materializes a shared in-flight result after fast subset unmount and remount`, async () => { + const deferred = createDeferred>() + const queryFn = vi.fn(() => deferred.promise) + const collection = createCollection( + queryCollectionOptions({ + id: `fast-subset-remount-test`, + queryClient, + queryKey: [`fast-subset-remount-test`], + queryFn, + getKey, + syncMode: `on-demand`, + }), + ) + const firstLiveQuery = createSubset(collection) + void firstLiveQuery.preload().catch(() => undefined) + await vi.waitFor(() => expect(queryFn).toHaveBeenCalledTimes(1)) + + await firstLiveQuery.cleanup() + const secondLiveQuery = createSubset(collection) + const secondPreload = secondLiveQuery.preload() + deferred.resolve([{ id: `1`, name: `Remounted item` }]) + await secondPreload + + expect(queryFn).toHaveBeenCalledTimes(1) + expect(stripVirtualProps(collection.get(`1`))).toEqual({ + id: `1`, + name: `Remounted item`, + }) + await secondLiveQuery.cleanup() + expect(collection.size).toBe(0) + await collection.cleanup() + }) + + it(`keeps invalidate-unsubscribe-resubscribe compatible while removing stale subset rows and observers`, async () => { + let items: Array = [{ id: `1`, name: `Initial item` }] + const queryFn = vi.fn(() => Promise.resolve(items)) + const collection = createCollection( + queryCollectionOptions({ + id: `subset-invalidation-remount-test`, + queryClient, + queryKey: [`subset-invalidation-remount-test`], + queryFn, + getKey, + syncMode: `on-demand`, + }), + ) + const firstLiveQuery = createSubset(collection) + await firstLiveQuery.preload() + const subsetQuery = queryClient.getQueryCache().findAll({ + queryKey: [`subset-invalidation-remount-test`], + })[0] + expect(subsetQuery).toBeDefined() + + items = [{ id: `1`, name: `Invalidated item` }] + await queryClient.invalidateQueries({ + queryKey: subsetQuery!.queryKey, + exact: true, + }) + await vi.waitFor(() => + expect(collection.get(`1`)?.name).toBe(`Invalidated item`), + ) + + await firstLiveQuery.cleanup() + expect(collection.size).toBe(0) + expect(subsetQuery!.getObserversCount()).toBe(0) + + items = [{ id: `1`, name: `Remounted item` }] + const secondLiveQuery = createSubset(collection) + await secondLiveQuery.preload() + await queryClient.invalidateQueries({ + queryKey: subsetQuery!.queryKey, + exact: true, + }) + await vi.waitFor(() => + expect(collection.get(`1`)?.name).toBe(`Remounted item`), + ) + + await secondLiveQuery.cleanup() + expect(collection.size).toBe(0) + expect( + queryClient + .getQueryCache() + .findAll({ + queryKey: [`subset-invalidation-remount-test`], + })[0] + ?.getObserversCount() ?? 0, + ).toBe(0) + await collection.cleanup() + }) + }) + it(`should maintain data consistency during rapid updates`, async () => { const queryKey = [`rapid-updates-test`] let updateCount = 0 From eff4005419ab389f1f6d274782a8bcc1d3cc5638 Mon Sep 17 00:00:00 2001 From: Kyle Mathews Date: Mon, 13 Jul 2026 14:06:36 -0600 Subject: [PATCH 2/4] test: cover query cleanup review findings --- .../query-db-collection/tests/query.test.ts | 68 +++++++++++++++---- 1 file changed, 56 insertions(+), 12 deletions(-) diff --git a/packages/query-db-collection/tests/query.test.ts b/packages/query-db-collection/tests/query.test.ts index e20c9cd98..d72a7a3da 100644 --- a/packages/query-db-collection/tests/query.test.ts +++ b/packages/query-db-collection/tests/query.test.ts @@ -2427,7 +2427,7 @@ describe(`QueryCollection`, () => { ).toBe(0) }) - it(`does not rematerialize a late fulfilled result after an in-flight subset is unloaded`, async () => { + it(`cleans listeners immediately and resolves the unloaded preload before its request settles`, async () => { const deferred = createDeferred>() const collection = createCollection( queryCollectionOptions({ @@ -2440,23 +2440,67 @@ describe(`QueryCollection`, () => { }), ) const liveQuery = createSubset(collection) - void liveQuery.preload().catch(() => undefined) + let preloadResolved = false + void liveQuery.preload().then(() => { + preloadResolved = true + }) await vi.waitFor(() => expect(queryClient.isFetching()).toBe(1)) await liveQuery.cleanup() + + const subsetQuery = queryClient.getQueryCache().findAll({ + queryKey: [`late-subset-result-test`], + })[0] + // This assertion runs while the request is unresolved and directly guards the + // ready-listener bookkeeping bug: unload must synchronously detach its observer. + expect(subsetQuery?.getObserversCount() ?? 0).toBe(0) + // Live-query cleanup resolves its preload even though Query Core is still fetching. + expect(preloadResolved).toBe(true) + deferred.resolve([{ id: `1`, name: `Late item` }]) - await flushPromises() + await vi.waitFor(() => expect(queryClient.isFetching()).toBe(0)) expect(collection.size).toBe(0) - expect(queryClient.isFetching()).toBe(0) - expect( - queryClient - .getQueryCache() - .findAll({ - queryKey: [`late-subset-result-test`], - })[0] - ?.getObserversCount() ?? 0, - ).toBe(0) + expect(preloadResolved).toBe(true) + expect(subsetQuery?.getObserversCount() ?? 0).toBe(0) + await collection.cleanup() + }) + + it(`preserves an active subset across cache removal and accepts a late notification from its detached observer`, async () => { + const queryKey = [`cache-removal-late-notification-test`] + const collection = createCollection( + queryCollectionOptions({ + id: `cache-removal-late-notification-test`, + queryClient, + queryKey, + queryFn: () => Promise.resolve([{ id: `1`, name: `Initial item` }]), + getKey, + syncMode: `on-demand`, + }), + ) + const liveQuery = createSubset(collection) + await liveQuery.preload() + const subsetQuery = queryClient.getQueryCache().findAll({ queryKey })[0] + expect(subsetQuery).toBeDefined() + + // A Query Core `removed` event can arrive before this collection's observer + // is detached. Existing semantics retain the active rows and observer. + queryClient.getQueryCache().remove(subsetQuery!) + expect(queryClient.getQueryCache().findAll({ queryKey })).toHaveLength( + 0, + ) + expect(collection.get(`1`)?.name).toBe(`Initial item`) + expect(subsetQuery!.getObserversCount()).toBe(1) + + // The retained observer can still notify after its query left the cache. + subsetQuery!.setData([{ id: `1`, name: `Late notification` }]) + await vi.waitFor(() => + expect(collection.get(`1`)?.name).toBe(`Late notification`), + ) + + await liveQuery.cleanup() + expect(subsetQuery!.getObserversCount()).toBe(0) + expect(collection.size).toBe(0) await collection.cleanup() }) From 438d96277e23031ebc80f1a8a43a7a7b08cc74d2 Mon Sep 17 00:00:00 2001 From: Kyle Mathews Date: Mon, 13 Jul 2026 15:08:55 -0600 Subject: [PATCH 3/4] Apply query readiness simplification --- packages/query-db-collection/src/query.ts | 77 +++++++++-------------- 1 file changed, 29 insertions(+), 48 deletions(-) diff --git a/packages/query-db-collection/src/query.ts b/packages/query-db-collection/src/query.ts index 18b59d63a..f104cba3b 100644 --- a/packages/query-db-collection/src/query.ts +++ b/packages/query-db-collection/src/query.ts @@ -1175,26 +1175,35 @@ export function queryCollectionOptions( } }) - const trackPendingReadyUnsubscribe = ( + const waitForQueryReady = ( + observer: QueryObserver, any, Array, Array, any>, hashedQueryKey: string, - unsubscribe: () => void, - ) => { - const pending = pendingReadyUnsubscribes.get(hashedQueryKey) ?? new Set() - pending.add(unsubscribe) - pendingReadyUnsubscribes.set(hashedQueryKey, pending) - } + ): Promise => + new Promise((resolve, reject) => { + const unsubscribe = observer.subscribe((result) => { + // Use a microtask in case `subscribe` is called synchronously, before `unsubscribe` is initialized + queueMicrotask(() => { + if (result.isSuccess || result.isError) { + unsubscribe() + const pending = pendingReadyUnsubscribes.get(hashedQueryKey) + pending?.delete(unsubscribe) + if (pending?.size === 0) { + pendingReadyUnsubscribes.delete(hashedQueryKey) + } - const releasePendingReadyUnsubscribe = ( - hashedQueryKey: string, - unsubscribe: () => void, - ) => { - unsubscribe() - const pending = pendingReadyUnsubscribes.get(hashedQueryKey) - pending?.delete(unsubscribe) - if (pending?.size === 0) { - pendingReadyUnsubscribes.delete(hashedQueryKey) - } - } + if (result.isSuccess) { + resolve() + } else { + reject(result.error) + } + } + }) + }) + const pending = + pendingReadyUnsubscribes.get(hashedQueryKey) ?? new Set() + pending.add(unsubscribe) + pendingReadyUnsubscribes.set(hashedQueryKey, pending) + }) const createQueryFromOpts = ( opts: LoadSubsetOptions = {}, @@ -1256,21 +1265,7 @@ export function queryCollectionOptions( } // Query is still loading, wait for the first result - return new Promise((resolve, reject) => { - const unsubscribe = observer.subscribe((result) => { - // Use a microtask in case `subscribe` is called synchronously, before `unsubscribe` is initialized - queueMicrotask(() => { - if (result.isSuccess) { - releasePendingReadyUnsubscribe(hashedQueryKey, unsubscribe) - resolve() - } else if (result.isError) { - releasePendingReadyUnsubscribe(hashedQueryKey, unsubscribe) - reject(result.error) - } - }) - }) - trackPendingReadyUnsubscribe(hashedQueryKey, unsubscribe) - }) + return waitForQueryReady(observer, hashedQueryKey) } } @@ -1340,21 +1335,7 @@ export function queryCollectionOptions( } // Create a promise that resolves when the query result is first available - const readyPromise = new Promise((resolve, reject) => { - const unsubscribe = localObserver.subscribe((result) => { - // Use a microtask in case `subscribe` is called synchronously, before `unsubscribe` is initialized - queueMicrotask(() => { - if (result.isSuccess) { - releasePendingReadyUnsubscribe(hashedQueryKey, unsubscribe) - resolve() - } else if (result.isError) { - releasePendingReadyUnsubscribe(hashedQueryKey, unsubscribe) - reject(result.error) - } - }) - }) - trackPendingReadyUnsubscribe(hashedQueryKey, unsubscribe) - }) + const readyPromise = waitForQueryReady(localObserver, hashedQueryKey) // If sync has started or there are subscribers to the collection, subscribe to the query straight away // This creates the main subscription that handles data updates From 9b08e3c3c4aebadd2b5284dbeac85aabe6c39495 Mon Sep 17 00:00:00 2001 From: Kyle Mathews Date: Mon, 13 Jul 2026 15:09:22 -0600 Subject: [PATCH 4/4] Add query cleanup changeset --- .changeset/characterize-query-cleanup.md | 5 +++++ 1 file changed, 5 insertions(+) create mode 100644 .changeset/characterize-query-cleanup.md diff --git a/.changeset/characterize-query-cleanup.md b/.changeset/characterize-query-cleanup.md new file mode 100644 index 000000000..9f41155f9 --- /dev/null +++ b/.changeset/characterize-query-cleanup.md @@ -0,0 +1,5 @@ +--- +'@tanstack/query-db-collection': patch +--- + +Fix temporary query readiness listeners so subset unload and collection cleanup release them correctly during in-flight requests.