Skip to content
File

Blob: worker/tests/helpers/sfu-fixture.ts

typescript161 lines
1import { z } from 'zod';
2import {
3 answerSchema,
4 descriptionSchema,
5 joinedSchema,
6 statusSchema,
7} from '@/shared/contracts/signaling.ts';
8import type { DeviceStart } from '@/shared/contracts/signaling.ts';
9 
10export 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};
14export const answer = { ...offer, type: 'answer' as const };
15export const track = { title: 'Test track', artist: 'Test artist', durationMs: 6000 };
16export const nowPlaying = { revision: 0, trackIndex: 0, trackCount: 3, track };
17export 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};
24export const device = { Authorization: `Bearer ${bindings.DEVICE_TOKEN}` };
25export const startedSchema = z.object({ generation: z.uuid(), sessionDescription: answerSchema });
26const 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();
56export type SfuCall = { path: string; input: z.infer<typeof inputSchema> };
57export type Responder = (call: SfuCall) => Response | undefined | Promise<Response | undefined>;
58 
59export 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 
105export 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}