Skip to content

Commit 30622cf

Browse files
Merge pull request #482 from corbitsdev/cl-7202-presence-client-error-handling
presence client: report request failures and rejoin after eviction (CL-7202)
2 parents 5ef58a2 + 912eab2 commit 30622cf

4 files changed

Lines changed: 669 additions & 29 deletions

File tree

‎packages/presence/src/client.test.ts‎

Lines changed: 363 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@ import { describe, expect, test } from "bun:test";
22
import * as Y from "yjs";
33
import {
44
connectPresence,
5+
type PresenceError,
56
type PresenceEventSourceLike,
67
type PresenceFetch,
78
type PresenceStreamEvent,
@@ -48,11 +49,35 @@ function fakeFetch(
4849
calls.push({ path: url, body });
4950
return Promise.resolve({
5051
ok: true,
52+
status: 200,
5153
json: () => Promise.resolve(joinResponse),
5254
});
5355
};
5456
}
5557

58+
/** A `PresenceFetch` whose response per operation is scripted by
59+
* `statusFor`, defaulting to 200 for anything unlisted — for exercising
60+
* failure and rejoin paths `fakeFetch`'s always-succeeds shape can't. */
61+
function scriptedFetch(
62+
calls: string[],
63+
statusFor: (path: string) => number | "reject",
64+
): PresenceFetch {
65+
return (url) => {
66+
calls.push(url);
67+
const status = statusFor(url);
68+
if (status === "reject") return Promise.reject(new Error("network down"));
69+
return Promise.resolve({
70+
ok: status >= 200 && status < 300,
71+
status,
72+
json: () => Promise.resolve({}),
73+
});
74+
};
75+
}
76+
77+
async function flushMicrotasks(): Promise<void> {
78+
for (let i = 0; i < 5; i++) await Promise.resolve();
79+
}
80+
5681
describe("connectPresence", () => {
5782
test("joins immediately and opens the room's SSE stream", () => {
5883
const calls: { path: string; body: unknown }[] = [];
@@ -104,13 +129,18 @@ describe("connectPresence", () => {
104129
handle.disconnect();
105130
});
106131

107-
test("publishCursor and publishTyping post heartbeats with the patch", () => {
132+
test("publishCursor and publishTyping post heartbeats with the patch", async () => {
108133
const calls: { path: string; body: unknown }[] = [];
109134
const handle = connectPresence({
110135
roomUrl: "/rooms/channel:chn_1",
111136
fetchImpl: fakeFetch(calls),
112137
openEventSource: () => new FakeEventSource(),
113138
});
139+
// Let the initial join's response settle before publishing: a patch
140+
// published before the room has confirmed the join rejoins instead
141+
// of heartbeating a membership the server doesn't have yet (CL-7202).
142+
await Promise.resolve();
143+
await Promise.resolve();
114144
calls.length = 0; // drop the initial join call
115145

116146
handle.publishCursor({ x: 5, y: 6, surfaceVersion: 1 });
@@ -307,3 +337,335 @@ describe("connectPresence: doc sync", () => {
307337
handle.disconnect();
308338
});
309339
});
340+
341+
// CL-7202: the client used to blind-post every join/heartbeat/leave/update
342+
// request (`.catch(() => undefined)`, no `response.ok` check anywhere, no
343+
// way for a caller to hear about a failure), and never rejoined after a
344+
// failed join or a heartbeat the server had already forgotten about (404
345+
// `not_joined`) — all while the SSE stream stayed open regardless, so the
346+
// UI kept reading as live.
347+
describe("connectPresence: error reporting and rejoin (CL-7202)", () => {
348+
const ROOM_URL = "/rooms/channel:chn_1";
349+
350+
test("a join request that never reaches the server is reported through onError, not swallowed", async () => {
351+
const calls: string[] = [];
352+
const errors: PresenceError[] = [];
353+
const handle = connectPresence({
354+
roomUrl: ROOM_URL,
355+
fetchImpl: scriptedFetch(calls, () => "reject"),
356+
openEventSource: () => new FakeEventSource(),
357+
});
358+
handle.onError((error) => errors.push(error));
359+
360+
await flushMicrotasks();
361+
362+
expect(calls).toEqual([`${ROOM_URL}/join`]);
363+
expect(errors).toEqual([{ operation: "join" }]);
364+
handle.disconnect();
365+
});
366+
367+
test("a non-ok join response is reported through onError with its status", async () => {
368+
const calls: string[] = [];
369+
const errors: PresenceError[] = [];
370+
const handle = connectPresence({
371+
roomUrl: ROOM_URL,
372+
fetchImpl: scriptedFetch(calls, () => 500),
373+
openEventSource: () => new FakeEventSource(),
374+
});
375+
handle.onError((error) => errors.push(error));
376+
377+
await flushMicrotasks();
378+
379+
expect(errors).toEqual([{ operation: "join", status: 500 }]);
380+
handle.disconnect();
381+
});
382+
383+
test("publishing before the room has confirmed the join rejoins instead of heartbeating a membership it doesn't have yet", async () => {
384+
let clock = 0;
385+
const calls: string[] = [];
386+
const handle = connectPresence({
387+
roomUrl: ROOM_URL,
388+
fetchImpl: scriptedFetch(calls, () => 500),
389+
openEventSource: () => new FakeEventSource(),
390+
now: () => clock,
391+
});
392+
393+
await flushMicrotasks();
394+
clock = 1_000; // past the first failure's backoff window
395+
handle.publishTyping(true);
396+
await flushMicrotasks();
397+
398+
expect(calls).toEqual([`${ROOM_URL}/join`, `${ROOM_URL}/join`]);
399+
handle.disconnect();
400+
});
401+
402+
test("a heartbeat 404 (self-eviction) triggers an automatic rejoin", async () => {
403+
const calls: string[] = [];
404+
let joinCount = 0;
405+
const handle = connectPresence({
406+
roomUrl: ROOM_URL,
407+
fetchImpl: scriptedFetch(calls, (path) => {
408+
if (path.endsWith("/join")) {
409+
joinCount += 1;
410+
return 200;
411+
}
412+
if (path.endsWith("/heartbeat")) return 404;
413+
return 200;
414+
}),
415+
openEventSource: () => new FakeEventSource(),
416+
});
417+
const errors: PresenceError[] = [];
418+
handle.onError((error) => errors.push(error));
419+
420+
await flushMicrotasks();
421+
expect(joinCount).toBe(1);
422+
423+
handle.publishCursor({ x: 1, y: 2, surfaceVersion: 1 });
424+
await flushMicrotasks();
425+
426+
expect(calls).toEqual([
427+
`${ROOM_URL}/join`,
428+
`${ROOM_URL}/heartbeat`,
429+
`${ROOM_URL}/join`,
430+
]);
431+
expect(errors).toEqual([{ operation: "heartbeat", status: 404 }]);
432+
expect(joinCount).toBe(2);
433+
handle.disconnect();
434+
});
435+
436+
test("a heartbeat succeeds normally once join has succeeded, without rejoining", async () => {
437+
const calls: string[] = [];
438+
const handle = connectPresence({
439+
roomUrl: ROOM_URL,
440+
fetchImpl: scriptedFetch(calls, () => 200),
441+
openEventSource: () => new FakeEventSource(),
442+
});
443+
444+
await flushMicrotasks();
445+
handle.publishTyping(true);
446+
await flushMicrotasks();
447+
448+
expect(calls).toEqual([`${ROOM_URL}/join`, `${ROOM_URL}/heartbeat`]);
449+
handle.disconnect();
450+
});
451+
452+
test("a failed doc update is reported through onError", async () => {
453+
const calls: string[] = [];
454+
const doc = new Y.Doc();
455+
const errors: PresenceError[] = [];
456+
const handle = connectPresence({
457+
roomUrl: ROOM_URL,
458+
fetchImpl: scriptedFetch(calls, (path) =>
459+
path.endsWith("/update") ? 413 : 200,
460+
),
461+
openEventSource: () => new FakeEventSource(),
462+
doc,
463+
});
464+
handle.onError((error) => errors.push(error));
465+
466+
await flushMicrotasks();
467+
doc.getText("content").insert(0, "hello");
468+
await flushMicrotasks();
469+
470+
expect(errors).toEqual([{ operation: "update", status: 413 }]);
471+
handle.disconnect();
472+
});
473+
474+
test("onError's unsubscribe stops further delivery", async () => {
475+
const calls: string[] = [];
476+
const errors: PresenceError[] = [];
477+
const handle = connectPresence({
478+
roomUrl: ROOM_URL,
479+
fetchImpl: scriptedFetch(calls, () => 500),
480+
openEventSource: () => new FakeEventSource(),
481+
});
482+
const unsubscribe = handle.onError((error) => errors.push(error));
483+
await flushMicrotasks();
484+
unsubscribe();
485+
486+
handle.publishTyping(true);
487+
await flushMicrotasks();
488+
489+
expect(errors).toEqual([{ operation: "join", status: 500 }]);
490+
handle.disconnect();
491+
});
492+
493+
test("a client stuck unable to join backs off instead of re-posting on every publish call", async () => {
494+
let clock = 0;
495+
const calls: string[] = [];
496+
const handle = connectPresence({
497+
roomUrl: ROOM_URL,
498+
fetchImpl: scriptedFetch(calls, () => 500),
499+
openEventSource: () => new FakeEventSource(),
500+
now: () => clock,
501+
});
502+
503+
await flushMicrotasks();
504+
expect(calls).toEqual([`${ROOM_URL}/join`]); // the initial attempt
505+
506+
// A caller publishing cursor moves in a tight loop while unjoined
507+
// must not turn into a `/join` per call.
508+
handle.publishCursor({ x: 1, y: 1, surfaceVersion: 1 });
509+
handle.publishCursor({ x: 2, y: 2, surfaceVersion: 1 });
510+
handle.publishCursor({ x: 3, y: 3, surfaceVersion: 1 });
511+
await flushMicrotasks();
512+
expect(calls).toEqual([`${ROOM_URL}/join`]);
513+
514+
// Still inside the backoff window after the first failure.
515+
clock = 999;
516+
handle.publishCursor({ x: 4, y: 4, surfaceVersion: 1 });
517+
await flushMicrotasks();
518+
expect(calls).toEqual([`${ROOM_URL}/join`]);
519+
520+
// Past the backoff window: one retry is allowed.
521+
clock = 1_000;
522+
handle.publishCursor({ x: 5, y: 5, surfaceVersion: 1 });
523+
await flushMicrotasks();
524+
expect(calls).toEqual([`${ROOM_URL}/join`, `${ROOM_URL}/join`]);
525+
526+
handle.disconnect();
527+
});
528+
529+
test("a join success resets the backoff, so a later failure streak starts from the base delay again", async () => {
530+
let clock = 0;
531+
const calls: string[] = [];
532+
let failNextJoin = true;
533+
let heartbeatStatus = 200;
534+
const handle = connectPresence({
535+
roomUrl: ROOM_URL,
536+
fetchImpl: scriptedFetch(calls, (path) => {
537+
if (path.endsWith("/join")) return failNextJoin ? 500 : 200;
538+
if (path.endsWith("/heartbeat")) return heartbeatStatus;
539+
return 200;
540+
}),
541+
openEventSource: () => new FakeEventSource(),
542+
now: () => clock,
543+
});
544+
545+
await flushMicrotasks(); // fails once (attempt 1): 1s backoff scheduled
546+
547+
clock = 1_000;
548+
failNextJoin = false;
549+
handle.publishCursor({ x: 1, y: 1, surfaceVersion: 1 });
550+
await flushMicrotasks(); // succeeds; joinFailureCount resets to 0
551+
552+
// Evict via a 404'd heartbeat and let the immediate rejoin attempt
553+
// fail too: if the reset above hadn't happened, this failure would
554+
// be attempt 3 (4s backoff) rather than a fresh attempt 1 (1s).
555+
heartbeatStatus = 404;
556+
failNextJoin = true;
557+
calls.length = 0;
558+
handle.publishTyping(true); // joined === true, so this sends a heartbeat
559+
await flushMicrotasks();
560+
expect(calls).toEqual([`${ROOM_URL}/heartbeat`, `${ROOM_URL}/join`]);
561+
562+
clock = 1_999; // short of a full second past the failure at clock 1_000
563+
handle.publishCursor({ x: 2, y: 2, surfaceVersion: 1 });
564+
await flushMicrotasks();
565+
expect(calls).toEqual([`${ROOM_URL}/heartbeat`, `${ROOM_URL}/join`]);
566+
567+
clock = 2_000; // a full second past the failure
568+
handle.publishCursor({ x: 3, y: 3, surfaceVersion: 1 });
569+
await flushMicrotasks();
570+
expect(calls).toEqual([
571+
`${ROOM_URL}/heartbeat`,
572+
`${ROOM_URL}/join`,
573+
`${ROOM_URL}/join`,
574+
]);
575+
576+
handle.disconnect();
577+
});
578+
579+
test("a doc update that fails to post is redelivered once the next heartbeat succeeds", async () => {
580+
const calls: { path: string; body: unknown }[] = [];
581+
const doc = new Y.Doc();
582+
const errors: PresenceError[] = [];
583+
let updateStatus = 500;
584+
const fetchImpl: PresenceFetch = (url, init) => {
585+
const body = JSON.parse(init.body) as unknown;
586+
calls.push({ path: url, body });
587+
const status = url.endsWith("/update") ? updateStatus : 200;
588+
return Promise.resolve({
589+
ok: status >= 200 && status < 300,
590+
status,
591+
json: () => Promise.resolve({}),
592+
});
593+
};
594+
const handle = connectPresence({
595+
roomUrl: ROOM_URL,
596+
fetchImpl,
597+
openEventSource: () => new FakeEventSource(),
598+
doc,
599+
});
600+
handle.onError((error) => errors.push(error));
601+
602+
await flushMicrotasks(); // join settles
603+
doc.getText("content").insert(0, "hello");
604+
await flushMicrotasks(); // the update POST fails and is queued
605+
606+
expect(errors).toEqual([{ operation: "update", status: 500 }]);
607+
const failedCall = calls.find((c) => c.path.endsWith("/update"));
608+
calls.length = 0;
609+
610+
// The next successful heartbeat redelivers the queued update with
611+
// the exact same payload — nothing about the failed edit is lost.
612+
updateStatus = 200;
613+
handle.publishTyping(true);
614+
await flushMicrotasks();
615+
616+
expect(calls.map((c) => c.path)).toEqual([
617+
`${ROOM_URL}/heartbeat`,
618+
`${ROOM_URL}/update`,
619+
]);
620+
expect(calls.find((c) => c.path.endsWith("/update"))?.body).toEqual(
621+
failedCall?.body,
622+
);
623+
// No second `onError` fire for the now-successful redelivery.
624+
expect(errors).toEqual([{ operation: "update", status: 500 }]);
625+
handle.disconnect();
626+
});
627+
628+
test("queued updates are redelivered in the order they were made", async () => {
629+
const calls: { path: string; body: unknown }[] = [];
630+
const doc = new Y.Doc();
631+
let updateStatus = 500;
632+
const fetchImpl: PresenceFetch = (url, init) => {
633+
const body = JSON.parse(init.body) as unknown;
634+
calls.push({ path: url, body });
635+
const status = url.endsWith("/update") ? updateStatus : 200;
636+
return Promise.resolve({
637+
ok: status >= 200 && status < 300,
638+
status,
639+
json: () => Promise.resolve({}),
640+
});
641+
};
642+
const handle = connectPresence({
643+
roomUrl: ROOM_URL,
644+
fetchImpl,
645+
openEventSource: () => new FakeEventSource(),
646+
doc,
647+
});
648+
649+
await flushMicrotasks();
650+
doc.getText("content").insert(0, "a");
651+
await flushMicrotasks();
652+
doc.getText("content").insert(1, "b");
653+
await flushMicrotasks();
654+
const [firstFailedUpdate, secondFailedUpdate] = calls
655+
.filter((c) => c.path.endsWith("/update"))
656+
.map((c) => c.body);
657+
calls.length = 0;
658+
659+
updateStatus = 200;
660+
handle.publishTyping(true);
661+
await flushMicrotasks(); // redelivers the first queued update
662+
handle.publishTyping(true);
663+
await flushMicrotasks(); // redelivers the second
664+
665+
const redeliveredUpdates = calls
666+
.filter((c) => c.path.endsWith("/update"))
667+
.map((c) => c.body);
668+
expect(redeliveredUpdates).toEqual([firstFailedUpdate, secondFailedUpdate]);
669+
handle.disconnect();
670+
});
671+
});

0 commit comments

Comments
 (0)