Part of the docs-migration chore (issue #708). Remove GitHub issue numbers from inline comments, option descriptions, and rustdoc — these are contextless noise for anyone reading the code without access to the original discussions. Replace with prose that captures the same rationale directly. No functional change. Build still clean (cargo check passes).
135 lines
5.7 KiB
JavaScript
135 lines
5.7 KiB
JavaScript
// SharedWorker that holds ONE EventSource per stream URL and fans
|
||
// every server-sent event out to every connected tab via MessagePort.
|
||
//
|
||
// Problem this solves: every dashboard / agent tab opens its own
|
||
// `EventSource('/dashboard/stream')`. Browsers cap concurrent
|
||
// connections per host (~6), and Firefox throttles / disconnects
|
||
// background tabs when many are open. The result: tabs silently
|
||
// drop the SSE, fall behind, and only catch up on focus.
|
||
//
|
||
// Centralising the connection in a SharedWorker means N tabs share
|
||
// ONE backend EventSource regardless of focus state — way under the
|
||
// per-host cap, immune to per-tab throttling, and the worker survives
|
||
// any individual tab being suspended.
|
||
//
|
||
// Wire protocol (port.postMessage payloads):
|
||
//
|
||
// tab → worker
|
||
// { kind: 'subscribe', url: '/dashboard/stream' }
|
||
// { kind: 'unsubscribe', url: '/dashboard/stream' }
|
||
//
|
||
// worker → tab
|
||
// { kind: 'open', url: '...' } relayed from EventSource.onopen,
|
||
// plus a synthetic open fired to
|
||
// a brand-new subscriber when the
|
||
// upstream is already OPEN — so
|
||
// the page's onStreamOpen still
|
||
// runs and triggers a snapshot
|
||
// re-sync after a reconnect gap.
|
||
// { kind: 'message', url: '...', data: '<raw SSE data string>' }
|
||
// { kind: 'error', url: '...' } relayed from EventSource.onerror.
|
||
// { kind: 'ping' } heartbeat — fired every
|
||
// PING_INTERVAL_MS to every
|
||
// connected port. The client's
|
||
// watchdog uses these as
|
||
// proof-of-life; silence past
|
||
// ~3× the interval triggers a
|
||
// re-subscribe on a fresh port
|
||
// (recovers from Firefox killing
|
||
// the SharedWorker out from
|
||
// under us, which it does under
|
||
// memory pressure with no native
|
||
// signal to the client).
|
||
//
|
||
// Subscriptions are tracked per (port, url): a single port can
|
||
// subscribe to multiple URLs (today only one is in use but the shape
|
||
// stays open for the per-agent /events/stream multiplexing follow-up).
|
||
// Unsubscribing the last port for a URL closes the EventSource so we
|
||
// don't keep idle streams open.
|
||
|
||
const streams = new Map();
|
||
// All currently-connected ports. Used by the heartbeat tick to fan
|
||
// pings out across every tab regardless of which URLs each port is
|
||
// subscribed to. Stays disjoint from per-stream `entry.ports` (which
|
||
// is URL-scoped); a port may be in `allPorts` without any active
|
||
// subscription (e.g. between a tab loading and its first subscribe).
|
||
const allPorts = new Set();
|
||
const PING_INTERVAL_MS = 30_000;
|
||
setInterval(() => {
|
||
for (const port of allPorts) {
|
||
try { port.postMessage({ kind: 'ping' }); }
|
||
catch { /* port dead — left in the Set; see onconnect's closing comment */ }
|
||
}
|
||
}, PING_INTERVAL_MS);
|
||
|
||
function getOrCreateStream(url) {
|
||
let entry = streams.get(url);
|
||
if (entry) return entry;
|
||
const es = new EventSource(url);
|
||
entry = { es, url, ports: new Set() };
|
||
es.onopen = () => {
|
||
for (const port of entry.ports) {
|
||
try { port.postMessage({ kind: 'open', url }); }
|
||
catch { /* port dead — cleanup happens on unsubscribe / next subscribe */ }
|
||
}
|
||
};
|
||
es.onmessage = (e) => {
|
||
for (const port of entry.ports) {
|
||
try { port.postMessage({ kind: 'message', url, data: e.data }); }
|
||
catch { /* same */ }
|
||
}
|
||
};
|
||
es.onerror = () => {
|
||
for (const port of entry.ports) {
|
||
try { port.postMessage({ kind: 'error', url }); }
|
||
catch { /* same */ }
|
||
}
|
||
};
|
||
streams.set(url, entry);
|
||
return entry;
|
||
}
|
||
|
||
function unsubscribe(port, url) {
|
||
const entry = streams.get(url);
|
||
if (!entry) return;
|
||
entry.ports.delete(port);
|
||
if (entry.ports.size === 0) {
|
||
entry.es.close();
|
||
streams.delete(url);
|
||
}
|
||
}
|
||
|
||
self.onconnect = (connectEvent) => {
|
||
const port = connectEvent.ports[0];
|
||
allPorts.add(port);
|
||
const subscribedUrls = new Set();
|
||
port.onmessage = (e) => {
|
||
const msg = e.data;
|
||
if (!msg || typeof msg.url !== 'string') return;
|
||
if (msg.kind === 'subscribe') {
|
||
if (subscribedUrls.has(msg.url)) return; // idempotent
|
||
const entry = getOrCreateStream(msg.url);
|
||
entry.ports.add(port);
|
||
subscribedUrls.add(msg.url);
|
||
// Synthetic open for late subscribers — the upstream EventSource
|
||
// may already be OPEN when this tab joins, in which case the
|
||
// native onopen has long since fired and won't fire again until
|
||
// the next reconnect. Hand the new tab the open event explicitly
|
||
// so its onStreamOpen handler runs.
|
||
if (entry.es.readyState === EventSource.OPEN) {
|
||
try { port.postMessage({ kind: 'open', url: msg.url }); }
|
||
catch { /* port dead immediately — give up */ }
|
||
}
|
||
} else if (msg.kind === 'unsubscribe') {
|
||
if (!subscribedUrls.has(msg.url)) return;
|
||
unsubscribe(port, msg.url);
|
||
subscribedUrls.delete(msg.url);
|
||
}
|
||
};
|
||
// A port has no explicit "disconnect" event in the SharedWorker API
|
||
// — tabs close, the GC eventually reclaims the port, but postMessage
|
||
// to a dead port throws which the senders above catch. We don't
|
||
// proactively prune ports on a timer because the cost is bounded
|
||
// (one dead Set entry per stale tab) and the next subscribe / catch
|
||
// catches it.
|
||
};
|