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. diff --git a/packages/query-db-collection/src/query.ts b/packages/query-db-collection/src/query.ts index 7b4578d9e..f104cba3b 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,36 @@ export function queryCollectionOptions( } }) + const waitForQueryReady = ( + observer: QueryObserver, any, Array, Array, any>, + hashedQueryKey: string, + ): 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) + } + + 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 = {}, queryFunction: typeof queryFn = queryFn, @@ -1234,20 +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) { - unsubscribe() - resolve() - } else if (result.isError) { - unsubscribe() - reject(result.error) - } - }) - }) - }) + return waitForQueryReady(observer, hashedQueryKey) } } @@ -1317,20 +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) { - unsubscribe() - resolve() - } else if (result.isError) { - unsubscribe() - reject(result.error) - } - }) - }) - }) + 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 @@ -1590,9 +1595,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 +1668,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..d72a7a3da 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,219 @@ 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(`cleans listeners immediately and resolves the unloaded preload before its request settles`, 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) + 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 vi.waitFor(() => expect(queryClient.isFetching()).toBe(0)) + + expect(collection.size).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() + }) + + 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