API index and search · Build metadata
Source snapshot
packages/client/src/ensure.ts
1// Query readiness for route loaders and intent prefetching. `ResultType` remains the2// server-authority axis; this module only decides when a caller has enough data to continue.34import { stableKey } from "./key.ts";5import type { AnyQuery } from "./query.ts";6import type { ColsMap } from "./schema.ts";7import type { Store } from "./store.ts";8import type { ResultType } from "./types.ts";910/** When an {@link QueryEnsureCache.ensure} call may resolve.11 *12 * - `complete` waits for the server-authoritative result (the default).13 * - `present` resolves as soon as the local view contains a result, while the remote retain keeps14 * revalidating in the background. An authoritative empty result also resolves it, so a real15 * not-found query never waits forever.16 *17 * `present` deliberately does not add a `partial` {@link ResultType}: a locally useful answer and18 * server authority are independent facts. While it resolves early, the view's result type remains19 * `unknown` until the server says otherwise. */20export type EnsureQueryUntil = "complete" | "present";2122export interface EnsureQueryOptions {23 /** Readiness policy. Defaults to `complete`. */24 until?: EnsureQueryUntil;25 /** Cancel this caller's wait. The shared query may stay retained for another waiter/prefetch. */26 signal?: AbortSignal;27}2829export interface QueryEnsureCacheOptions {30 /** Keep a completed preload alive for this long so the destination can adopt its rows. */31 releaseDelayMs?: number;32 /** Bound retained preloads. In-flight waits are never evicted. */33 maxEntries?: number;34}3536/** The structural surface consumed by framework adapters such as `@rindle/tanstack`. */37export interface QueryEnsurer {38 ensure(query: AnyQuery, options?: EnsureQueryOptions): Promise<void>;39}4041interface EnsureView {42 readonly data: unknown;43 readonly resultType: ResultType;44 subscribe(listener: () => void): () => void;45 destroy(): void;46}4748interface EnsureWaiter {49 until: EnsureQueryUntil;50 resolve: () => void;51 reject: (reason: Error) => void;52 signal?: AbortSignal;53 onAbort?: () => void;54}5556interface EnsureEntry {57 key: string;58 view: EnsureView;59 waiters: Set<EnsureWaiter>;60 unsubscribe: () => void;61 releaseTimer?: ReturnType<typeof setTimeout>;62 lastUsedAt: number;63}6465const DEFAULT_RELEASE_DELAY_MS = 10_000;66const DEFAULT_MAX_ENTRIES = 32;6768/**69 * Deduplicates route/intent preloads and holds their live view through the navigation handoff.70 *71 * The cache materializes the real named query rather than retaining sync coverage alone: that is72 * what makes `until: "present"` observable for overlapping queries already satisfied by local73 * normalized rows. Once the server marks the query complete, the entry is kept briefly for a74 * destination component to take its own retain, then released automatically.75 */76export class QueryEnsureCache<S extends ColsMap> implements QueryEnsurer {77 private readonly store: Store<S>;78 private readonly releaseDelayMs: number;79 private readonly maxEntries: number;80 private readonly entries = new Map<string, EnsureEntry>();81 private closed = false;8283 constructor(store: Store<S>, options: QueryEnsureCacheOptions = {}) {84 this.store = store;85 this.releaseDelayMs = Math.max(0, options.releaseDelayMs ?? DEFAULT_RELEASE_DELAY_MS);86 this.maxEntries = Math.max(1, options.maxEntries ?? DEFAULT_MAX_ENTRIES);87 }8889 /**90 * Ensure a named query is retained, resolving according to `options.until`.91 *92 * Concurrent calls for the same `(name, args, AST)` share one materialized view and one remote93 * subscription. A `present` call can resolve from local rows while a concurrent `complete` call94 * continues waiting for server authority.95 */96 async ensure<Q extends AnyQuery>(query: Q, options: EnsureQueryOptions = {}): Promise<void> {97 if (this.closed) throw new Error("QueryEnsureCache.ensure: cache is closed.");98 if (options.signal?.aborted) throw abortError();99 if (typeof query.name !== "string") {100 throw new Error("QueryEnsureCache.ensure: preloading requires a named query.");101 }102103 const key = stableKey({ ast: query.ast(), remote: { name: query.name, args: query.args } });104 let entry = this.entries.get(key);105 if (!entry) entry = this.createEntry(key, query);106 this.touch(entry);107108 const until = options.until ?? "complete";109 if (entry.view.resultType === "error") {110 const error = this.queryError();111 this.dispose(entry, error);112 throw error;113 }114 if (this.isReady(entry, until)) {115 if (entry.view.resultType === "complete") this.scheduleRelease(entry);116 this.trim();117 return;118 }119120 const ready = new Promise<void>((resolve, reject) => {121 const waiter: EnsureWaiter = { until, resolve, reject, signal: options.signal };122 if (options.signal) {123 waiter.onAbort = () => {124 if (!entry!.waiters.delete(waiter)) return;125 reject(abortError());126 };127 options.signal.addEventListener("abort", waiter.onAbort, { once: true });128 }129 entry!.waiters.add(waiter);130 });131 this.trim();132 return ready;133 }134135 /** Release every retained preload and reject outstanding waits. Idempotent. */136 close(): void {137 if (this.closed) return;138 this.closed = true;139 const error = new Error("QueryEnsureCache: closed before the query became ready.");140 for (const entry of [...this.entries.values()]) this.dispose(entry, error);141 }142143 /** Number of retained query entries (primarily useful for tests/devtools). */144 size(): number {145 return this.entries.size;146 }147148 private createEntry(key: string, query: AnyQuery): EnsureEntry {149 const view = this.store.materialize(query) as unknown as EnsureView;150 const entry: EnsureEntry = {151 key,152 view,153 waiters: new Set(),154 unsubscribe: () => {},155 lastUsedAt: Date.now(),156 };157 this.entries.set(key, entry);158 entry.unsubscribe = view.subscribe(() => this.inspect(entry));159 return entry;160 }161162 private inspect(entry: EnsureEntry): void {163 if (this.entries.get(entry.key) !== entry) return;164 if (entry.view.resultType === "error") {165 this.dispose(entry, this.queryError());166 return;167 }168169 for (const waiter of [...entry.waiters]) {170 if (!this.isReady(entry, waiter.until)) continue;171 entry.waiters.delete(waiter);172 this.detachAbort(waiter);173 waiter.resolve();174 }175176 // A present waiter may have continued while authority was still unknown. Keep the retain alive177 // through that revalidation; only a completed query starts the handoff/expiry clock.178 if (entry.view.resultType === "complete" && entry.waiters.size === 0) {179 this.scheduleRelease(entry);180 this.trim();181 }182 }183184 private isReady(entry: EnsureEntry, until: EnsureQueryUntil): boolean {185 if (entry.view.resultType === "complete") return true;186 return until === "present" && hasLocalResult(entry.view.data);187 }188189 private touch(entry: EnsureEntry): void {190 entry.lastUsedAt = Date.now();191 if (entry.releaseTimer !== undefined) {192 clearTimeout(entry.releaseTimer);193 entry.releaseTimer = undefined;194 }195 }196197 private scheduleRelease(entry: EnsureEntry): void {198 if (this.entries.get(entry.key) !== entry || entry.waiters.size > 0) return;199 if (entry.releaseTimer !== undefined) clearTimeout(entry.releaseTimer);200 // Even a configured zero delay crosses a task boundary. Some local-first backends publish201 // `complete` immediately before folding the authoritative catch-up batch; synchronous teardown202 // here would unregister the view in the middle of that delivery.203 entry.releaseTimer = setTimeout(() => this.dispose(entry), this.releaseDelayMs);204 (entry.releaseTimer as { unref?: () => void }).unref?.();205 }206207 private trim(): void {208 if (this.entries.size <= this.maxEntries) return;209 const idle = [...this.entries.values()]210 .filter((entry) => entry.waiters.size === 0)211 .sort((left, right) => left.lastUsedAt - right.lastUsedAt);212 while (this.entries.size > this.maxEntries && idle.length > 0) this.dispose(idle.shift()!);213 }214215 private dispose(entry: EnsureEntry, error?: Error): void {216 if (this.entries.get(entry.key) !== entry) return;217 this.entries.delete(entry.key);218 if (entry.releaseTimer !== undefined) clearTimeout(entry.releaseTimer);219 entry.unsubscribe();220 entry.view.destroy();221 if (error) {222 for (const waiter of entry.waiters) {223 this.detachAbort(waiter);224 waiter.reject(error);225 }226 }227 entry.waiters.clear();228 }229230 private queryError(): Error {231 return new Error("QueryEnsureCache.ensure: query entered the error result state.");232 }233234 private detachAbort(waiter: EnsureWaiter): void {235 if (waiter.signal && waiter.onAbort) waiter.signal.removeEventListener("abort", waiter.onAbort);236 }237}238239/** A plural query is locally present when it has at least one row; a `.one()` query when non-null. */240function hasLocalResult(data: unknown): boolean {241 return Array.isArray(data) ? data.length > 0 : data !== null && data !== undefined;242}243244function abortError(): Error {245 const error = new Error("QueryEnsureCache.ensure: wait was aborted.");246 error.name = "AbortError";247 return error;248}249