fix(web): harden cross-tab SSE fallback

Stop coordinated election and leadership work when BroadcastChannel posting fails, so tabs degrade cleanly to per-subscriber fallback without stale resync or heartbeat side effects. Expand election coverage for the edge cases called out in the coordination plan.
This commit is contained in:
Bryan Helmkamp 2026-05-04 15:24:25 -04:00
parent ade721ae65
commit 38726666af
No known key found for this signature in database
2 changed files with 189 additions and 15 deletions

View file

@ -40,6 +40,7 @@ class FakeEventSource {
class FakeBroadcastChannel implements BroadcastChannelLike {
static channels = new Set<FakeBroadcastChannel>();
static muted = false;
static throwOnTypes = new Set<CrossTabSseMessage["type"]>();
onmessage: ((event: { data: unknown }) => void) | null = null;
closed = false;
@ -49,6 +50,9 @@ class FakeBroadcastChannel implements BroadcastChannelLike {
}
postMessage(message: CrossTabSseMessage) {
if (FakeBroadcastChannel.throwOnTypes.has(message.type)) {
throw new Error(`postMessage failed for ${message.type}`);
}
if (FakeBroadcastChannel.muted) return;
const recipients = [...FakeBroadcastChannel.channels].filter(
(channel) => channel !== this && !channel.closed && channel.name === this.name,
@ -61,6 +65,15 @@ class FakeBroadcastChannel implements BroadcastChannelLike {
});
}
static broadcastExternal(message: CrossTabSseMessage) {
queueMicrotask(() => {
for (const channel of FakeBroadcastChannel.channels) {
if (channel.closed) continue;
channel.onmessage?.({ data: { ...message } });
}
});
}
close() {
this.closed = true;
FakeBroadcastChannel.channels.delete(this);
@ -72,6 +85,7 @@ class FakeBroadcastChannel implements BroadcastChannelLike {
}
FakeBroadcastChannel.channels.clear();
FakeBroadcastChannel.muted = false;
FakeBroadcastChannel.throwOnTypes.clear();
}
}
@ -166,6 +180,40 @@ describe("subscribeToCrossTabSse", () => {
expect(keysByTab.get("c")).toEqual(["event"]);
});
test("board and run subscriptions coexist on the same global stream", async () => {
const harness = newHarness();
const coordinator = harness.createTab("a");
const boardKeys: string[] = [];
const runKeys: string[] = [];
subscribeForEvent(coordinator, {
subscriptionKey: "board",
keys: boardKeys,
resolveInvalidation: (payload) => ({
keys: payload.event === "run.running" ? ["board"] : [],
}),
resyncKeys: () => ["board-resync"],
});
subscribeForEvent(coordinator, {
subscriptionKey: "run:run-1",
keys: runKeys,
resolveInvalidation: (payload) => ({
keys: payload.event === "run.running" && payload.run_id === "run-1" ? ["run"] : [],
}),
resyncKeys: () => ["run-resync"],
});
await waitFor(() => harness.openSources().length === 1);
boardKeys.length = 0;
runKeys.length = 0;
harness.openSources()[0].emit(runEvent({ id: "evt-coexist", runId: "run-1", seq: 1 }));
expect(boardKeys).toEqual(["board"]);
expect(runKeys).toEqual(["run"]);
expect(harness.openSources().map((source) => source.url)).toEqual(["/api/v1/attach"]);
});
test("dedupes duplicate event ids until TTL or max-size eviction", async () => {
const harness = newHarness();
const keys: string[] = [];
@ -218,7 +266,8 @@ describe("subscribeToCrossTabSse", () => {
subscribeForRunEvent(harness.createTab("b", "visible"), []);
subscribeForRunEvent(harness.createTab("a", "visible"), []);
await waitFor(() => harness.openSources().length === 1 && harness.openSources()[0].owner === "a");
await waitFor(() => harness.openSources().length === 1 && harness.openSources()[0].owner !== "z");
expect(harness.openSources().map((source) => source.owner)).toEqual(["a"]);
});
test("a lower lexical follower does not preempt a fresh visible leader", async () => {
@ -250,6 +299,46 @@ describe("subscribeToCrossTabSse", () => {
expect(followerKeys).toContain("resync");
});
test("simultaneous stale leader elections resolve to the lexical winner", async () => {
const harness = newHarness();
subscribeForRunEvent(harness.createTab("z"), []);
await waitFor(() => harness.openSources().length === 1 && harness.openSources()[0].owner === "z");
subscribeForRunEvent(harness.createTab("b"), []);
subscribeForRunEvent(harness.createTab("a"), []);
await sleep(TEST_TIMING.heartbeatMs * 2);
harness.coordinators.get("z")?.close();
harness.now += TEST_TIMING.leaderStaleMs + TEST_TIMING.heartbeatMs + 1;
await waitFor(() => harness.openSources().length === 1 && harness.openSources()[0].owner === "a");
});
test("hidden leader ignores candidates for old observed leadership", async () => {
const harness = newHarness();
subscribeForRunEvent(harness.createTab("z", "hidden"), []);
await waitFor(() => harness.openSources().length === 1 && harness.openSources()[0].owner === "z");
const hiddenSource = harness.openSources()[0];
FakeBroadcastChannel.broadcastExternal({
type: "candidate",
version: 1,
tabId: "ghost",
sentAt: harness.now,
candidateId: "ghost",
candidateGeneration: 1,
visibility: "visible",
observedLeaderId: "z",
observedGeneration: 0,
reason: "hidden-leader",
});
await sleep(TEST_TIMING.electionJitterMs * 2);
expect(hiddenSource.closed).toBe(false);
expect(harness.openSources().map((source) => source.owner)).toEqual(["z"]);
});
test("same-generation split brain converges to the higher-priority visible leader", async () => {
const harness = newHarness();
FakeBroadcastChannel.muted = true;
@ -278,6 +367,29 @@ describe("subscribeToCrossTabSse", () => {
expect(keys).toEqual([]);
});
test("old leader heartbeats are ignored after takeover", async () => {
const harness = newHarness();
subscribeForRunEvent(harness.createTab("z", "hidden"), []);
await waitFor(() => harness.openSources().length === 1 && harness.openSources()[0].owner === "z");
subscribeForRunEvent(harness.createTab("a", "visible"), []);
await waitFor(() => harness.openSources().length === 1 && harness.openSources()[0].owner === "a");
FakeBroadcastChannel.broadcastExternal({
type: "heartbeat",
version: 1,
tabId: "z",
sentAt: harness.now,
leaderId: "z",
generation: 1,
visibility: "hidden",
});
await sleep(TEST_TIMING.heartbeatMs * 2);
expect(harness.openSources().map((source) => source.owner)).toEqual(["a"]);
});
test("last unsubscribe closes the leader source and releases leadership", async () => {
const harness = newHarness();
const cleanup = subscribeForRunEvent(harness.createTab("a"), []);
@ -320,6 +432,41 @@ describe("subscribeToCrossTabSse", () => {
expect(fallbackStarted).toBe(1);
expect(fallbackStopped).toBe(1);
});
test("postMessage failure after initialization degrades to fallback without coordinated resync", async () => {
const harness = newHarness();
const coordinator = harness.createTab("a");
const keys: string[] = [];
let fallbackStarted = 0;
let fallbackStopped = 0;
FakeBroadcastChannel.throwOnTypes.add("leader-changed");
const cleanup = subscribeToCrossTabSse<EventPayload>({
coordinator,
subscriptionKey: "throwing-channel",
mutate: ((key: string) => {
keys.push(key);
return Promise.resolve();
}) as MutateFn,
resolveInvalidation: () => ({ keys: ["event"] }),
resyncKeys: () => ["resync"],
fallbackSubscribe: () => {
fallbackStarted += 1;
return () => {
fallbackStopped += 1;
};
},
debounceMs: 0,
});
await waitFor(() => fallbackStarted === 1);
expect(harness.openSources()).toEqual([]);
expect(keys).toEqual([]);
cleanup();
expect(fallbackStopped).toBe(1);
});
});
function newHarness() {
@ -329,17 +476,39 @@ function newHarness() {
}
function subscribeForRunEvent(coordinator: CrossTabSseCoordinator, keys: string[]) {
return subscribeToCrossTabSse<EventPayload>({
coordinator,
return subscribeForEvent(coordinator, {
subscriptionKey: "run-feed",
mutate: ((key: string) => {
keys.push(key);
return Promise.resolve();
}) as MutateFn,
keys,
resolveInvalidation: (payload) => ({
keys: payload.event === "run.running" ? ["event"] : [],
}),
resyncKeys: () => ["resync"],
});
}
function subscribeForEvent(
coordinator: CrossTabSseCoordinator,
{
subscriptionKey,
keys,
resolveInvalidation,
resyncKeys,
}: {
subscriptionKey: string;
keys: string[];
resolveInvalidation: (payload: EventPayload) => { keys: string[] };
resyncKeys: () => string[];
},
) {
return subscribeToCrossTabSse<EventPayload>({
coordinator,
subscriptionKey,
mutate: ((key: string) => {
keys.push(key);
return Promise.resolve();
}) as MutateFn,
resolveInvalidation,
resyncKeys,
fallbackSubscribe: () => {
throw new Error("fallback should not be used");
},

View file

@ -594,7 +594,7 @@ export class CrossTabSseCoordinator {
this.ownCandidate = candidate;
this.candidates.set(candidateKey(candidate), candidate);
this.post(candidate);
if (!this.post(candidate)) return;
this.candidateTimer = setTimeout(() => {
this.completeCandidacy(candidate);
}, this.timing.electionJitterMs);
@ -651,7 +651,7 @@ export class CrossTabSseCoordinator {
this.handleLeaderEventSourceMessage(message.data);
};
this.post({
const announced = this.post({
type: "leader-changed",
version: MESSAGE_VERSION,
tabId: this.tabId,
@ -660,6 +660,7 @@ export class CrossTabSseCoordinator {
generation,
visibility: this.currentVisibility(),
});
if (!announced) return;
this.startHeartbeat();
this.resyncAll();
}
@ -761,14 +762,14 @@ export class CrossTabSseCoordinator {
if (this.heartbeatTimer) {
clearInterval(this.heartbeatTimer);
}
this.sendHeartbeat();
if (!this.sendHeartbeat()) return;
this.heartbeatTimer = setInterval(() => {
this.sendHeartbeat();
}, this.timing.heartbeatMs);
}
private sendHeartbeat() {
if (!this.isLeader) return;
private sendHeartbeat(): boolean {
if (!this.isLeader) return false;
const visibility = this.currentVisibility();
this.leader = {
leaderId: this.tabId,
@ -776,7 +777,7 @@ export class CrossTabSseCoordinator {
visibility,
lastSeen: this.now(),
};
this.post({
return this.post({
type: "heartbeat",
version: MESSAGE_VERSION,
tabId: this.tabId,
@ -806,7 +807,7 @@ export class CrossTabSseCoordinator {
this.leader = null;
if (broadcast) {
this.post({
const released = this.post({
type: "release",
version: MESSAGE_VERSION,
tabId: this.tabId,
@ -814,6 +815,7 @@ export class CrossTabSseCoordinator {
leaderId: this.tabId,
generation,
});
if (!released) return;
this.post({
type: "resync",
version: MESSAGE_VERSION,
@ -888,7 +890,10 @@ export class CrossTabSseCoordinator {
}
private post(message: CrossTabSseMessage): boolean {
if (!this.channel) return false;
if (!this.channel) {
this.degradeToFallback();
return false;
}
try {
this.channel.postMessage(message);
return true;