Rindle

API index and search · Build metadata

Source snapshot

packages/remote/src/backend.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// The RemoteBackend: a `@rindle/client` `Backend` over a network transport. It owns the2// epoch/seq/gap protocol (via {@link Subscriber}) and emits only CLEAN, in-order3// `hello`/`snapshot`/`batch` `ChangeEvent`s upward — so the same `Store`/`ArrayView` that4// drives the local backends drives this one, never seeing seq/epoch (WASM-CLIENT-DESIGN.md §2.2).5//6// On a gap (or epoch/schema drift) the backend re-hydrates INTERNALLY: it re-subscribes, the7// server re-registers the query under a NEW epoch and replies with a fresh hello + snapshot,8// and the `ArrayView` resets in place — the caller's materialized view reference survives.910import { Store } from "@rindle/client";11import type { Backend, BackendDevObserver, ChangeEvent, ColsMap, Mutation, QueryId, RemoteQuery, Schema } from "@rindle/client";1213import { ProtocolError, Subscriber } from "./protocol.ts";14import type { Batch, Hello, ServerMsg } from "./protocol.ts";15import {16  defaultSubscribeTarget,17  isThenable,18  subscribeMessage,19  type RawMutationSender,20  type SubscribeResolver,21} from "./subscribe.ts";22import { WsTransport } from "./transport.ts";23import { retryDelayMs } from "./query-error.ts";24import type { Transport } from "./transport.ts";2526interface QState {27  remote: RemoteQuery;28  subscriber: Subscriber | null;29  /** The epoch of the current subscription (0 before the first hello). */30  epoch: number;31  /** True between sending a re-subscribe and receiving its hello (so a second gap is ignored). */32  resubscribing: boolean;33  /** Monotonic token that cancels stale async lease resolutions. */34  subscribeTicket: number;35  /** Pending retryable-error re-subscribe timer (FOLLOWER-LAG-SHED §6.3), if any. */36  retryTimer: ReturnType<typeof setTimeout> | undefined;37  /** Consecutive retryable errors without a successful hello — the backoff exponent. */38  retryAttempt: number;39}4041export interface RemoteBackendOptions {42  /** Resolve the upstream subscribe target. Defaults to embedded-server `{name,args}`. */43  resolveSubscribe?: SubscribeResolver;44  /** Override raw authoritative writes, e.g. POST them to an app API server. */45  sendMutation?: RawMutationSender;46}4748export class RemoteBackend implements Backend {49  private readonly transport: Transport;50  private readonly resolveSubscribe: SubscribeResolver;51  private readonly sendMutation?: RawMutationSender;52  private handler: (qid: QueryId, ev: ChangeEvent) => void = () => {};53  private readonly devObservers = new Set<BackendDevObserver>();54  private readonly subs = new Map<QueryId, QState>();5556  constructor(transport: Transport, opts: RemoteBackendOptions = {}) {57    this.transport = transport;58    this.resolveSubscribe = opts.resolveSubscribe ?? defaultSubscribeTarget;59    this.sendMutation = opts.sendMutation;60    this.transport.onMessage((msg) => this.onServerMsg(msg));61  }6263  registerQuery(qid: QueryId, _ast: unknown, remote?: RemoteQuery): void {64    if (!remote) {65      console.error(`[rindle-remote] flat query ${qid} has no named remote identity`);66      return;67    }68    this.subs.set(qid, { remote, subscriber: null, epoch: 0, resubscribing: false, subscribeTicket: 0, retryTimer: undefined, retryAttempt: 0 });69    this.subscribe(qid, remote);70  }7172  unregisterQuery(qid: QueryId): void {73    const s = this.subs.get(qid);74    if (s?.retryTimer !== undefined) clearTimeout(s.retryTimer);75    this.subs.delete(qid);76    this.transport.send({ t: "unsubscribe", queryId: qid });77  }7879  /** Eventually-consistent: the mutation is sent; the resulting batches arrive async on the80   *  stream. The promise resolves once sent (the §2.3 write asymmetry, named not leaked). */81  mutate(mutations: Mutation[]): Promise<void> {82    if (this.sendMutation) return Promise.resolve(this.sendMutation(mutations));83    this.transport.send({ t: "mutate", mutations });84    return Promise.resolve();85  }8687  onEvent(handler: (qid: QueryId, ev: ChangeEvent) => void): void {88    this.handler = handler;89  }9091  __attachDevtoolsServerDeltas(observer: BackendDevObserver): () => void {92    this.devObservers.add(observer);93    return () => {94      this.devObservers.delete(observer);95    };96  }9798  // --- internals ---------------------------------------------------------------99100  private subscribe(qid: QueryId, remote: RemoteQuery): void {101    const s = this.subs.get(qid);102    if (!s) return;103    // A fresh subscribe (gap recovery, the retry timer itself) supersedes any scheduled104    // retryable-error retry — never leave two subscribe paths racing for one query.105    if (s.retryTimer !== undefined) {106      clearTimeout(s.retryTimer);107      s.retryTimer = undefined;108    }109    const request = { queryId: qid, remote, mode: "flat" as const };110    const ticket = ++s.subscribeTicket;111    const send = (target: ReturnType<typeof defaultSubscribeTarget>) => {112      const cur = this.subs.get(qid);113      if (cur !== s || cur.subscribeTicket !== ticket) return;114      this.transport.send(subscribeMessage(request, target));115    };116    const fail = (err: unknown) => {117      const cur = this.subs.get(qid);118      if (cur !== s || cur.subscribeTicket !== ticket) return;119      s.resubscribing = false;120      console.error(`[rindle-remote] query ${qid} subscribe resolution failed: ${String((err as Error)?.message ?? err)}`);121    };122    try {123      const target = this.resolveSubscribe(request);124      if (isThenable(target)) void target.then(send, fail);125      else send(target);126    } catch (err) {127      fail(err);128    }129  }130131  private onServerMsg(msg: ServerMsg): void {132    if (msg.t === "queryError") {133      this.onQueryError(msg.queryId, msg);134      return;135    }136    // This backend is flat-only; it ignores normalized frames (`nhello`/`nbatch`).137    if (msg.t !== "hello" && msg.t !== "batch") return;138    const s = this.subs.get(msg.queryId);139    if (!s) return; // unsubscribed / unknown query140    if (msg.t === "hello") this.openSubscriber(msg.queryId, s, msg.hello);141    else this.applyBatch(msg.queryId, s, msg.batch);142  }143144  /** Route a `queryError` by its 101 §5 classification: retryable ⇒ keep the QState (and the145   *  rows already folded downstream — 101 §6) and re-subscribe after a jittered backoff146   *  honoring `retryAfterMs`; terminal (or pre-classification servers) ⇒ drop the147   *  subscription, exactly as before (FOLLOWER-LAG-SHED §6.3). */148  private onQueryError(qid: QueryId, err: { message: string; code?: string; retryable?: boolean; retryAfterMs?: number }): void {149    const s = this.subs.get(qid);150    if (!s) return;151    if (err.retryable !== true) {152      if (s.retryTimer !== undefined) clearTimeout(s.retryTimer);153      this.subs.delete(qid);154      console.error(`[rindle-remote] query ${qid} subscription rejected: ${err.message}`);155      return;156    }157    if (s.retryTimer !== undefined) return; // a retry is already scheduled — don't stack them158    s.subscriber = null; // stop validating the dead epoch; recovery is a fresh seq-0 hydrate159    s.resubscribing = true;160    const delay = retryDelayMs(s.retryAttempt++, err.retryAfterMs);161    console.warn(162      `[rindle-remote] query ${qid} ${err.code ?? "error"} (retryable): re-subscribing in ${delay}ms: ${err.message}`,163    );164    const timer = setTimeout(() => {165      const cur = this.subs.get(qid);166      if (cur !== s) return;167      s.retryTimer = undefined;168      this.subscribe(qid, s.remote);169    }, delay);170    (timer as { unref?: () => void }).unref?.();171    s.retryTimer = timer;172  }173174  private openSubscriber(qid: QueryId, s: QState, hello: Hello): void {175    try {176      // The Subscriber ctor validates the comparator + fingerprint and emits the `hello`177      // ChangeEvent (→ the Store resets the view to this schema).178      s.subscriber = new Subscriber(hello, (ev) => this.emitServerEvent(qid, ev));179      s.epoch = hello.epoch;180      s.resubscribing = false;181      s.retryAttempt = 0; // a successful hello resets the retryable-error backoff182    } catch (e) {183      // A comparator/schema mismatch at hello is unrecoverable (a code-contract divergence) —184      // leave the view pending and report it; do not loop.185      s.subscriber = null;186      console.error(`[rindle-remote] query ${qid} subscription rejected: ${(e as Error).message}`);187    }188  }189190  private applyBatch(qid: QueryId, s: QState, batch: Batch): void {191    if (!s.subscriber) return; // no hello yet (or mid re-hydrate)192    if (batch.epoch < s.epoch) return; // a stale batch from a superseded epoch — drop193    try {194      s.subscriber.apply(batch);195    } catch (e) {196      if (!(e instanceof ProtocolError)) throw e;197      if (s.resubscribing) return; // already recovering198      // Gap / drift → re-hydrate under a new epoch (the server bumps it on re-subscribe).199      s.resubscribing = true;200      s.subscriber = null;201      this.subscribe(qid, s.remote);202    }203  }204205  private emitServerEvent(qid: QueryId, ev: ChangeEvent): void {206    this.handler(qid, ev);207    if (this.devObservers.size) {208      for (const o of this.devObservers) o.onServerDelta?.(qid, { format: "flat", event: ev });209    }210  }211}212213/** Convenience: a `Store` backed by a remote server, over a ws URL or a custom transport. */214export function createRemoteStore<S extends ColsMap>(215  schema: Schema<S>,216  urlOrTransport: string | Transport,217  opts: RemoteBackendOptions = {},218): Store<S> {219  const transport = typeof urlOrTransport === "string" ? new WsTransport(urlOrTransport) : urlOrTransport;220  return new Store(schema, new RemoteBackend(transport, opts));221}222