API index and search · Build metadata
Source snapshot
packages/remote/src/backend.ts
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