File
Blob: worker/src/server/robot-room.ts
| 1 | import { startupMetadata, checkStartupMetadata, heartbeatMetadata } from './publisher-metadata.ts'; |
| 2 | import type { NowPlaying } from '@/shared/contracts/track.ts'; |
| 3 | import { DurableObject } from 'cloudflare:workers'; |
| 4 | import { answerSchema, offerSchema } from '@/shared/contracts/signaling.ts'; |
| 5 | import { CHANNEL_PROFILES, channelListSchema } from '@/shared/contracts/channels.ts'; |
| 6 | import type { |
| 7 | Description, |
| 8 | DeviceStart, |
| 9 | RoomStatus, |
| 10 | ViewerKey, |
| 11 | } from '@/shared/contracts/signaling.ts'; |
| 12 | import { digest } from './crypto.ts'; |
| 13 | import { demand, OperationError } from './rpc.ts'; |
| 14 | import type { RpcResult } from './rpc.ts'; |
| 15 | import type { Publisher, RoomState, SessionResources, Viewer } from './room-state.ts'; |
| 16 | import { SfuClient } from './sfu.ts'; |
| 17 | import type { AllocationReceipt, SfuResult } from './sfu.ts'; |
| 18 | |
| 19 | /** One physical board, its publisher, its viewers, and one renewable controller. */ |
| 20 | export 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 | } |