diff --git a/frontend/packages/dashboard/build.mjs b/frontend/packages/dashboard/build.mjs index e8cf87b3..b18ebf06 100644 --- a/frontend/packages/dashboard/build.mjs +++ b/frontend/packages/dashboard/build.mjs @@ -45,29 +45,6 @@ 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. `format: 'iife'` -// matches the classic-script load (`new SharedWorker(url, name)` with -// no `{ type: 'module' }`); Firefox is the #448 target and module -// SharedWorker support there is patchy, so keeping the worker as a -// classic script + IIFE bundle is the compatible default. argus nit -// on #453: if a future contributor adds an `import` to this bundle, -// the IIFE format will surface it as a build error rather than -// silently shipping broken code. -await build({ - entryPoints: [src('stream-worker.js')], - outdir: staticDir(''), - bundle: true, - format: 'iife', - platform: 'browser', - target: ['es2022'], - sourcemap: true, - logLevel: 'info', -}); - // Bundle the CSS — esbuild resolves @import including the package // re-exports from @hive/shared. await build({ diff --git a/frontend/packages/dashboard/src/app.js b/frontend/packages/dashboard/src/app.js index 7f55eb38..3a597a8d 100644 --- a/frontend/packages/dashboard/src/app.js +++ b/frontend/packages/dashboard/src/app.js @@ -18,7 +18,6 @@ import { fmtAgeSecs, Panel, NOTIF, makePathLink, appendText, appendLinkified, - openStream, } from './common.js'; // mdNode (in common.js) reads `window.marked` for the markdown side @@ -2070,14 +2069,7 @@ window.marked = marked; rebuild_queue_changed: applyRebuildQueueChanged, }; (function bindDashboardStream() { - // #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'); + const es = new EventSource('/dashboard/stream'); es.onmessage = (e) => { let ev; try { ev = JSON.parse(e.data); } catch { return; } diff --git a/frontend/packages/dashboard/src/common.js b/frontend/packages/dashboard/src/common.js index 8998ff46..e0870207 100644 --- a/frontend/packages/dashboard/src/common.js +++ b/frontend/packages/dashboard/src/common.js @@ -58,154 +58,6 @@ 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 port per page, reused by all openStream calls on -// that page. Invalidated on `pagehide` so a bfcache restore picks up -// a fresh port (the cached one may have been collected if all other -// tabs closed while this page was frozen — argus nit on #453). -let _sharedPort = null; -function makeSharedPort() { - if (typeof SharedWorker === 'undefined') return null; - try { - const sw = new SharedWorker(SHARED_WORKER_PATH, SHARED_WORKER_NAME); - sw.port.start(); - return sw.port; - } catch (err) { - console.warn('SharedWorker unavailable, falling back to direct EventSource:', err); - return null; - } -} -function getSharedPort() { - if (!_sharedPort) _sharedPort = makeSharedPort(); - return _sharedPort; -} - -// Registry of live subscriptions on this page. Keyed by url so a -// second openStream call for the same URL (would only happen on a -// hypothetical multi-consumer page) attaches to the existing route -// rather than overlapping. Each entry caches the route function so -// bfcache-restore re-bind can re-attach it to the fresh port. -// -// Today's pages only call openStream once with one URL; the registry -// shape just keeps the bfcache-restore path correct if that changes -// (e.g. /index.html later subscribing to two streams). -const _activeSubs = new Map(); - -// One-shot wiring of the page-wide lifecycle hooks: on bfcache -// freeze (`pagehide { persisted: true }`) we unsubscribe so the -// worker can close the upstream when the last live subscriber -// leaves; on bfcache restore (`pageshow { persisted: true }`) we -// invalidate the cached port (it may be dead if all other tabs -// closed during the freeze) and re-attach every active subscription -// to a fresh port. argus nit on #453: without this, the consumer's -// onmessage stays bound but no events flow after a bfcache restore. -let _lifecycleBound = false; -function bindLifecycleOnce() { - if (_lifecycleBound) return; - _lifecycleBound = true; - window.addEventListener('pagehide', () => { - if (!_sharedPort) return; - for (const url of _activeSubs.keys()) { - try { _sharedPort.postMessage({ kind: 'unsubscribe', url }); } - catch { /* port dead — worker side already cleaned up */ } - } - // Drop port routes too; the bfcache-restore path will re-add - // them on a fresh port. Leaving stale routes on a dead port - // would just keep a closure alive without cost, but cleaning - // up keeps the registry shape honest. - for (const sub of _activeSubs.values()) { - try { _sharedPort.removeEventListener('message', sub.route); } - catch { /* same */ } - } - _sharedPort = null; - }); - window.addEventListener('pageshow', (ev) => { - if (!ev.persisted) return; // cold load — openStream just bound listeners - if (!_activeSubs.size) return; - const port = getSharedPort(); - if (!port) return; // SharedWorker really gone; fallback already in place - for (const [url, sub] of _activeSubs) { - sub.target.readyState = 0; // CONNECTING — the worker will fire 'open' - port.addEventListener('message', sub.route); - try { port.postMessage({ kind: 'subscribe', url }); } - catch { /* port dead immediately — skip */ } - } - }); -} - -export function openStream(url) { - const port = getSharedPort(); - if (!port) return new EventSource(url); - bindLifecycleOnce(); - - // 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() { - const p = _sharedPort; - if (p) { - try { p.postMessage({ kind: 'unsubscribe', url }); } - catch { /* port dead */ } - try { p.removeEventListener('message', route); } - catch { /* same */ } - } - _activeSubs.delete(url); - }, - }; - 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); } - } - } - }; - _activeSubs.set(url, { target, route }); - port.addEventListener('message', route); - port.postMessage({ kind: 'subscribe', url }); - return target; -} - // ─── side panel ───────────────────────────────────────────────────────── // Singleton drawer that swipes in from the right. Long content // (file previews, approval diffs, journald logs, applied config) diff --git a/frontend/packages/dashboard/src/flow.js b/frontend/packages/dashboard/src/flow.js index 41487071..18f0a0f7 100644 --- a/frontend/packages/dashboard/src/flow.js +++ b/frontend/packages/dashboard/src/flow.js @@ -16,7 +16,6 @@ import { $, el, Panel, NOTIF, appendLinkified, - openStream, } from './common.js'; (() => { @@ -186,10 +185,6 @@ 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, '✓'), diff --git a/frontend/packages/dashboard/src/stream-worker.js b/frontend/packages/dashboard/src/stream-worker.js deleted file mode 100644 index 3d7af676..00000000 --- a/frontend/packages/dashboard/src/stream-worker.js +++ /dev/null @@ -1,110 +0,0 @@ -// 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: '' } -// { 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. -}; diff --git a/frontend/packages/shared/src/terminal.js b/frontend/packages/shared/src/terminal.js index b47b08c2..f7677805 100644 --- a/frontend/packages/shared/src/terminal.js +++ b/frontend/packages/shared/src/terminal.js @@ -318,15 +318,7 @@ export function create(opts) { let live = false; let buffered = []; - // #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); + const es = new EventSource(opts.streamUrl); es.onmessage = (e) => { let ev; try { ev = JSON.parse(e.data); } @@ -339,10 +331,7 @@ export function create(opts) { } }; es.onerror = () => { - // 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…]'); + if (es.readyState === EventSource.CONNECTING) row('note', '[reconnecting…]'); else row('note', '[disconnected]'); }; es.onopen = () => {