diff --git a/.changeset/fix-query-cache-scope-invalidation.md b/.changeset/fix-query-cache-scope-invalidation.md new file mode 100644 index 0000000000..4b9a61ed68 --- /dev/null +++ b/.changeset/fix-query-cache-scope-invalidation.md @@ -0,0 +1,5 @@ +--- +'@tanstack/query-db-collection': patch +--- + +Keep on-demand Query cache ownership and post-write readiness isolated across collections, co-owners, deferred cleanup, errors, and custom query hashes. Active enabled scopes revalidate from post-write provider results, while inactive collection-owned entries are removed without disturbing unrelated or foreign-observed Queries. diff --git a/docs/collections/query-collection.md b/docs/collections/query-collection.md index 25b5f565de..8d2f52869c 100644 --- a/docs/collections/query-collection.md +++ b/docs/collections/query-collection.md @@ -340,9 +340,9 @@ derived Query cache entries instead. If a stale initial response triggers a fetch, the initial rows remain available while it is in flight. A successful response reconciles them through the normal -row ownership pipeline; an error retains the initial rows. Direct writes use the -same Query cache-patching rules as fetched data, and a later successful server -response may reconcile or replace those writes. +row ownership pipeline; an error retains the initial rows. Direct writes patch +the eager Query cache in place. On-demand direct writes revalidate scoped +entries as described below. ### Selecting Rows from Wrapped Responses @@ -382,7 +382,7 @@ preserving the envelope in the Query cache. This differs from TanStack Query's observer-level `select`: query-db-collection uses this option to bridge Query's response object into DB's normalized row store. -Direct write utilities such as `writeInsert`, `writeUpdate`, and `writeDelete` make a best-effort attempt to update the matching row array inside wrapped Query cache entries while preserving wrapper metadata. +In eager mode, direct write utilities such as `writeInsert`, `writeUpdate`, and `writeDelete` make a best-effort attempt to update the matching row array inside wrapped Query cache entries while preserving wrapper metadata. In on-demand mode, they revalidate active scoped queries and remove inactive or disabled cache entries instead of patching them with the full collection snapshot. This works automatically for simple wrappers such as: @@ -599,7 +599,7 @@ The collection provides these utility methods via `collection.utils`: ## Direct Writes -Direct writes are intended for scenarios where the normal query/mutation flow doesn't fit your needs. They allow you to write directly to the synced data store, bypassing the optimistic update system and query refetch mechanism. +Direct writes are intended for scenarios where the normal query/mutation flow doesn't fit your needs. They write directly to the synced data store and bypass the optimistic update system. Their Query cache behavior depends on the collection's sync mode. ### Understanding the Data Stores @@ -615,7 +615,7 @@ Normal collection operations (insert, update, delete) create optimistic mutation - Rolled back automatically if the server request fails - Replaced with server data when the query refetches -Direct writes bypass this system entirely and write directly to the synced data store, making them ideal for handling real-time updates from alternative sources. +Direct writes bypass this system entirely and write directly to the synced data store, making them useful for handling real-time updates from alternative sources. Active on-demand queries still refetch so each scoped cache remains authoritative for its own request. ### When to Use Direct Writes @@ -654,9 +654,9 @@ These operations: - Write directly to the synced data store - Do NOT create optimistic mutations -- Do NOT trigger automatic query refetches -- Update the TanStack Query cache immediately - Are immediately visible in the UI +- In eager mode, update the full-result TanStack Query cache in place without refetching +- In on-demand mode, refetch active enabled queries and remove inactive or disabled cache entries ### Batch Operations @@ -933,13 +933,15 @@ This pattern allows you to: ### Direct Writes and Query Sync -Direct writes update the collection immediately and also update the TanStack Query cache. However, they do not prevent the normal query sync behavior. If your `queryFn` returns data that conflicts with your direct writes, the query data will take precedence. +Direct writes update the collection immediately. In eager mode, they also patch the full-result TanStack Query cache in place. + +In on-demand mode, each Query cache entry may represent a different predicate, order, limit, or offset. A full collection snapshot cannot safely replace those scoped results. Direct writes therefore refetch active enabled queries and remove inactive or disabled entries. A successful `queryFn` result remains authoritative and may reconcile or replace a direct write. To handle this properly: -1. Use `{ refetch: false }` in your persistence handlers when using direct writes -2. Set appropriate `staleTime` to prevent unnecessary refetches -3. Design your `queryFn` to be aware of incremental updates (e.g., only fetch new data) +1. Use `{ refetch: false }` in persistence handlers to avoid the handler's additional refetch after a direct write. On-demand cache revalidation still runs. +2. Make sure an on-demand `queryFn` returns the current server result for its pushed-down predicate, order, limit, and offset. +3. Use eager mode when direct writes must update one complete cached result without a network request. ## Complete Direct Write API Reference diff --git a/docs/reference/query-db-collection/interfaces/QueryCollectionUtils.md b/docs/reference/query-db-collection/interfaces/QueryCollectionUtils.md index 2596da3674..5f2877d303 100644 --- a/docs/reference/query-db-collection/interfaces/QueryCollectionUtils.md +++ b/docs/reference/query-db-collection/interfaces/QueryCollectionUtils.md @@ -3,10 +3,11 @@ id: QueryCollectionUtils title: QueryCollectionUtils --- -Defined in: [packages/query-db-collection/src/query.ts:261](https://github.com/TanStack/db/blob/main/packages/query-db-collection/src/query.ts#L261) +Defined in: [packages/query-db-collection/src/query.ts:262](https://github.com/TanStack/db/blob/main/packages/query-db-collection/src/query.ts#L262) Utility methods available on Query Collections for direct writes and manual operations. -Direct writes bypass the normal query/mutation flow and write directly to the synced data store. +Direct writes bypass optimistic mutations and write to the synced data store. +Eager collections patch Query cache; on-demand collections revalidate scoped entries. ## Extends @@ -185,7 +186,7 @@ writeBatch: (callback) => void; Defined in: [packages/query-db-collection/src/query.ts:278](https://github.com/TanStack/db/blob/main/packages/query-db-collection/src/query.ts#L278) -Execute multiple write operations as a single atomic batch to the synced data store +Execute direct writes as one atomic batch, then update or revalidate the Query cache #### Parameters @@ -207,7 +208,7 @@ writeDelete: (keys) => void; Defined in: [packages/query-db-collection/src/query.ts:274](https://github.com/TanStack/db/blob/main/packages/query-db-collection/src/query.ts#L274) -Delete one or more items directly from the synced data store without triggering a query refetch or optimistic update +Delete items without an optimistic update. On-demand queries revalidate their scoped cache entries. #### Parameters @@ -229,7 +230,7 @@ writeInsert: (data) => void; Defined in: [packages/query-db-collection/src/query.ts:270](https://github.com/TanStack/db/blob/main/packages/query-db-collection/src/query.ts#L270) -Insert one or more items directly into the synced data store without triggering a query refetch or optimistic update +Insert items without an optimistic update. On-demand queries revalidate their scoped cache entries. #### Parameters @@ -251,7 +252,7 @@ writeUpdate: (updates) => void; Defined in: [packages/query-db-collection/src/query.ts:272](https://github.com/TanStack/db/blob/main/packages/query-db-collection/src/query.ts#L272) -Update one or more items directly in the synced data store without triggering a query refetch or optimistic update +Update items without an optimistic update. On-demand queries revalidate their scoped cache entries. #### Parameters @@ -273,7 +274,7 @@ writeUpsert: (data) => void; Defined in: [packages/query-db-collection/src/query.ts:276](https://github.com/TanStack/db/blob/main/packages/query-db-collection/src/query.ts#L276) -Insert or update one or more items directly in the synced data store without triggering a query refetch or optimistic update +Insert or update items without an optimistic update. On-demand queries revalidate their scoped cache entries. #### Parameters diff --git a/packages/query-db-collection/src/manual-sync.ts b/packages/query-db-collection/src/manual-sync.ts index ab61dd5eb6..5ba48dff38 100644 --- a/packages/query-db-collection/src/manual-sync.ts +++ b/packages/query-db-collection/src/manual-sync.ts @@ -52,7 +52,7 @@ export interface SyncContext< * Handles both direct array caches and wrapped response formats (when `select` is used). * If not provided, falls back to directly setting the cache with the raw array. */ - updateCacheData?: (items: Array) => void + updateCacheData?: (getItems: () => Array) => void } interface NormalizedOperation< @@ -221,12 +221,16 @@ export function performWriteOperations< ctx.commit() // Update query cache after successful commit - const updatedData = Array.from(ctx.collection._state.syncedData.values()) if (ctx.updateCacheData) { - ctx.updateCacheData(updatedData) + ctx.updateCacheData(() => + Array.from(ctx.collection._state.syncedData.values()), + ) } else { // Fallback: directly set the cache with raw array (for non-Query Collection consumers) - ctx.queryClient.setQueryData(ctx.queryKey, updatedData) + ctx.queryClient.setQueryData( + ctx.queryKey, + Array.from(ctx.collection._state.syncedData.values()), + ) } } diff --git a/packages/query-db-collection/src/query.ts b/packages/query-db-collection/src/query.ts index cc797f59bb..0bb0020275 100644 --- a/packages/query-db-collection/src/query.ts +++ b/packages/query-db-collection/src/query.ts @@ -1,4 +1,4 @@ -import { QueryObserver, hashKey } from '@tanstack/query-core' +import { QueryObserver, hashKey, partialMatchKey } from '@tanstack/query-core' import { LoadSubsetOperationAbortedError, deepEquals, @@ -28,6 +28,7 @@ import type { } from '@tanstack/db' import type { FetchStatus, + Query, QueryClient, QueryFunctionContext, QueryKey, @@ -252,7 +253,8 @@ export type RefetchFn = (opts?: { /** * Utility methods available on Query Collections for direct writes and manual operations. - * Direct writes bypass the normal query/mutation flow and write directly to the synced data store. + * Direct writes bypass optimistic mutations and write to the synced data store. + * Eager collections patch Query cache; on-demand collections revalidate scoped entries. * @template TItem - The type of items stored in the collection * @template TKey - The type of the item keys * @template TInsertInput - The type accepted for insert operations @@ -266,15 +268,15 @@ export interface QueryCollectionUtils< > extends UtilsRecord { /** Manually trigger a refetch of the query */ refetch: RefetchFn - /** Insert one or more items directly into the synced data store without triggering a query refetch or optimistic update */ + /** Insert items without an optimistic update. On-demand queries revalidate their scoped cache entries. */ writeInsert: (data: TInsertInput | Array) => void - /** Update one or more items directly in the synced data store without triggering a query refetch or optimistic update */ + /** Update items without an optimistic update. On-demand queries revalidate their scoped cache entries. */ writeUpdate: (updates: Partial | Array>) => void - /** Delete one or more items directly from the synced data store without triggering a query refetch or optimistic update */ + /** Delete items without an optimistic update. On-demand queries revalidate their scoped cache entries. */ writeDelete: (keys: TKey | Array) => void - /** Insert or update one or more items directly in the synced data store without triggering a query refetch or optimistic update */ + /** Insert or update items without an optimistic update. On-demand queries revalidate their scoped cache entries. */ writeUpsert: (data: Partial | Array>) => void - /** Execute multiple write operations as a single atomic batch to the synced data store */ + /** Execute direct writes as one atomic batch, then update or revalidate the Query cache */ writeBatch: (callback: () => void) => void // Query Observer State (getters) @@ -332,6 +334,15 @@ type PersistedQueryRetentionEntry = const QUERY_COLLECTION_GC_PREFIX = `queryCollection:gc:` +type AnyQuery = Query + +let nextQueryCollectionFetchStart = 0 +const queryCollectionFetchActionStarts = new WeakMap() +const queryCollectionCurrentFetchStarts = new WeakMap() +const queryCollectionSuccessfulFetchStarts = new WeakMap() +const queryCollectionRequiredFetchStarts = new WeakMap() +const queryCollectionCacheOwners = new WeakMap>() + type PersistedScannedRowForQuery = { key: string | number value: TItem @@ -791,6 +802,83 @@ export function queryCollectionOptions( >(), } + // Query-cache ownership is scoped to this sync generation and keyed by the + // actual Query object. Weak membership survives subset unload without + // retaining entries after Query Core garbage-collects them. + let ownedCacheQueries = new WeakSet() + const cacheOwnerToken = {} + let trackedCacheQueries: Set + const logicalHashesByQuery = new WeakMap>() + + // Manual writes require a successful fetch which started after the write. + // Observe Query Core's fetch/success actions so foreign query functions and + // initialPromise fetches count, while cancelled/reverted requests do not. + const requiredFetchStarts = new Map() + const postWriteRefetchGenerations = new Map() + + const trackCacheQuery = (query: AnyQuery, logicalHash?: string): void => { + trackedCacheQueries.add(query) + if (logicalHash !== undefined) { + const hashes = logicalHashesByQuery.get(query) ?? new Set() + hashes.add(logicalHash) + logicalHashesByQuery.set(query, hashes) + } + } + + const trackOwnedCacheQuery = (query: AnyQuery, logicalHash: string): void => { + ownedCacheQueries.add(query) + trackCacheQuery(query, logicalHash) + const owners = queryCollectionCacheOwners.get(query) ?? new Set() + owners.add(cacheOwnerToken) + queryCollectionCacheOwners.set(query, owners) + } + + const getLogicalHashes = (query: AnyQuery): Set => + logicalHashesByQuery.get(query) ?? new Set([hashKey(query.queryKey)]) + + const hasPostWriteAuthority = ( + hashedQueryKey: string, + query: AnyQuery, + ): boolean => { + const localRequiredStart = requiredFetchStarts.get(hashedQueryKey) + const sharedRequiredStart = queryCollectionRequiredFetchStarts.get(query) + const requiredStart = Math.max( + localRequiredStart ?? 0, + sharedRequiredStart ?? 0, + ) + return ( + (localRequiredStart === undefined && sharedRequiredStart === undefined) || + (queryCollectionSuccessfulFetchStarts.get(query) ?? 0) > requiredStart + ) + } + + const requirePostWriteAuthority = ( + hashedQueryKey: string, + query: AnyQuery, + ): number => { + requiredFetchStarts.set(hashedQueryKey, nextQueryCollectionFetchStart) + queryCollectionRequiredFetchStarts.set( + query, + Math.max( + queryCollectionRequiredFetchStarts.get(query) ?? 0, + nextQueryCollectionFetchStart, + ), + ) + const generation = + (postWriteRefetchGenerations.get(hashedQueryKey) ?? 0) + 1 + postWriteRefetchGenerations.set(hashedQueryKey, generation) + return generation + } + + const isObserverEnabled = ( + observer: QueryObserver, any, Array, Array, any>, + ): boolean => { + const observerEnabled = observer.options.enabled + return typeof observerEnabled === `function` + ? observerEnabled(observer.getCurrentQuery()) !== false + : observerEnabled !== false + } + // hashedQueryKey → queryKey const hashToQueryKey = new Map() @@ -804,6 +892,10 @@ export function queryCollectionOptions( // queryKey → QueryObserver's unsubscribe function const unsubscribes = new Map void>() const pendingReadyUnsubscribes = new Map void>>() + const manualWriteSnapshots = new Map< + string, + { data: unknown; dataUpdateCount: number } + >() // queryKey → reference count (how many loadSubset calls are active) // Reference counting for QueryObserver lifecycle management @@ -886,6 +978,10 @@ export function queryCollectionOptions( } const internalSync: SyncConfig[`sync`] = (params) => { + // Rebuild on every start so caches created while sync was stopped are owned. + trackedCacheQueries = new Set( + queryClient.getQueryCache().findAll({ queryKey: baseKey }), + ) const { begin, write, commit, markReady, markError, collection, metadata } = params const persistedMetadata = metadata as @@ -1280,9 +1376,12 @@ export function queryCollectionOptions( const unsubscribe = observer.subscribe((result) => { // Use a microtask in case `subscribe` is called synchronously, before `unsubscribe` is initialized queueMicrotask(() => { + const query = observer.getCurrentQuery() if ( - (result.isSuccess && !collection.deferDataRefresh) || - result.isError + (result.isSuccess && + hasPostWriteAuthority(hashedQueryKey, query) && + !collection.deferDataRefresh) || + (result.isError && !result.isFetching) ) { unsubscribe() const pending = pendingReadyUnsubscribes.get(hashedQueryKey) @@ -1309,6 +1408,15 @@ export function queryCollectionOptions( pendingReadyUnsubscribes.set(hashedQueryKey, pending) }) + const waitForQueryReadyAndApplied = ( + observer: QueryObserver, any, Array, Array, any>, + hashedQueryKey: string, + ): Promise => + waitForQueryReady(observer, hashedQueryKey).then(() => { + const settlement = getResultApplicationSettlement(hashedQueryKey) + return settlement === true ? undefined : settlement + }) + const createQueryFromOpts = ( opts: LoadSubsetOptions = {}, queryFunction: typeof queryFn = queryFn, @@ -1361,22 +1469,19 @@ export function queryCollectionOptions( const observer = state.observers.get(hashedQueryKey)! const currentResult = observer.getCurrentResult() - if (currentResult.isSuccess) { + if ( + currentResult.isSuccess && + hasPostWriteAuthority(hashedQueryKey, observer.getCurrentQuery()) + ) { if (collection.deferDataRefresh) { - return waitForQueryReady(observer, hashedQueryKey).then(() => { - const settlement = getResultApplicationSettlement(hashedQueryKey) - return settlement === true ? undefined : settlement - }) + return waitForQueryReadyAndApplied(observer, hashedQueryKey) } return getResultApplicationSettlement(hashedQueryKey) - } else if (currentResult.isError) { + } else if (currentResult.isError && !currentResult.isFetching) { // Error already occurred, reject immediately return Promise.reject(currentResult.error) } else { - return waitForQueryReady(observer, hashedQueryKey).then(() => { - const settlement = getResultApplicationSettlement(hashedQueryKey) - return settlement === true ? undefined : settlement - }) + return waitForQueryReadyAndApplied(observer, hashedQueryKey) } } @@ -1414,6 +1519,8 @@ export function queryCollectionOptions( Array, any >(queryClient, observerOptions) + const localQuery = localObserver.getCurrentQuery() + trackOwnedCacheQuery(localQuery, hashedQueryKey) const resolvedQueryGcTime = queryClient.getQueryCache().find({ queryKey: key, exact: true, @@ -1443,17 +1550,30 @@ export function queryCollectionOptions( if (syncStarted || collection.subscriberCount > 0) { subscribeToQuery(localObserver, hashedQueryKey) } + const currentResult = localObserver.getCurrentResult() + if (currentResult.isError && !currentResult.isFetching) { + return Promise.reject(currentResult.error) + } + if ( + !currentResult.isSuccess || + !hasPostWriteAuthority( + hashedQueryKey, + localObserver.getCurrentQuery(), + ) + ) { + return waitForQueryReadyAndApplied(localObserver, hashedQueryKey) + } if (collection.deferDataRefresh) { - return waitForQueryReady(localObserver, hashedQueryKey).then(() => { - const settlement = getResultApplicationSettlement(hashedQueryKey) - return settlement === true ? undefined : settlement - }) + return waitForQueryReadyAndApplied(localObserver, hashedQueryKey) } return getResultApplicationSettlement(hashedQueryKey) } // Create a promise that resolves when the query result is first available - const readyPromise = waitForQueryReady(localObserver, hashedQueryKey) + const readyPromise = waitForQueryReadyAndApplied( + 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 @@ -1461,10 +1581,7 @@ export function queryCollectionOptions( subscribeToQuery(localObserver, hashedQueryKey) } - return readyPromise.then(() => { - const settlement = getResultApplicationSettlement(hashedQueryKey) - return settlement === true ? undefined : settlement - }) + return readyPromise } type UpdateHandler = Parameters[0] @@ -1715,6 +1832,31 @@ export function queryCollectionOptions( const makeQueryResultHandler = (queryKey: QueryKey) => { const hashedQueryKey = hashKey(queryKey) const handleQueryResult: UpdateHandler = (result) => { + const observer = state.observers.get(hashedQueryKey) + if (observer) { + const query = observer.getCurrentQuery() + trackOwnedCacheQuery(query, hashedQueryKey) + if (result.isSuccess) { + if (!hasPostWriteAuthority(hashedQueryKey, query)) { + // Query observers are notified before Query Cache subscribers. + // Recheck after the cache success action records fetch authority. + queueMicrotask(() => { + const currentObserver = state.observers.get(hashedQueryKey) + if ( + currentObserver === observer && + hasPostWriteAuthority( + hashedQueryKey, + currentObserver.getCurrentQuery(), + ) + ) { + handleQueryResult(currentObserver.getCurrentResult()) + } + }) + return + } + requiredFetchStarts.delete(hashedQueryKey) + } + } if (result.isSuccess) { // Error state follows observer notification order, not the later // publication time of a queued successful result. @@ -1736,6 +1878,21 @@ export function queryCollectionOptions( return } + const manualWriteSnapshot = manualWriteSnapshots.get(hashedQueryKey) + if (manualWriteSnapshot) { + const currentQuery = state.observers + .get(hashedQueryKey) + ?.getCurrentQuery() + if ( + currentQuery?.state.dataUpdateCount === + manualWriteSnapshot.dataUpdateCount && + currentQuery.state.data === manualWriteSnapshot.data + ) { + return + } + manualWriteSnapshots.delete(hashedQueryKey) + } + if (retainedQueriesPendingRevalidation.has(hashedQueryKey)) { const query = queryClient.getQueryCache().find({ queryKey, @@ -1770,7 +1927,14 @@ export function queryCollectionOptions( applySuccessfulResult(queryKey, result, undefined, signal), ) } - } else if (result.isError) { + } else { + // A reset/recreation can reuse a dataUpdateCount. Retire the old + // snapshot marker on the intervening non-success notification so a + // fresh authoritative result cannot be mistaken for stale data. + manualWriteSnapshots.delete(hashedQueryKey) + } + + if (result.isError) { const isNewError = result.errorUpdatedAt !== state.lastErrorUpdatedAt || result.error !== state.lastError @@ -1909,6 +2073,7 @@ export function queryCollectionOptions( cancelPersistedRetentionExpiry(hashedQueryKey) retainedQueriesPendingRevalidation.delete(hashedQueryKey) invalidatePendingResultApplication(hashedQueryKey) + manualWriteSnapshots.delete(hashedQueryKey) const nextOwnersByRow = removeQueryOwnership(hashedQueryKey) const rowsToDelete: Array = [] @@ -2001,6 +2166,7 @@ export function queryCollectionOptions( persistedMetadata?.row.scanPersisted ) { invalidatePendingResultApplication(hashedQueryKey) + manualWriteSnapshots.delete(hashedQueryKey) begin() metadata.collection.set( `${QUERY_COLLECTION_GC_PREFIX}${hashedQueryKey}`, @@ -2046,10 +2212,36 @@ export function queryCollectionOptions( const unsubscribeQueryCache = queryClient .getQueryCache() .subscribe((event) => { + if ( + event.type === `added` && + partialMatchKey(event.query.queryKey, baseKey) + ) { + trackCacheQuery(event.query) + } + + if (event.type === `updated`) { + if (event.action.type === `fetch`) { + let fetchStart = queryCollectionFetchActionStarts.get(event.action) + if (fetchStart === undefined) { + fetchStart = ++nextQueryCollectionFetchStart + queryCollectionFetchActionStarts.set(event.action, fetchStart) + } + queryCollectionCurrentFetchStarts.set(event.query, fetchStart) + } else if (event.action.type === `success` && !event.action.manual) { + const fetchStart = queryCollectionCurrentFetchStarts.get( + event.query, + ) + if (fetchStart !== undefined) { + queryCollectionSuccessfulFetchStarts.set(event.query, fetchStart) + } + } + } + // Ownership uses our stable key, not the Query client's optional // custom cache hash function. const hashedKey = hashKey(event.query.queryKey) if (event.type === `removed`) { + trackedCacheQueries.delete(event.query) // Only cleanup if this is OUR query (we track it) if (hashToQueryKey.has(hashedKey)) { if (syncMode === `eager`) { @@ -2092,9 +2284,23 @@ export function queryCollectionOptions( // Removing a Query destroys it and synchronously cancels its retryer. // Finish this before a later collection sync can create a replacement. - queryClient.removeQueries({ - predicate: (query) => allHashedKeys.has(hashKey(query.queryKey)), - }) + for (const query of [...trackedCacheQueries]) { + const belongsToCleanup = [...getLogicalHashes(query)].some((hash) => + allHashedKeys.has(hash), + ) + if ( + belongsToCleanup && + (syncMode === `eager` || + (ownedCacheQueries.has(query) && query.getObserversCount() === 0)) + ) { + queryClient.getQueryCache().remove(query) + } + } + for (const query of trackedCacheQueries) { + queryCollectionCacheOwners.get(query)?.delete(cacheOwnerToken) + } + trackedCacheQueries.clear() + ownedCacheQueries = new WeakSet() } /** @@ -2254,20 +2460,173 @@ export function queryCollectionOptions( } /** - * Updates the query cache with new items for ALL query keys matching this collection, - * including stale/inactive cache entries from destroyed observers. + * Updates Query cache state after a manual write. * - * This prevents ghost items: when an observer is destroyed but gcTime > 0, TanStack Query - * keeps the cached data. If syncedData changes (via writeDelete/writeInsert/writeUpdate) - * after the observer is destroyed, the stale cache becomes inconsistent. When a new observer - * later picks up this stale cache, makeQueryResultHandler would create spurious sync - * operations (re-inserting deleted items, reverting updated values, etc). - * - * By updating all cache entries (active and stale), we ensure the cache always reflects - * the current syncedData state. + * An on-demand cache entry belongs to one exact queryFn result and may be + * predicate-, order-, or window-scoped. A normalized collection snapshot + * cannot preserve that shape. Revalidate actively owned, enabled observers + * and remove every other scoped entry so a later owner fetches it again. + * Eager collections retain their single full-result cache patch. */ - const updateCacheData = (items: Array): void => { - const allCached = queryClient.getQueryCache().findAll({ queryKey: baseKey }) + const updateCacheData = (getItems: () => Array): void => { + if (syncMode === `on-demand`) { + const deferredRefresh = writeContext?.collection.deferDataRefresh + const revalidatingQueries = new Set() + const ownObserverCounts = new Map() + + for (const observer of state.observers.values()) { + const query = observer.getCurrentQuery() + ownObserverCounts.set(query, (ownObserverCounts.get(query) ?? 0) + 1) + } + + const refetchTrackedQuery = async ( + query: AnyQuery, + logicalHashes: Set, + generations: Map, + ): Promise => { + try { + await query.fetch(undefined, { cancelRefetch: false }) + } catch { + // A failed post-write refetch is terminal for this attempt. + return + } + + const stillNeedsAuthority = [...logicalHashes].some( + (hashedQueryKey) => + postWriteRefetchGenerations.get(hashedQueryKey) === + generations.get(hashedQueryKey) && + !hasPostWriteAuthority(hashedQueryKey, query), + ) + if ( + stillNeedsAuthority && + queryClient.getQueryCache().get(query.queryHash) === query && + !query.isDisabled() + ) { + try { + // The first call may have reused a request which began before the + // write. One bounded follow-up then establishes authority. + await query.fetch(undefined, { cancelRefetch: false }) + } catch { + // The Query result owns error publication. + } + } + } + + for (const [hashedQueryKey, observer] of state.observers) { + if ((queryRefCounts.get(hashedQueryKey) ?? 0) <= 0) { + continue + } + + const query = observer.getCurrentQuery() + if (!isObserverEnabled(observer)) { + continue + } + + const generation = requirePostWriteAuthority(hashedQueryKey, query) + revalidatingQueries.add(query) + const ownedAtSchedule = ownedCacheQueries.has(query) + query.invalidate() + manualWriteSnapshots.set(hashedQueryKey, { + data: query.state.data, + dataUpdateCount: query.state.dataUpdateCount, + }) + + const retireProtectedQuery = () => { + if (!ownedAtSchedule) return + if (queryClient.getQueryCache().get(query.queryHash) !== query) return + if (query.getObserversCount() > 0) { + query.invalidate() + if (!query.isDisabled()) { + const logicalHashes = new Set([hashedQueryKey]) + const generations = new Map([[hashedQueryKey, generation]]) + void refetchTrackedQuery(query, logicalHashes, generations) + } + } else { + queryClient.getQueryCache().remove(query) + } + } + + const refetchObserver = async () => { + if ( + state.observers.get(hashedQueryKey) !== observer || + (queryRefCounts.get(hashedQueryKey) ?? 0) <= 0 + ) { + retireProtectedQuery() + return + } + + if ( + hasPostWriteAuthority(hashedQueryKey, observer.getCurrentQuery()) + ) { + return + } + + const result = await observer.refetch().catch(() => undefined) + if ( + result?.isError || + postWriteRefetchGenerations.get(hashedQueryKey) !== generation + ) { + return + } + + const currentQuery = observer.getCurrentQuery() + if (hasPostWriteAuthority(hashedQueryKey, currentQuery)) return + if ( + state.observers.get(hashedQueryKey) === observer && + (queryRefCounts.get(hashedQueryKey) ?? 0) > 0 && + isObserverEnabled(observer) + ) { + // A cancellation/revert or a reused pre-write request can resolve + // without authority. Retry once; errors end the attempt. + await observer.refetch().catch(() => undefined) + } else { + retireProtectedQuery() + } + } + + if (deferredRefresh) { + void deferredRefresh.then(refetchObserver, refetchObserver) + } else { + void refetchObserver() + } + } + + for (const query of [...trackedCacheQueries]) { + if (revalidatingQueries.has(query)) continue + + const ownedByAnotherCollection = [ + ...(queryCollectionCacheOwners.get(query) ?? []), + ].some((owner) => owner !== cacheOwnerToken) + if (ownedByAnotherCollection) continue + + const ownObservers = ownObserverCounts.get(query) ?? 0 + if (query.getObserversCount() > ownObservers) { + // A disabled collection observer must not authorize a foreign + // observer to refetch on the collection's behalf. + if (ownObservers > 0) continue + + const logicalHashes = getLogicalHashes(query) + const generations = new Map() + for (const hashedQueryKey of logicalHashes) { + generations.set( + hashedQueryKey, + requirePostWriteAuthority(hashedQueryKey, query), + ) + } + query.invalidate() + if (!query.isDisabled()) { + void refetchTrackedQuery(query, logicalHashes, generations) + } + continue + } + + queryClient.getQueryCache().remove(query) + } + return + } + + const items = getItems() + const allCached = [...trackedCacheQueries] if (allCached.length > 0) { for (const query of allCached) { @@ -2289,7 +2648,7 @@ export function queryCollectionOptions( begin: () => void write: (message: Omit, `key`>) => void commit: () => SyncAppliedReceipt - updateCacheData?: (items: Array) => void + updateCacheData?: (getItems: () => Array) => void } | null = null // Enhanced internalSync that captures write functions for manual use diff --git a/packages/query-db-collection/tests/ownership-lifecycle.oracle.test.ts b/packages/query-db-collection/tests/ownership-lifecycle.oracle.test.ts index 0ef3821b44..e2f325c736 100644 --- a/packages/query-db-collection/tests/ownership-lifecycle.oracle.test.ts +++ b/packages/query-db-collection/tests/ownership-lifecycle.oracle.test.ts @@ -1,5 +1,16 @@ -import { QueryClient, hashKey, isCancelledError } from '@tanstack/query-core' -import { IR, createCollection, eq, getLoadSubsetDemandKey } from '@tanstack/db' +import { + QueryClient, + QueryObserver, + hashKey, + isCancelledError, +} from '@tanstack/query-core' +import { + IR, + createCollection, + createLiveQueryCollection, + eq, + getLoadSubsetDemandKey, +} from '@tanstack/db' import { afterEach, describe, expect, it, vi } from 'vitest' import { createDeferred } from '../../db/src/deferred.js' import { persistedCollectionOptions } from '../../db-sqlite-persistence-core/src/index.js' @@ -799,6 +810,129 @@ describe(`query collection ownership lifecycle`, () => { expect(rows(collection)).toEqual([]) }) + it(`keeps an inactive scoped cache isolated from another scope's manual write`, async () => { + const id = `inactive-scoped-cache-isolation` + const queryClient = createQueryClient() + const sourceRows = new Map([ + [`1`, { id: `1`, category: `A`, name: `Category A` }], + [`2`, { id: `2`, category: `B`, name: `Category B` }], + ]) + const expectedByCategory = new Map>() + const providerCalls: Array = [] + let manualWrites = 0 + let remountReached = false + let comparisonCount = 0 + + const recomputeCategory = (category: string): Array => + Array.from(sourceRows.values()) + .filter((row) => row.category === category) + .map((row) => row.id) + .sort() + + expectedByCategory.set(`A`, recomputeCategory(`A`)) + expectedByCategory.set(`B`, recomputeCategory(`B`)) + + const queryFn = vi.fn((context: QueryFunctionContext) => { + const where = context.meta?.loadSubsetOptions?.where + const operands = where?.type === `func` ? where.args : [] + const categoryRef = operands.find( + (operand) => + operand.type === `ref` && + operand.path.length === 1 && + operand.path[0] === `category`, + ) + const categoryValue = operands.find( + (operand) => + operand.type === `val` && typeof operand.value === `string`, + ) + if ( + where?.type !== `func` || + where.name !== `eq` || + operands.length !== 2 || + categoryRef?.type !== `ref` || + categoryValue?.type !== `val` || + typeof categoryValue.value !== `string` + ) { + throw new Error(`Category fixture received an unsupported request`) + } + + const category = categoryValue.value + providerCalls.push(category) + return Promise.resolve( + Array.from(sourceRows.values(), (row) => structuredClone(row)).filter( + (row) => row.category === category, + ), + ) + }) + const collection = createCollection( + queryCollectionOptions({ + id, + queryClient, + queryKey: [id], + queryFn, + getKey: (item) => item.id, + syncMode: `on-demand`, + startSync: true, + }), + ) + cleanups.push(async () => { + await collection.cleanup() + queryClient.clear() + }) + const createCategoryQuery = (category: string) => + createLiveQueryCollection({ + query: (query) => + query + .from({ item: collection }) + .where(({ item }) => eq(item.category, category)), + }) + + type CategoryQuery = ReturnType + let inactiveCategoryA: CategoryQuery | undefined + let activeCategoryB: CategoryQuery | undefined + let remountedCategoryA: CategoryQuery | undefined + + try { + inactiveCategoryA = createCategoryQuery(`A`) + await inactiveCategoryA.preload() + expect(rows(inactiveCategoryA)).toEqual(expectedByCategory.get(`A`)) + await inactiveCategoryA.cleanup() + inactiveCategoryA = undefined + + activeCategoryB = createCategoryQuery(`B`) + await activeCategoryB.preload() + expect(rows(activeCategoryB)).toEqual(expectedByCategory.get(`B`)) + + const added = { id: `3`, category: `B`, name: `New category B` } + sourceRows.set(added.id, structuredClone(added)) + expectedByCategory.set(`B`, recomputeCategory(`B`)) + manualWrites++ + collection.utils.writeUpsert(structuredClone(added)) + const cacheEntriesAfterManualWrite = queryClient + .getQueryCache() + .findAll({ queryKey: [id] }).length + + remountedCategoryA = createCategoryQuery(`A`) + await remountedCategoryA.preload() + remountReached = true + + expect(providerCalls.slice(0, 2)).toEqual([`A`, `B`]) + expect(manualWrites).toBe(1) + expect(remountReached).toBe(true) + comparisonCount++ + expect(rows(remountedCategoryA)).toEqual(expectedByCategory.get(`A`)) + expect(comparisonCount).toBe(1) + expect(cacheEntriesAfterManualWrite).toBe(1) + expect(providerCalls.filter((category) => category === `A`)).toHaveLength( + 2, + ) + } finally { + await remountedCategoryA?.cleanup() + await activeCategoryB?.cleanup() + await inactiveCategoryA?.cleanup() + } + }) + it(`emits metadata for every owner of rows shared by overlapping queries`, async () => { const metadata: MetadataRecorder = { rows: new Map(), writes: [] } const { collection } = createOwnershipFixture({ @@ -865,4 +999,286 @@ describe(`query collection ownership lifecycle`, () => { expect(setupCalls).toBe(1) expect(persistedOwners(metadata.rows, shared.id)).toEqual([queryHash]) }) + + it(`stops after one failed post-write refetch`, async () => { + const failedRefetch = createDeferred>() + const consoleError = vi.spyOn(console, `error`).mockImplementation(() => {}) + const { collection, queryFn } = createOwnershipFixture({ + id: `failed-post-write-refetch`, + results: [[shared], failedRefetch.promise], + }) + + try { + await collection._sync.loadSubset({}) + collection.utils.writeUpdate({ ...shared, name: `Manual` }) + await vi.waitFor(() => expect(queryFn).toHaveBeenCalledTimes(2)) + + failedRefetch.reject(new Error(`Controlled refetch failure`)) + for (let turn = 0; turn < 30; turn++) await Promise.resolve() + + expect(queryFn).toHaveBeenCalledTimes(2) + expect(consoleError).toHaveBeenCalledTimes(1) + } finally { + consoleError.mockRestore() + } + }) + + it(`keeps overlapping post-write refetches bounded`, async () => { + const id = `overlapping-post-write-refetches` + const queryClient = createQueryClient() + const firstWriteResult = createDeferred>() + const secondWriteResult = createDeferred>() + const firstWriteStarted = createDeferred() + const secondWriteStarted = createDeferred() + let starts = 0 + let aborts = 0 + const queryFn = vi.fn((context: QueryFunctionContext) => { + starts++ + if (starts === 1) return Promise.resolve([shared]) + context.signal.addEventListener(`abort`, () => aborts++, { once: true }) + if (starts === 2) { + firstWriteStarted.resolve() + return firstWriteResult.promise + } + secondWriteStarted.resolve() + return secondWriteResult.promise + }) + const collection = createCollection( + queryCollectionOptions({ + id, + queryClient, + queryKey: [id], + queryFn, + getKey: (item) => item.id, + syncMode: `on-demand`, + startSync: true, + }), + ) + cleanups.push(async () => { + firstWriteResult.resolve([]) + secondWriteResult.resolve([]) + await collection.cleanup() + queryClient.clear() + }) + + await collection._sync.loadSubset({}) + collection.utils.writeUpdate({ ...shared, name: `First write` }) + await firstWriteStarted.promise + collection.utils.writeUpdate({ ...shared, name: `Second write` }) + await secondWriteStarted.promise + for (let turn = 0; turn < 30; turn++) await Promise.resolve() + + expect({ starts, aborts }).toEqual({ starts: 3, aborts: 1 }) + secondWriteResult.resolve([{ ...shared, name: `Authoritative` }]) + await vi.waitFor(() => { + expect(collection.get(shared.id)?.name).toBe(`Authoritative`) + }) + }) + + it(`accepts a successful foreign fetch as post-write authority`, async () => { + const id = `foreign-post-write-authority` + const queryClient = createQueryClient() + const siblingOptions = categorySubset(`detail`) + const siblingKey = [id, getLoadSubsetDemandKey(siblingOptions)] + queryClient.setQueryData(siblingKey, [detailOnly]) + + const foreignResult = { ...detailOnly, name: `Foreign authority` } + const foreignQueryFn = vi.fn(() => Promise.resolve([foreignResult])) + const foreignObserver = new QueryObserver(queryClient, { + queryKey: siblingKey, + queryFn: foreignQueryFn, + staleTime: Number.POSITIVE_INFINITY, + retry: false, + }) + const unsubscribeForeign = foreignObserver.subscribe(() => {}) + const queryFn = vi.fn(() => Promise.resolve([shared])) + const collection = createCollection( + queryCollectionOptions({ + id, + queryClient, + queryKey: [id], + queryFn, + getKey: (item) => item.id, + syncMode: `on-demand`, + startSync: true, + }), + ) + cleanups.push(async () => { + unsubscribeForeign() + await collection.cleanup() + queryClient.clear() + }) + + await collection._sync.loadSubset({}) + collection.utils.writeUpdate({ ...shared, name: `Manual` }) + await vi.waitFor(() => expect(foreignQueryFn).toHaveBeenCalledTimes(1)) + + let settled = false + const load = collection._sync.loadSubset(siblingOptions) + void Promise.resolve(load === true ? undefined : load).then( + () => { + settled = true + }, + () => { + settled = true + }, + ) + for (let turn = 0; turn < 30; turn++) await Promise.resolve() + + expect(settled).toBe(true) + expect(foreignObserver.getCurrentResult().data).toEqual([foreignResult]) + }) + + it(`uses the replacement Query when the cache is cleared before refetch`, async () => { + const barrier = createDeferred() + const thirdFetch = createDeferred>() + const authoritative = { ...shared, name: `Authoritative` } + const { collection, queryClient, queryFn } = createOwnershipFixture({ + id: `clear-before-post-write-refetch`, + results: [[shared], [authoritative], thirdFetch.promise], + }) + const barrierCompletion = barrier.promise.then(() => { + collection.deferDataRefresh = null + }) + + try { + await collection._sync.loadSubset({}) + collection.deferDataRefresh = barrier.promise + collection.utils.writeUpdate({ ...shared, name: `Manual` }) + queryClient.clear() + + barrier.resolve() + await barrierCompletion + await vi.waitFor(() => expect(queryFn).toHaveBeenCalledTimes(2)) + for (let turn = 0; turn < 30; turn++) await Promise.resolve() + + expect(queryFn).toHaveBeenCalledTimes(2) + expect(collection.get(shared.id)?.name).toBe(`Authoritative`) + } finally { + barrier.resolve() + collection.deferDataRefresh = null + thirdFetch.resolve([]) + } + }) + + it(`requires another fetch after cancellation reverts a raw result`, async () => { + const id = `cancelled-post-write-refetch` + const queryClient = createQueryClient() + const cancelledResult = createDeferred>() + const cancelledStarted = createDeferred() + const authoritativeStarted = createDeferred() + const authoritative = { ...shared, name: `Authoritative` } + let call = 0 + const queryFn = vi.fn((context: QueryFunctionContext) => { + call++ + if (call === 1) return Promise.resolve([shared]) + if (call === 2) { + void context.signal.aborted + cancelledStarted.resolve() + return cancelledResult.promise + } + authoritativeStarted.resolve() + return Promise.resolve([authoritative]) + }) + const collection = createCollection( + queryCollectionOptions({ + id, + queryClient, + queryKey: [id], + queryFn, + getKey: (item) => item.id, + syncMode: `on-demand`, + startSync: true, + }), + ) + cleanups.push(async () => { + cancelledResult.resolve([]) + await collection.cleanup() + queryClient.clear() + }) + + await collection._sync.loadSubset({}) + collection.utils.writeUpdate({ ...shared, name: `Manual` }) + await cancelledStarted.promise + + const query = queryClient.getQueryCache().find({ + queryKey: [id], + exact: true, + })! + await query.cancel({ revert: true }) + cancelledResult.resolve([{ ...shared, name: `Cancelled raw result` }]) + + await authoritativeStarted.promise + await vi.waitFor(() => { + expect(queryFn).toHaveBeenCalledTimes(3) + expect(collection.get(shared.id)?.name).toBe(`Authoritative`) + }) + }) + + it(`removes unobserved sibling scopes without scanning or snapshots`, async () => { + const id = `unobserved-sibling-manual-write` + const queryClient = createQueryClient() + const siblingKey = [id, `prefetched-sibling`] + queryClient.setQueryData(siblingKey, [shared]) + const queryFn = vi.fn(() => Promise.resolve([shared])) + const collection = createCollection( + queryCollectionOptions({ + id, + queryClient, + queryKey: [id], + queryFn, + getKey: (item) => item.id, + syncMode: `on-demand`, + startSync: true, + }), + ) + cleanups.push(async () => { + await collection.cleanup() + queryClient.clear() + }) + + await collection._sync.loadSubset({}) + const findAll = vi.spyOn(queryClient.getQueryCache(), `findAll`) + const values = vi.spyOn(collection._state.syncedData, `values`) + + collection.utils.writeDelete(shared.id) + + expect(queryClient.getQueryData(siblingKey)).toBeUndefined() + expect(findAll).not.toHaveBeenCalled() + expect(values).not.toHaveBeenCalled() + }) }) + +it.each([false, true])( + `evicts sibling caches created before sync starts, restart=%s`, + async (restart) => { + const queryClient = createQueryClient() + const id = `delayed-sync-cache` + const collection = createCollection( + queryCollectionOptions({ + id, + queryClient, + queryKey: [id], + queryFn: () => Promise.resolve([shared]), + getKey: (item) => item.id, + syncMode: `on-demand`, + startSync: false, + }), + ) + cleanups.push(async () => { + await collection.cleanup() + queryClient.clear() + }) + if (restart) { + collection.startSyncImmediate() + await collection._sync.loadSubset({}) + await collection.cleanup() + } + const siblingKey = [id, `prefetched-sibling`] + queryClient.setQueryData(siblingKey, [shared]) + collection.startSyncImmediate() + await collection._sync.loadSubset({}) + collection.utils.writeDelete(shared.id) + expect(queryClient.getQueryData(siblingKey)).toBeUndefined() + }, +) diff --git a/packages/query-db-collection/tests/query.test.ts b/packages/query-db-collection/tests/query.test.ts index 11b3f8db60..60b1086e9b 100644 --- a/packages/query-db-collection/tests/query.test.ts +++ b/packages/query-db-collection/tests/query.test.ts @@ -1,6 +1,7 @@ import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' import { QueryClient, + QueryObserver, dehydrate, focusManager, hashKey, @@ -8232,17 +8233,19 @@ describe(`QueryCollection`, () => { }) }) - describe(`On-demand collection directWrite cache update`, () => { - it(`should update query cache for all active query keys when using writeUpdate with computed queryKey`, async () => { - // Ensures writeUpdate on on-demand collections with computed query keys - // updates all active cache keys to prevent data loss on remount + describe(`On-demand collection directWrite cache revalidation`, () => { + it(`should revalidate an active computed queryKey after writeUpdate`, async () => { + // Ensures writeUpdate on on-demand collections revalidates the active + // computed query key so the authoritative result survives a remount. - const items: Array = [ + const serverItems: Array = [ { id: `1`, name: `Item 1`, category: `A` }, { id: `2`, name: `Item 2`, category: `A` }, ] - const queryFn = vi.fn().mockResolvedValue(items) + const queryFn = vi.fn(() => + Promise.resolve(serverItems.map((item) => structuredClone(item))), + ) // Use a custom queryClient with longer gcTime to prevent cache from being removed const customQueryClient = new QueryClient({ @@ -8282,61 +8285,67 @@ describe(`QueryCollection`, () => { .where(({ item }) => eq(item.category, `A`)) .select(({ item }) => ({ id: item.id, name: item.name })), }) + let query2: typeof query1 | undefined - await query1.preload() + try { + await query1.preload() - // Wait for data to load - await vi.waitFor(() => { - expect(collection.size).toBe(2) - }) + // Wait for data to load + await vi.waitFor(() => { + expect(collection.size).toBe(2) + }) - // Perform a direct write update - collection.utils.writeUpdate({ id: `1`, name: `Updated Item 1` }) + serverItems[0] = { ...serverItems[0]!, name: `Updated Item 1` } + collection.utils.writeUpdate({ id: `1`, name: `Updated Item 1` }) - // Verify the collection reflects the update - expect(collection.get(`1`)?.name).toBe(`Updated Item 1`) + // Verify the collection reflects the update + expect(collection.get(`1`)?.name).toBe(`Updated Item 1`) - // IMPORTANT: Simulate remount by cleaning up and recreating the live query - // This is where the bug manifests - the updated data should persist - await query1.cleanup() - await flushPromises() + await vi.waitFor(() => expect(queryFn).toHaveBeenCalledTimes(2)) - // Recreate the same live query (simulating component remount) - const query2 = createLiveQueryCollection({ - query: (q) => - q - .from({ item: collection }) - .where(({ item }) => eq(item.category, `A`)) - .select(({ item }) => ({ id: item.id, name: item.name })), - }) + // IMPORTANT: Simulate remount by cleaning up and recreating the live query + // This is where the bug manifests - the updated data should persist + await query1.cleanup() + await flushPromises() - await query2.preload() + // Recreate the same live query (simulating component remount) + query2 = createLiveQueryCollection({ + query: (q) => + q + .from({ item: collection }) + .where(({ item }) => eq(item.category, `A`)) + .select(({ item }) => ({ id: item.id, name: item.name })), + }) - // Wait for data to be available - await vi.waitFor(() => { - expect(collection.size).toBe(2) - }) + await query2.preload() - // BUG ASSERTION: After remount, the updated data should persist - // With the bug, this will fail because writeUpdate updated the wrong cache key - // and on remount, the stale cached data is loaded instead - expect(collection.get(`1`)?.name).toBe(`Updated Item 1`) + // Wait for data to be available + await vi.waitFor(() => { + expect(collection.size).toBe(2) + }) - // Cleanup - await query2.cleanup() - customQueryClient.clear() + // After remount, the authoritative updated data should persist. + expect(collection.get(`1`)?.name).toBe(`Updated Item 1`) + } finally { + await query2?.cleanup() + await query1.cleanup() + await collection.cleanup() + customQueryClient.clear() + } }) - it(`should update query cache for static queryKey with where clause in on-demand mode`, async () => { + it(`should revalidate a scoped static queryKey after writeUpdate`, async () => { // Scenario: static queryKey + on-demand mode + where clause // The where clause causes a computed query key to be generated - const items: Array = [ + const serverItems: Array = [ { id: `1`, name: `Item 1`, category: `A` }, { id: `2`, name: `Item 2`, category: `A` }, ] - const queryFn = vi.fn().mockResolvedValue(items) + const queryFn = vi.fn(() => + Promise.resolve(serverItems.map((item) => structuredClone(item))), + ) const customQueryClient = new QueryClient({ defaultOptions: { @@ -8369,53 +8378,62 @@ describe(`QueryCollection`, () => { .where(({ item }) => eq(item.category, `A`)) .select(({ item }) => ({ id: item.id, name: item.name })), }) + let query2: typeof query1 | undefined - await query1.preload() + try { + await query1.preload() - await vi.waitFor(() => { - expect(collection.size).toBe(2) - }) + await vi.waitFor(() => { + expect(collection.size).toBe(2) + }) - // Perform a direct write update - collection.utils.writeUpdate({ id: `1`, name: `Updated Item 1` }) + serverItems[0] = { ...serverItems[0]!, name: `Updated Item 1` } + collection.utils.writeUpdate({ id: `1`, name: `Updated Item 1` }) - expect(collection.get(`1`)?.name).toBe(`Updated Item 1`) + expect(collection.get(`1`)?.name).toBe(`Updated Item 1`) - // Simulate remount - await query1.cleanup() - await flushPromises() + await vi.waitFor(() => expect(queryFn).toHaveBeenCalledTimes(2)) - const query2 = createLiveQueryCollection({ - query: (q) => - q - .from({ item: collection }) - .where(({ item }) => eq(item.category, `A`)) - .select(({ item }) => ({ id: item.id, name: item.name })), - }) + // Simulate remount + await query1.cleanup() + await flushPromises() - await query2.preload() + query2 = createLiveQueryCollection({ + query: (q) => + q + .from({ item: collection }) + .where(({ item }) => eq(item.category, `A`)) + .select(({ item }) => ({ id: item.id, name: item.name })), + }) - await vi.waitFor(() => { - expect(collection.size).toBe(2) - }) + await query2.preload() - // After remount, the updated data should persist - expect(collection.get(`1`)?.name).toBe(`Updated Item 1`) + await vi.waitFor(() => { + expect(collection.size).toBe(2) + }) - await query2.cleanup() - customQueryClient.clear() + // After remount, the updated data should persist + expect(collection.get(`1`)?.name).toBe(`Updated Item 1`) + } finally { + await query2?.cleanup() + await query1.cleanup() + await collection.cleanup() + customQueryClient.clear() + } }) - it(`should update query cache for function queryKey that returns constant value in on-demand mode`, async () => { + it(`should revalidate a constant function queryKey after writeUpdate`, async () => { // Scenario: function queryKey that returns same value // This creates an undefined entry in the cache - const items: Array = [ + const serverItems: Array = [ { id: `1`, name: `Item 1` }, { id: `2`, name: `Item 2` }, ] - const queryFn = vi.fn().mockResolvedValue(items) + const queryFn = vi.fn(() => + Promise.resolve(serverItems.map((item) => structuredClone(item))), + ) const customQueryClient = new QueryClient({ defaultOptions: { @@ -8443,38 +8461,1497 @@ describe(`QueryCollection`, () => { const query1 = createLiveQueryCollection({ query: (q) => q.from({ item: collection }).select(({ item }) => item), }) + let query2: typeof query1 | undefined - await query1.preload() + try { + await query1.preload() - await vi.waitFor(() => { - expect(collection.size).toBe(2) - }) + await vi.waitFor(() => { + expect(collection.size).toBe(2) + }) - // Perform a direct write update - collection.utils.writeUpdate({ id: `1`, name: `Updated Item 1` }) + serverItems[0] = { ...serverItems[0]!, name: `Updated Item 1` } + collection.utils.writeUpdate({ id: `1`, name: `Updated Item 1` }) - expect(collection.get(`1`)?.name).toBe(`Updated Item 1`) + expect(collection.get(`1`)?.name).toBe(`Updated Item 1`) - // Simulate remount - await query1.cleanup() - await flushPromises() + await vi.waitFor(() => expect(queryFn).toHaveBeenCalledTimes(2)) - const query2 = createLiveQueryCollection({ - query: (q) => q.from({ item: collection }).select(({ item }) => item), + // Simulate remount + await query1.cleanup() + await flushPromises() + + query2 = createLiveQueryCollection({ + query: (q) => q.from({ item: collection }).select(({ item }) => item), + }) + + await query2.preload() + + await vi.waitFor(() => { + expect(collection.size).toBe(2) + }) + + // After remount, the updated data should persist + expect(collection.get(`1`)?.name).toBe(`Updated Item 1`) + } finally { + await query2?.cleanup() + await query1.cleanup() + await collection.cleanup() + customQueryClient.clear() + } + }) + + it.each([false, true])( + `keeps a manual update visible until an active scoped refetch settles with custom hash %s`, + async (customHash) => { + const initial = { id: `1`, name: `Initial`, category: `A` } + const authoritative = { + id: `1`, + name: `Authoritative`, + category: `A`, + } + const refetchResult = createDeferred>() + const customQueryClient = new QueryClient({ + defaultOptions: { + queries: { + gcTime: Number.POSITIVE_INFINITY, + staleTime: Number.POSITIVE_INFINITY, + retry: false, + queryKeyHashFn: customHash + ? (key) => `custom:${hashKey(key)}` + : undefined, + }, + }, + }) + const queryFn = vi + .fn<() => Promise>>() + .mockResolvedValueOnce([initial]) + .mockImplementationOnce(() => refetchResult.promise) + const collection = createCollection( + queryCollectionOptions({ + id: `active-scoped-manual-write-${customHash}`, + queryClient: customQueryClient, + queryKey: [`active-scoped-manual-write`], + queryFn, + getKey: (item) => item.id, + syncMode: `on-demand`, + startSync: true, + }), + ) + const active = createLiveQueryCollection({ + query: (query) => + query + .from({ item: collection }) + .where(({ item }) => eq(item.category, `A`)), + }) + + try { + await active.preload() + const scopedQuery = customQueryClient.getQueryCache().getAll()[0]! + expect(scopedQuery.state.dataUpdateCount).toBe(1) + + collection.utils.writeUpdate({ id: `1`, name: `Manual` }) + expect(collection.get(`1`)?.name).toBe(`Manual`) + await vi.waitFor(() => expect(queryFn).toHaveBeenCalledTimes(2)) + + expect(collection.get(`1`)?.name).toBe(`Manual`) + expect(scopedQuery.state.data).toEqual([initial]) + + refetchResult.resolve([authoritative]) + await vi.waitFor(() => { + expect(collection.get(`1`)?.name).toBe(`Authoritative`) + expect(scopedQuery.state.dataUpdateCount).toBe(2) + }) + } finally { + refetchResult.resolve([authoritative]) + await active.cleanup() + await collection.cleanup() + customQueryClient.clear() + } + }, + ) + + it(`keeps active scoped queries isolated across collections sharing a QueryClient base key`, async () => { + const baseKey = [`shared-query-client-ownership`] + const serverA: Array = [ + { id: `a1`, name: `A one`, category: `A` }, + ] + let serverB: Array = [ + { id: `b1`, name: `B one`, category: `B` }, + ] + const queryA = vi.fn(() => + Promise.resolve(serverA.map((row) => structuredClone(row))), + ) + const queryB = vi.fn(() => + Promise.resolve(serverB.map((row) => structuredClone(row))), + ) + const sharedQueryClient = new QueryClient({ + defaultOptions: { + queries: { + gcTime: Number.POSITIVE_INFINITY, + staleTime: Number.POSITIVE_INFINITY, + retry: false, + }, + }, + }) + const collectionA = createCollection( + queryCollectionOptions({ + id: `shared-query-client-owner-a`, + queryClient: sharedQueryClient, + queryKey: baseKey, + queryFn: queryA, + getKey: (row) => row.id, + syncMode: `on-demand`, + startSync: true, + }), + ) + const collectionB = createCollection( + queryCollectionOptions({ + id: `shared-query-client-owner-b`, + queryClient: sharedQueryClient, + queryKey: baseKey, + queryFn: queryB, + getKey: (row) => row.id, + syncMode: `on-demand`, + startSync: true, + }), + ) + const activeA = createLiveQueryCollection({ + query: (query) => + query + .from({ row: collectionA }) + .where(({ row }) => eq(row.category, `A`)), + }) + const activeB = createLiveQueryCollection({ + query: (query) => + query + .from({ row: collectionB }) + .where(({ row }) => eq(row.category, `B`)), }) - await query2.preload() + try { + await Promise.all([activeA.preload(), activeB.preload()]) + expect(queryA).toHaveBeenCalledTimes(1) + expect(queryB).toHaveBeenCalledTimes(1) + + const cachedQueries = sharedQueryClient.getQueryCache().getAll() + const aQueryBefore = cachedQueries.find((query) => + (query.state.data as Array | undefined)?.some( + (row) => row.id === `a1`, + ), + ) + const bQueryBefore = cachedQueries.find((query) => + (query.state.data as Array | undefined)?.some( + (row) => row.id === `b1`, + ), + ) + expect(aQueryBefore).toBeDefined() + expect(bQueryBefore).toBeDefined() + expect(aQueryBefore!.queryKey).not.toEqual(bQueryBefore!.queryKey) - await vi.waitFor(() => { - expect(collection.size).toBe(2) + const bQueryKey = bQueryBefore!.queryKey + expect(activeB.toArray.map((row) => row.id)).toEqual([`b1`]) + + collectionA.utils.writeInsert({ + id: `a-outside`, + name: `Outside A scope`, + category: `outside`, + }) + await vi.waitFor(() => expect(queryA).toHaveBeenCalledTimes(2)) + + const afterAWrite = { + queryPresent: + sharedQueryClient + .getQueryCache() + .find({ queryKey: bQueryKey, exact: true }) !== undefined, + queryCalls: queryB.mock.calls.length, + materializedRows: activeB.toArray.map((row) => row.id), + } + + serverB = [...serverB, { id: `b2`, name: `B two`, category: `B` }] + await sharedQueryClient.invalidateQueries({ + queryKey: bQueryKey, + exact: true, + }) + await flushPromises() + + const afterBInvalidation = { + queryPresent: + sharedQueryClient + .getQueryCache() + .find({ queryKey: bQueryKey, exact: true }) !== undefined, + queryCalls: queryB.mock.calls.length, + materializedRows: activeB.toArray.map((row) => row.id), + } + + expect({ afterAWrite, afterBInvalidation }).toEqual({ + afterAWrite: { + queryPresent: true, + queryCalls: 1, + materializedRows: [`b1`], + }, + afterBInvalidation: { + queryPresent: true, + queryCalls: 2, + materializedRows: [`b1`, `b2`], + }, + }) + } finally { + await activeA.cleanup() + await activeB.cleanup() + await collectionA.cleanup() + await collectionB.cleanup() + sharedQueryClient.clear() + } + }) + + it(`does not let a foreign enabled observer authorize a disabled collection refetch`, async () => { + const queryKey = [`disabled-collection-foreign-observer`] + const initial: CategorisedItem = { + id: `1`, + name: `Initial`, + category: `A`, + } + const collectionQueryFn = vi.fn(() => + Promise.resolve([structuredClone(initial)]), + ) + const foreignQueryFn = vi.fn(() => + Promise.resolve([structuredClone(initial)]), + ) + const customQueryClient = new QueryClient({ + defaultOptions: { + queries: { + gcTime: Number.POSITIVE_INFINITY, + staleTime: Number.POSITIVE_INFINITY, + retry: false, + }, + }, + }) + customQueryClient.setQueryData(queryKey, [initial]) + + const collection = createCollection( + queryCollectionOptions({ + id: `disabled-collection-foreign-observer`, + queryClient: customQueryClient, + queryKey: () => queryKey, + queryFn: collectionQueryFn, + enabled: false, + getKey: (row) => row.id, + syncMode: `on-demand`, + startSync: true, + }), + ) + const active = createLiveQueryCollection({ + query: (query) => + query + .from({ row: collection }) + .where(({ row }) => eq(row.category, `A`)), }) + const foreignObserver = new QueryObserver(customQueryClient, { + queryKey, + queryFn: foreignQueryFn, + enabled: true, + staleTime: Number.POSITIVE_INFINITY, + retry: false, + }) + let unsubscribeForeign: (() => void) | undefined + + try { + await active.preload() + expect(collectionQueryFn).not.toHaveBeenCalled() + expect(active.toArray.map((row) => row.name)).toEqual([`Initial`]) - // After remount, the updated data should persist - expect(collection.get(`1`)?.name).toBe(`Updated Item 1`) + unsubscribeForeign = foreignObserver.subscribe(() => {}) + const sharedQuery = customQueryClient.getQueryCache().find({ + queryKey, + exact: true, + })! - await query2.cleanup() - customQueryClient.clear() + expect(sharedQuery.getObserversCount()).toBe(2) + expect(sharedQuery.isDisabled()).toBe(false) + expect(foreignQueryFn).not.toHaveBeenCalled() + + collection.utils.writeInsert({ + id: `outside`, + name: `Outside`, + category: `B`, + }) + + for (let turn = 0; turn < 20; turn++) await Promise.resolve() + + expect(collectionQueryFn).not.toHaveBeenCalled() + expect(foreignQueryFn).not.toHaveBeenCalled() + } finally { + unsubscribeForeign?.() + await active.cleanup() + await collection.cleanup() + customQueryClient.clear() + } }) + + it(`keeps a foreign-observed Query reachable after deferred owner release`, async () => { + const barrier = createDeferred() + const queryKey = [`deferred-foreign-owner-release`] + let serverRows: Array = [ + { id: `1`, name: `Initial`, category: `A` }, + ] + const queryFn = vi.fn(() => + Promise.resolve(serverRows.map((row) => structuredClone(row))), + ) + const customQueryClient = new QueryClient({ + defaultOptions: { + queries: { + gcTime: Number.POSITIVE_INFINITY, + staleTime: Number.POSITIVE_INFINITY, + retry: false, + }, + }, + }) + const collection = createCollection( + queryCollectionOptions({ + id: `deferred-foreign-owner-release`, + queryClient: customQueryClient, + queryKey: () => queryKey, + queryFn, + getKey: (row) => row.id, + syncMode: `on-demand`, + startSync: true, + }), + ) + const active = createLiveQueryCollection({ + query: (query) => + query + .from({ row: collection }) + .where(({ row }) => eq(row.category, `A`)), + }) + const foreignObserver = new QueryObserver(customQueryClient, { + queryKey, + queryFn, + enabled: true, + staleTime: Number.POSITIVE_INFINITY, + retry: false, + }) + let unsubscribeForeign: (() => void) | undefined + const barrierCompletion = barrier.promise.then(() => { + collection.deferDataRefresh = null + }) + + try { + await active.preload() + expect(queryFn).toHaveBeenCalledTimes(1) + + unsubscribeForeign = foreignObserver.subscribe(() => {}) + const sharedQuery = customQueryClient.getQueryCache().find({ + queryKey, + exact: true, + })! + expect(sharedQuery.getObserversCount()).toBe(2) + + collection.deferDataRefresh = barrier.promise + serverRows = [{ id: `1`, name: `Authoritative`, category: `A` }] + collection.utils.writeUpdate({ id: `1`, name: `Manual` }) + + for (let turn = 0; turn < 20; turn++) await Promise.resolve() + expect(queryFn).toHaveBeenCalledTimes(1) + + await active.cleanup() + expect(sharedQuery.getObserversCount()).toBe(1) + expect( + customQueryClient.getQueryCache().find({ queryKey, exact: true }), + ).toBe(sharedQuery) + + barrier.resolve() + await barrierCompletion + for (let turn = 0; turn < 20; turn++) await Promise.resolve() + + expect( + customQueryClient.getQueryCache().find({ queryKey, exact: true }), + ).toBe(sharedQuery) + expect(queryFn).toHaveBeenCalledTimes(2) + expect( + foreignObserver.getCurrentResult().data?.map((row) => row.name), + ).toEqual([`Authoritative`]) + } finally { + barrier.resolve() + collection.deferDataRefresh = null + unsubscribeForeign?.() + await active.cleanup() + await collection.cleanup() + customQueryClient.clear() + } + }) + + it(`evicts an owned custom-hash alias outside the base-key prefix`, async () => { + const foreignKey = [`custom-hash-foreign-retained-key`] + const baseKey = [`custom-hash-owned-base-key`] + let enabled = false + const initial: CategorisedItem = { + id: `1`, + name: `Stale`, + category: `A`, + } + const authoritative: CategorisedItem = { + id: `1`, + name: `Authoritative`, + category: `A`, + } + const queryFn = vi.fn(() => + Promise.resolve([structuredClone(authoritative)]), + ) + const customQueryClient = new QueryClient({ + defaultOptions: { + queries: { + gcTime: Number.POSITIVE_INFINITY, + staleTime: Number.POSITIVE_INFINITY, + retry: false, + queryKeyHashFn: () => `forced-custom-hash-alias`, + }, + }, + }) + customQueryClient.setQueryData(foreignKey, [initial]) + + const aliasedQuery = customQueryClient.getQueryCache().getAll()[0]! + expect(aliasedQuery.queryKey).toEqual(foreignKey) + + const collection = createCollection( + queryCollectionOptions({ + id: `custom-hash-owned-alias`, + queryClient: customQueryClient, + queryKey: () => baseKey, + queryFn, + enabled: () => enabled, + getKey: (row) => row.id, + syncMode: `on-demand`, + startSync: true, + }), + ) + const createActive = () => + createLiveQueryCollection({ + query: (query) => + query + .from({ row: collection }) + .where(({ row }) => eq(row.category, `A`)), + }) + const first = createActive() + let remounted: ReturnType | undefined + + try { + await first.preload() + expect(queryFn).not.toHaveBeenCalled() + expect(first.toArray.map((row) => row.name)).toEqual([`Stale`]) + + collection.utils.writeInsert({ + id: `outside`, + name: `Outside`, + category: `B`, + }) + for (let turn = 0; turn < 20; turn++) await Promise.resolve() + + expect( + customQueryClient.getQueryCache().getAll().includes(aliasedQuery), + ).toBe(false) + + enabled = true + await first.cleanup() + + remounted = createActive() + await remounted.preload() + + expect(queryFn).toHaveBeenCalledTimes(1) + expect(remounted.toArray.map((row) => row.name)).toEqual([ + `Authoritative`, + ]) + } finally { + await remounted?.cleanup() + await first.cleanup() + await collection.cleanup() + customQueryClient.clear() + } + }) + + it(`keeps an exact shared Query reachable when one owning collection is cleaned up`, async () => { + const baseKey = [`shared-query-cleanup-ownership`] + let serverRows: Array = [ + { id: `1`, name: `First`, category: `A` }, + ] + const queryFn = vi.fn(() => + Promise.resolve(serverRows.map((row) => structuredClone(row))), + ) + const sharedQueryClient = new QueryClient({ + defaultOptions: { + queries: { + gcTime: Number.POSITIVE_INFINITY, + staleTime: Number.POSITIVE_INFINITY, + retry: false, + }, + }, + }) + const createSharedCollection = (id: string) => + createCollection( + queryCollectionOptions({ + id, + queryClient: sharedQueryClient, + queryKey: baseKey, + queryFn, + getKey: (row) => row.id, + syncMode: `on-demand`, + startSync: true, + }), + ) + const collectionA = createSharedCollection(`shared-query-cleanup-a`) + const collectionB = createSharedCollection(`shared-query-cleanup-b`) + const createActive = (collection: typeof collectionA) => + createLiveQueryCollection({ + query: (query) => + query + .from({ row: collection }) + .where(({ row }) => eq(row.category, `A`)), + }) + const activeA = createActive(collectionA) + const activeB = createActive(collectionB) + + try { + await Promise.all([activeA.preload(), activeB.preload()]) + const sharedQuery = sharedQueryClient.getQueryCache().getAll()[0]! + const sharedQueryKey = sharedQuery.queryKey + expect(sharedQuery.getObserversCount()).toBe(2) + + await activeA.cleanup() + await collectionA.cleanup() + + expect( + sharedQueryClient + .getQueryCache() + .find({ queryKey: sharedQueryKey, exact: true }), + ).toBe(sharedQuery) + expect(sharedQuery.getObserversCount()).toBe(1) + + serverRows = [...serverRows, { id: `2`, name: `Second`, category: `A` }] + await sharedQueryClient.invalidateQueries({ + queryKey: sharedQueryKey, + exact: true, + }) + + await vi.waitFor(() => { + expect(activeB.toArray.map((row) => row.id)).toEqual([`1`, `2`]) + }) + } finally { + await activeA.cleanup() + await activeB.cleanup() + await collectionA.cleanup() + await collectionB.cleanup() + sharedQueryClient.clear() + } + }) + + it(`does not ready a remounted co-owner from a shared pre-write result`, async () => { + const refetchResult = createDeferred>() + const refetchStarted = createDeferred() + const queryFn = vi + .fn<() => Promise>>() + .mockResolvedValueOnce([{ id: `1`, name: `Initial`, category: `A` }]) + .mockImplementationOnce(() => { + refetchStarted.resolve() + return refetchResult.promise + }) + const sharedQueryClient = new QueryClient({ + defaultOptions: { + queries: { + gcTime: Number.POSITIVE_INFINITY, + staleTime: Number.POSITIVE_INFINITY, + retry: false, + }, + }, + }) + const createSharedCollection = (id: string) => + createCollection( + queryCollectionOptions({ + id, + queryClient: sharedQueryClient, + queryKey: [`shared-pre-write-result`], + queryFn, + getKey: (row) => row.id, + syncMode: `on-demand`, + startSync: true, + }), + ) + const collectionA = createSharedCollection(`shared-write-owner-a`) + const collectionB = createSharedCollection(`shared-write-owner-b`) + const createActive = (collection: typeof collectionA) => + createLiveQueryCollection({ + query: (query) => + query + .from({ row: collection }) + .where(({ row }) => eq(row.category, `A`)), + }) + const activeA = createActive(collectionA) + const firstB = createActive(collectionB) + let remountedB: ReturnType | undefined + + try { + await Promise.all([activeA.preload(), firstB.preload()]) + + collectionA.utils.writeUpdate({ id: `1`, name: `Manual` }) + await refetchStarted.promise + await firstB.cleanup() + expect(collectionB.size).toBe(0) + + remountedB = createActive(collectionB) + let preloadSettled = false + let rowsAtSettlement: Array | undefined + const preload = remountedB.preload().then(() => { + rowsAtSettlement = remountedB!.toArray.map((row) => row.name) + preloadSettled = true + }) + await flushPromises() + + expect(queryFn).toHaveBeenCalledTimes(2) + expect(sharedQueryClient.isFetching()).toBe(1) + expect(remountedB.toArray).toEqual([]) + expect(preloadSettled).toBe(false) + + refetchResult.resolve([ + { id: `1`, name: `Authoritative`, category: `A` }, + ]) + await preload + + expect(rowsAtSettlement).toEqual([`Authoritative`]) + expect(remountedB.toArray.map((row) => row.name)).toEqual([ + `Authoritative`, + ]) + } finally { + refetchResult.resolve([]) + await remountedB?.cleanup() + await firstB.cleanup() + await activeA.cleanup() + await collectionA.cleanup() + await collectionB.cleanup() + sharedQueryClient.clear() + } + }) + + it(`evicts a protected scoped Query when its owner unloads before deferred revalidation`, async () => { + const barrier = createDeferred() + let serverRows: Array = [ + { id: `1`, name: `Initial`, category: `A` }, + ] + const queryFn = vi.fn(() => + Promise.resolve(serverRows.map((row) => structuredClone(row))), + ) + const customQueryClient = new QueryClient({ + defaultOptions: { + queries: { + gcTime: Number.POSITIVE_INFINITY, + staleTime: Number.POSITIVE_INFINITY, + retry: false, + }, + }, + }) + const collection = createCollection( + queryCollectionOptions({ + id: `deferred-owner-unload`, + queryClient: customQueryClient, + queryKey: [`deferred-owner-unload`], + queryFn, + getKey: (row) => row.id, + syncMode: `on-demand`, + startSync: true, + }), + ) + const createActive = () => + createLiveQueryCollection({ + query: (query) => + query + .from({ row: collection }) + .where(({ row }) => eq(row.category, `A`)), + }) + const first = createActive() + let remounted: ReturnType | undefined + + const barrierCompletion = barrier.promise.then(() => { + collection.deferDataRefresh = null + }) + + try { + await first.preload() + collection.deferDataRefresh = barrier.promise + serverRows = [{ id: `1`, name: `Authoritative`, category: `A` }] + collection.utils.writeUpdate({ id: `1`, name: `Manual` }) + await first.cleanup() + + barrier.resolve() + await barrierCompletion + + expect(customQueryClient.getQueryCache().getAll()).toEqual([]) + + remounted = createActive() + await remounted.preload() + expect(queryFn).toHaveBeenCalledTimes(2) + expect(remounted.toArray.map((row) => row.name)).toEqual([ + `Authoritative`, + ]) + } finally { + barrier.resolve() + collection.deferDataRefresh = null + await remounted?.cleanup() + await first.cleanup() + await collection.cleanup() + customQueryClient.clear() + } + }) + + it(`requires post-write authority when an initial request is reused`, async () => { + const initialResult = createDeferred>() + const postWriteResult = createDeferred>() + const initialStarted = createDeferred() + const postWriteStarted = createDeferred() + let call = 0 + const queryFn = vi.fn(() => { + call++ + if (call === 1) { + initialStarted.resolve() + return initialResult.promise + } + postWriteStarted.resolve() + return postWriteResult.promise + }) + const customQueryClient = new QueryClient({ + defaultOptions: { + queries: { + gcTime: Number.POSITIVE_INFINITY, + staleTime: Number.POSITIVE_INFINITY, + retry: false, + }, + }, + }) + const collection = createCollection( + queryCollectionOptions({ + id: `pre-write-request-authority`, + queryClient: customQueryClient, + queryKey: [`pre-write-request-authority`], + queryFn, + getKey: (row) => row.id, + syncMode: `on-demand`, + startSync: true, + }), + ) + const active = createLiveQueryCollection({ + query: (query) => + query + .from({ row: collection }) + .where(({ row }) => eq(row.category, `A`)), + }) + + try { + collection.utils.writeInsert({ + id: `1`, + name: `Existing`, + category: `A`, + }) + + let preloadSettled = false + let rowsAtSettlement: Array | undefined + const preload = active.preload().then(() => { + rowsAtSettlement = active.toArray.map((row) => row.name) + preloadSettled = true + }) + await initialStarted.promise + + collection.utils.writeDelete(`1`) + expect(collection.has(`1`)).toBe(false) + + initialResult.resolve([{ id: `1`, name: `Stale`, category: `A` }]) + for (let turn = 0; turn < 30; turn++) await Promise.resolve() + + expect({ + providerCalls: queryFn.mock.calls.length, + preloadSettled, + materializedName: collection.get(`1`)?.name, + rowsAtSettlement, + }).toEqual({ + providerCalls: 2, + preloadSettled: false, + materializedName: undefined, + rowsAtSettlement: undefined, + }) + + await postWriteStarted.promise + postWriteResult.resolve([ + { id: `1`, name: `Authoritative`, category: `A` }, + ]) + await preload + + expect(rowsAtSettlement).toEqual([`Authoritative`]) + } finally { + initialResult.resolve([]) + postWriteResult.resolve([]) + await active.cleanup() + await collection.cleanup() + customQueryClient.clear() + } + }) + + it(`accepts replacement data after shared initialData reset repeats its update count`, async () => { + const firstFetch = createDeferred>() + const writeRefetch = createDeferred>() + const resetFetch = createDeferred>() + const firstStarted = createDeferred() + const writeRefetchStarted = createDeferred() + const resetFetchStarted = createDeferred() + const replacementPublished = createDeferred() + const queryKey = [`shared-initial-data-reset`] + let call = 0 + const queryFn = vi.fn((context: { signal: AbortSignal }) => { + void context.signal.aborted + call++ + if (call === 1) { + firstStarted.resolve() + return firstFetch.promise + } + if (call === 2) { + writeRefetchStarted.resolve() + return writeRefetch.promise + } + resetFetchStarted.resolve() + return resetFetch.promise + }) + const customQueryClient = new QueryClient({ + defaultOptions: { + queries: { + gcTime: Number.POSITIVE_INFINITY, + staleTime: Number.POSITIVE_INFINITY, + retry: false, + }, + }, + }) + const outside = new QueryObserver>( + customQueryClient, + { + queryKey, + queryFn, + initialData: [{ id: `1`, name: `Seed`, category: `A` }], + staleTime: Number.POSITIVE_INFINITY, + retry: false, + }, + ) + const unsubscribeOutside = outside.subscribe((result) => { + if (result.data?.[0]?.name === `Replacement`) { + replacementPublished.resolve() + } + }) + const collection = createCollection( + queryCollectionOptions({ + id: `shared-initial-data-reset`, + queryClient: customQueryClient, + queryKey: () => queryKey, + queryFn, + getKey: (row) => row.id, + syncMode: `on-demand`, + startSync: true, + }), + ) + const active = createLiveQueryCollection({ + query: (query) => query.from({ row: collection }), + }) + + try { + await active.preload() + expect(active.toArray.map((row) => row.name)).toEqual([`Seed`]) + + const firstRefetch = collection.utils.refetch({ throwOnError: true }) + await firstStarted.promise + firstFetch.resolve([{ id: `1`, name: `Initial`, category: `A` }]) + await firstRefetch + + expect(active.toArray.map((row) => row.name)).toEqual([`Initial`]) + expect( + customQueryClient.getQueryCache().find({ queryKey, exact: true }) + ?.state.dataUpdateCount, + ).toBe(1) + + collection.utils.writeUpdate({ id: `1`, name: `Manual` }) + await writeRefetchStarted.promise + expect(active.toArray.map((row) => row.name)).toEqual([`Manual`]) + + const reset = customQueryClient.resetQueries({ + queryKey, + exact: true, + }) + await resetFetchStarted.promise + resetFetch.resolve([{ id: `1`, name: `Replacement`, category: `A` }]) + await reset + await replacementPublished.promise + for (let turn = 0; turn < 30; turn++) await Promise.resolve() + + expect({ + providerCalls: queryFn.mock.calls.length, + cachedName: (customQueryClient.getQueryData>( + queryKey, + ) ?? [])[0]?.name, + cachedUpdateCount: customQueryClient.getQueryCache().find({ + queryKey, + exact: true, + })?.state.dataUpdateCount, + materializedNames: active.toArray.map((row) => row.name), + }).toEqual({ + providerCalls: 3, + cachedName: `Replacement`, + cachedUpdateCount: 1, + materializedNames: [`Replacement`], + }) + } finally { + firstFetch.resolve([]) + writeRefetch.resolve([]) + resetFetch.resolve([]) + unsubscribeOutside() + await active.cleanup() + await collection.cleanup() + customQueryClient.clear() + } + }) + + it(`revalidates a preserved shared Query after deferred owner cleanup`, async () => { + const barrier = createDeferred() + const secondResult = createDeferred>() + const secondStarted = createDeferred() + const outsideUpdated = createDeferred() + const queryKey = [`deferred-cleanup-shared-query`] + let call = 0 + const queryFn = vi.fn(() => { + call++ + if (call === 1) { + return Promise.resolve([{ id: `1`, name: `Initial`, category: `A` }]) + } + secondStarted.resolve() + return secondResult.promise + }) + const customQueryClient = new QueryClient({ + defaultOptions: { + queries: { + gcTime: Number.POSITIVE_INFINITY, + staleTime: Number.POSITIVE_INFINITY, + retry: false, + }, + }, + }) + const collection = createCollection( + queryCollectionOptions({ + id: `deferred-cleanup-shared-query`, + queryClient: customQueryClient, + queryKey: () => queryKey, + queryFn, + getKey: (row) => row.id, + syncMode: `on-demand`, + startSync: true, + }), + ) + const active = createLiveQueryCollection({ + query: (query) => query.from({ row: collection }), + }) + let unsubscribeOutside: (() => void) | undefined + const barrierCompletion = barrier.promise.then(() => { + collection.deferDataRefresh = null + }) + + try { + await active.preload() + const outside = new QueryObserver>( + customQueryClient, + { + queryKey, + queryFn, + staleTime: Number.POSITIVE_INFINITY, + retry: false, + }, + ) + unsubscribeOutside = outside.subscribe((result) => { + if (result.data?.[0]?.name === `Authoritative`) { + outsideUpdated.resolve() + } + }) + + expect(outside.getCurrentResult().data?.[0]?.name).toBe(`Initial`) + expect( + customQueryClient + .getQueryCache() + .find({ queryKey, exact: true }) + ?.getObserversCount(), + ).toBe(2) + + collection.deferDataRefresh = barrier.promise + collection.utils.writeUpdate({ id: `1`, name: `Manual` }) + await active.cleanup() + await collection.cleanup() + + const preserved = customQueryClient.getQueryCache().find({ + queryKey, + exact: true, + }) + expect(preserved).toBeDefined() + expect(preserved?.getObserversCount()).toBe(1) + + barrier.resolve() + await barrierCompletion + for (let turn = 0; turn < 30; turn++) await Promise.resolve() + + expect(queryFn).toHaveBeenCalledTimes(2) + await secondStarted.promise + secondResult.resolve([ + { id: `1`, name: `Authoritative`, category: `A` }, + ]) + await outsideUpdated.promise + + expect(outside.getCurrentResult().data?.[0]?.name).toBe(`Authoritative`) + } finally { + barrier.resolve() + collection.deferDataRefresh = null + secondResult.resolve([]) + unsubscribeOutside?.() + await active.cleanup() + await collection.cleanup() + customQueryClient.clear() + } + }) + + it.each([undefined, `static`] as const)( + `replenishes an active limited query after a manual delete with staleTime %s`, + async (staleTime) => { + let serverRows: Array = [ + { id: `1`, name: `First`, category: `A` }, + { id: `2`, name: `Second`, category: `A` }, + { id: `3`, name: `Third`, category: `A` }, + ] + const customQueryClient = new QueryClient({ + defaultOptions: { queries: { retry: false } }, + }) + const queryFn = vi.fn((context: QueryFunctionContext) => { + const limit = context.meta?.loadSubsetOptions?.limit + return Promise.resolve( + serverRows.slice(0, limit).map((row) => structuredClone(row)), + ) + }) + const collection = createCollection( + queryCollectionOptions({ + id: `limited-manual-write-${String(staleTime)}`, + queryClient: customQueryClient, + queryKey: [`limited-manual-write-${String(staleTime)}`], + queryFn, + getKey: (item) => item.id, + syncMode: `on-demand`, + staleTime, + autoIndex: `eager`, + defaultIndexType: BTreeIndex, + }), + ) + const active = createLiveQueryCollection({ + query: (query) => + query + .from({ item: collection }) + .orderBy(({ item }) => item.id, `asc`) + .limit(2), + }) + + try { + await active.preload() + expect(active.toArray.map((row) => row.id)).toEqual([`1`, `2`]) + + serverRows = serverRows.filter((row) => row.id !== `1`) + collection.utils.writeDelete(`1`) + + await vi.waitFor(() => { + expect(queryFn.mock.calls.length).toBeGreaterThan(1) + expect(active.toArray.map((row) => row.id)).toEqual([`2`, `3`]) + }) + } finally { + await active.cleanup() + await collection.cleanup() + customQueryClient.clear() + } + }, + ) + + it(`does not settle an in-flight initial scope from a manual write`, async () => { + const initialResult = createDeferred>() + const customQueryClient = new QueryClient({ + defaultOptions: { queries: { retry: false } }, + }) + const queryFn = vi.fn(() => initialResult.promise) + const collection = createCollection( + queryCollectionOptions({ + id: `in-flight-manual-write`, + queryClient: customQueryClient, + queryKey: [`in-flight-manual-write`], + queryFn, + getKey: (item) => item.id, + syncMode: `on-demand`, + startSync: true, + }), + ) + const active = createLiveQueryCollection({ + query: (query) => + query + .from({ item: collection }) + .where(({ item }) => eq(item.category, `A`)), + }) + let preloadSettled = false + + try { + const preload = active.preload().then(() => { + preloadSettled = true + }) + await vi.waitFor(() => expect(queryFn).toHaveBeenCalledTimes(1)) + + collection.utils.writeInsert({ + id: `outside`, + name: `Outside`, + category: `B`, + }) + await flushPromises() + + expect(preloadSettled).toBe(false) + expect( + customQueryClient.getQueryCache().getAll()[0]?.state.data, + ).toBeUndefined() + + initialResult.resolve([{ id: `1`, name: `First`, category: `A` }]) + await preload + expect(active.toArray.map((row) => row.id)).toEqual([`1`]) + } finally { + initialResult.resolve([]) + await active.cleanup() + await collection.cleanup() + customQueryClient.clear() + } + }) + + it(`removes a disabled cache entry and permits explicit recovery`, async () => { + const queryKey = [`disabled-manual-write`] + const customQueryClient = new QueryClient({ + defaultOptions: { queries: { retry: false } }, + }) + customQueryClient.setQueryData(queryKey, [ + { id: `1`, name: `First`, category: `A` }, + ]) + const queryFn = vi.fn(() => + Promise.resolve([{ id: `1`, name: `First`, category: `A` }]), + ) + const collection = createCollection( + queryCollectionOptions({ + id: `disabled-manual-write`, + queryClient: customQueryClient, + queryKey: () => queryKey, + queryFn, + enabled: false, + getKey: (item) => item.id, + syncMode: `on-demand`, + startSync: true, + }), + ) + const active = createLiveQueryCollection({ + query: (query) => + query + .from({ item: collection }) + .where(({ item }) => eq(item.category, `A`)), + }) + + try { + await active.preload() + expect(queryFn).not.toHaveBeenCalled() + + collection.utils.writeInsert({ + id: `outside`, + name: `Outside`, + category: `B`, + }) + expect(customQueryClient.getQueryCache().findAll({ queryKey })).toEqual( + [], + ) + + await collection.utils.refetch({ throwOnError: true }) + expect(queryFn).toHaveBeenCalledTimes(1) + expect(active.toArray.map((row) => row.id)).toEqual([`1`]) + } finally { + await active.cleanup() + await collection.cleanup() + customQueryClient.clear() + } + }) + + it(`accepts a fresh result whose update count repeats after reset`, async () => { + const heldRefetch = createDeferred>() + const replacement = [ + { id: `1`, name: `First`, category: `A` }, + { id: `2`, name: `Second`, category: `A` }, + ] + const queryKey = [`reset-manual-write`] + const customQueryClient = new QueryClient({ + defaultOptions: { queries: { retry: false } }, + }) + const queryFn = vi + .fn<() => Promise>>() + .mockResolvedValueOnce([replacement[0]!]) + .mockImplementationOnce(() => heldRefetch.promise) + .mockResolvedValueOnce(replacement) + const collection = createCollection( + queryCollectionOptions({ + id: `reset-manual-write`, + queryClient: customQueryClient, + queryKey: () => queryKey, + queryFn, + getKey: (item) => item.id, + syncMode: `on-demand`, + startSync: true, + }), + ) + const active = createLiveQueryCollection({ + query: (query) => + query + .from({ item: collection }) + .where(({ item }) => eq(item.category, `A`)), + }) + + try { + await active.preload() + expect( + customQueryClient.getQueryCache().find({ queryKey, exact: true }) + ?.state.dataUpdateCount, + ).toBe(1) + + collection.utils.writeInsert({ + id: `outside`, + name: `Outside`, + category: `B`, + }) + await vi.waitFor(() => expect(queryFn).toHaveBeenCalledTimes(2)) + + await customQueryClient.resetQueries({ queryKey, exact: true }) + expect(queryFn).toHaveBeenCalledTimes(3) + expect( + customQueryClient.getQueryCache().find({ queryKey, exact: true }) + ?.state.dataUpdateCount, + ).toBe(1) + expect(active.toArray.map((row) => row.id)).toEqual([`1`, `2`]) + + heldRefetch.resolve([ + { id: `obsolete`, name: `Obsolete`, category: `A` }, + ]) + await flushPromises() + expect(active.toArray.map((row) => row.id)).toEqual([`1`, `2`]) + } finally { + heldRefetch.resolve([]) + await active.cleanup() + await collection.cleanup() + customQueryClient.clear() + } + }) + + it(`waits for recovery rows when remounting after a write-triggered refetch error`, async () => { + const initialResult = createDeferred>() + const failedRefetch = createDeferred>() + const recoveryResult = createDeferred>() + const initialStarted = createDeferred() + const failedRefetchStarted = createDeferred() + const recoveryStarted = createDeferred() + const refetchErrored = createDeferred() + const recoverySucceeded = createDeferred() + const failure = new Error(`write-triggered refetch failed`) + const queryKey = [`error-remount-recovery`] + const customQueryClient = new QueryClient({ + defaultOptions: { + queries: { + gcTime: Number.POSITIVE_INFINITY, + staleTime: Number.POSITIVE_INFINITY, + retry: false, + }, + }, + }) + + let queryCall = 0 + const queryFn = vi.fn(() => { + queryCall++ + + if (queryCall === 1) { + initialStarted.resolve() + return initialResult.promise + } + if (queryCall === 2) { + failedRefetchStarted.resolve() + return failedRefetch.promise + } + + recoveryStarted.resolve() + return recoveryResult.promise + }) + + const unsubscribeCache = customQueryClient + .getQueryCache() + .subscribe((event) => { + if (event.query.queryKey[0] !== queryKey[0]) return + + if (event.query.state.status === `error`) { + refetchErrored.resolve() + } + if ( + queryCall === 3 && + event.query.state.status === `success` && + event.query.state.fetchStatus === `idle` + ) { + recoverySucceeded.resolve() + } + }) + + const collection = createCollection( + queryCollectionOptions({ + id: `error-remount-recovery`, + queryClient: customQueryClient, + queryKey, + queryFn, + getKey: (item) => item.id, + syncMode: `on-demand`, + startSync: true, + }), + ) + const createActive = () => + createLiveQueryCollection({ + query: (query) => + query + .from({ item: collection }) + .where(({ item }) => eq(item.category, `A`)), + }) + + const first = createActive() + let remounted: ReturnType | undefined + + try { + const initialPreload = first.preload() + await initialStarted.promise + initialResult.resolve([{ id: `1`, name: `Initial`, category: `A` }]) + await initialPreload + expect(first.toArray.map((row) => row.name)).toEqual([`Initial`]) + + collection.utils.writeUpdate({ id: `1`, name: `Manual` }) + await failedRefetchStarted.promise + failedRefetch.reject(failure) + await refetchErrored.promise + expect(collection.utils.lastError).toBe(failure) + + await first.cleanup() + expect(collection.size).toBe(0) + + const current = createActive() + remounted = current + let preloadSettled = false + let rowsAtPreloadSettlement: Array | undefined + const remountPreload = current.preload().then(() => { + rowsAtPreloadSettlement = current.toArray.map((row) => row.name) + preloadSettled = true + }) + + await recoveryStarted.promise + + // Drain deterministic promise work after the recovery-start signal. + // Recovery remains controlled and unresolved throughout this checkpoint. + for (let turn = 0; turn < 20; turn++) { + await Promise.resolve() + } + + expect(queryFn).toHaveBeenCalledTimes(3) + expect(customQueryClient.isFetching()).toBe(1) + expect(current.toArray).toEqual([]) + expect(preloadSettled).toBe(false) + expect(rowsAtPreloadSettlement).toBeUndefined() + + recoveryResult.resolve([{ id: `1`, name: `Recovered`, category: `A` }]) + await recoverySucceeded.promise + await remountPreload + + // Readiness must be observed only after recovery rows are published. + expect(rowsAtPreloadSettlement).toEqual([`Recovered`]) + expect(current.toArray.map((row) => row.name)).toEqual([`Recovered`]) + } finally { + initialResult.resolve([]) + failedRefetch.reject(failure) + recoveryResult.resolve([]) + unsubscribeCache() + await remounted?.cleanup() + await first.cleanup() + await collection.cleanup() + customQueryClient.clear() + } + }) + + it.each([`resolve`, `reject`] as const)( + `revalidates after a deferred refresh barrier %s`, + async (outcome) => { + const barrier = createDeferred() + const failure = new Error(`deferred refresh failed`) + const serverRows: Array = [ + { id: `1`, name: `First`, category: `A` }, + ] + const customQueryClient = new QueryClient({ + defaultOptions: { queries: { retry: false } }, + }) + const queryFn = vi.fn(() => + Promise.resolve(serverRows.map((row) => structuredClone(row))), + ) + const collection = createCollection( + queryCollectionOptions({ + id: `deferred-manual-write-${outcome}`, + queryClient: customQueryClient, + queryKey: [`deferred-manual-write-${outcome}`], + queryFn, + getKey: (item) => item.id, + syncMode: `on-demand`, + startSync: true, + }), + ) + const active = createLiveQueryCollection({ + query: (query) => + query + .from({ item: collection }) + .where(({ item }) => eq(item.category, `A`)), + }) + let observedFailure: unknown + const barrierSettlement = barrier.promise.then( + () => { + collection.deferDataRefresh = null + }, + (error: unknown) => { + observedFailure = error + collection.deferDataRefresh = null + }, + ) + + try { + await active.preload() + collection.deferDataRefresh = barrier.promise + serverRows.push({ id: `2`, name: `Second`, category: `A` }) + collection.utils.writeInsert({ + id: `outside`, + name: `Outside`, + category: `B`, + }) + await flushPromises() + expect(queryFn).toHaveBeenCalledTimes(1) + + if (outcome === `resolve`) barrier.resolve() + else barrier.reject(failure) + await barrierSettlement + + await vi.waitFor(() => { + expect(queryFn).toHaveBeenCalledTimes(2) + expect(active.toArray.map((row) => row.id)).toEqual([`1`, `2`]) + }) + if (outcome === `reject`) expect(observedFailure).toBe(failure) + } finally { + barrier.resolve() + collection.deferDataRefresh = null + await active.cleanup() + await collection.cleanup() + customQueryClient.clear() + } + }, + ) }) describe(`rows from external sync sources`, () => {