argus pointed out the inline catch comment claimed dead ports get cleaned up on next subscribe, but nothing actually prunes the allPorts Set on subscribe — the honest answer is the one already at the bottom of onconnect: dead entries are left in the Set, the bound cost is acceptable, and postMessage's throw is the ambient signal we use. Point at that comment instead of repeating a wrong description.
136 lines
5.7 KiB
JavaScript
136 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 (#448 — mara: "firefox disconnects bc of too
|
||
// many tabs"): 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' } #515: 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.
|
||
};
|