| 1 | // Timing broadcast hub. |
| 2 | // |
| 3 | // The browser tab that owns the audio is the *producer*: it pushes the current |
| 4 | // source code and, a few times a second, the window of upcoming events (haps) |
| 5 | // with their character ranges in that source. Everyone else is a *viewer* — |
| 6 | // the terminal client, another browser, whatever — and just receives. |
| 7 | // |
| 8 | // Viewers animate from their own clock rather than from a per-frame feed: each |
| 9 | // `haps` message carries `now` (the cycle position at send time) and `cps`, so |
| 10 | // a viewer can advance the playhead locally at whatever framerate it likes and |
| 11 | // stay smooth across a laggy link. See tui/strudel-tui.py for a consumer. |
| 12 | import { WebSocketServer } from 'ws'; |
| 13 | |
| 14 | // The last thing each kind of state message said, replayed to anyone who joins |
| 15 | // mid-stream so a fresh viewer isn't staring at a blank screen until the next |
| 16 | // eval. `haps` is deliberately absent — stale event windows are worse than none. |
| 17 | const snapshot = { |
| 18 | code: null, |
| 19 | transport: null, |
| 20 | chat: null, |
| 21 | }; |
| 22 | |
| 23 | // Viewers are receive-only with exactly one exception: they may ask Claude for |
| 24 | // a pattern, which is the point of the terminal's prompt line. Nothing else |
| 25 | // they send is relayed, so a viewer cannot drive the editor, the transport, or |
| 26 | // anyone else's screen. |
| 27 | const MAX_PROMPT = 4000; |
| 28 | |
| 29 | const viewers = new Set(); |
| 30 | |
| 31 | // Exactly one producer is live at a time; any other open tab waits its turn in |
| 32 | // `standby` and stays silent. The incumbent keeps the slot rather than being |
| 33 | // kicked out by whoever connected last — an earlier version handed the slot to |
| 34 | // the newest tab and closed the old one, which reconnected 500ms later and |
| 35 | // took it back, so two open tabs swapped the terminal's view twice a second |
| 36 | // forever. A slot that only changes hands when it is actually free cannot |
| 37 | // oscillate. |
| 38 | let producer = null; |
| 39 | let producerSeen = 0; |
| 40 | const standby = new Set(); |
| 41 | |
| 42 | // How long an incumbent may go quiet before a newcomer may take over. The |
| 43 | // producer ticks every 500ms, so this is several missed beats — long enough |
| 44 | // not to trip on a slow frame, short enough that reloading a page doesn't |
| 45 | // leave the new tab waiting on a socket the old one never cleanly closed. |
| 46 | const STALE_MS = 3000; |
| 47 | |
| 48 | function send(ws, msg) { |
| 49 | if (ws.readyState === 1) ws.send(typeof msg === 'string' ? msg : JSON.stringify(msg)); |
| 50 | } |
| 51 | |
| 52 | function fanout(raw) { |
| 53 | for (const ws of viewers) send(ws, raw); |
| 54 | } |
| 55 | |
| 56 | function announceSource() { |
| 57 | fanout(JSON.stringify({ |
| 58 | type: 'source', |
| 59 | connected: Boolean(producer), |
| 60 | standby: standby.size, |
| 61 | })); |
| 62 | } |
| 63 | |
| 64 | // Hand the slot to `ws`. Whatever the previous producer said described *its* |
| 65 | // session, so the snapshot goes with it — keeping it would pair a dead tab's |
| 66 | // transcript with the new tab's code, which reads as real and isn't. |
| 67 | function promote(ws) { |
| 68 | producer = ws; |
| 69 | producerSeen = Date.now(); |
| 70 | standby.delete(ws); |
| 71 | snapshot.code = snapshot.transport = snapshot.chat = null; |
| 72 | send(ws, { type: 'active' }); |
| 73 | announceSource(); |
| 74 | } |
| 75 | |
| 76 | export function attachBroadcast(server, { path = '/ws' } = {}) { |
| 77 | const wss = new WebSocketServer({ server, path }); |
| 78 | |
| 79 | wss.on('connection', (ws, req) => { |
| 80 | // Role comes from the query string so it is known before the first message: |
| 81 | // /ws?role=producer from the app, anything else is a viewer. |
| 82 | const role = new URL(req.url, 'http://x').searchParams.get('role') === 'producer' |
| 83 | ? 'producer' |
| 84 | : 'viewer'; |
| 85 | ws.__role = role; |
| 86 | |
| 87 | if (role === 'producer') { |
| 88 | const held = producer && producer !== ws && producer.readyState === 1 |
| 89 | && Date.now() - producerSeen < STALE_MS; |
| 90 | if (held) { |
| 91 | // Someone else is live. Wait, quietly — `standby` tells the page to |
| 92 | // stop sending, and it gets promoted if the incumbent goes away. |
| 93 | standby.add(ws); |
| 94 | send(ws, { type: 'standby' }); |
| 95 | announceSource(); |
| 96 | } else { |
| 97 | // The incumbent is gone or has stopped talking (a tab closed without a |
| 98 | // clean close, usually). Park it and take over. |
| 99 | if (producer && producer !== ws) { |
| 100 | standby.add(producer); |
| 101 | send(producer, { type: 'standby' }); |
| 102 | } |
| 103 | promote(ws); |
| 104 | } |
| 105 | } else { |
| 106 | viewers.add(ws); |
| 107 | send(ws, { |
| 108 | type: 'hello', |
| 109 | viewers: viewers.size, |
| 110 | source: Boolean(producer), |
| 111 | standby: standby.size, |
| 112 | }); |
| 113 | if (snapshot.code) send(ws, snapshot.code); |
| 114 | if (snapshot.transport) send(ws, snapshot.transport); |
| 115 | if (snapshot.chat) send(ws, snapshot.chat); |
| 116 | // The snapshot only covers what this hub has seen the producer say. Ask |
| 117 | // it to re-announce everything as well, so a viewer joining a session |
| 118 | // that has been quiet — or one whose state this hub predates — gets a |
| 119 | // full picture instead of waiting for the next change. |
| 120 | if (producer?.readyState === 1) send(producer, { type: 'refresh' }); |
| 121 | } |
| 122 | |
| 123 | ws.on('message', (data) => { |
| 124 | const raw = data.toString(); |
| 125 | let msg; |
| 126 | try { |
| 127 | msg = JSON.parse(raw); |
| 128 | } catch { |
| 129 | return; |
| 130 | } |
| 131 | |
| 132 | if (role !== 'producer') { |
| 133 | // See MAX_PROMPT above: a prompt is the only thing a viewer may send. |
| 134 | if (msg.type !== 'prompt' || typeof msg.text !== 'string') return; |
| 135 | const text = msg.text.slice(0, MAX_PROMPT).trim(); |
| 136 | if (!text) return; |
| 137 | if (producer?.readyState === 1) send(producer, { type: 'prompt', text }); |
| 138 | else send(ws, { type: 'notice', text: 'no browser is connected to send that to' }); |
| 139 | return; |
| 140 | } |
| 141 | |
| 142 | // A tab that is on standby has been told to stop, but may have had a |
| 143 | // message already in flight. Only the live producer reaches the viewers. |
| 144 | if (ws !== producer) return; |
| 145 | producerSeen = Date.now(); |
| 146 | |
| 147 | // Keep the join-snapshot current, but relay the original bytes — no |
| 148 | // reserialization of the hot `haps` path. |
| 149 | if (msg.type === 'code') snapshot.code = raw; |
| 150 | else if (msg.type === 'transport') snapshot.transport = raw; |
| 151 | else if (msg.type === 'chat') snapshot.chat = raw; |
| 152 | fanout(raw); |
| 153 | }); |
| 154 | |
| 155 | const drop = () => { |
| 156 | viewers.delete(ws); |
| 157 | standby.delete(ws); |
| 158 | if (producer !== ws) { |
| 159 | if (role === 'producer') announceSource(); // the standby count changed |
| 160 | return; |
| 161 | } |
| 162 | producer = null; |
| 163 | snapshot.transport = null; |
| 164 | // Hand the slot to whoever is still waiting — most recent first, since |
| 165 | // that is the tab most likely to be the one in use. |
| 166 | const next = [...standby].reverse().find((c) => c.readyState === 1); |
| 167 | if (next) promote(next); |
| 168 | else announceSource(); |
| 169 | }; |
| 170 | |
| 171 | ws.on('close', drop); |
| 172 | ws.on('error', drop); |
| 173 | }); |
| 174 | |
| 175 | // Drop half-open connections (laptop closed, tunnel died) so `viewers` and |
| 176 | // the producer slot reflect reality. |
| 177 | const heartbeat = setInterval(() => { |
| 178 | for (const ws of wss.clients) { |
| 179 | if (ws.__alive === false) { |
| 180 | ws.terminate(); |
| 181 | continue; |
| 182 | } |
| 183 | ws.__alive = false; |
| 184 | ws.ping(); |
| 185 | } |
| 186 | }, 30_000); |
| 187 | wss.on('connection', (ws) => { |
| 188 | ws.__alive = true; |
| 189 | ws.on('pong', () => { |
| 190 | ws.__alive = true; |
| 191 | }); |
| 192 | }); |
| 193 | wss.on('close', () => clearInterval(heartbeat)); |
| 194 | |
| 195 | return wss; |
| 196 | } |