Skip to content
File

Blob: worker/src/client/radio/session.ts

typescript450 lines
1import { acceptMetadata, matchesPlayback } from './playback.ts';
2import type { NowPlaying } from '@/shared/contracts/track.ts';
3import { radioApi, ApiError, errorText } from './api.ts';
4import type { RadioApi } from './api.ts';
5import { commandId, parseRobot, parseSpectrum } from './protocol.ts';
6import type { Command } from '@/shared/contracts/robot.ts';
7import type { Joined } from '@/shared/contracts/signaling.ts';
8import { initialState } from './session-state.ts';
9import type { IssueSource, SessionState, SignalBuffer } from './session-state.ts';
10import { 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 */
15export 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}