From 480bdbc8de76994982716e91130c42631a4cca37 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 04:55:47 -0700 Subject: [PATCH 1/4] Add tests for presence client error reporting and rejoin (CL-7202) Reproduces CL-7202: connectPresence blind-posts every join/heartbeat/ leave/update request, never checks response.ok, exposes no way for a caller to hear about a failure, and never rejoins after a failed join or a heartbeat the server has already forgotten about (404 not_joined) -- all while its EventSource stays open regardless, so the UI keeps reading as live. --- packages/presence/src/client.test.ts | 182 ++++++++++++++++++++++++++- 1 file changed, 181 insertions(+), 1 deletion(-) diff --git a/packages/presence/src/client.test.ts b/packages/presence/src/client.test.ts index 652e9cbbd..e4256f7b0 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,153 @@ 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 () => { + const calls: string[] = []; + const handle = connectPresence({ + roomUrl: ROOM_URL, + fetchImpl: scriptedFetch(calls, () => 500), + openEventSource: () => new FakeEventSource(), + }); + + await flushMicrotasks(); + 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(); + }); +}); From 20ee4b3040698766dc5871cf7d0bf445bde482b5 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 04:55:53 -0700 Subject: [PATCH 2/4] presence client: report request failures and rejoin after eviction (CL-7202) connectPresence blind-posted every join/heartbeat/leave/update request and never checked response.ok, so a failed join or a heartbeat 404'd by the server (not_joined) went unnoticed while the SSE stream stayed open regardless -- the UI kept reading as live with no signal that membership had been lost. PresenceHandle now exposes onError, firing for every request that never reached the server or came back non-2xx. Heartbeats are gated behind a joined flag: a heartbeat is only sent once join has actually succeeded, and a 404 response flips that flag back off and triggers an immediate rejoin -- the same self-healing a publishCursor/ publishTyping call gets if it fires before the initial join settles. --- packages/presence/src/client.ts | 142 ++++++++++++++++++++++++++------ 1 file changed, 119 insertions(+), 23 deletions(-) diff --git a/packages/presence/src/client.ts b/packages/presence/src/client.ts index a159dd219..133fef545 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`. */ @@ -91,6 +108,8 @@ 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; } @@ -181,21 +200,96 @@ export function connectPresence( const doc = options.doc; 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; 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. + */ + function doJoin(): void { + if (joinInFlight) return; + joinInFlight = true; + void post("join", { displayName: options.displayName }).then((response) => { + joinInFlight = false; + if (disconnected) return; + if (response === undefined) { + reportError("join"); + return; + } + if (!response.ok) { + reportError("join", response.status); + return; + } + joined = true; + 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(); + } + } + }); + } + function applyRemoteUpdate(base64Update: string): void { if (doc === undefined) return; try { @@ -211,24 +305,19 @@ export function connectPresence( if (doc !== undefined) { onLocalDocUpdate = (update, origin) => { if (disconnected || origin === REMOTE_UPDATE_ORIGIN) return; - void post("update", { update: encodeBase64(update) }); + void post("update", { update: encodeBase64(update) }).then((response) => { + if (disconnected) return; + if (response === undefined) { + reportError("update"); + } else if (!response.ok) { + reportError("update", response.status); + } + }); }; 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 +337,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 +359,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", {}); }, }; From 06198c127a347cb1dc35ba096cc6ae575060abdc Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 05:39:45 -0700 Subject: [PATCH 3/4] Back off repeated join failures in the presence client A caller publishing cursor/typing updates while unjoined called doJoin on every publish; with no delay between attempts, a client stuck unable to join (or a room that keeps evicting it) would re-POST /join as fast as each failed attempt resolved. Joins now back off exponentially after a failure and reset the moment one succeeds. --- packages/presence/src/client.test.ts | 89 ++++++++++++++++++++++++++++ packages/presence/src/client.ts | 38 +++++++++++- 2 files changed, 126 insertions(+), 1 deletion(-) diff --git a/packages/presence/src/client.test.ts b/packages/presence/src/client.test.ts index e4256f7b0..84738acd4 100644 --- a/packages/presence/src/client.test.ts +++ b/packages/presence/src/client.test.ts @@ -381,14 +381,17 @@ describe("connectPresence: error reporting and rejoin (CL-7202)", () => { }); 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(); @@ -486,4 +489,90 @@ describe("connectPresence: error reporting and rejoin (CL-7202)", () => { 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(); + }); }); diff --git a/packages/presence/src/client.ts b/packages/presence/src/client.ts index 133fef545..051fe670b 100644 --- a/packages/presence/src/client.ts +++ b/packages/presence/src/client.ts @@ -82,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; } /** @@ -116,6 +118,22 @@ export interface PresenceHandle { 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); @@ -198,6 +216,7 @@ 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>(); @@ -213,6 +232,13 @@ export function connectPresence( // 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; const notify = (members: readonly PresenceState[]) => { latestMembers = members; @@ -236,22 +262,32 @@ export function connectPresence( * 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. + * (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; if (doc === undefined) return; void response.json().then((body) => { From 912eab282597e27ce6dd351992ea923df306efcf Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 05:54:20 -0700 Subject: [PATCH 4/4] Retry a doc update that failed to post instead of dropping it A local Yjs edit whose /update POST failed was only ever surfaced through onError; the edit itself was gone for good, so a dropped update meant silent, permanent divergence from the server's doc for a collaborative document. Failed updates now queue and are redelivered in order on the next successful join or heartbeat -- safe because Yjs updates are idempotent against a doc that has already applied them. --- packages/presence/src/client.test.ts | 93 ++++++++++++++++++++++++++++ packages/presence/src/client.ts | 66 +++++++++++++++++--- 2 files changed, 151 insertions(+), 8 deletions(-) diff --git a/packages/presence/src/client.test.ts b/packages/presence/src/client.test.ts index 84738acd4..c864caf4e 100644 --- a/packages/presence/src/client.test.ts +++ b/packages/presence/src/client.test.ts @@ -575,4 +575,97 @@ describe("connectPresence: error reporting and rejoin (CL-7202)", () => { 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 051fe670b..0dbb6ecde 100644 --- a/packages/presence/src/client.ts +++ b/packages/presence/src/client.ts @@ -239,6 +239,13 @@ export function connectPresence( // `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; @@ -289,6 +296,7 @@ export function connectPresence( } joinFailureCount = 0; joined = true; + flushPendingUpdates(); if (doc === undefined) return; void response.json().then((body) => { if (disconnected) return; @@ -322,7 +330,9 @@ export function connectPresence( joined = false; doJoin(); } + return; } + flushPendingUpdates(); }); } @@ -336,19 +346,59 @@ 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) }).then((response) => { - if (disconnected) return; - if (response === undefined) { - reportError("update"); - } else if (!response.ok) { - reportError("update", response.status); - } - }); + void postUpdate(encodeBase64(update), false); }; doc.on("update", onLocalDocUpdate); }