Rindle

API index and search · Build metadata

Source snapshot

packages/client/src/ensure.ts

Source revision 05d0bf2c2e56 · build details
Source revision: 05d0bf2c2e56.
TypeScript input SHA-256: aabe6cfcc4172b870d5e272142958e9ea8d8784c2aa23133156e5a7ee633318e
Generated 2026-09-04T23:58:25.590Z with TypeScript 6.0.3. Public TypeScript checks and declaration emit passed. Package runtime tests are separate.
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