diff --git a/apps/fabro-web/app/lib/board-events.test.tsx b/apps/fabro-web/app/lib/board-events.test.tsx index 853d300b9..9b398fdc1 100644 --- a/apps/fabro-web/app/lib/board-events.test.tsx +++ b/apps/fabro-web/app/lib/board-events.test.tsx @@ -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"); +} diff --git a/apps/fabro-web/app/lib/board-events.ts b/apps/fabro-web/app/lib/board-events.ts index ef30009b1..a182d9ad2 100644 --- a/apps/fabro-web/app/lib/board-events.ts +++ b/apps/fabro-web/app/lib/board-events.ts @@ -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({ - subscriptions, + return subscribeToCrossTabSse({ + 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({ + 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(); diff --git a/apps/fabro-web/app/lib/cross-tab-sse.test.ts b/apps/fabro-web/app/lib/cross-tab-sse.test.ts new file mode 100644 index 000000000..32cff001c --- /dev/null +++ b/apps/fabro-web/app/lib/cross-tab-sse.test.ts @@ -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(); + 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(); + readonly visibility = new Map(); + readonly visibilityHandlers = new Map 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(); + + 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({ + 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({ + 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) { + for (const keys of keysByTab.values()) { + keys.length = 0; + } +} diff --git a/apps/fabro-web/app/lib/cross-tab-sse.ts b/apps/fabro-web/app/lib/cross-tab-sse.ts new file mode 100644 index 000000000..255c4c7ce --- /dev/null +++ b/apps/fabro-web/app/lib/cross-tab-sse.ts @@ -0,0 +1,1175 @@ +import { queryKeys } from "./query-keys"; +import { + createBrowserEventSource, + type EventInvalidation, + type EventPayload, + type EventSourceLike, + type MutateFn, +} from "./sse"; +import { isRecord } from "./unknown"; + +export const CROSS_TAB_SSE_CHANNEL = "fabro:sse:v1"; +export const HEARTBEAT_MS = 1000; +export const LEADER_STALE_MS = 4000; +export const ELECTION_JITTER_MS = 150; + +const MESSAGE_VERSION = 1 as const; +const EVENT_DEDUPE_TTL_MS = 5 * 60 * 1000; +const EVENT_DEDUPE_MAX = 1000; + +type TabVisibility = "visible" | "hidden"; +type CandidateReason = "hidden-leader" | "stale-leader" | "release" | "no-leader"; + +interface BaseMessage { + type: string; + version: typeof MESSAGE_VERSION; + tabId: string; + sentAt: number; +} + +interface HelloMessage extends BaseMessage { + type: "hello"; +} + +interface HeartbeatMessage extends BaseMessage { + type: "heartbeat"; + leaderId: string; + generation: number; + visibility: TabVisibility; +} + +interface CandidateMessage extends BaseMessage { + type: "candidate"; + candidateId: string; + candidateGeneration: number; + visibility: TabVisibility; + observedLeaderId: string | null; + observedGeneration: number; + reason: CandidateReason; +} + +interface LeaderChangedMessage extends BaseMessage { + type: "leader-changed"; + leaderId: string; + generation: number; + visibility: TabVisibility; +} + +interface ReleaseMessage extends BaseMessage { + type: "release"; + leaderId: string; + generation: number; +} + +interface ResyncMessage extends BaseMessage { + type: "resync"; + leaderId: string | null; + generation: number; + reason: CandidateReason; +} + +interface EventMessage extends BaseMessage { + type: "event"; + leaderId: string; + generation: number; + payload: EventPayload; +} + +export type CrossTabSseMessage = + | HelloMessage + | HeartbeatMessage + | CandidateMessage + | LeaderChangedMessage + | ReleaseMessage + | ResyncMessage + | EventMessage; + +export interface BroadcastChannelLike { + onmessage: ((event: { data: unknown }) => void) | null; + postMessage(message: CrossTabSseMessage): void; + close(): void; +} + +interface TimingOptions { + heartbeatMs: number; + leaderStaleMs: number; + electionJitterMs: number; +} + +export interface CrossTabSseCoordinatorOptions { + tabId?: string; + channelFactory?: (name: string) => BroadcastChannelLike; + eventSourceFactory?: (url: string) => EventSourceLike; + getVisibility?: () => TabVisibility; + addVisibilityChangeListener?: (handler: () => void) => () => void; + addPagehideListener?: (handler: () => void) => () => void; + now?: () => number; + timing?: Partial; +} + +interface SubscribeOptions { + subscriptionKey: string; + mutate: MutateFn; + resolveInvalidation: (payload: TPayload) => EventInvalidation; + resyncKeys: () => string[]; + fallbackSubscribe: () => () => void; + eventSourceFactory?: (url: string) => EventSourceLike; + debounceMs?: number; +} + +export interface SubscribeToCrossTabSseOptions + extends SubscribeOptions { + coordinator?: CrossTabSseCoordinator; +} + +interface FallbackEntry { + count: number; + subscribe: () => () => void; + cleanup?: () => void; +} + +interface LocalSubscription { + refcount: number; + mutators: Map; + fallbacks: Map; + pendingKeys: Set; + debounceTimer: ReturnType | null; + debounceMs: number; + resolveInvalidation: (payload: EventPayload) => EventInvalidation; + resyncKeys: () => string[]; +} + +interface LeaderState { + leaderId: string; + generation: number; + visibility: TabVisibility; + lastSeen: number; +} + +class RecentEventCache { + private readonly seen = new Map(); + + constructor( + private readonly maxSize: number, + private readonly ttlMs: number, + ) {} + + remember(key: string | undefined, now: number): boolean { + if (!key) return true; + this.prune(now); + if (this.seen.has(key)) return false; + + this.seen.set(key, now); + while (this.seen.size > this.maxSize) { + const oldest = this.seen.keys().next().value; + if (oldest === undefined) break; + this.seen.delete(oldest); + } + return true; + } + + private prune(now: number) { + for (const [key, seenAt] of this.seen) { + if (now - seenAt > this.ttlMs) { + this.seen.delete(key); + } + } + } +} + +export class CrossTabSseCoordinator { + readonly tabId: string; + + private readonly channelFactory: (name: string) => BroadcastChannelLike; + private readonly getVisibility: () => TabVisibility; + private readonly addVisibilityChangeListener: (handler: () => void) => () => void; + private readonly addPagehideListener: (handler: () => void) => () => void; + private readonly now: () => number; + private readonly timing: TimingOptions; + private readonly recentEvents = new RecentEventCache(EVENT_DEDUPE_MAX, EVENT_DEDUPE_TTL_MS); + private readonly subscriptions = new Map(); + private readonly candidates = new Map(); + + private sourceFactory: (url: string) => EventSourceLike; + private channel: BroadcastChannelLike | null = null; + private source: EventSourceLike | null = null; + private initialized = false; + private coordinationUnavailable = false; + private fallbackMode = false; + private degradingToFallback = false; + private isLeader = false; + private leader: LeaderState | null = null; + private generation = 0; + private ownCandidate: CandidateMessage | null = null; + private candidateTimer: ReturnType | null = null; + private noLeaderTimer: ReturnType | null = null; + private heartbeatTimer: ReturnType | null = null; + private leaderCheckTimer: ReturnType | null = null; + private removeVisibilityListener: (() => void) | null = null; + private removePagehideListener: (() => void) | null = null; + private sourceFactoryLocked = false; + + constructor(options: CrossTabSseCoordinatorOptions = {}) { + this.tabId = options.tabId ?? createTabId(); + this.channelFactory = options.channelFactory ?? createBrowserBroadcastChannel; + this.sourceFactory = options.eventSourceFactory ?? createBrowserEventSource; + this.getVisibility = options.getVisibility ?? getBrowserVisibility; + this.addVisibilityChangeListener = + options.addVisibilityChangeListener ?? addBrowserVisibilityChangeListener; + this.addPagehideListener = options.addPagehideListener ?? addBrowserPagehideListener; + this.now = options.now ?? Date.now; + this.timing = { + heartbeatMs: options.timing?.heartbeatMs ?? HEARTBEAT_MS, + leaderStaleMs: options.timing?.leaderStaleMs ?? LEADER_STALE_MS, + electionJitterMs: options.timing?.electionJitterMs ?? ELECTION_JITTER_MS, + }; + } + + subscribe(options: SubscribeOptions): () => void { + if (options.eventSourceFactory && !this.source && !this.sourceFactoryLocked) { + this.sourceFactory = options.eventSourceFactory; + this.sourceFactoryLocked = true; + } + + if (this.coordinationUnavailable) { + return options.fallbackSubscribe(); + } + + if (!this.initialized && !this.initialize()) { + this.coordinationUnavailable = true; + return options.fallbackSubscribe(); + } + + const subscription = this.addLocalSubscription(options); + if (this.fallbackMode) { + this.startFallbacksFor(subscription); + } else { + this.ensureLeadershipProgress(); + } + + let active = true; + return () => { + if (!active) return; + active = false; + this.removeLocalSubscription(options.subscriptionKey, options.mutate); + }; + } + + close() { + this.releaseLeadership({ broadcast: false, resync: false }); + this.clearCandidate(); + this.clearNoLeaderTimer(); + this.shutdownTimersAndChannel(); + this.closeFallbacks(); + this.subscriptions.clear(); + this.leader = null; + this.initialized = false; + this.fallbackMode = false; + this.sourceFactoryLocked = false; + } + + private initialize(): boolean { + try { + this.channel = this.channelFactory(CROSS_TAB_SSE_CHANNEL); + } catch { + this.channel = null; + return false; + } + + this.channel.onmessage = (event) => this.handleMessage(event.data); + this.initialized = true; + this.removeVisibilityListener = this.addVisibilityChangeListener(() => { + this.handleVisibilityChange(); + }); + this.removePagehideListener = this.addPagehideListener(() => { + this.handlePagehide(); + }); + this.leaderCheckTimer = setInterval(() => { + this.checkLeaderFreshness(); + }, this.timing.heartbeatMs); + + return this.post({ type: "hello", version: MESSAGE_VERSION, tabId: this.tabId, sentAt: this.now() }); + } + + private addLocalSubscription( + options: SubscribeOptions, + ): LocalSubscription { + let subscription = this.subscriptions.get(options.subscriptionKey); + if (!subscription) { + subscription = { + refcount: 0, + mutators: new Map(), + fallbacks: new Map(), + pendingKeys: new Set(), + debounceTimer: null, + debounceMs: options.debounceMs ?? 300, + resolveInvalidation: options.resolveInvalidation as (payload: EventPayload) => EventInvalidation, + resyncKeys: options.resyncKeys, + }; + this.subscriptions.set(options.subscriptionKey, subscription); + } else { + subscription.resolveInvalidation = + options.resolveInvalidation as (payload: EventPayload) => EventInvalidation; + subscription.resyncKeys = options.resyncKeys; + subscription.debounceMs = options.debounceMs ?? subscription.debounceMs; + } + + subscription.refcount += 1; + subscription.mutators.set( + options.mutate, + (subscription.mutators.get(options.mutate) ?? 0) + 1, + ); + + const fallback = subscription.fallbacks.get(options.mutate); + if (fallback) { + fallback.count += 1; + fallback.subscribe = options.fallbackSubscribe; + } else { + subscription.fallbacks.set(options.mutate, { + count: 1, + subscribe: options.fallbackSubscribe, + }); + } + + return subscription; + } + + private removeLocalSubscription(subscriptionKey: string, mutate: MutateFn) { + const subscription = this.subscriptions.get(subscriptionKey); + if (!subscription) return; + + const mutateCount = subscription.mutators.get(mutate) ?? 0; + if (mutateCount <= 1) { + subscription.mutators.delete(mutate); + } else { + subscription.mutators.set(mutate, mutateCount - 1); + } + + const fallback = subscription.fallbacks.get(mutate); + if (fallback) { + fallback.count -= 1; + if (fallback.count <= 0) { + fallback.cleanup?.(); + subscription.fallbacks.delete(mutate); + } + } + + subscription.refcount -= 1; + if (subscription.refcount <= 0) { + if (subscription.debounceTimer) { + clearTimeout(subscription.debounceTimer); + } + this.subscriptions.delete(subscriptionKey); + } + + if (this.subscriptions.size === 0) { + this.releaseLeadership({ broadcast: true, resync: false }); + this.clearCandidate(); + this.clearNoLeaderTimer(); + this.shutdownTimersAndChannel(); + this.initialized = false; + this.fallbackMode = false; + this.sourceFactoryLocked = false; + } + } + + private handleMessage(data: unknown) { + const message = parseMessage(data); + if (!message || message.tabId === this.tabId) return; + + switch (message.type) { + case "hello": + if (this.isLeader) this.sendHeartbeat(); + break; + case "heartbeat": + this.handleLeaderAnnouncement(message, { resyncOnChange: true }); + break; + case "candidate": + this.handleCandidate(message); + break; + case "leader-changed": + this.handleLeaderAnnouncement(message, { resyncOnChange: true }); + break; + case "release": + this.handleRelease(message); + break; + case "resync": + this.handleResync(message); + break; + case "event": + this.handleBroadcastEvent(message); + break; + } + } + + private handleLeaderAnnouncement( + message: HeartbeatMessage | LeaderChangedMessage, + { resyncOnChange }: { resyncOnChange: boolean }, + ) { + const incoming: LeaderState = { + leaderId: message.leaderId, + generation: message.generation, + visibility: message.visibility, + lastSeen: this.now(), + }; + + if (this.isLeader && incoming.leaderId !== this.tabId) { + const own: LeaderState = { + leaderId: this.tabId, + generation: this.generation, + visibility: this.currentVisibility(), + lastSeen: this.now(), + }; + if ( + incoming.generation > own.generation || + (incoming.generation === own.generation && leaderHasHigherPriority(incoming, own)) + ) { + this.releaseLeadership({ broadcast: false, resync: true }); + } else { + return; + } + } + + const previous = this.leader; + if (this.ownCandidate && incoming.generation === this.ownCandidate.candidateGeneration) { + const sawIncomingCandidate = this.candidates.has( + `${incoming.generation}:${incoming.leaderId}`, + ); + const shouldAcceptFreshVisibleLeader = + this.ownCandidate.reason === "no-leader" && + incoming.visibility === "visible" && + !sawIncomingCandidate; + + if ( + !shouldAcceptFreshVisibleLeader && + !leaderHasHigherPriority(incoming, leaderStateForCandidate(this.ownCandidate, this.now())) + ) { + return; + } + } + + if (!this.shouldAcceptLeader(incoming)) return; + + this.leader = incoming; + this.clearNoLeaderTimer(); + this.generation = Math.max(this.generation, incoming.generation); + if (this.ownCandidate && incoming.generation >= this.ownCandidate.candidateGeneration) { + this.clearCandidate(); + } + + const changed = + !previous || + previous.leaderId !== incoming.leaderId || + previous.generation !== incoming.generation; + + if (changed && resyncOnChange) { + this.resyncAll(); + } + + if (incoming.visibility === "hidden" && this.currentVisibility() === "visible") { + this.enterCandidacy("hidden-leader", incoming); + } + } + + private handleCandidate(message: CandidateMessage) { + this.candidates.set(candidateKey(message), message); + + if ( + this.isLeader && + message.observedLeaderId === this.tabId && + message.observedGeneration >= this.generation + ) { + this.releaseLeadership({ broadcast: true, resync: true }); + } + + if ( + this.ownCandidate && + message.candidateGeneration === this.ownCandidate.candidateGeneration && + candidateHasHigherPriority(message, this.ownCandidate) + ) { + this.clearCandidate(); + } + } + + private handleRelease(message: ReleaseMessage) { + const current = this.leader; + if ( + current && + message.leaderId === current.leaderId && + message.generation >= current.generation + ) { + this.leader = null; + this.generation = Math.max(this.generation, message.generation); + this.enterCandidacy("release", { + leaderId: message.leaderId, + generation: message.generation, + visibility: "hidden", + lastSeen: this.now(), + }); + } + } + + private handleResync(message: ResyncMessage) { + if (this.leader && message.generation < this.leader.generation) return; + this.resyncAll(); + } + + private handleBroadcastEvent(message: EventMessage) { + if (!this.isCurrentLeader(message.leaderId, message.generation)) return; + if (!this.recentEvents.remember(eventDedupeKey(message.payload), this.now())) return; + this.dispatchPayload(message.payload); + } + + private ensureLeadershipProgress() { + if (this.subscriptions.size === 0 || this.isLeader || this.ownCandidate) return; + + if (!this.leader) { + this.scheduleNoLeaderCandidacy(); + return; + } + + if (this.leader.visibility === "hidden" && this.currentVisibility() === "visible") { + this.enterCandidacy("hidden-leader", this.leader); + } + } + + private checkLeaderFreshness() { + if (this.subscriptions.size === 0 || this.fallbackMode) return; + if (this.isLeader) return; + + const current = this.leader; + if (!current) { + this.scheduleNoLeaderCandidacy(); + return; + } + + if (this.now() - current.lastSeen > this.timing.leaderStaleMs) { + this.leader = null; + this.generation = Math.max(this.generation, current.generation); + this.resyncAll(); + this.enterCandidacy("stale-leader", current); + return; + } + + if (current.visibility === "hidden" && this.currentVisibility() === "visible") { + this.enterCandidacy("hidden-leader", current); + } + } + + private enterCandidacy(reason: CandidateReason, observedLeader: LeaderState | null = this.leader) { + if (this.subscriptions.size === 0 || this.fallbackMode) return; + + if ( + reason === "no-leader" && + this.leader && + this.leader.visibility === "visible" && + this.now() - this.leader.lastSeen <= this.timing.leaderStaleMs + ) { + return; + } + + const observedGeneration = observedLeader?.generation ?? this.generation; + const candidateGeneration = observedGeneration + 1; + if ( + this.ownCandidate && + this.ownCandidate.candidateGeneration >= candidateGeneration + ) { + return; + } + + this.clearCandidate(); + this.clearNoLeaderTimer(); + const candidate: CandidateMessage = { + type: "candidate", + version: MESSAGE_VERSION, + tabId: this.tabId, + sentAt: this.now(), + candidateId: this.tabId, + candidateGeneration, + visibility: this.currentVisibility(), + observedLeaderId: observedLeader?.leaderId ?? null, + observedGeneration, + reason, + }; + + this.ownCandidate = candidate; + this.candidates.set(candidateKey(candidate), candidate); + this.post(candidate); + this.candidateTimer = setTimeout(() => { + this.completeCandidacy(candidate); + }, this.timing.electionJitterMs); + } + + private completeCandidacy(candidate: CandidateMessage) { + if (this.ownCandidate !== candidate || this.fallbackMode) return; + + for (const other of this.candidates.values()) { + if ( + other.candidateGeneration === candidate.candidateGeneration && + candidateHasHigherPriority(other, candidate) + ) { + this.clearCandidate(); + return; + } + } + + if ( + this.leader && + this.now() - this.leader.lastSeen <= this.timing.leaderStaleMs + ) { + if (this.leader.generation > candidate.candidateGeneration) { + this.clearCandidate(); + return; + } + if ( + this.leader.generation === candidate.candidateGeneration && + leaderHasHigherPriority(this.leader, leaderStateForCandidate(candidate, this.now())) + ) { + this.clearCandidate(); + return; + } + } + + this.becomeLeader(candidate.candidateGeneration); + } + + private becomeLeader(generation: number) { + this.clearCandidate(); + this.closeSource(); + this.isLeader = true; + this.generation = generation; + this.leader = { + leaderId: this.tabId, + generation, + visibility: this.currentVisibility(), + lastSeen: this.now(), + }; + + const source = this.sourceFactory(queryKeys.system.attach()); + this.source = source; + source.onmessage = (message) => { + this.handleLeaderEventSourceMessage(message.data); + }; + + this.post({ + type: "leader-changed", + version: MESSAGE_VERSION, + tabId: this.tabId, + sentAt: this.now(), + leaderId: this.tabId, + generation, + visibility: this.currentVisibility(), + }); + this.startHeartbeat(); + this.resyncAll(); + } + + private handleLeaderEventSourceMessage(data: string) { + if (!this.isLeader) return; + + let payload: EventPayload; + try { + payload = JSON.parse(data) as EventPayload; + } catch { + return; + } + + if (!this.recentEvents.remember(eventDedupeKey(payload), this.now())) return; + this.dispatchPayload(payload); + this.post({ + type: "event", + version: MESSAGE_VERSION, + tabId: this.tabId, + sentAt: this.now(), + leaderId: this.tabId, + generation: this.generation, + payload, + }); + } + + private dispatchPayload(payload: EventPayload) { + for (const subscription of this.subscriptions.values()) { + const invalidation = subscription.resolveInvalidation(payload); + this.queueInvalidations(subscription, invalidation.keys, { + immediate: invalidation.immediate, + }); + } + } + + private queueInvalidations( + subscription: LocalSubscription, + keys: string[], + { immediate = false }: { immediate?: boolean } = {}, + ) { + if (keys.length === 0) return; + for (const key of keys) { + subscription.pendingKeys.add(key); + } + + if (immediate || subscription.debounceMs <= 0) { + this.flushInvalidations(subscription); + return; + } + + if (subscription.debounceTimer) { + clearTimeout(subscription.debounceTimer); + } + subscription.debounceTimer = setTimeout(() => { + subscription.debounceTimer = null; + this.flushInvalidations(subscription); + }, subscription.debounceMs); + } + + private flushInvalidations(subscription: LocalSubscription) { + if (subscription.pendingKeys.size === 0) return; + const keys = [...subscription.pendingKeys]; + subscription.pendingKeys.clear(); + + for (const mutator of subscription.mutators.keys()) { + for (const key of keys) { + void mutator(key); + } + } + } + + private resyncAll() { + for (const subscription of this.subscriptions.values()) { + this.queueInvalidations(subscription, subscription.resyncKeys(), { immediate: true }); + } + } + + private handleVisibilityChange() { + if (this.isLeader) { + this.sendHeartbeat(); + } + + if (this.currentVisibility() === "visible") { + this.resyncAll(); + if (this.leader?.visibility === "hidden") { + this.enterCandidacy("hidden-leader", this.leader); + } + } + } + + private handlePagehide() { + this.releaseLeadership({ broadcast: true, resync: false }); + this.clearCandidate(); + this.clearNoLeaderTimer(); + } + + private startHeartbeat() { + if (this.heartbeatTimer) { + clearInterval(this.heartbeatTimer); + } + this.sendHeartbeat(); + this.heartbeatTimer = setInterval(() => { + this.sendHeartbeat(); + }, this.timing.heartbeatMs); + } + + private sendHeartbeat() { + if (!this.isLeader) return; + const visibility = this.currentVisibility(); + this.leader = { + leaderId: this.tabId, + generation: this.generation, + visibility, + lastSeen: this.now(), + }; + this.post({ + type: "heartbeat", + version: MESSAGE_VERSION, + tabId: this.tabId, + sentAt: this.now(), + leaderId: this.tabId, + generation: this.generation, + visibility, + }); + } + + private releaseLeadership({ + broadcast, + resync, + }: { + broadcast: boolean; + resync: boolean; + }) { + if (!this.isLeader && !this.source) return; + + const generation = this.generation; + this.closeSource(); + this.isLeader = false; + if (this.heartbeatTimer) { + clearInterval(this.heartbeatTimer); + this.heartbeatTimer = null; + } + this.leader = null; + + if (broadcast) { + this.post({ + type: "release", + version: MESSAGE_VERSION, + tabId: this.tabId, + sentAt: this.now(), + leaderId: this.tabId, + generation, + }); + this.post({ + type: "resync", + version: MESSAGE_VERSION, + tabId: this.tabId, + sentAt: this.now(), + leaderId: this.tabId, + generation, + reason: "release", + }); + } + if (resync) this.resyncAll(); + } + + private closeSource() { + if (!this.source) return; + this.source.close(); + this.source = null; + } + + private clearCandidate() { + if (this.candidateTimer) { + clearTimeout(this.candidateTimer); + this.candidateTimer = null; + } + this.ownCandidate = null; + } + + private scheduleNoLeaderCandidacy() { + if ( + this.noLeaderTimer || + this.ownCandidate || + this.isLeader || + this.leader || + this.subscriptions.size === 0 || + this.fallbackMode + ) { + return; + } + + this.noLeaderTimer = setTimeout(() => { + this.noLeaderTimer = null; + if (!this.leader && !this.isLeader) { + this.enterCandidacy("no-leader"); + } + }, this.timing.electionJitterMs); + } + + private clearNoLeaderTimer() { + if (!this.noLeaderTimer) return; + clearTimeout(this.noLeaderTimer); + this.noLeaderTimer = null; + } + + private shutdownTimersAndChannel() { + this.clearNoLeaderTimer(); + if (this.heartbeatTimer) { + clearInterval(this.heartbeatTimer); + this.heartbeatTimer = null; + } + if (this.leaderCheckTimer) { + clearInterval(this.leaderCheckTimer); + this.leaderCheckTimer = null; + } + this.removeVisibilityListener?.(); + this.removeVisibilityListener = null; + this.removePagehideListener?.(); + this.removePagehideListener = null; + if (this.channel) { + this.channel.close(); + this.channel = null; + } + } + + private post(message: CrossTabSseMessage): boolean { + if (!this.channel) return false; + try { + this.channel.postMessage(message); + return true; + } catch { + this.degradeToFallback(); + return false; + } + } + + private degradeToFallback() { + if (this.degradingToFallback || this.fallbackMode) return; + this.degradingToFallback = true; + this.releaseLeadership({ broadcast: false, resync: false }); + this.clearCandidate(); + this.shutdownTimersAndChannel(); + this.initialized = false; + this.fallbackMode = true; + this.coordinationUnavailable = true; + for (const subscription of this.subscriptions.values()) { + this.startFallbacksFor(subscription); + } + this.degradingToFallback = false; + } + + private startFallbacksFor(subscription: LocalSubscription) { + for (const fallback of subscription.fallbacks.values()) { + if (!fallback.cleanup) { + fallback.cleanup = fallback.subscribe(); + } + } + } + + private closeFallbacks() { + for (const subscription of this.subscriptions.values()) { + for (const fallback of subscription.fallbacks.values()) { + fallback.cleanup?.(); + fallback.cleanup = undefined; + } + } + } + + private shouldAcceptLeader(incoming: LeaderState): boolean { + if (!this.leader) return true; + if (incoming.generation > this.leader.generation) return true; + if (incoming.generation < this.leader.generation) return false; + if (incoming.leaderId === this.leader.leaderId) return true; + return leaderHasHigherPriority(incoming, this.leader); + } + + private isCurrentLeader(leaderId: string, generation: number): boolean { + return Boolean( + this.leader && + this.leader.leaderId === leaderId && + this.leader.generation === generation, + ); + } + + private currentVisibility(): TabVisibility { + return this.getVisibility() === "hidden" ? "hidden" : "visible"; + } +} + +const defaultCoordinator = new CrossTabSseCoordinator(); + +export function createCrossTabSseCoordinator(options: CrossTabSseCoordinatorOptions = {}) { + return new CrossTabSseCoordinator(options); +} + +export function subscribeToCrossTabSse({ + coordinator = defaultCoordinator, + ...options +}: SubscribeToCrossTabSseOptions): () => void { + return coordinator.subscribe(options); +} + +function createBrowserBroadcastChannel(name: string): BroadcastChannelLike { + if (typeof BroadcastChannel === "undefined") { + throw new Error("BroadcastChannel is unavailable"); + } + return new BroadcastChannel(name); +} + +function createTabId(): string { + if (typeof crypto !== "undefined" && typeof crypto.randomUUID === "function") { + return crypto.randomUUID(); + } + return `tab-${Date.now().toString(36)}-${Math.random().toString(36).slice(2)}`; +} + +function getBrowserVisibility(): TabVisibility { + if (typeof document === "undefined") return "visible"; + return document.visibilityState === "visible" ? "visible" : "hidden"; +} + +function addBrowserVisibilityChangeListener(handler: () => void): () => void { + if (typeof document === "undefined") return () => {}; + document.addEventListener("visibilitychange", handler); + return () => document.removeEventListener("visibilitychange", handler); +} + +function addBrowserPagehideListener(handler: () => void): () => void { + if (typeof window === "undefined") return () => {}; + window.addEventListener("pagehide", handler); + return () => window.removeEventListener("pagehide", handler); +} + +function parseMessage(data: unknown): CrossTabSseMessage | undefined { + if (!isRecord(data)) return undefined; + if (data.version !== MESSAGE_VERSION) return undefined; + const type = data.type; + const tabId = data.tabId; + const sentAt = data.sentAt; + if (typeof type !== "string") return undefined; + if (typeof tabId !== "string") return undefined; + if (typeof sentAt !== "number") return undefined; + const base = { tabId, sentAt }; + + switch (type) { + case "hello": + return baseMessage(base, "hello"); + case "heartbeat": { + const { leaderId, generation, visibility } = data; + if (typeof leaderId === "string" && typeof generation === "number" && isVisibility(visibility)) { + return { + ...baseMessage(base, "heartbeat"), + leaderId, + generation, + visibility, + }; + } + return undefined; + } + case "candidate": { + const { + candidateId, + candidateGeneration, + visibility, + observedLeaderId, + observedGeneration, + reason, + } = data; + if ( + typeof candidateId === "string" && + typeof candidateGeneration === "number" && + isVisibility(visibility) && + (typeof observedLeaderId === "string" || observedLeaderId === null) && + typeof observedGeneration === "number" && + isCandidateReason(reason) + ) { + const normalizedObservedLeaderId = + typeof observedLeaderId === "string" ? observedLeaderId : null; + return { + ...baseMessage(base, "candidate"), + candidateId, + candidateGeneration, + visibility, + observedLeaderId: normalizedObservedLeaderId, + observedGeneration, + reason, + }; + } + return undefined; + } + case "leader-changed": { + const { leaderId, generation, visibility } = data; + if (typeof leaderId === "string" && typeof generation === "number" && isVisibility(visibility)) { + return { + ...baseMessage(base, "leader-changed"), + leaderId, + generation, + visibility, + }; + } + return undefined; + } + case "release": { + const { leaderId, generation } = data; + if (typeof leaderId === "string" && typeof generation === "number") { + return { + ...baseMessage(base, "release"), + leaderId, + generation, + }; + } + return undefined; + } + case "resync": { + const { leaderId, generation, reason } = data; + if ( + (typeof leaderId === "string" || leaderId === null) && + typeof generation === "number" && + isCandidateReason(reason) + ) { + const normalizedLeaderId = typeof leaderId === "string" ? leaderId : null; + return { + ...baseMessage(base, "resync"), + leaderId: normalizedLeaderId, + generation, + reason, + }; + } + return undefined; + } + case "event": { + const { leaderId, generation, payload } = data; + if (typeof leaderId === "string" && typeof generation === "number" && isRecord(payload)) { + return { + ...baseMessage(base, "event"), + leaderId, + generation, + payload, + }; + } + return undefined; + } + default: + return undefined; + } +} + +function baseMessage( + data: { + tabId: string; + sentAt: number; + }, + type: TType, +) { + return { + type, + version: MESSAGE_VERSION, + tabId: data.tabId, + sentAt: data.sentAt, + }; +} + +function isVisibility(value: unknown): value is TabVisibility { + return value === "visible" || value === "hidden"; +} + +function isCandidateReason(value: unknown): value is CandidateReason { + return ( + value === "hidden-leader" || + value === "stale-leader" || + value === "release" || + value === "no-leader" + ); +} + +function candidateHasHigherPriority(candidate: CandidateMessage, other: CandidateMessage): boolean { + if (candidate.visibility !== other.visibility) return candidate.visibility === "visible"; + return candidate.candidateId < other.candidateId; +} + +function leaderHasHigherPriority(candidate: LeaderState, other: LeaderState): boolean { + if (candidate.visibility !== other.visibility) return candidate.visibility === "visible"; + return candidate.leaderId < other.leaderId; +} + +function leaderStateForCandidate(candidate: CandidateMessage, lastSeen: number): LeaderState { + return { + leaderId: candidate.candidateId, + generation: candidate.candidateGeneration, + visibility: candidate.visibility, + lastSeen, + }; +} + +function candidateKey(candidate: CandidateMessage): string { + return `${candidate.candidateGeneration}:${candidate.candidateId}`; +} + +function eventDedupeKey(payload: EventPayload): string | undefined { + if (typeof payload.id === "string" && payload.id.length > 0) { + return payload.id; + } + + const runId = typeof payload.run_id === "string" ? payload.run_id : undefined; + const seq = typeof payload.seq === "number" ? payload.seq : undefined; + const event = typeof payload.event === "string" ? payload.event : undefined; + if (runId && seq != null && event) { + return `${runId}:${seq}:${event}`; + } + return undefined; +} diff --git a/apps/fabro-web/app/lib/run-events.test.tsx b/apps/fabro-web/app/lib/run-events.test.tsx index 14acf4f23..db35a757d 100644 --- a/apps/fabro-web/app/lib/run-events.test.tsx +++ b/apps/fabro-web/app/lib/run-events.test.tsx @@ -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"); +} diff --git a/apps/fabro-web/app/lib/run-events.ts b/apps/fabro-web/app/lib/run-events.ts index 924264af2..80394da48 100644 --- a/apps/fabro-web/app/lib/run-events.ts +++ b/apps/fabro-web/app/lib/run-events.ts @@ -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; } +interface RunEventOptions { + debounceMs?: number; + coordinator?: CrossTabSseCoordinator; +} + const subscriptions = new Map(); 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({ - subscriptions, - subscriptionKey: runId, - url: queryKeys.runs.attach(runId), + return subscribeToCrossTabSse({ + 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({ + 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;