feat(web): coordinate SSE subscriptions across tabs

Elect a single browser tab to own the global attach stream and broadcast run events to sibling tabs. Keep the existing per-tab EventSource path as the fallback when cross-tab coordination is unavailable.
This commit is contained in:
Bryan Helmkamp 2026-05-04 14:54:39 -04:00
parent f39e512990
commit ade721ae65
No known key found for this signature in database
6 changed files with 1852 additions and 36 deletions

View file

@ -4,6 +4,10 @@ import {
shouldRefreshBoardForEvent,
subscribeToBoardEvents,
} from "./board-events";
import {
createCrossTabSseCoordinator,
type BroadcastChannelLike,
} from "./cross-tab-sse";
import { queryKeys } from "./query-keys";
type MessageHandler = ((event: { data: string }) => void) | null;
@ -21,6 +25,14 @@ class FakeEventSource {
}
}
class FakeBroadcastChannel implements BroadcastChannelLike {
onmessage: ((event: { data: unknown }) => void) | null = null;
postMessage() {}
close() {}
}
describe("shouldRefreshBoardForEvent", () => {
test("refreshes board for run and interview status changes only", () => {
expect(shouldRefreshBoardForEvent("run.running")).toBe(true);
@ -31,10 +43,11 @@ describe("shouldRefreshBoardForEvent", () => {
});
describe("subscribeToBoardEvents", () => {
test("shares one source and invalidates the board runs key", () => {
test("coordinated mode shares one global source and invalidates the board runs key", async () => {
const source = new FakeEventSource();
const created: string[] = [];
const keys: string[] = [];
const coordinator = createCoordinator();
const mutate = (key: string) => {
keys.push(key);
return Promise.resolve();
@ -43,10 +56,13 @@ describe("subscribeToBoardEvents", () => {
const firstCleanup = subscribeToBoardEvents(mutate, (url) => {
created.push(url);
return source;
}, { debounceMs: 0 });
}, { debounceMs: 0, coordinator });
const secondCleanup = subscribeToBoardEvents(mutate, () => {
throw new Error("source should be reused");
}, { debounceMs: 0 });
}, { debounceMs: 0, coordinator });
await waitFor(() => created.length === 1);
keys.length = 0;
source.emit({ event: "run.running" });
@ -57,5 +73,67 @@ describe("subscribeToBoardEvents", () => {
expect(source.closed).toBe(false);
secondCleanup();
expect(source.closed).toBe(true);
coordinator.close();
});
test("fallback mode preserves the existing shared board EventSource", () => {
const source = new FakeEventSource();
const created: string[] = [];
const keys: string[] = [];
const coordinator = createFallbackCoordinator();
const mutate = (key: string) => {
keys.push(key);
return Promise.resolve();
};
const firstCleanup = subscribeToBoardEvents(mutate, (url) => {
created.push(url);
return source;
}, { debounceMs: 0, coordinator });
const secondCleanup = subscribeToBoardEvents(mutate, () => {
throw new Error("source should be reused");
}, { debounceMs: 0, coordinator });
source.emit({ event: "run.running" });
expect(created).toEqual(["/api/v1/attach"]);
expect(keys).toEqual([queryKeys.boards.runs()]);
firstCleanup();
expect(source.closed).toBe(false);
secondCleanup();
expect(source.closed).toBe(true);
coordinator.close();
});
});
function createCoordinator() {
return createCrossTabSseCoordinator({
tabId: "board-test",
channelFactory: () => new FakeBroadcastChannel(),
addVisibilityChangeListener: () => () => {},
addPagehideListener: () => () => {},
timing: {
heartbeatMs: 10,
leaderStaleMs: 50,
electionJitterMs: 0,
},
});
}
function createFallbackCoordinator() {
return createCrossTabSseCoordinator({
channelFactory: () => {
throw new Error("BroadcastChannel unavailable");
},
});
}
async function waitFor(condition: () => boolean, timeoutMs = 200) {
const deadline = Date.now() + timeoutMs;
while (Date.now() < deadline) {
if (condition()) return;
await new Promise((resolve) => setTimeout(resolve, 2));
}
throw new Error("condition did not become true before timeout");
}

View file

@ -1,6 +1,10 @@
import { useEffect } from "react";
import { useSWRConfig } from "swr";
import {
subscribeToCrossTabSse,
type CrossTabSseCoordinator,
} from "./cross-tab-sse";
import { queryKeys } from "./query-keys";
import {
createBrowserEventSource,
@ -11,6 +15,11 @@ import {
type SharedEventSubscription,
} from "./sse";
interface BoardEventOptions {
debounceMs?: number;
coordinator?: CrossTabSseCoordinator;
}
const BOARD_STATUS_EVENTS = new Set([
"run.submitted",
"run.queued",
@ -41,23 +50,37 @@ export function shouldRefreshBoardForEvent(event: string) {
export function subscribeToBoardEvents(
mutate: MutateFn,
eventSourceFactory: (url: string) => EventSourceLike = createBrowserEventSource,
{ debounceMs = 500 }: { debounceMs?: number } = {},
{ debounceMs = 500, coordinator }: BoardEventOptions = {},
): () => void {
return subscribeToSharedEventSource<EventPayload>({
subscriptions,
return subscribeToCrossTabSse<EventPayload>({
coordinator,
subscriptionKey: BOARD_SUBSCRIPTION_KEY,
url: queryKeys.system.attach(),
mutate,
eventSourceFactory,
debounceMs,
resolveInvalidation: (payload) => ({
keys: payload.event && shouldRefreshBoardForEvent(payload.event)
? [queryKeys.boards.runs()]
: [],
}),
resyncKeys: () => [queryKeys.boards.runs()],
resolveInvalidation: boardInvalidation,
fallbackSubscribe: () =>
subscribeToSharedEventSource<EventPayload>({
subscriptions,
subscriptionKey: BOARD_SUBSCRIPTION_KEY,
url: queryKeys.system.attach(),
mutate,
eventSourceFactory,
debounceMs,
resolveInvalidation: boardInvalidation,
}),
});
}
function boardInvalidation(payload: EventPayload) {
return {
keys: payload.event && shouldRefreshBoardForEvent(payload.event)
? [queryKeys.boards.runs()]
: [],
};
}
export function useBoardEvents() {
const { mutate } = useSWRConfig();

View file

@ -0,0 +1,385 @@
import { afterEach, describe, expect, test } from "bun:test";
import {
CROSS_TAB_SSE_CHANNEL,
createCrossTabSseCoordinator,
subscribeToCrossTabSse,
type BroadcastChannelLike,
type CrossTabSseCoordinator,
type CrossTabSseMessage,
} from "./cross-tab-sse";
import type { EventPayload, MutateFn } from "./sse";
type MessageHandler = ((event: { data: string }) => void) | null;
type TabVisibility = "visible" | "hidden";
const TEST_TIMING = {
heartbeatMs: 10,
leaderStaleMs: 35,
electionJitterMs: 5,
};
class FakeEventSource {
onmessage: MessageHandler = null;
closed = false;
constructor(
readonly url: string,
readonly owner: string,
) {}
emit(payload: unknown) {
this.onmessage?.({ data: JSON.stringify(payload) });
}
close() {
this.closed = true;
}
}
class FakeBroadcastChannel implements BroadcastChannelLike {
static channels = new Set<FakeBroadcastChannel>();
static muted = false;
onmessage: ((event: { data: unknown }) => void) | null = null;
closed = false;
constructor(readonly name: string) {
FakeBroadcastChannel.channels.add(this);
}
postMessage(message: CrossTabSseMessage) {
if (FakeBroadcastChannel.muted) return;
const recipients = [...FakeBroadcastChannel.channels].filter(
(channel) => channel !== this && !channel.closed && channel.name === this.name,
);
queueMicrotask(() => {
for (const channel of recipients) {
if (channel.closed) continue;
channel.onmessage?.({ data: { ...message } });
}
});
}
close() {
this.closed = true;
FakeBroadcastChannel.channels.delete(this);
}
static reset() {
for (const channel of FakeBroadcastChannel.channels) {
channel.closed = true;
}
FakeBroadcastChannel.channels.clear();
FakeBroadcastChannel.muted = false;
}
}
class Harness {
readonly sources: FakeEventSource[] = [];
readonly coordinators = new Map<string, CrossTabSseCoordinator>();
readonly visibility = new Map<string, TabVisibility>();
readonly visibilityHandlers = new Map<string, () => void>();
now = 1000;
createTab(tabId: string, visibility: TabVisibility = "visible") {
this.visibility.set(tabId, visibility);
const coordinator = createCrossTabSseCoordinator({
tabId,
channelFactory: (name) => new FakeBroadcastChannel(name),
eventSourceFactory: (url) => {
const source = new FakeEventSource(url, tabId);
this.sources.push(source);
return source;
},
getVisibility: () => this.visibility.get(tabId) ?? "visible",
addVisibilityChangeListener: (handler) => {
this.visibilityHandlers.set(tabId, handler);
return () => this.visibilityHandlers.delete(tabId);
},
addPagehideListener: () => () => {},
now: () => this.now,
timing: TEST_TIMING,
});
this.coordinators.set(tabId, coordinator);
return coordinator;
}
setVisibility(tabId: string, visibility: TabVisibility) {
this.visibility.set(tabId, visibility);
this.visibilityHandlers.get(tabId)?.();
}
openSources() {
return this.sources.filter((source) => !source.closed);
}
close() {
for (const coordinator of this.coordinators.values()) {
coordinator.close();
}
}
}
const harnesses: Harness[] = [];
afterEach(() => {
for (const harness of harnesses.splice(0)) {
harness.close();
}
FakeBroadcastChannel.reset();
});
describe("subscribeToCrossTabSse", () => {
test("opens one leader-owned global EventSource and keeps followers passive", async () => {
const harness = newHarness();
const cleanups = ["a", "b", "c"].map((tabId) => {
const coordinator = harness.createTab(tabId);
return subscribeForRunEvent(coordinator, []);
});
await waitFor(() => harness.openSources().length === 1);
expect(harness.openSources().map((source) => source.url)).toEqual(["/api/v1/attach"]);
expect([...FakeBroadcastChannel.channels].every((channel) => channel.name === CROSS_TAB_SSE_CHANNEL)).toBe(true);
cleanups.forEach((cleanup) => cleanup());
});
test("leader broadcasts events to all local subscribers", async () => {
const harness = newHarness();
const keysByTab = new Map<string, string[]>();
for (const tabId of ["a", "b", "c"]) {
keysByTab.set(tabId, []);
subscribeForRunEvent(harness.createTab(tabId), keysByTab.get(tabId)!);
}
await waitFor(() => harness.openSources().length === 1);
clearRecordedKeys(keysByTab);
harness.openSources()[0].emit(runEvent({ id: "evt-1", runId: "run-1", seq: 1 }));
await waitFor(() => [...keysByTab.values()].every((keys) => keys.length === 1));
expect(keysByTab.get("a")).toEqual(["event"]);
expect(keysByTab.get("b")).toEqual(["event"]);
expect(keysByTab.get("c")).toEqual(["event"]);
});
test("dedupes duplicate event ids until TTL or max-size eviction", async () => {
const harness = newHarness();
const keys: string[] = [];
subscribeForRunEvent(harness.createTab("a"), keys);
await waitFor(() => harness.openSources().length === 1);
keys.length = 0;
const source = harness.openSources()[0];
source.emit(runEvent({ id: "evt-dup", runId: "run-1", seq: 1 }));
source.emit(runEvent({ id: "evt-dup", runId: "run-1", seq: 1 }));
expect(keys).toEqual(["event"]);
harness.now += 5 * 60 * 1000 + 1;
source.emit(runEvent({ id: "evt-dup", runId: "run-1", seq: 1 }));
expect(keys).toEqual(["event", "event"]);
keys.length = 0;
for (let i = 0; i < 1001; i += 1) {
source.emit(runEvent({ id: `evt-${i}`, runId: "run-1", seq: i + 2 }));
}
source.emit(runEvent({ id: "evt-0", runId: "run-1", seq: 2 }));
expect(keys).toHaveLength(1002);
});
test("visible followers take over from a fresh hidden leader and resync", async () => {
const harness = newHarness();
const hiddenKeys: string[] = [];
const visibleKeys: string[] = [];
subscribeForRunEvent(harness.createTab("z", "hidden"), hiddenKeys);
await waitFor(() => harness.openSources().length === 1);
const hiddenSource = harness.openSources()[0];
subscribeForRunEvent(harness.createTab("a", "visible"), visibleKeys);
await waitFor(() => harness.openSources().length === 1 && harness.openSources()[0].owner === "a");
expect(hiddenSource.closed).toBe(true);
expect(visibleKeys).toContain("resync");
});
test("visible candidates racing for the same hidden leader resolve lexically", async () => {
const harness = newHarness();
subscribeForRunEvent(harness.createTab("z", "hidden"), []);
await waitFor(() => harness.openSources().length === 1 && harness.openSources()[0].owner === "z");
subscribeForRunEvent(harness.createTab("b", "visible"), []);
subscribeForRunEvent(harness.createTab("a", "visible"), []);
await waitFor(() => harness.openSources().length === 1 && harness.openSources()[0].owner === "a");
});
test("a lower lexical follower does not preempt a fresh visible leader", async () => {
const harness = newHarness();
subscribeForRunEvent(harness.createTab("z", "visible"), []);
await waitFor(() => harness.openSources().length === 1 && harness.openSources()[0].owner === "z");
subscribeForRunEvent(harness.createTab("a", "visible"), []);
await sleep(TEST_TIMING.electionJitterMs * 4);
expect(harness.openSources().map((source) => source.owner)).toEqual(["z"]);
});
test("stale leader detection opens a new leader source and resyncs followers", async () => {
const harness = newHarness();
const followerKeys: string[] = [];
subscribeForRunEvent(harness.createTab("a"), []);
subscribeForRunEvent(harness.createTab("b"), followerKeys);
await waitFor(() => harness.openSources().length === 1);
const staleLeader = harness.openSources()[0];
harness.coordinators.get(staleLeader.owner)?.close();
harness.now += TEST_TIMING.leaderStaleMs + TEST_TIMING.heartbeatMs + 1;
await waitFor(() => harness.openSources().length === 1 && harness.openSources()[0].owner !== staleLeader.owner);
expect(followerKeys).toContain("resync");
});
test("same-generation split brain converges to the higher-priority visible leader", async () => {
const harness = newHarness();
FakeBroadcastChannel.muted = true;
subscribeForRunEvent(harness.createTab("b"), []);
subscribeForRunEvent(harness.createTab("a"), []);
await waitFor(() => harness.openSources().length === 2);
FakeBroadcastChannel.muted = false;
await waitFor(() => harness.openSources().length === 1 && harness.openSources()[0].owner === "a");
});
test("old leader events are ignored after takeover", async () => {
const harness = newHarness();
const keys: string[] = [];
subscribeForRunEvent(harness.createTab("z", "hidden"), []);
await waitFor(() => harness.openSources().length === 1);
const oldSource = harness.openSources()[0];
subscribeForRunEvent(harness.createTab("a", "visible"), keys);
await waitFor(() => harness.openSources().length === 1 && harness.openSources()[0].owner === "a");
keys.length = 0;
oldSource.emit(runEvent({ id: "evt-old", runId: "run-1", seq: 1 }));
expect(keys).toEqual([]);
});
test("last unsubscribe closes the leader source and releases leadership", async () => {
const harness = newHarness();
const cleanup = subscribeForRunEvent(harness.createTab("a"), []);
await waitFor(() => harness.openSources().length === 1);
const source = harness.openSources()[0];
cleanup();
expect(source.closed).toBe(true);
expect(harness.openSources()).toEqual([]);
});
test("missing BroadcastChannel uses subscriber fallback", () => {
const coordinator = createCrossTabSseCoordinator({
channelFactory: () => {
throw new Error("no channel");
},
});
let fallbackStarted = 0;
let fallbackStopped = 0;
const cleanup = subscribeToCrossTabSse<EventPayload>({
coordinator,
subscriptionKey: "fallback",
mutate: (() => Promise.resolve()) as MutateFn,
resolveInvalidation: () => ({ keys: [] }),
resyncKeys: () => [],
fallbackSubscribe: () => {
fallbackStarted += 1;
return () => {
fallbackStopped += 1;
};
},
debounceMs: 0,
});
cleanup();
expect(fallbackStarted).toBe(1);
expect(fallbackStopped).toBe(1);
});
});
function newHarness() {
const harness = new Harness();
harnesses.push(harness);
return harness;
}
function subscribeForRunEvent(coordinator: CrossTabSseCoordinator, keys: string[]) {
return subscribeToCrossTabSse<EventPayload>({
coordinator,
subscriptionKey: "run-feed",
mutate: ((key: string) => {
keys.push(key);
return Promise.resolve();
}) as MutateFn,
resolveInvalidation: (payload) => ({
keys: payload.event === "run.running" ? ["event"] : [],
}),
resyncKeys: () => ["resync"],
fallbackSubscribe: () => {
throw new Error("fallback should not be used");
},
debounceMs: 0,
});
}
function runEvent({
id,
runId,
seq,
}: {
id: string;
runId: string;
seq: number;
}) {
return {
id,
seq,
run_id: runId,
event: "run.running",
ts: "2026-05-04T12:00:00.000Z",
};
}
async function waitFor(condition: () => boolean, timeoutMs = 500) {
const deadline = Date.now() + timeoutMs;
while (Date.now() < deadline) {
if (condition()) return;
await sleep(2);
}
throw new Error("condition did not become true before timeout");
}
function sleep(ms: number) {
return new Promise((resolve) => setTimeout(resolve, ms));
}
function clearRecordedKeys(keysByTab: Map<string, string[]>) {
for (const keys of keysByTab.values()) {
keys.length = 0;
}
}

File diff suppressed because it is too large Load diff

View file

@ -4,6 +4,10 @@ import {
queryKeysForRunEvent,
subscribeToRunEvents,
} from "./run-events";
import {
createCrossTabSseCoordinator,
type BroadcastChannelLike,
} from "./cross-tab-sse";
import { queryKeys } from "./query-keys";
type MessageHandler = ((event: { data: string }) => void) | null;
@ -25,6 +29,14 @@ class FakeEventSource {
}
}
class FakeBroadcastChannel implements BroadcastChannelLike {
onmessage: ((event: { data: unknown }) => void) | null = null;
postMessage() {}
close() {}
}
describe("queryKeysForRunEvent", () => {
test("terminal events invalidate run-scoped resources", () => {
expect(queryKeysForRunEvent("run-1", "run.completed")).toEqual([
@ -39,10 +51,74 @@ describe("queryKeysForRunEvent", () => {
});
describe("subscribeToRunEvents", () => {
test("refcounts shared sources and keeps mutators active until final unsubscribe", () => {
test("coordinated mode uses the global attach stream and filters by run_id", async () => {
const source = new FakeEventSource();
const created: string[] = [];
const keys: string[] = [];
const coordinator = createCoordinator();
const cleanup = subscribeToRunEvents(
"run-coordinated",
(key) => {
keys.push(key);
return Promise.resolve();
},
(url) => {
created.push(url);
return source;
},
{ debounceMs: 0, coordinator },
);
await waitFor(() => created.length === 1);
keys.length = 0;
source.emit({ event: "checkpoint.completed", run_id: "other-run" });
source.emit({ event: "checkpoint.completed", run_id: "run-coordinated" });
expect(created).toEqual(["/api/v1/attach"]);
expect(keys).toEqual([queryKeys.runs.files("run-coordinated")]);
cleanup();
coordinator.close();
});
test("coordinated terminal events invalidate without closing the global stream", async () => {
const source = new FakeEventSource();
const keys: string[] = [];
const coordinator = createCoordinator();
const cleanup = subscribeToRunEvents(
"run-terminal",
(key) => {
keys.push(key);
return Promise.resolve();
},
() => source,
{ debounceMs: 0, coordinator },
);
await waitFor(() => source.onmessage !== null);
keys.length = 0;
source.emit({ event: "run.failed", run_id: "run-terminal" });
expect(source.closed).toBe(false);
expect(keys).toContain(queryKeys.runs.files("run-terminal"));
expect(keys).toContain(queryKeys.runs.billing("run-terminal"));
keys.length = 0;
source.emit({ event: "run.archived", run_id: "run-terminal" });
expect(source.closed).toBe(false);
expect(keys).toEqual([queryKeys.runs.detail("run-terminal")]);
cleanup();
coordinator.close();
});
test("fallback refcounts run-scoped sources and keeps mutators active until final unsubscribe", () => {
const source = new FakeEventSource();
const created: string[] = [];
const keys: string[] = [];
const coordinator = createFallbackCoordinator();
const mutate = (key: string) => {
keys.push(key);
return Promise.resolve();
@ -51,10 +127,10 @@ describe("subscribeToRunEvents", () => {
const firstCleanup = subscribeToRunEvents("run-refcount", mutate, (url) => {
created.push(url);
return source;
}, { debounceMs: 0 });
}, { debounceMs: 0, coordinator });
const secondCleanup = subscribeToRunEvents("run-refcount", mutate, () => {
throw new Error("source should be reused");
}, { debounceMs: 0 });
}, { debounceMs: 0, coordinator });
expect(created).toEqual(["/api/v1/runs/run-refcount/attach"]);
@ -66,11 +142,13 @@ describe("subscribeToRunEvents", () => {
secondCleanup();
expect(source.closed).toBe(true);
coordinator.close();
});
test("terminal events close the source after invalidating keys", () => {
test("fallback terminal events close the source after invalidating keys", () => {
const source = new FakeEventSource();
const keys: string[] = [];
const coordinator = createFallbackCoordinator();
const cleanup = subscribeToRunEvents(
"run-terminal",
(key) => {
@ -78,7 +156,7 @@ describe("subscribeToRunEvents", () => {
return Promise.resolve();
},
() => source,
{ debounceMs: 0 },
{ debounceMs: 0, coordinator },
);
source.emit({ event: "run.failed" });
@ -88,13 +166,15 @@ describe("subscribeToRunEvents", () => {
expect(keys).toContain(queryKeys.runs.billing("run-terminal"));
cleanup();
coordinator.close();
});
test("malformed events are ignored and StrictMode-style cleanup does not underflow", () => {
test("fallback malformed events are ignored and StrictMode-style cleanup does not underflow", () => {
const firstSource = new FakeEventSource();
const secondSource = new FakeEventSource();
const sources = [firstSource, secondSource];
const keys: string[] = [];
const coordinator = createFallbackCoordinator();
const firstCleanup = subscribeToRunEvents(
"run-strict",
@ -103,7 +183,7 @@ describe("subscribeToRunEvents", () => {
return Promise.resolve();
},
() => sources.shift()!,
{ debounceMs: 0 },
{ debounceMs: 0, coordinator },
);
firstSource.emitRaw("{broken");
firstCleanup();
@ -115,12 +195,44 @@ describe("subscribeToRunEvents", () => {
return Promise.resolve();
},
() => sources.shift()!,
{ debounceMs: 0 },
{ debounceMs: 0, coordinator },
);
secondCleanup();
expect(keys).toEqual([]);
expect(firstSource.closed).toBe(true);
expect(secondSource.closed).toBe(true);
coordinator.close();
});
});
function createCoordinator() {
return createCrossTabSseCoordinator({
tabId: "run-test",
channelFactory: () => new FakeBroadcastChannel(),
addVisibilityChangeListener: () => () => {},
addPagehideListener: () => () => {},
timing: {
heartbeatMs: 10,
leaderStaleMs: 50,
electionJitterMs: 0,
},
});
}
function createFallbackCoordinator() {
return createCrossTabSseCoordinator({
channelFactory: () => {
throw new Error("BroadcastChannel unavailable");
},
});
}
async function waitFor(condition: () => boolean, timeoutMs = 200) {
const deadline = Date.now() + timeoutMs;
while (Date.now() < deadline) {
if (condition()) return;
await new Promise((resolve) => setTimeout(resolve, 2));
}
throw new Error("condition did not become true before timeout");
}

View file

@ -1,6 +1,10 @@
import { useEffect } from "react";
import { useSWRConfig } from "swr";
import {
subscribeToCrossTabSse,
type CrossTabSseCoordinator,
} from "./cross-tab-sse";
import { queryKeys } from "./query-keys";
import {
createBrowserEventSource,
@ -13,10 +17,16 @@ import {
interface RunEventPayload extends EventPayload {
event?: string;
run_id?: string;
node_id?: string;
properties?: Record<string, unknown>;
}
interface RunEventOptions {
debounceMs?: number;
coordinator?: CrossTabSseCoordinator;
}
const subscriptions = new Map<string, SharedEventSubscription>();
const TERMINAL_EVENTS = new Set(["run.completed", "run.failed"]);
@ -104,31 +114,64 @@ export function subscribeToRunEvents(
runId: string,
mutate: MutateFn,
eventSourceFactory: (url: string) => EventSourceLike = createBrowserEventSource,
{ debounceMs = 300 }: { debounceMs?: number } = {},
{ debounceMs = 300, coordinator }: RunEventOptions = {},
): () => void {
return subscribeToSharedEventSource<RunEventPayload>({
subscriptions,
subscriptionKey: runId,
url: queryKeys.runs.attach(runId),
return subscribeToCrossTabSse<RunEventPayload>({
coordinator,
subscriptionKey: `run:${runId}`,
mutate,
eventSourceFactory,
debounceMs,
resyncKeys: () => resyncKeysForRun(runId),
resolveInvalidation: (payload) => {
const event = payload.event;
if (!event) return { keys: [] };
const stageId = stageIdFromPayload(payload);
const keys = queryKeysForRunEvent(runId, event, stageId);
const terminal = TERMINAL_EVENTS.has(event);
return {
keys,
close: terminal,
immediate: terminal,
};
if (payload.run_id !== runId) return { keys: [] };
return runInvalidation(runId, payload, { closeOnTerminal: false });
},
fallbackSubscribe: () =>
subscribeToSharedEventSource<RunEventPayload>({
subscriptions,
subscriptionKey: runId,
url: queryKeys.runs.attach(runId),
mutate,
eventSourceFactory,
debounceMs,
resolveInvalidation: (payload) =>
runInvalidation(runId, payload, { closeOnTerminal: true }),
}),
});
}
function runInvalidation(
runId: string,
payload: RunEventPayload,
{ closeOnTerminal }: { closeOnTerminal: boolean },
) {
const event = payload.event;
if (!event) return { keys: [] };
const stageId = stageIdFromPayload(payload);
const keys = queryKeysForRunEvent(runId, event, stageId);
const terminal = TERMINAL_EVENTS.has(event);
return {
keys,
close: closeOnTerminal && terminal,
immediate: terminal,
};
}
function resyncKeysForRun(runId: string) {
return [
queryKeys.runs.detail(runId),
queryKeys.runs.files(runId),
queryKeys.runs.billing(runId),
queryKeys.runs.stages(runId),
queryKeys.runs.events(runId, 1000),
queryKeys.runs.graph(runId, "LR"),
queryKeys.runs.graph(runId, "TB"),
queryKeys.runs.questions(runId, 25, 0),
];
}
function stageIdFromPayload(payload: RunEventPayload): string | undefined {
if (typeof payload.node_id === "string") return payload.node_id;
const nodeId = payload.properties?.node_id;