dashboard: SharedWorker for SSE multiplexing (closes #448)

mara on #448: "firefox disconnects bc of too many tabs. needs bg
service worker". picked SharedWorker over full Service Worker:
smaller change, addresses the actual problem (shared connection
across tabs), no offline-cache scope creep.

architecture

per-tab `new EventSource('/dashboard/stream')` replaced with a
SharedWorker-backed facade. one SharedWorker instance per origin
holds ONE upstream EventSource and fans every server-sent event
out to every connected tab via MessagePort. N hyperhive tabs now
share ONE backend connection, immune to Firefox's per-tab SSE
throttling under many-open-tabs pressure.

wire protocol (port.postMessage):

  tab → worker
    { kind: 'subscribe',   url: '/dashboard/stream' }
    { kind: 'unsubscribe', url: '/dashboard/stream' }

  worker → tab
    { kind: 'open',    url }
    { kind: 'message', url, data: '<raw SSE data>' }
    { kind: 'error',   url }

subscription tracking is per (port, url). a late subscriber that
joins after the upstream is already OPEN gets a synthetic 'open'
event so its onStreamOpen handler still runs (triggers the
snapshot re-sync that recovers events lost during the join gap).
unsubscribing the last port for a URL closes the upstream
EventSource so we don't leak idle streams.

files

- frontend/packages/dashboard/src/stream-worker.js: new — the
  worker. multi-URL multiplexing via Map<url, {es, ports}>.
- frontend/packages/dashboard/src/common.js: new exported helper
  openStream(url) — returns an EventSource-shaped facade backed
  by the SharedWorker. graceful fallback to direct EventSource
  when SharedWorker is unavailable.
- frontend/packages/dashboard/src/app.js: replaces the inline
  new EventSource('/dashboard/stream') with openStream.
- frontend/packages/dashboard/src/flow.js: passes
  streamFactory: openStream to termCreate so the broker
  terminal's SSE goes through the worker too.
- frontend/packages/shared/src/terminal.js: accepts an optional
  streamFactory(url) option. default unchanged — non-dashboard
  consumers (per-agent UI) keep using direct EventSource.
- frontend/packages/dashboard/build.mjs: new esbuild entry for
  stream-worker.js → dist/static/stream-worker.js (separate
  bundle because SharedWorker scripts run in a different global
  scope and can't be inlined into app.js).

scope kept tight

- per-agent UI's /events/stream stays on direct EventSource. the
  agent UI's tab count per agent is typically 1; SharedWorker
  helps when you have N tabs hitting the SAME stream and the
  per-agent stream URLs differ. if mara wants the agent UI to
  share its workers too it's a separate small PR.
- no offline-cache, no push notifications — those need full
  Service Worker; explicit non-goal here per the design Q.

validation

- npm run build --workspace=@hive/dashboard clean.
- stream-worker.js bundle: 1.8 kb.
- app.js: 154 kb → 158 kb. flow.js: 29.9 kb → 32 kb.
- browser smoke test isn't possible from inside iris's container;
  the EventSource-shaped facade preserves the exact onmessage /
  onopen / onerror surface the existing IIFE consumers use.
This commit is contained in:
iris 2026-05-26 01:15:43 +02:00 committed by Mara
commit 4504f9ede3
6 changed files with 246 additions and 3 deletions

View file

@ -45,6 +45,24 @@ await build({
logLevel: 'info',
});
// Stream-worker entry (#448). Lives in a separate bundle: SharedWorker
// scripts run in a different global (`self` is the worker scope, no
// `window`) so they can't be inlined into app.js / flow.js. Output is
// at `static/stream-worker.js`; common.js's `openStream` references
// `/static/stream-worker.js` as the SharedWorker URL. ES-module-shaped
// (so a future `import` from within the worker can pull in shared
// utilities), browser target — webworker target subset.
await build({
entryPoints: [src('stream-worker.js')],
outdir: staticDir(''),
bundle: true,
format: 'esm',
platform: 'browser',
target: ['es2022'],
sourcemap: true,
logLevel: 'info',
});
// Bundle the CSS — esbuild resolves @import including the package
// re-exports from @hive/shared.
await build({

View file

@ -18,6 +18,7 @@ import {
fmtAgeSecs,
Panel, NOTIF,
makePathLink, appendText, appendLinkified,
openStream,
} from './common.js';
// mdNode (in common.js) reads `window.marked` for the markdown side
@ -2069,7 +2070,14 @@ window.marked = marked;
rebuild_queue_changed: applyRebuildQueueChanged,
};
(function bindDashboardStream() {
const es = new EventSource('/dashboard/stream');
// #448: route the EventSource through a SharedWorker so all open
// hyperhive tabs share ONE backend SSE connection. Survives
// Firefox's per-tab connection throttling under many-open-tabs
// pressure (the actual mara symptom). `openStream` returns an
// EventSource-shaped facade so the rest of this IIFE is unchanged;
// graceful fallback to direct `new EventSource` when SharedWorker
// isn't supported.
const es = openStream('/dashboard/stream');
es.onmessage = (e) => {
let ev;
try { ev = JSON.parse(e.data); } catch { return; }

View file

@ -58,6 +58,97 @@ export const form = (action, btnClass, btnLabel, confirmMsg, extra = {}, opts =
// tied to its caller, so they don't generalise cleanly. We can lift
// them when a second consumer needs the same shape.
// ─── shared-worker SSE pipe (#448) ──────────────────────────────────────
// Returns an EventSource-shaped object backed by a SharedWorker that
// holds ONE upstream `new EventSource(url)` and fans events out to
// every connected tab. Replaces direct `new EventSource(url)` at the
// dashboard's two consumer sites (app.js inline + flow.js via
// terminal.js's `streamFactory` option) so N hyperhive tabs share
// ONE backend connection — way under the browser's per-host
// connection cap, immune to per-tab throttling that drops the SSE
// when Firefox suspends background tabs.
//
// Graceful fallback to direct EventSource on environments without
// SharedWorker (some embedded browsers, some Safari versions). The
// per-tab connection cost is the same as today — no regression.
//
// The page consumer uses the returned object like a regular
// EventSource: assign `onmessage` / `onopen` / `onerror`. `.close()`
// tells the worker to drop the subscription; the worker closes the
// upstream EventSource when the last subscriber leaves.
const SHARED_WORKER_PATH = '/static/stream-worker.js';
const SHARED_WORKER_NAME = 'hyperhive-stream';
// One SharedWorker instance per page, reused by all openStream calls
// on that page. Lazy — pages with no streams don't spawn the worker.
let _sharedPort = null;
function getSharedPort() {
if (_sharedPort) return _sharedPort;
if (typeof SharedWorker === 'undefined') return null;
try {
const sw = new SharedWorker(SHARED_WORKER_PATH, SHARED_WORKER_NAME);
_sharedPort = sw.port;
_sharedPort.start();
return _sharedPort;
} catch (err) {
console.warn('SharedWorker unavailable, falling back to direct EventSource:', err);
return null;
}
}
export function openStream(url) {
const port = getSharedPort();
if (!port) return new EventSource(url);
// Build an EventSource-shaped facade so consumer code is unchanged.
// `target.onmessage` / `onopen` / `onerror` are assigned by the
// consumer; the routing function below forwards events received
// from the worker (filtered by url, since one port can multiplex
// multiple subscriptions).
const target = {
onmessage: null,
onopen: null,
onerror: null,
readyState: 0, // CONNECTING
close() {
try { port.postMessage({ kind: 'unsubscribe', url }); }
catch { /* port dead */ }
port.removeEventListener('message', route);
},
};
const route = (e) => {
const m = e.data;
if (!m || m.url !== url) return;
if (m.kind === 'open') {
target.readyState = 1; // OPEN
if (target.onopen) {
try { target.onopen({ target }); }
catch (err) { console.error('openStream onopen threw', err); }
}
} else if (m.kind === 'message') {
if (target.onmessage) {
try { target.onmessage({ data: m.data, target }); }
catch (err) { console.error('openStream onmessage threw', err); }
}
} else if (m.kind === 'error') {
if (target.onerror) {
try { target.onerror({ target }); }
catch (err) { console.error('openStream onerror threw', err); }
}
}
};
port.addEventListener('message', route);
port.postMessage({ kind: 'subscribe', url });
// Drop the subscription when the tab unloads so the worker can
// close the upstream EventSource when the last subscriber leaves.
// `pagehide` fires for both real unloads and bfcache transitions.
window.addEventListener('pagehide', () => {
try { port.postMessage({ kind: 'unsubscribe', url }); }
catch { /* port dead — worker side already cleaned up */ }
});
return target;
}
// ─── side panel ─────────────────────────────────────────────────────────
// Singleton drawer that swipes in from the right. Long content
// (file previews, approval diffs, journald logs, applied config)

View file

@ -16,6 +16,7 @@ import {
$, el,
Panel, NOTIF,
appendLinkified,
openStream,
} from './common.js';
(() => {
@ -185,6 +186,10 @@ import {
pillAnchor: flowMain,
historyUrl: '/dashboard/history',
streamUrl: '/dashboard/stream',
// #448: route through the SharedWorker so this page's SSE shares
// a single backend connection with /index.html (and any other
// open hyperhive tab).
streamFactory: openStream,
renderers: {
sent: (ev, api) => renderMsg(ev, api, '→'),
delivered: (ev, api) => renderMsg(ev, api, '✓'),

View file

@ -0,0 +1,110 @@
// 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.
//
// 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();
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];
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.
};

View file

@ -318,7 +318,15 @@ export function create(opts) {
let live = false;
let buffered = [];
const es = new EventSource(opts.streamUrl);
// #448: callers can supply a `streamFactory(url)` that returns an
// EventSource-shaped object (must expose onmessage/onopen/onerror
// + .close()). The dashboard pages pass a SharedWorker-backed
// factory so all open hyperhive tabs share ONE upstream SSE
// connection. Default keeps the direct `new EventSource(url)`
// behaviour so non-dashboard consumers (per-agent UI) are unchanged.
const es = opts.streamFactory
? opts.streamFactory(opts.streamUrl)
: new EventSource(opts.streamUrl);
es.onmessage = (e) => {
let ev;
try { ev = JSON.parse(e.data); }
@ -331,7 +339,10 @@ export function create(opts) {
}
};
es.onerror = () => {
if (es.readyState === EventSource.CONNECTING) row('note', '[reconnecting…]');
// SharedWorker-backed facades expose `readyState` mirroring the
// upstream EventSource state; the native EventSource exposes the
// same. Either way the CONNECTING vs. closed distinction works.
if (es.readyState === 0 /* CONNECTING */) row('note', '[reconnecting…]');
else row('note', '[disconnected]');
};
es.onopen = () => {