Skip to content
File

Blob: worker/src/server/sfu.ts

typescript148 lines
1import { z } from 'zod';
2import { descriptionSchema } from '@/shared/contracts/signaling.ts';
3import { demand, OperationError } from './rpc.ts';
4 
5const resultSchema = z.object({
6 sessionId: z.string().optional(),
7 sessionDescription: descriptionSchema.optional(),
8 requiresImmediateRenegotiation: z.boolean().optional(),
9 errorCode: z.string().optional(),
10 dataChannels: z
11 .array(
12 z.object({
13 id: z.number().int().optional(),
14 dataChannelName: z.string().optional(),
15 ordered: z.boolean().optional(),
16 maxRetransmits: z.number().int().optional(),
17 errorCode: z.string().optional(),
18 }),
19 )
20 .optional(),
21 tracks: z
22 .array(
23 z.object({
24 mid: z.string().optional(),
25 trackName: z.string().optional(),
26 errorCode: z.string().optional(),
27 }),
28 )
29 .optional(),
30});
31export type SfuResult = z.infer<typeof resultSchema>;
32 
33// A malformed sibling or SDP must not hide a valid allocation receipt.
34const receiptsSchema = z
35 .object({
36 dataChannels: z
37 .array(z.object({ id: z.number().int().min(0).max(65534).optional() }).catch({}))
38 .catch([]),
39 tracks: z.array(z.object({ mid: z.string().min(1).max(128).optional() }).catch({})).catch([]),
40 })
41 .catch({ dataChannels: [], tracks: [] });
42export type AllocationReceipt = { channelIds: number[]; mids: string[] };
43 
44export class SfuClient {
45 private env: Pick<Env, 'REALTIME_APP_ID' | 'REALTIME_APP_TOKEN'>;
46 constructor(env: Pick<Env, 'REALTIME_APP_ID' | 'REALTIME_APP_TOKEN'>) {
47 this.env = env;
48 }
49 async call(
50 path: string,
51 input?: unknown,
52 method = 'POST',
53 goneIsClosed = false,
54 ): Promise<SfuResult> {
55 return this.request(path, input, method, goneIsClosed);
56 }
57 async allocate(
58 path: string,
59 input: unknown,
60 retain: (receipt: AllocationReceipt) => Promise<void>,
61 ): Promise<SfuResult> {
62 return this.request(path, input, 'POST', false, retain);
63 }
64 async closeTracks(sessionId: string, mids: string[]): Promise<void> {
65 await this.call(
66 `/sessions/${sessionId}/tracks/close`,
67 { tracks: mids.map((mid) => ({ mid })), force: true },
68 'PUT',
69 );
70 }
71 private async request(
72 path: string,
73 input: unknown,
74 method: string,
75 goneIsClosed: boolean,
76 retain?: (receipt: AllocationReceipt) => Promise<void>,
77 ): Promise<SfuResult> {
78 let response: Response;
79 try {
80 response = await fetch(
81 `https://rtc.live.cloudflare.com/v1/apps/${encodeURIComponent(this.env.REALTIME_APP_ID)}${path}`,
82 {
83 method,
84 headers: {
85 Authorization: `Bearer ${this.env.REALTIME_APP_TOKEN}`,
86 'Content-Type': 'application/json',
87 },
88 ...(input === undefined ? {} : { body: JSON.stringify(input) }),
89 signal: AbortSignal.timeout(15000),
90 },
91 );
92 } catch {
93 throw new OperationError(502, 'Could not reach the SFU. Please reconnect.');
94 }
95 const closing = path.endsWith('/close');
96 if ((closing || goneIsClosed) && [404, 410].includes(response.status)) return {};
97 let payload: unknown;
98 try {
99 payload = await response.json();
100 } catch {
101 throw new OperationError(502, 'The SFU returned an unreadable response.');
102 }
103 if (retain) {
104 const receipts = receiptsSchema.parse(payload);
105 await retain({
106 channelIds: [
107 ...new Set(receipts.dataChannels.flatMap(({ id }) => (id === undefined ? [] : [id]))),
108 ],
109 mids: [...new Set(receipts.tracks.flatMap(({ mid }) => (mid === undefined ? [] : [mid])))],
110 });
111 }
112 const parsed = resultSchema.safeParse(payload);
113 demand(parsed.success, 502, 'The SFU returned an invalid response.');
114 const value = parsed.data;
115 // Explicit absence satisfies cleanup whether reported for the request or an item.
116 const failed = (item: { errorCode?: string }) =>
117 item.errorCode && !(closing && item.errorCode === 'close_track_error');
118 if (
119 !response.ok ||
120 failed(value) ||
121 value.tracks?.some(failed) ||
122 value.dataChannels?.some(failed)
123 ) {
124 console.warn(
125 JSON.stringify({
126 phase: 'sfu',
127 operation: path.split('/').slice(-2).join('/'),
128 status: response.status,
129 errorCode: value.errorCode,
130 trackErrorCodes: value.tracks?.flatMap((item) =>
131 item.errorCode ? [item.errorCode] : [],
132 ),
133 dataChannelErrorCodes: value.dataChannels?.flatMap((item) =>
134 item.errorCode ? [item.errorCode] : [],
135 ),
136 }),
137 );
138 }
139 demand(response.ok && !failed(value), 502, `SFU operation failed (${response.status}).`);
140 demand(
141 !value.dataChannels?.some(failed) && !value.tracks?.some(failed),
142 502,
143 'The SFU could not configure a channel or track.',
144 );
145 return value;
146 }
147}