File
Blob: worker/src/client/radio/session.ts
| 1 | import { acceptMetadata, matchesPlayback } from './playback.ts'; |
| 2 | import type { NowPlaying } from '@/shared/contracts/track.ts'; |
| 3 | import { radioApi, ApiError, errorText } from './api.ts'; |
| 4 | import type { RadioApi } from './api.ts'; |
| 5 | import { commandId, parseRobot, parseSpectrum } from './protocol.ts'; |
| 6 | import type { Command } from '@/shared/contracts/robot.ts'; |
| 7 | import type { Joined } from '@/shared/contracts/signaling.ts'; |
| 8 | import { initialState } from './session-state.ts'; |
| 9 | import type { IssueSource, SessionState, SignalBuffer } from './session-state.ts'; |
| 10 | import { acknowledge, gatheredAnswer } from './webrtc.ts'; |
| 11 | |
| 12 | /** Owns a viewer's peer, browser audio, timers, and command acknowledgments. |
| 13 | * React subscribes to snapshots; the spectrum canvas reads a separate sample buffer. |
| 14 | */ |
| 15 | export class RadioSession { |
| 16 | readonly signal: SignalBuffer = { at: 0, count: 0 }; |
| 17 | private state = initialState(); |
| 18 | private listeners = new Set<() => void>(); |
| 19 | private api: RadioApi; |
| 20 | private peer?: RTCPeerConnection; |
| 21 | private member?: Joined; |
| 22 | private robot?: RTCDataChannel; |
| 23 | private audio?: HTMLAudioElement; |
| 24 | private context?: AudioContext; |
| 25 | private source?: MediaElementAudioSourceNode; |
| 26 | private epoch = 0; |
| 27 | private pending?: { id: string; at: number; timer: ReturnType<typeof setTimeout> }; |
| 28 | private heartbeatBusy = false; |
| 29 | private statusBusy = false; |
| 30 | private statsBusy = false; |
| 31 | |
| 32 | constructor(api: RadioApi = radioApi) { |
| 33 | this.api = api; |
| 34 | } |
| 35 | getSnapshot = (): SessionState => this.state; |
| 36 | subscribe = (listener: () => void): (() => void) => { |
| 37 | this.listeners.add(listener); |
| 38 | return () => { |
| 39 | this.listeners.delete(listener); |
| 40 | }; |
| 41 | }; |
| 42 | private update(patch: Partial<SessionState>) { |
| 43 | this.state = { ...this.state, ...patch }; |
| 44 | for (const listener of this.listeners) listener(); |
| 45 | } |
| 46 | private issue(source: IssueSource, message?: string) { |
| 47 | this.update({ issues: { ...this.state.issues, [source]: message } }); |
| 48 | } |
| 49 | get held(): boolean { |
| 50 | return ( |
| 51 | this.state.phase === 'connected' && |
| 52 | !!this.member && |
| 53 | this.state.room.controller === this.member.id |
| 54 | ); |
| 55 | } |
| 56 | |
| 57 | start(audio: HTMLAudioElement): () => void { |
| 58 | this.audio = audio; |
| 59 | audio.volume = this.state.volume; |
| 60 | void this.refreshStatus(); |
| 61 | const statusTimer = setInterval(() => { |
| 62 | if (!document.hidden) void this.refreshStatus(); |
| 63 | }, 5000); |
| 64 | const heartbeatTimer = setInterval(() => { |
| 65 | void this.heartbeat(); |
| 66 | }, 5000); |
| 67 | const clockTimer = setInterval(() => { |
| 68 | this.update({ now: performance.now(), spectrumCount: this.signal.count }); |
| 69 | void this.collectStats(); |
| 70 | }, 1000); |
| 71 | const pagehide = () => { |
| 72 | if (this.member) |
| 73 | navigator.sendBeacon( |
| 74 | `/api/viewers/${this.member.id}/leave`, |
| 75 | new Blob([JSON.stringify({ viewerToken: this.member.viewerToken })], { |
| 76 | type: 'application/json', |
| 77 | }), |
| 78 | ); |
| 79 | if (this.peer || this.member || this.state.phase !== 'idle') void this.disconnect(); |
| 80 | }; |
| 81 | window.addEventListener('pagehide', pagehide); |
| 82 | return () => { |
| 83 | clearInterval(statusTimer); |
| 84 | clearInterval(heartbeatTimer); |
| 85 | clearInterval(clockTimer); |
| 86 | window.removeEventListener('pagehide', pagehide); |
| 87 | if (this.peer || this.member || this.state.phase !== 'idle') void this.disconnect(); |
| 88 | // A StrictMode effect replay can reuse the same media element. Keep its |
| 89 | // one MediaElementSource and suspend the context until the next gesture. |
| 90 | void this.context?.suspend(); |
| 91 | }; |
| 92 | } |
| 93 | |
| 94 | async refreshStatus(): Promise<void> { |
| 95 | if (this.statusBusy) return; |
| 96 | this.statusBusy = true; |
| 97 | try { |
| 98 | const room = await this.api.getStatus(); |
| 99 | const changed = room.generation !== this.state.room.generation; |
| 100 | if (changed) { |
| 101 | this.signal.frame = undefined; |
| 102 | this.signal.at = 0; |
| 103 | this.update({ nowPlaying: undefined, paused: undefined }); |
| 104 | } |
| 105 | this.issue('status'); |
| 106 | this.update({ room, auth: 'ready' }); |
| 107 | if (room.nowPlaying) this.acceptNowPlaying(room.nowPlaying); |
| 108 | if (this.member && (!room.online || room.generation !== this.member.generation)) { |
| 109 | await this.disconnect('Board connection changed. Start listening when it is online.'); |
| 110 | } else if (this.state.phase === 'idle') { |
| 111 | this.update({ |
| 112 | message: room.online |
| 113 | ? 'Board online. Start listening.' |
| 114 | : 'Board offline. You can still explore the diagram.', |
| 115 | }); |
| 116 | } |
| 117 | } catch (error) { |
| 118 | if (error instanceof ApiError && error.status === 401) { |
| 119 | this.issue('status'); |
| 120 | this.update({ auth: 'required' }); |
| 121 | if (this.member) await this.disconnect('Enter the viewer password to reconnect.'); |
| 122 | } else { |
| 123 | this.issue('status', 'Radio service unavailable. Retrying shortly.'); |
| 124 | } |
| 125 | } finally { |
| 126 | this.statusBusy = false; |
| 127 | } |
| 128 | } |
| 129 | |
| 130 | async login(password: string): Promise<void> { |
| 131 | await this.api.login(password); |
| 132 | this.issue('connection'); |
| 133 | this.update({ auth: 'ready' }); |
| 134 | await this.refreshStatus(); |
| 135 | } |
| 136 | |
| 137 | async connect(): Promise<void> { |
| 138 | if (this.state.phase !== 'idle' || !this.audio) return; |
| 139 | const epoch = ++this.epoch; |
| 140 | const current = () => { |
| 141 | if (epoch !== this.epoch) throw new Error('Connection cancelled.'); |
| 142 | }; |
| 143 | this.issue('connection'); |
| 144 | this.update({ phase: 'connecting', message: 'Opening data channels…' }); |
| 145 | try { |
| 146 | // Resume within the user gesture. There is no microphone or camera track. |
| 147 | if (typeof AudioContext !== 'undefined') { |
| 148 | this.context ??= new AudioContext(); |
| 149 | if (!this.source) { |
| 150 | this.source = this.context.createMediaElementSource(this.audio); |
| 151 | this.source.connect(this.context.destination); |
| 152 | } |
| 153 | await this.context.resume(); |
| 154 | } |
| 155 | current(); |
| 156 | const peer = new RTCPeerConnection({ iceServers: [], bundlePolicy: 'max-bundle' }); |
| 157 | this.peer = peer; |
| 158 | peer.ontrack = ({ track }) => { |
| 159 | if (this.peer !== peer || track.kind !== 'audio' || !this.audio) return; |
| 160 | this.audio.srcObject = new MediaStream([track]); |
| 161 | void this.audio.play().catch(() => { |
| 162 | if (this.peer === peer) this.update({ audioBlocked: true }); |
| 163 | }); |
| 164 | }; |
| 165 | peer.onconnectionstatechange = () => { |
| 166 | if (this.peer !== peer) return; |
| 167 | if (peer.connectionState === 'failed') |
| 168 | void this.disconnect('Connection closed. Start listening again.'); |
| 169 | else if (peer.connectionState === 'disconnected') |
| 170 | this.update({ message: 'Connection interrupted. Waiting for the signal…' }); |
| 171 | else if (peer.connectionState === 'connected' && this.state.phase === 'connected') |
| 172 | this.update({ message: 'Connected to the board.' }); |
| 173 | }; |
| 174 | const joined = await this.api.joinViewer(); |
| 175 | if (epoch !== this.epoch) { |
| 176 | await this.leave(joined); |
| 177 | return; |
| 178 | } |
| 179 | this.member = joined; |
| 180 | this.update({ viewerId: joined.id }); |
| 181 | await peer.setRemoteDescription(joined.sessionDescription); |
| 182 | current(); |
| 183 | const answer = await gatheredAnswer(peer); |
| 184 | current(); |
| 185 | const response = await this.api.answerViewer(joined, answer); |
| 186 | current(); |
| 187 | const channels = response.channels.map((config) => { |
| 188 | const channel = peer.createDataChannel(config.dataChannelName, { |
| 189 | negotiated: true, |
| 190 | id: config.id, |
| 191 | ordered: config.ordered, |
| 192 | ...('maxRetransmits' in config ? { maxRetransmits: config.maxRetransmits } : {}), |
| 193 | }); |
| 194 | channel.binaryType = 'arraybuffer'; |
| 195 | channel.onmessage = (event) => { |
| 196 | if (this.peer === peer) this.receive(config.dataChannelName, event.data); |
| 197 | }; |
| 198 | channel.onclose = () => { |
| 199 | if (this.peer === peer && this.state.phase === 'connected') |
| 200 | void this.disconnect('Data channel closed. Start listening again.'); |
| 201 | }; |
| 202 | if (config.dataChannelName === 'robot') this.robot = channel; |
| 203 | return channel; |
| 204 | }); |
| 205 | await Promise.all(channels.map(acknowledge)); |
| 206 | current(); |
| 207 | this.update({ message: 'Connecting audio…' }); |
| 208 | const remote = await this.api.subscribeAudio(joined); |
| 209 | current(); |
| 210 | await peer.setRemoteDescription(remote.sessionDescription); |
| 211 | current(); |
| 212 | const audioAnswer = await gatheredAnswer(peer); |
| 213 | current(); |
| 214 | await this.api.renegotiateViewer(joined, audioAnswer); |
| 215 | current(); |
| 216 | this.update({ phase: 'connected', message: 'Connected to the board.' }); |
| 217 | await this.refreshStatus(); |
| 218 | } catch (error) { |
| 219 | if (epoch === this.epoch) { |
| 220 | await this.disconnect(); |
| 221 | this.issue('connection', errorText(error)); |
| 222 | } |
| 223 | } |
| 224 | } |
| 225 | |
| 226 | async disconnect(message = 'Disconnected. Shared playback is unchanged.'): Promise<void> { |
| 227 | ++this.epoch; |
| 228 | const leaving = this.member; |
| 229 | this.member = undefined; |
| 230 | const peer = this.peer; |
| 231 | this.peer = undefined; |
| 232 | this.robot = undefined; |
| 233 | peer?.close(); |
| 234 | if (this.audio) { |
| 235 | this.audio.pause(); |
| 236 | this.audio.srcObject = null; |
| 237 | } |
| 238 | clearTimeout(this.pending?.timer); |
| 239 | this.pending = undefined; |
| 240 | this.signal.frame = undefined; |
| 241 | this.signal.at = 0; |
| 242 | this.signal.count = 0; |
| 243 | const initial = initialState(); |
| 244 | this.update({ |
| 245 | phase: 'idle', |
| 246 | viewerId: undefined, |
| 247 | paused: undefined, |
| 248 | issues: { status: this.state.issues.status, connection: this.state.issues.connection }, |
| 249 | telemetry: undefined, |
| 250 | led: undefined, |
| 251 | telemetryAt: 0, |
| 252 | received: 0, |
| 253 | spectrumCount: 0, |
| 254 | audio: initial.audio, |
| 255 | ack: initial.ack, |
| 256 | busy: false, |
| 257 | audioBlocked: false, |
| 258 | message, |
| 259 | }); |
| 260 | if (leaving) await this.leave(leaving); |
| 261 | } |
| 262 | private async leave(member: Joined) { |
| 263 | try { |
| 264 | await this.api.leaveViewer(member); |
| 265 | } catch { |
| 266 | /* Server membership expires after inactivity. */ |
| 267 | } |
| 268 | } |
| 269 | private async heartbeat() { |
| 270 | const member = this.member; |
| 271 | if (!member || this.state.phase !== 'connected' || this.heartbeatBusy) return; |
| 272 | this.heartbeatBusy = true; |
| 273 | try { |
| 274 | const result = await this.api.heartbeatViewer(member); |
| 275 | if (this.member === member) { |
| 276 | this.issue('heartbeat'); |
| 277 | this.update({ room: { ...this.state.room, controller: result.controller } }); |
| 278 | } |
| 279 | } catch (error) { |
| 280 | if (this.member !== member) return; |
| 281 | this.issue('heartbeat', 'Connection renewal failed. Controls are unavailable.'); |
| 282 | this.update({ room: { ...this.state.room, controller: null } }); |
| 283 | if (error instanceof ApiError && [401, 403, 409].includes(error.status)) |
| 284 | await this.disconnect('Session expired. Start listening again.'); |
| 285 | } finally { |
| 286 | this.heartbeatBusy = false; |
| 287 | } |
| 288 | } |
| 289 | |
| 290 | async toggleControl(): Promise<void> { |
| 291 | const member = this.member; |
| 292 | if (!member || this.state.phase !== 'connected' || this.state.busy) return; |
| 293 | this.issue('control'); |
| 294 | this.update({ busy: true }); |
| 295 | try { |
| 296 | const result = this.held |
| 297 | ? await this.api.releaseControl(member) |
| 298 | : await this.api.claimControl(member); |
| 299 | if (this.member === member) |
| 300 | this.update({ room: { ...this.state.room, controller: result.controller } }); |
| 301 | } catch (error) { |
| 302 | if (this.member === member) this.issue('control', errorText(error)); |
| 303 | } finally { |
| 304 | if (this.member === member) this.update({ busy: false }); |
| 305 | } |
| 306 | } |
| 307 | |
| 308 | send(command: Command): void { |
| 309 | if (!this.held || this.robot?.readyState !== 'open' || this.pending) return; |
| 310 | if (this.robot.bufferedAmount > 4096) { |
| 311 | this.update({ |
| 312 | ack: { |
| 313 | state: 'error', |
| 314 | text: 'Connection busy. Try again shortly.', |
| 315 | at: performance.now(), |
| 316 | }, |
| 317 | }); |
| 318 | return; |
| 319 | } |
| 320 | const id = commandId(), |
| 321 | at = performance.now(); |
| 322 | const timer = setTimeout(() => { |
| 323 | this.pending = undefined; |
| 324 | this.update({ |
| 325 | ack: { |
| 326 | state: 'error', |
| 327 | text: 'No board confirmation. Try again.', |
| 328 | at: performance.now(), |
| 329 | }, |
| 330 | }); |
| 331 | }, 5000); |
| 332 | this.pending = { id, at, timer }; |
| 333 | this.update({ ack: { state: 'pending', text: 'Waiting for board confirmation…', at } }); |
| 334 | try { |
| 335 | this.robot.send(JSON.stringify({ ...command, command_id: id })); |
| 336 | } catch { |
| 337 | clearTimeout(timer); |
| 338 | this.pending = undefined; |
| 339 | this.update({ |
| 340 | ack: { state: 'error', text: 'Command not sent. Reconnect and try again.', at }, |
| 341 | }); |
| 342 | } |
| 343 | } |
| 344 | |
| 345 | private acceptNowPlaying(next: NowPlaying) { |
| 346 | const previous = this.state.nowPlaying; |
| 347 | if (acceptMetadata(previous, next) === previous) return; |
| 348 | this.signal.frame = undefined; |
| 349 | this.signal.at = 0; |
| 350 | this.update({ |
| 351 | nowPlaying: next, |
| 352 | }); |
| 353 | } |
| 354 | |
| 355 | private receive(channel: 'robot' | 'spectrum', data: unknown) { |
| 356 | const now = performance.now(); |
| 357 | if (channel === 'spectrum') { |
| 358 | const frame = parseSpectrum(data, this.signal.frame?.pts); |
| 359 | if (frame && matchesPlayback(frame.revision, this.state.nowPlaying)) { |
| 360 | this.signal.frame = frame; |
| 361 | this.signal.at = now; |
| 362 | ++this.signal.count; |
| 363 | } |
| 364 | return; |
| 365 | } |
| 366 | const message = parseRobot(data); |
| 367 | if (!message) return; |
| 368 | if (message.event === 'nowPlaying') { |
| 369 | this.acceptNowPlaying(message.nowPlaying); |
| 370 | } else if (message.event === 'telemetry') { |
| 371 | if (!matchesPlayback(message.playbackRevision, this.state.nowPlaying)) return; |
| 372 | this.update({ |
| 373 | telemetry: message, |
| 374 | paused: message.paused, |
| 375 | led: message.led, |
| 376 | telemetryAt: now, |
| 377 | now, |
| 378 | received: this.state.received + 1, |
| 379 | }); |
| 380 | } else if (this.pending?.id === message.command_id) { |
| 381 | const rtt = Math.round(now - this.pending.at); |
| 382 | clearTimeout(this.pending.timer); |
| 383 | this.pending = undefined; |
| 384 | this.update({ |
| 385 | led: message.led, |
| 386 | paused: message.paused, |
| 387 | ack: { |
| 388 | state: message.result === 0 ? 'confirmed' : 'error', |
| 389 | text: |
| 390 | message.result === 0 |
| 391 | ? `Board confirmed · ${rtt} ms round trip` |
| 392 | : 'The board could not apply the command.', |
| 393 | at: now, |
| 394 | rtt, |
| 395 | }, |
| 396 | }); |
| 397 | } |
| 398 | } |
| 399 | |
| 400 | async toggleAudio(): Promise<void> { |
| 401 | if (!this.audio || this.state.phase !== 'connected') return; |
| 402 | try { |
| 403 | if (this.state.audioBlocked || this.audio.paused) { |
| 404 | await this.context?.resume(); |
| 405 | await this.audio.play(); |
| 406 | this.audio.muted = false; |
| 407 | } else this.audio.muted = !this.audio.muted; |
| 408 | this.update({ muted: this.audio.muted, audioBlocked: false }); |
| 409 | } catch { |
| 410 | this.update({ audioBlocked: true }); |
| 411 | } |
| 412 | } |
| 413 | setVolume(volume: number) { |
| 414 | const clamped = Math.max(0, Math.min(1, volume)); |
| 415 | if (this.audio) this.audio.volume = clamped; |
| 416 | this.update({ volume: clamped }); |
| 417 | } |
| 418 | private async collectStats() { |
| 419 | const peer = this.peer; |
| 420 | if (!peer || this.state.phase !== 'connected' || this.statsBusy) return; |
| 421 | this.statsBusy = true; |
| 422 | try { |
| 423 | const report = await peer.getStats(); |
| 424 | if (this.peer !== peer) return; |
| 425 | const now = performance.now(); |
| 426 | report.forEach((stat) => { |
| 427 | if (stat.type !== 'inbound-rtp' || stat.kind !== 'audio') return; |
| 428 | const previous = this.state.audio, |
| 429 | bytes = Number(stat.bytesReceived ?? 0), |
| 430 | packets = Number(stat.packetsReceived ?? 0); |
| 431 | this.update({ |
| 432 | audio: { |
| 433 | bytes, |
| 434 | packets, |
| 435 | lost: Number(stat.packetsLost ?? 0), |
| 436 | kbps: previous.at |
| 437 | ? Math.max(0, ((bytes - previous.bytes) * 8) / (now - previous.at)) |
| 438 | : 0, |
| 439 | at: packets > previous.packets ? now : previous.at, |
| 440 | }, |
| 441 | }); |
| 442 | }); |
| 443 | } catch { |
| 444 | /* Browser statistics are optional; media and control keep working. */ |
| 445 | } finally { |
| 446 | this.statsBusy = false; |
| 447 | } |
| 448 | } |
| 449 | } |