File
Blob: worker/tests/helpers/sfu-fixture.ts
| 1 | import { z } from 'zod'; |
| 2 | import { |
| 3 | answerSchema, |
| 4 | descriptionSchema, |
| 5 | joinedSchema, |
| 6 | statusSchema, |
| 7 | } from '@/shared/contracts/signaling.ts'; |
| 8 | import type { DeviceStart } from '@/shared/contracts/signaling.ts'; |
| 9 | |
| 10 | export const offer = { |
| 11 | type: 'offer' as const, |
| 12 | sdp: 'v=0\r\nm=audio 9 UDP/TLS/RTP/SAVPF 111\r\na=mid:0\r\na=sendonly\r\n', |
| 13 | }; |
| 14 | export const answer = { ...offer, type: 'answer' as const }; |
| 15 | export const track = { title: 'Test track', artist: 'Test artist', durationMs: 6000 }; |
| 16 | export const nowPlaying = { revision: 0, trackIndex: 0, trackCount: 3, track }; |
| 17 | export const bindings = { |
| 18 | REALTIME_APP_ID: 'test-app', |
| 19 | REALTIME_APP_TOKEN: 'test-sfu-token', |
| 20 | DEVICE_TOKEN: 'test-device-token', |
| 21 | VIEWER_PASSWORD: 'test-viewer-password', |
| 22 | ROBOT_NAME: 'test-board', |
| 23 | }; |
| 24 | export const device = { Authorization: `Bearer ${bindings.DEVICE_TOKEN}` }; |
| 25 | export const startedSchema = z.object({ generation: z.uuid(), sessionDescription: answerSchema }); |
| 26 | const inputSchema = z |
| 27 | .object({ |
| 28 | sessionDescription: descriptionSchema.optional(), |
| 29 | tracks: z |
| 30 | .array( |
| 31 | z |
| 32 | .object({ |
| 33 | location: z.enum(['local', 'remote']).optional(), |
| 34 | mid: z.string().optional(), |
| 35 | trackName: z.string().optional(), |
| 36 | sessionId: z.string().optional(), |
| 37 | }) |
| 38 | .passthrough(), |
| 39 | ) |
| 40 | .default([]), |
| 41 | dataChannels: z |
| 42 | .array( |
| 43 | z |
| 44 | .object({ |
| 45 | id: z.number().optional(), |
| 46 | location: z.enum(['local', 'remote']).optional(), |
| 47 | sessionId: z.string().optional(), |
| 48 | dataChannelName: z.string().optional(), |
| 49 | canReply: z.boolean().optional(), |
| 50 | }) |
| 51 | .passthrough(), |
| 52 | ) |
| 53 | .default([]), |
| 54 | }) |
| 55 | .passthrough(); |
| 56 | export type SfuCall = { path: string; input: z.infer<typeof inputSchema> }; |
| 57 | export type Responder = (call: SfuCall) => Response | undefined | Promise<Response | undefined>; |
| 58 | |
| 59 | export function createSfu(respond: Responder = () => undefined) { |
| 60 | const calls: SfuCall[] = []; |
| 61 | let allocated = 0; |
| 62 | return { |
| 63 | calls, |
| 64 | allocations: () => allocated, |
| 65 | async fetch(request: Request): Promise<Response> { |
| 66 | const url = new URL(request.url); |
| 67 | if ( |
| 68 | url.hostname !== 'rtc.live.cloudflare.com' || |
| 69 | !url.pathname.startsWith('/v1/apps/test-app/') || |
| 70 | request.headers.get('Authorization') !== `Bearer ${bindings.REALTIME_APP_TOKEN}` |
| 71 | ) { |
| 72 | throw new Error('Unexpected outbound target or test credentials'); |
| 73 | } |
| 74 | const text = await request.text(); |
| 75 | const input = inputSchema.parse(text ? JSON.parse(text) : {}); |
| 76 | const path = url.pathname, |
| 77 | call = { path, input }; |
| 78 | calls.push(call); |
| 79 | const response = await respond(call); |
| 80 | if (response) return response; |
| 81 | if (path.endsWith('/sessions/new')) |
| 82 | return Response.json({ sessionId: `session-${++allocated}` }); |
| 83 | if (path.endsWith('/tracks/new')) |
| 84 | return Response.json({ |
| 85 | sessionDescription: input.tracks[0].location === 'local' ? answer : offer, |
| 86 | tracks: [{ mid: '0', trackName: 'music' }], |
| 87 | }); |
| 88 | if (path.endsWith('/datachannels/establish')) |
| 89 | return Response.json({ sessionDescription: offer, requiresImmediateRenegotiation: true }); |
| 90 | if (path.endsWith('/datachannels/new')) |
| 91 | return Response.json({ |
| 92 | dataChannels: input.dataChannels.map((c, i) => ({ |
| 93 | dataChannelName: c.dataChannelName, |
| 94 | id: c.location === 'local' ? 2 + i * 2 : 1 + i * 2, |
| 95 | })), |
| 96 | }); |
| 97 | if (path.endsWith('/datachannels/update')) |
| 98 | return Response.json({ dataChannels: input.dataChannels.map((c) => ({ ...c, id: 1 })) }); |
| 99 | if (path.endsWith('/renegotiate') || path.endsWith('/close')) return Response.json({}); |
| 100 | throw new Error(`Unexpected test SFU operation: ${path}`); |
| 101 | }, |
| 102 | }; |
| 103 | } |
| 104 | |
| 105 | export function createRadioClient(dispatch: (request: Request) => Promise<Response>) { |
| 106 | let cookie: string | undefined; |
| 107 | const call = async (path: string, input?: unknown, headers: Record<string, string> = {}) => { |
| 108 | const response = await dispatch( |
| 109 | new Request(`https://radio.example/api${path}`, { |
| 110 | method: input === undefined ? 'GET' : 'POST', |
| 111 | headers: { |
| 112 | 'Content-Type': 'application/json', |
| 113 | ...(cookie ? { Cookie: cookie } : {}), |
| 114 | ...headers, |
| 115 | }, |
| 116 | ...(input === undefined |
| 117 | ? {} |
| 118 | : { body: typeof input === 'string' ? input : JSON.stringify(input) }), |
| 119 | }), |
| 120 | ); |
| 121 | // Drain runtime-backed bodies even when a case only inspects its status. |
| 122 | return new Response(await response.arrayBuffer(), { |
| 123 | status: response.status, |
| 124 | headers: response.headers, |
| 125 | }); |
| 126 | }; |
| 127 | async function successful(path: string, input: unknown, headers: Record<string, string> = {}) { |
| 128 | const response = await call(path, input, headers); |
| 129 | if (response.status !== 200) |
| 130 | throw new Error(`${path} returned ${response.status}: ${await response.text()}`); |
| 131 | return response; |
| 132 | } |
| 133 | return { |
| 134 | call, |
| 135 | status: async () => statusSchema.parse(await (await call('/status')).json()), |
| 136 | async login() { |
| 137 | const response = await successful('/login', { password: bindings.VIEWER_PASSWORD }); |
| 138 | const header = response.headers.get('Set-Cookie'); |
| 139 | if (!header) throw new Error('Missing test login cookie'); |
| 140 | cookie = header.split(';')[0]; |
| 141 | return header; |
| 142 | }, |
| 143 | async start(fields: Partial<DeviceStart> = { track }) { |
| 144 | const input: DeviceStart = { sessionDescription: offer, bootId: 'a'.repeat(32), ...fields }; |
| 145 | const publisher = startedSchema.parse( |
| 146 | await (await successful('/device/start', input, device)).json(), |
| 147 | ); |
| 148 | const identity = { generation: publisher.generation }; |
| 149 | await successful('/device/channels', identity, device); |
| 150 | await successful('/device/ready', identity, device); |
| 151 | return { input, publisher, identity }; |
| 152 | }, |
| 153 | async viewer() { |
| 154 | const viewer = joinedSchema.parse(await (await successful('/viewers', {})).json()); |
| 155 | const owner = { 'X-Viewer-Token': viewer.viewerToken }; |
| 156 | await successful(`/viewers/${viewer.id}/answer`, { sessionDescription: answer }, owner); |
| 157 | return { ...viewer, owner }; |
| 158 | }, |
| 159 | }; |
| 160 | } |