Skip to content
File

Blob: worker/src/server/robot-room.ts

typescript470 lines
1import { startupMetadata, checkStartupMetadata, heartbeatMetadata } from './publisher-metadata.ts';
2import type { NowPlaying } from '@/shared/contracts/track.ts';
3import { DurableObject } from 'cloudflare:workers';
4import { answerSchema, offerSchema } from '@/shared/contracts/signaling.ts';
5import { CHANNEL_PROFILES, channelListSchema } from '@/shared/contracts/channels.ts';
6import type {
7 Description,
8 DeviceStart,
9 RoomStatus,
10 ViewerKey,
11} from '@/shared/contracts/signaling.ts';
12import { digest } from './crypto.ts';
13import { demand, OperationError } from './rpc.ts';
14import type { RpcResult } from './rpc.ts';
15import type { Publisher, RoomState, SessionResources, Viewer } from './room-state.ts';
16import { SfuClient } from './sfu.ts';
17import type { AllocationReceipt, SfuResult } from './sfu.ts';
18 
19/** One physical board, its publisher, its viewers, and one renewable controller. */
20export class RobotRoom extends DurableObject<Env> {
21 private state: RoomState = { viewers: {} };
22 private tail: Promise<void> = Promise.resolve();
23 private sfu: SfuClient;
24 
25 constructor(ctx: DurableObjectState, env: Env) {
26 super(ctx, env);
27 this.sfu = new SfuClient(env);
28 ctx.blockConcurrencyWhile(async () => {
29 this.state = (await ctx.storage.get<RoomState>('room')) ?? { viewers: {} };
30 });
31 }
32 private serialize<T>(fn: () => Promise<T>): Promise<T> {
33 const result = this.tail.then(fn);
34 this.tail = result.then(
35 () => undefined,
36 () => undefined,
37 );
38 return result;
39 }
40 private run<T>(fn: () => Promise<T>): Promise<RpcResult<T>> {
41 return this.serialize(async () => {
42 try {
43 return { ok: true, value: await fn() };
44 } catch (error) {
45 if (!(error instanceof OperationError)) throw error;
46 return { ok: false, status: error.status, error: error.message };
47 } finally {
48 await this.save();
49 }
50 });
51 }
52 private async save(): Promise<void> {
53 await this.ctx.storage.put('room', this.state);
54 if (this.state.publisher || Object.keys(this.state.viewers).length) {
55 const next = Date.now() + 10000,
56 scheduled = await this.ctx.storage.getAlarm();
57 // Polling must not push an already scheduled cleanup farther into the future.
58 if (!scheduled || scheduled > next) await this.ctx.storage.setAlarm(next);
59 } else await this.ctx.storage.deleteAlarm();
60 }
61 private publisher(): Publisher {
62 const p = this.state.publisher;
63 demand(
64 p?.ready && Date.now() - p.seen < 25000,
65 409,
66 'The board is offline. Try again when it is online.',
67 );
68 return p;
69 }
70 private viewer(key: ViewerKey): Viewer {
71 const v = this.state.viewers[key.id];
72 demand(v && key.token === v.token, 403, 'This browser does not own that viewer session.');
73 demand(
74 !v.closing && v.generation === this.state.publisher?.generation,
75 409,
76 'The S3 session changed. Reconnect to the radio.',
77 );
78 v.seen = Date.now();
79 return v;
80 }
81 private device(generation: string): Publisher {
82 const p = this.state.publisher;
83 demand(p && generation === p.generation, 409, 'Device generation is no longer current.');
84 p.seen = Date.now();
85 return p;
86 }
87 private channels(result: SfuResult) {
88 const parsed = channelListSchema.safeParse(result.dataChannels);
89 demand(parsed.success, 502, 'The SFU did not return both data channels.');
90 return parsed.data;
91 }
92 private async permission(v: Viewer, enabled: boolean): Promise<void> {
93 const p = this.state.publisher;
94 if (!p || !v.channels.length) return;
95 await this.sfu.call(
96 `/sessions/${v.sessionId}/datachannels/update`,
97 {
98 dataChannels: [
99 {
100 location: 'remote',
101 sessionId: p.sessionId,
102 dataChannelName: 'robot',
103 canReply: enabled,
104 },
105 ],
106 },
107 'PUT',
108 !enabled,
109 );
110 }
111 private async retain(session: SessionResources, receipt: AllocationReceipt): Promise<void> {
112 if (receipt.channelIds.length)
113 session.pendingChannels = [
114 ...new Set([...(session.pendingChannels ?? []), ...receipt.channelIds]),
115 ];
116 if (receipt.mids.length)
117 session.pendingMids = [...new Set([...(session.pendingMids ?? []), ...receipt.mids])];
118 if (receipt.channelIds.length || receipt.mids.length) await this.save();
119 }
120 private async closePending(session: SessionResources): Promise<boolean> {
121 let ok = true;
122 if (session.pendingChannels?.length) {
123 try {
124 await this.sfu.call(
125 `/sessions/${session.sessionId}/datachannels/close`,
126 {
127 dataChannels: session.pendingChannels.map((id) => ({ id })),
128 },
129 'PUT',
130 );
131 delete session.pendingChannels;
132 } catch {
133 ok = false;
134 }
135 }
136 if (session.pendingMids?.length) {
137 try {
138 await this.sfu.closeTracks(session.sessionId, session.pendingMids);
139 delete session.pendingMids;
140 } catch {
141 ok = false;
142 }
143 }
144 return ok;
145 }
146 private async openChannels(session: SessionResources, dataChannels: unknown) {
147 if (!session.channels.length) {
148 demand(
149 await this.closePending(session),
150 502,
151 'Previous channel allocations are still closing. Try again.',
152 );
153 const result = await this.sfu.allocate(
154 `/sessions/${session.sessionId}/datachannels/new`,
155 { dataChannels },
156 (receipt) => this.retain(session, receipt),
157 );
158 session.channels = this.channels(result);
159 delete session.pendingChannels;
160 }
161 return { channels: session.channels };
162 }
163 private async close(session: SessionResources): Promise<boolean> {
164 let ok = await this.closePending(session);
165 if (session.channels.length) {
166 try {
167 await this.sfu.call(
168 `/sessions/${session.sessionId}/datachannels/close`,
169 { dataChannels: session.channels.map(({ id }) => ({ id })) },
170 'PUT',
171 );
172 session.channels = [];
173 } catch {
174 ok = false;
175 }
176 }
177 if (session.mid) {
178 try {
179 await this.sfu.closeTracks(session.sessionId, [session.mid]);
180 delete session.mid;
181 } catch {
182 ok = false;
183 }
184 }
185 return ok;
186 }
187 private async expire(): Promise<void> {
188 const p = this.state.publisher;
189 if (!p || Date.now() - p.seen > 90000) {
190 // The entire expired generation is disposable; let its SFU transports time out.
191 this.state = { viewers: {} };
192 return;
193 }
194 const controller = this.state.controller;
195 if (controller && controller.until <= Date.now()) {
196 const v = this.state.viewers[controller.id];
197 if (v) await this.permission(v, false);
198 delete this.state.controller;
199 }
200 for (const [id, v] of Object.entries(this.state.viewers)) {
201 if (v.closing || Date.now() - v.seen > 45000) {
202 v.closing = true;
203 if (await this.close(v)) delete this.state.viewers[id];
204 }
205 }
206 }
207 async alarm(): Promise<void> {
208 await this.serialize(async () => {
209 try {
210 await this.expire();
211 } finally {
212 await this.save();
213 }
214 });
215 }
216 
217 getStatus(): Promise<RpcResult<RoomStatus>> {
218 return this.run(async () => {
219 const p = this.state.publisher,
220 online = !!p?.ready && Date.now() - p.seen < 25000;
221 return {
222 online,
223 viewers: Object.values(this.state.viewers).filter(
224 (v) => !v.closing && Date.now() - v.seen < 45000,
225 ).length,
226 controller:
227 this.state.controller && this.state.controller.until > Date.now()
228 ? this.state.controller.id
229 : null,
230 generation: online ? p!.generation : null,
231 track: online ? (p!.track ?? null) : null,
232 nowPlaying: online ? (p!.nowPlaying ?? null) : null,
233 };
234 });
235 }
236 startDevice(input: DeviceStart) {
237 return this.run(async () => {
238 const offer = input.sessionDescription,
239 bootId = input.bootId;
240 const metadata = startupMetadata(input);
241 const offerHash = bootId ? await digest(offer.sdp) : undefined;
242 const existing = this.state.publisher;
243 if (bootId && existing?.bootId === bootId) {
244 demand(
245 existing.offerHash === offerHash,
246 409,
247 'A boot identifier cannot be reused with another offer.',
248 );
249 checkStartupMetadata(existing, metadata);
250 demand(
251 existing.answer,
252 409,
253 'This startup did not complete. Restart the device to recover.',
254 );
255 existing.seen = Date.now();
256 return { generation: existing.generation, sessionDescription: existing.answer };
257 }
258 const audio = offer.sdp.split(/\r?\nm=/).find((section) => section.startsWith('audio '));
259 const mid = audio?.match(/(?:^|\n)a=mid:([^\r\n]+)/)?.[1];
260 demand(
261 mid && /^[A-Za-z0-9_-]{1,32}$/.test(mid),
262 400,
263 'The S3 offer must contain an audio track.',
264 );
265 // A replacement boot owns a new transport. Retire the old generation before
266 // allocating, even if setup fails; its SFU resources expire independently.
267 this.state = { viewers: {} };
268 await this.save();
269 const session = await this.sfu.call('/sessions/new');
270 demand(session.sessionId, 502, 'The SFU did not create a publisher session.');
271 const p: Publisher = {
272 sessionId: session.sessionId,
273 generation: crypto.randomUUID(),
274 mid,
275 ready: false,
276 seen: Date.now(),
277 channels: [],
278 bootId,
279 offerHash,
280 ...metadata,
281 initialNowPlaying: input.nowPlaying,
282 };
283 this.state.publisher = p;
284 const published = await this.sfu.allocate(
285 `/sessions/${p.sessionId}/tracks/new`,
286 {
287 sessionDescription: offer,
288 tracks: [{ location: 'local', mid, trackName: 'music' }],
289 },
290 (receipt) => this.retain(p, receipt),
291 );
292 const answer = answerSchema.safeParse(published.sessionDescription);
293 demand(answer.success, 502, 'The SFU did not return a publisher answer.');
294 p.answer = answer.data;
295 p.pendingMids = p.pendingMids?.filter((value) => value !== mid);
296 return { generation: p.generation, sessionDescription: p.answer };
297 });
298 }
299 createDeviceChannels(generation: string) {
300 return this.run(async () => {
301 const p = this.device(generation);
302 return this.openChannels(
303 p,
304 CHANNEL_PROFILES.map((profile) => ({ ...profile, location: 'local' })),
305 );
306 });
307 }
308 deviceReady(generation: string) {
309 return this.run(async () => {
310 const p = this.device(generation);
311 demand(p.channels.length === 2, 409, 'Device channels are missing.');
312 p.ready = true;
313 return { ok: true as const };
314 });
315 }
316 deviceHeartbeat(value: string | { generation: string; nowPlaying?: NowPlaying }) {
317 return this.run(async () => {
318 // Accept the earlier RPC shape during a compatible Worker rollout.
319 const input = typeof value === 'string' ? { generation: value } : value;
320 const p = this.device(input.generation);
321 if (input.nowPlaying) {
322 p.nowPlaying = heartbeatMetadata(p.nowPlaying, input.nowPlaying);
323 p.track = p.nowPlaying.track ?? undefined;
324 }
325 return { ok: true as const };
326 });
327 }
328 
329 joinViewer() {
330 return this.run(async () => {
331 const p = this.publisher();
332 await this.expire();
333 demand(
334 Object.keys(this.state.viewers).length < 8,
335 429,
336 'All 8 listener spots are in use. Try again later.',
337 );
338 const session = await this.sfu.call('/sessions/new');
339 demand(session.sessionId, 502, 'The SFU did not create a viewer session.');
340 const id = crypto.randomUUID(),
341 token = crypto.randomUUID();
342 this.state.viewers[id] = {
343 sessionId: session.sessionId,
344 token,
345 generation: p.generation,
346 seen: Date.now(),
347 channels: [],
348 };
349 const transport = await this.sfu.call(
350 `/sessions/${session.sessionId}/datachannels/establish`,
351 { dataChannel: { location: 'remote', dataChannelName: 'server-events' } },
352 );
353 const offer = offerSchema.safeParse(transport.sessionDescription);
354 demand(offer.success, 502, 'The SFU did not return a viewer offer.');
355 return { id, viewerToken: token, generation: p.generation, sessionDescription: offer.data };
356 });
357 }
358 answerViewer(key: ViewerKey, answer: Description) {
359 return this.run(async () => {
360 const v = this.viewer(key),
361 p = this.publisher();
362 await this.sfu.call(
363 `/sessions/${v.sessionId}/renegotiate`,
364 { sessionDescription: answer },
365 'PUT',
366 );
367 return this.openChannels(
368 v,
369 CHANNEL_PROFILES.map((profile) => ({
370 ...profile,
371 location: 'remote',
372 sessionId: p.sessionId,
373 waitForAck: true,
374 canReply: false,
375 })),
376 );
377 });
378 }
379 subscribeAudio(key: ViewerKey) {
380 return this.run(async () => {
381 const v = this.viewer(key),
382 p = this.publisher();
383 demand(!v.mid, 409, 'Audio is already subscribed.');
384 demand(
385 await this.closePending(v),
386 502,
387 'Previous audio allocations are still closing. Try again.',
388 );
389 const result = await this.sfu.allocate(
390 `/sessions/${v.sessionId}/tracks/new`,
391 {
392 tracks: [{ location: 'remote', sessionId: p.sessionId, trackName: 'music' }],
393 },
394 (receipt) => this.retain(v, receipt),
395 );
396 const mid = result.tracks?.[0]?.mid;
397 const offer = offerSchema.safeParse(result.sessionDescription);
398 demand(offer.success && mid, 502, 'The SFU did not return an audio offer.');
399 v.mid = mid;
400 v.pendingMids = v.pendingMids?.filter((value) => value !== mid);
401 return { sessionDescription: offer.data };
402 });
403 }
404 renegotiateViewer(key: ViewerKey, answer: Description) {
405 return this.run(async () => {
406 const v = this.viewer(key);
407 this.publisher();
408 await this.sfu.call(
409 `/sessions/${v.sessionId}/renegotiate`,
410 { sessionDescription: answer },
411 'PUT',
412 );
413 return { ok: true as const };
414 });
415 }
416 heartbeatViewer(key: ViewerKey) {
417 return this.run(async () => {
418 this.viewer(key);
419 this.publisher();
420 if (this.state.controller?.id === key.id && this.state.controller.until > Date.now())
421 this.state.controller.until = Date.now() + 15000;
422 return {
423 controller:
424 this.state.controller && this.state.controller.until > Date.now()
425 ? this.state.controller.id
426 : null,
427 };
428 });
429 }
430 claimControl(key: ViewerKey) {
431 return this.run(async () => {
432 const v = this.viewer(key);
433 this.publisher();
434 await this.expire();
435 demand(
436 !this.state.controller || this.state.controller.id === key.id,
437 409,
438 'Another listener has control. Try again when it is available.',
439 );
440 demand(v.channels.length === 2, 409, 'Wait for the data channels to open.');
441 await this.permission(v, true);
442 this.state.controller = { id: key.id, until: Date.now() + 15000 };
443 return { controller: key.id };
444 });
445 }
446 releaseControl(key: ViewerKey) {
447 return this.run(async () => {
448 const v = this.viewer(key);
449 this.publisher();
450 if (this.state.controller?.id === key.id) {
451 await this.permission(v, false);
452 delete this.state.controller;
453 }
454 return { controller: this.state.controller?.id ?? null };
455 });
456 }
457 leaveViewer(key: ViewerKey) {
458 return this.run(async () => {
459 const v = this.viewer(key);
460 if (this.state.controller?.id === key.id) {
461 await this.permission(v, false);
462 delete this.state.controller;
463 }
464 v.closing = true;
465 if (await this.close(v)) delete this.state.viewers[key.id];
466 return { ok: true as const };
467 });
468 }
469}