diff --git a/packages/presence/src/client.test.ts b/packages/presence/src/client.test.ts index 652e9cbbd..c864caf4e 100644 --- a/packages/presence/src/client.test.ts +++ b/packages/presence/src/client.test.ts @@ -2,6 +2,7 @@ import { describe, expect, test } from "bun:test"; import * as Y from "yjs"; import { connectPresence, + type PresenceError, type PresenceEventSourceLike, type PresenceFetch, type PresenceStreamEvent, @@ -48,11 +49,35 @@ function fakeFetch( calls.push({ path: url, body }); return Promise.resolve({ ok: true, + status: 200, json: () => Promise.resolve(joinResponse), }); }; } +/** A `PresenceFetch` whose response per operation is scripted by + * `statusFor`, defaulting to 200 for anything unlisted — for exercising + * failure and rejoin paths `fakeFetch`'s always-succeeds shape can't. */ +function scriptedFetch( + calls: string[], + statusFor: (path: string) => number | "reject", +): PresenceFetch { + return (url) => { + calls.push(url); + const status = statusFor(url); + if (status === "reject") return Promise.reject(new Error("network down")); + return Promise.resolve({ + ok: status >= 200 && status < 300, + status, + json: () => Promise.resolve({}), + }); + }; +} + +async function flushMicrotasks(): Promise { + for (let i = 0; i < 5; i++) await Promise.resolve(); +} + describe("connectPresence", () => { test("joins immediately and opens the room's SSE stream", () => { const calls: { path: string; body: unknown }[] = []; @@ -104,13 +129,18 @@ describe("connectPresence", () => { handle.disconnect(); }); - test("publishCursor and publishTyping post heartbeats with the patch", () => { + test("publishCursor and publishTyping post heartbeats with the patch", async () => { const calls: { path: string; body: unknown }[] = []; const handle = connectPresence({ roomUrl: "/rooms/channel:chn_1", fetchImpl: fakeFetch(calls), openEventSource: () => new FakeEventSource(), }); + // Let the initial join's response settle before publishing: a patch + // published before the room has confirmed the join rejoins instead + // of heartbeating a membership the server doesn't have yet (CL-7202). + await Promise.resolve(); + await Promise.resolve(); calls.length = 0; // drop the initial join call handle.publishCursor({ x: 5, y: 6, surfaceVersion: 1 }); @@ -307,3 +337,335 @@ describe("connectPresence: doc sync", () => { handle.disconnect(); }); }); + +// CL-7202: the client used to blind-post every join/heartbeat/leave/update +// request (`.catch(() => undefined)`, no `response.ok` check anywhere, no +// way for a caller to hear about a failure), and never rejoined after a +// failed join or a heartbeat the server had already forgotten about (404 +// `not_joined`) — all while the SSE stream stayed open regardless, so the +// UI kept reading as live. +describe("connectPresence: error reporting and rejoin (CL-7202)", () => { + const ROOM_URL = "/rooms/channel:chn_1"; + + test("a join request that never reaches the server is reported through onError, not swallowed", async () => { + const calls: string[] = []; + const errors: PresenceError[] = []; + const handle = connectPresence({ + roomUrl: ROOM_URL, + fetchImpl: scriptedFetch(calls, () => "reject"), + openEventSource: () => new FakeEventSource(), + }); + handle.onError((error) => errors.push(error)); + + await flushMicrotasks(); + + expect(calls).toEqual([`${ROOM_URL}/join`]); + expect(errors).toEqual([{ operation: "join" }]); + handle.disconnect(); + }); + + test("a non-ok join response is reported through onError with its status", async () => { + const calls: string[] = []; + const errors: PresenceError[] = []; + const handle = connectPresence({ + roomUrl: ROOM_URL, + fetchImpl: scriptedFetch(calls, () => 500), + openEventSource: () => new FakeEventSource(), + }); + handle.onError((error) => errors.push(error)); + + await flushMicrotasks(); + + expect(errors).toEqual([{ operation: "join", status: 500 }]); + handle.disconnect(); + }); + + test("publishing before the room has confirmed the join rejoins instead of heartbeating a membership it doesn't have yet", async () => { + let clock = 0; + const calls: string[] = []; + const handle = connectPresence({ + roomUrl: ROOM_URL, + fetchImpl: scriptedFetch(calls, () => 500), + openEventSource: () => new FakeEventSource(), + now: () => clock, + }); + + await flushMicrotasks(); + clock = 1_000; // past the first failure's backoff window + handle.publishTyping(true); + await flushMicrotasks(); + + expect(calls).toEqual([`${ROOM_URL}/join`, `${ROOM_URL}/join`]); + handle.disconnect(); + }); + + test("a heartbeat 404 (self-eviction) triggers an automatic rejoin", async () => { + const calls: string[] = []; + let joinCount = 0; + const handle = connectPresence({ + roomUrl: ROOM_URL, + fetchImpl: scriptedFetch(calls, (path) => { + if (path.endsWith("/join")) { + joinCount += 1; + return 200; + } + if (path.endsWith("/heartbeat")) return 404; + return 200; + }), + openEventSource: () => new FakeEventSource(), + }); + const errors: PresenceError[] = []; + handle.onError((error) => errors.push(error)); + + await flushMicrotasks(); + expect(joinCount).toBe(1); + + handle.publishCursor({ x: 1, y: 2, surfaceVersion: 1 }); + await flushMicrotasks(); + + expect(calls).toEqual([ + `${ROOM_URL}/join`, + `${ROOM_URL}/heartbeat`, + `${ROOM_URL}/join`, + ]); + expect(errors).toEqual([{ operation: "heartbeat", status: 404 }]); + expect(joinCount).toBe(2); + handle.disconnect(); + }); + + test("a heartbeat succeeds normally once join has succeeded, without rejoining", async () => { + const calls: string[] = []; + const handle = connectPresence({ + roomUrl: ROOM_URL, + fetchImpl: scriptedFetch(calls, () => 200), + openEventSource: () => new FakeEventSource(), + }); + + await flushMicrotasks(); + handle.publishTyping(true); + await flushMicrotasks(); + + expect(calls).toEqual([`${ROOM_URL}/join`, `${ROOM_URL}/heartbeat`]); + handle.disconnect(); + }); + + test("a failed doc update is reported through onError", async () => { + const calls: string[] = []; + const doc = new Y.Doc(); + const errors: PresenceError[] = []; + const handle = connectPresence({ + roomUrl: ROOM_URL, + fetchImpl: scriptedFetch(calls, (path) => + path.endsWith("/update") ? 413 : 200, + ), + openEventSource: () => new FakeEventSource(), + doc, + }); + handle.onError((error) => errors.push(error)); + + await flushMicrotasks(); + doc.getText("content").insert(0, "hello"); + await flushMicrotasks(); + + expect(errors).toEqual([{ operation: "update", status: 413 }]); + handle.disconnect(); + }); + + test("onError's unsubscribe stops further delivery", async () => { + const calls: string[] = []; + const errors: PresenceError[] = []; + const handle = connectPresence({ + roomUrl: ROOM_URL, + fetchImpl: scriptedFetch(calls, () => 500), + openEventSource: () => new FakeEventSource(), + }); + const unsubscribe = handle.onError((error) => errors.push(error)); + await flushMicrotasks(); + unsubscribe(); + + handle.publishTyping(true); + await flushMicrotasks(); + + expect(errors).toEqual([{ operation: "join", status: 500 }]); + handle.disconnect(); + }); + + test("a client stuck unable to join backs off instead of re-posting on every publish call", async () => { + let clock = 0; + const calls: string[] = []; + const handle = connectPresence({ + roomUrl: ROOM_URL, + fetchImpl: scriptedFetch(calls, () => 500), + openEventSource: () => new FakeEventSource(), + now: () => clock, + }); + + await flushMicrotasks(); + expect(calls).toEqual([`${ROOM_URL}/join`]); // the initial attempt + + // A caller publishing cursor moves in a tight loop while unjoined + // must not turn into a `/join` per call. + handle.publishCursor({ x: 1, y: 1, surfaceVersion: 1 }); + handle.publishCursor({ x: 2, y: 2, surfaceVersion: 1 }); + handle.publishCursor({ x: 3, y: 3, surfaceVersion: 1 }); + await flushMicrotasks(); + expect(calls).toEqual([`${ROOM_URL}/join`]); + + // Still inside the backoff window after the first failure. + clock = 999; + handle.publishCursor({ x: 4, y: 4, surfaceVersion: 1 }); + await flushMicrotasks(); + expect(calls).toEqual([`${ROOM_URL}/join`]); + + // Past the backoff window: one retry is allowed. + clock = 1_000; + handle.publishCursor({ x: 5, y: 5, surfaceVersion: 1 }); + await flushMicrotasks(); + expect(calls).toEqual([`${ROOM_URL}/join`, `${ROOM_URL}/join`]); + + handle.disconnect(); + }); + + test("a join success resets the backoff, so a later failure streak starts from the base delay again", async () => { + let clock = 0; + const calls: string[] = []; + let failNextJoin = true; + let heartbeatStatus = 200; + const handle = connectPresence({ + roomUrl: ROOM_URL, + fetchImpl: scriptedFetch(calls, (path) => { + if (path.endsWith("/join")) return failNextJoin ? 500 : 200; + if (path.endsWith("/heartbeat")) return heartbeatStatus; + return 200; + }), + openEventSource: () => new FakeEventSource(), + now: () => clock, + }); + + await flushMicrotasks(); // fails once (attempt 1): 1s backoff scheduled + + clock = 1_000; + failNextJoin = false; + handle.publishCursor({ x: 1, y: 1, surfaceVersion: 1 }); + await flushMicrotasks(); // succeeds; joinFailureCount resets to 0 + + // Evict via a 404'd heartbeat and let the immediate rejoin attempt + // fail too: if the reset above hadn't happened, this failure would + // be attempt 3 (4s backoff) rather than a fresh attempt 1 (1s). + heartbeatStatus = 404; + failNextJoin = true; + calls.length = 0; + handle.publishTyping(true); // joined === true, so this sends a heartbeat + await flushMicrotasks(); + expect(calls).toEqual([`${ROOM_URL}/heartbeat`, `${ROOM_URL}/join`]); + + clock = 1_999; // short of a full second past the failure at clock 1_000 + handle.publishCursor({ x: 2, y: 2, surfaceVersion: 1 }); + await flushMicrotasks(); + expect(calls).toEqual([`${ROOM_URL}/heartbeat`, `${ROOM_URL}/join`]); + + clock = 2_000; // a full second past the failure + handle.publishCursor({ x: 3, y: 3, surfaceVersion: 1 }); + await flushMicrotasks(); + expect(calls).toEqual([ + `${ROOM_URL}/heartbeat`, + `${ROOM_URL}/join`, + `${ROOM_URL}/join`, + ]); + + handle.disconnect(); + }); + + test("a doc update that fails to post is redelivered once the next heartbeat succeeds", async () => { + const calls: { path: string; body: unknown }[] = []; + const doc = new Y.Doc(); + const errors: PresenceError[] = []; + let updateStatus = 500; + const fetchImpl: PresenceFetch = (url, init) => { + const body = JSON.parse(init.body) as unknown; + calls.push({ path: url, body }); + const status = url.endsWith("/update") ? updateStatus : 200; + return Promise.resolve({ + ok: status >= 200 && status < 300, + status, + json: () => Promise.resolve({}), + }); + }; + const handle = connectPresence({ + roomUrl: ROOM_URL, + fetchImpl, + openEventSource: () => new FakeEventSource(), + doc, + }); + handle.onError((error) => errors.push(error)); + + await flushMicrotasks(); // join settles + doc.getText("content").insert(0, "hello"); + await flushMicrotasks(); // the update POST fails and is queued + + expect(errors).toEqual([{ operation: "update", status: 500 }]); + const failedCall = calls.find((c) => c.path.endsWith("/update")); + calls.length = 0; + + // The next successful heartbeat redelivers the queued update with + // the exact same payload — nothing about the failed edit is lost. + updateStatus = 200; + handle.publishTyping(true); + await flushMicrotasks(); + + expect(calls.map((c) => c.path)).toEqual([ + `${ROOM_URL}/heartbeat`, + `${ROOM_URL}/update`, + ]); + expect(calls.find((c) => c.path.endsWith("/update"))?.body).toEqual( + failedCall?.body, + ); + // No second `onError` fire for the now-successful redelivery. + expect(errors).toEqual([{ operation: "update", status: 500 }]); + handle.disconnect(); + }); + + test("queued updates are redelivered in the order they were made", async () => { + const calls: { path: string; body: unknown }[] = []; + const doc = new Y.Doc(); + let updateStatus = 500; + const fetchImpl: PresenceFetch = (url, init) => { + const body = JSON.parse(init.body) as unknown; + calls.push({ path: url, body }); + const status = url.endsWith("/update") ? updateStatus : 200; + return Promise.resolve({ + ok: status >= 200 && status < 300, + status, + json: () => Promise.resolve({}), + }); + }; + const handle = connectPresence({ + roomUrl: ROOM_URL, + fetchImpl, + openEventSource: () => new FakeEventSource(), + doc, + }); + + await flushMicrotasks(); + doc.getText("content").insert(0, "a"); + await flushMicrotasks(); + doc.getText("content").insert(1, "b"); + await flushMicrotasks(); + const [firstFailedUpdate, secondFailedUpdate] = calls + .filter((c) => c.path.endsWith("/update")) + .map((c) => c.body); + calls.length = 0; + + updateStatus = 200; + handle.publishTyping(true); + await flushMicrotasks(); // redelivers the first queued update + handle.publishTyping(true); + await flushMicrotasks(); // redelivers the second + + const redeliveredUpdates = calls + .filter((c) => c.path.endsWith("/update")) + .map((c) => c.body); + expect(redeliveredUpdates).toEqual([firstFailedUpdate, secondFailedUpdate]); + handle.disconnect(); + }); +}); diff --git a/packages/presence/src/client.ts b/packages/presence/src/client.ts index a159dd219..0dbb6ecde 100644 --- a/packages/presence/src/client.ts +++ b/packages/presence/src/client.ts @@ -39,7 +39,24 @@ export interface PresenceEventSourceLike { export type PresenceFetch = ( url: string, init: { method: string; headers: Record; body: string }, -) => Promise<{ ok: boolean; json: () => Promise }>; +) => Promise<{ ok: boolean; status: number; json: () => Promise }>; + +/** The four requests this module ever POSTs — also each one's URL path + * segment under `roomUrl`. */ +export type PresenceOperation = "join" | "heartbeat" | "leave" | "update"; + +/** + * Reported through `PresenceHandle.onError` for every join/heartbeat/leave/ + * update request that never reached the server (a rejected `fetch`, no + * `status`) or came back non-2xx (`status` set). This is the only signal a + * consuming UI has that "connected" is a lie — the SSE stream opens and + * stays open regardless of whether join or any later heartbeat actually + * succeeded server-side. + */ +export interface PresenceError { + readonly operation: PresenceOperation; + readonly status?: number; +} export interface PresenceClientOptions { /** The room's base URL, e.g. `/api/tenants/tnt_1/presence/rooms/channel:chn_1`. */ @@ -65,6 +82,8 @@ export interface PresenceClientOptions { * from anything it did locally. */ readonly onSaved?: (info: { version: number; savedAt: number }) => void; + /** Injectable clock, for deterministic tests of rejoin backoff. */ + readonly now?: () => number; } /** @@ -91,12 +110,30 @@ export interface PresenceHandle { publishTyping(typing: boolean): void; /** Subscribes to every room snapshot; fires once with whatever has been received so far. */ subscribe(listener: (members: readonly PresenceState[]) => void): () => void; + /** Subscribes to every failed join/heartbeat/leave/update request. */ + onError(listener: (error: PresenceError) => void): () => void; /** Leaves the room and tears down the stream/heartbeat timer. */ disconnect(): void; } const DEFAULT_HEARTBEAT_INTERVAL_MS = 15_000; +// Exponential backoff for repeated join failures: without it, a client +// whose join keeps failing (server down, or a caller publishing cursor +// updates on every mouse move while unjoined) would re-POST `/join` as +// fast as each failed attempt resolves. Mirrors +// `packages/chat-ui/src/use-workbench-stream.ts`'s `backoffDelayMs` shape +// (not imported — that package depends on this one, not the reverse). +const REJOIN_BASE_DELAY_MS = 1_000; +const REJOIN_MAX_DELAY_MS = 30_000; + +function rejoinDelayMs(failureCount: number): number { + return Math.min( + REJOIN_BASE_DELAY_MS * 2 ** (failureCount - 1), + REJOIN_MAX_DELAY_MS, + ); +} + function parseMembers(data: string): readonly PresenceState[] { try { const parsed: unknown = JSON.parse(data); @@ -179,23 +216,126 @@ export function connectPresence( const heartbeatIntervalMs = options.heartbeatIntervalMs ?? DEFAULT_HEARTBEAT_INTERVAL_MS; const doc = options.doc; + const now = options.now ?? Date.now; const listeners = new Set<(members: readonly PresenceState[]) => void>(); + const errorListeners = new Set<(error: PresenceError) => void>(); let latestMembers: readonly PresenceState[] = []; let disconnected = false; + // Whether the last join attempt is known to have succeeded server-side. + // A heartbeat answered 404 (`not_joined`) — this room evicted us, most + // often ordinary heartbeat jitter racing the server's own timeout sweep + // — flips this back to `false` so the next heartbeat tick rejoins + // instead of heartbeating a membership that no longer exists. + let joined = false; + // Guards against a `publishCursor`/`publishTyping` call landing while + // the initial (or a rejoin) `join` request is still in flight from + // firing a second, redundant one. + let joinInFlight = false; + // Consecutive join failures since the last success — drives + // `rejoinDelayMs`'s backoff and resets to 0 the moment a join succeeds. + let joinFailureCount = 0; + // The earliest time `doJoin` will send another request while + // `joinFailureCount > 0` — the actual throttle; `joinFailureCount` by + // itself only picks the delay. + let nextJoinAttemptAt = 0; + // Local doc edits (base64-encoded) whose `/update` POST failed and + // hasn't yet been redelivered. Safe to replay in order once the + // connection is healthy again: a Yjs update is idempotent against a + // doc that has already applied it, so redelivering never double-edits + // — without this queue a dropped update is gone for good (see CL-7202). + let pendingUpdates: string[] = []; + let flushingPendingUpdates = false; const notify = (members: readonly PresenceState[]) => { latestMembers = members; for (const listener of listeners) listener(members); }; - const post = (path: string, body: unknown) => - fetchImpl(`${options.roomUrl}/${path}`, { + const reportError = (operation: PresenceOperation, status?: number) => { + const error: PresenceError = + status === undefined ? { operation } : { operation, status }; + for (const listener of errorListeners) listener(error); + }; + + const post = (operation: PresenceOperation, body: unknown) => + fetchImpl(`${options.roomUrl}/${operation}`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify(body ?? {}), }).catch(() => undefined); + /** + * Joins (or rejoins) the room. Called once at connect time and again + * whenever a heartbeat discovers the server no longer considers us + * joined — the only two situations this room's membership needs + * (re-)establishing. Backs off after repeated failures rather than + * retrying on every call: `publishCursor`/`publishTyping` land far more + * often than the heartbeat timer ticks, so without the backoff gate a + * still-unjoined client publishing cursor moves would re-POST `/join` + * as fast as each failed attempt resolves. + */ + function doJoin(): void { + if (joinInFlight) return; + if (joinFailureCount > 0 && now() < nextJoinAttemptAt) return; + joinInFlight = true; + void post("join", { displayName: options.displayName }).then((response) => { + joinInFlight = false; + if (disconnected) return; + if (response === undefined) { + joinFailureCount += 1; + nextJoinAttemptAt = now() + rejoinDelayMs(joinFailureCount); + reportError("join"); + return; + } + if (!response.ok) { + joinFailureCount += 1; + nextJoinAttemptAt = now() + rejoinDelayMs(joinFailureCount); + reportError("join", response.status); + return; + } + joinFailureCount = 0; + joined = true; + flushPendingUpdates(); + if (doc === undefined) return; + void response.json().then((body) => { + if (disconnected) return; + const docUpdate = docUpdateFromJoinResponse(body); + if (docUpdate !== undefined) applyRemoteUpdate(docUpdate); + }); + }); + } + + /** + * Sends a heartbeat (optionally carrying a cursor/typing patch) if the + * room still considers us joined, rejoining first if it doesn't — the + * self-healing half of CL-7202: a 404'd heartbeat now recovers instead + * of leaving the client heartbeating into the void while its own SSE + * stream keeps reading as live. + */ + function sendHeartbeat(patch: PresenceStatePatch): void { + if (!joined) { + doJoin(); + return; + } + void post("heartbeat", patch).then((response) => { + if (disconnected) return; + if (response === undefined) { + reportError("heartbeat"); + return; + } + if (!response.ok) { + reportError("heartbeat", response.status); + if (response.status === 404) { + joined = false; + doJoin(); + } + return; + } + flushPendingUpdates(); + }); + } + function applyRemoteUpdate(base64Update: string): void { if (doc === undefined) return; try { @@ -206,29 +346,64 @@ export function connectPresence( } } + /** + * POSTs one already-encoded doc update. `isRetry` distinguishes a fresh + * local edit (queued into `pendingUpdates` on failure) from a replay out + * of that queue (left in place on failure — `flushPendingUpdates` is + * already the thing re-adding it by not removing it). + */ + function postUpdate( + base64Update: string, + isRetry: boolean, + ): Promise { + return post("update", { update: base64Update }).then((response) => { + if (disconnected) return false; + if (response === undefined) { + reportError("update"); + if (!isRetry) pendingUpdates.push(base64Update); + return false; + } + if (!response.ok) { + reportError("update", response.status); + if (!isRetry) pendingUpdates.push(base64Update); + return false; + } + return true; + }); + } + + /** + * Redelivers queued failed updates in order, one at a time, stopping at + * the first failure so a later edit can never leapfrog an earlier one + * still waiting to land. Called whenever a join or heartbeat succeeds — + * the client's only signals that the connection is healthy again. + */ + function flushPendingUpdates(): void { + if (disconnected || flushingPendingUpdates || pendingUpdates.length === 0) { + return; + } + const [next, ...rest] = pendingUpdates; + if (next === undefined) return; + flushingPendingUpdates = true; + void postUpdate(next, true).then((succeeded) => { + flushingPendingUpdates = false; + if (!succeeded) return; + pendingUpdates = rest; + flushPendingUpdates(); + }); + } + let onLocalDocUpdate: ((update: Uint8Array, origin: unknown) => void) | undefined; if (doc !== undefined) { onLocalDocUpdate = (update, origin) => { if (disconnected || origin === REMOTE_UPDATE_ORIGIN) return; - void post("update", { update: encodeBase64(update) }); + void postUpdate(encodeBase64(update), false); }; doc.on("update", onLocalDocUpdate); } - void post("join", { displayName: options.displayName }) - .then((response) => { - if (doc === undefined || response === undefined || !response.ok) { - return undefined; - } - return response.json(); - }) - .then((body) => { - if (body === undefined) return; - const docUpdate = docUpdateFromJoinResponse(body); - if (docUpdate !== undefined) applyRemoteUpdate(docUpdate); - }) - .catch(() => undefined); + doJoin(); const source = openEventSource(`${options.roomUrl}/stream`); source.addEventListener("presence.state", (event) => { @@ -248,20 +423,18 @@ export function connectPresence( const heartbeatTimer = setInterval(() => { if (disconnected) return; - void post("heartbeat", {}); + sendHeartbeat({}); }, heartbeatIntervalMs); return { publishCursor(cursor) { if (disconnected) return; - const patch: PresenceStatePatch = { cursor }; - void post("heartbeat", patch); + sendHeartbeat({ cursor }); }, publishTyping(typing) { if (disconnected) return; - const patch: PresenceStatePatch = { typing }; - void post("heartbeat", patch); + sendHeartbeat({ typing }); }, subscribe(listener) { @@ -272,18 +445,27 @@ export function connectPresence( }; }, + onError(listener) { + errorListeners.add(listener); + return () => { + errorListeners.delete(listener); + }; + }, + disconnect() { if (disconnected) return; disconnected = true; clearInterval(heartbeatTimer); source.close(); listeners.clear(); + errorListeners.clear(); if (doc !== undefined && onLocalDocUpdate !== undefined) { doc.off("update", onLocalDocUpdate); } // Best-effort: the room drops this principal on its own heartbeat // timeout even if this never arrives (page unload racing the - // request), so a failed leave is not a correctness bug. + // request), so a failed leave is not reported as a `PresenceError` + // — there is no listener left to hear it by the time it would fire. void post("leave", {}); }, };