main
js 285 lines 7.7 KB
Raw
1 /**
2 * Copyright (c) Meta Platforms, Inc. and affiliates.
3 *
4 * This source code is licensed under the MIT license found in the
5 * LICENSE file in the root directory of this source tree.
6 *
7 * @flow
8 */
9
10 import type {Thenable} from 'shared/ReactTypes.js';
11
12 import type {
13 DebugChannel,
14 DebugChannelCallback,
15 FindSourceMapURLCallback,
16 Response as FlightResponse,
17 } from 'react-client/src/ReactFlightClient';
18
19 import type {ReactServerValue} from 'react-client/src/ReactFlightReplyClient';
20
21 import {
22 createResponse,
23 createStreamState,
24 getRoot,
25 reportGlobalError,
26 processBinaryChunk,
27 processStringChunk,
28 close,
29 injectIntoDevTools,
30 } from 'react-client/src/ReactFlightClient';
31
32 import {processReply} from 'react-client/src/ReactFlightReplyClient';
33
34 export {
35 createServerReference,
36 registerServerReference,
37 } from 'react-client/src/ReactFlightReplyClient';
38
39 import type {TemporaryReferenceSet} from 'react-client/src/ReactFlightTemporaryReferences';
40
41 export {createTemporaryReferenceSet} from 'react-client/src/ReactFlightTemporaryReferences';
42
43 export type {TemporaryReferenceSet};
44
45 type CallServerCallback = <A, T>(string, args: A) => Promise<T>;
46
47 export type Options = {
48 moduleBaseURL?: string,
49 callServer?: CallServerCallback,
50 debugChannel?: {writable?: WritableStream, readable?: ReadableStream, ...},
51 temporaryReferences?: TemporaryReferenceSet,
52 unstable_allowPartialStream?: boolean,
53 findSourceMapURL?: FindSourceMapURLCallback,
54 replayConsoleLogs?: boolean,
55 environmentName?: string,
56 startTime?: number,
57 endTime?: number,
58 };
59
60 function createDebugCallbackFromWritableStream(
61 debugWritable: WritableStream,
62 ): DebugChannelCallback {
63 const textEncoder = new TextEncoder();
64 const writer = debugWritable.getWriter();
65 return message => {
66 if (message === '') {
67 writer.close();
68 } else {
69 // Note: It's important that this function doesn't close over the Response object or it can't be GC:ed.
70 // Therefore, we can't report errors from this write back to the Response object.
71 if (__DEV__) {
72 writer.write(textEncoder.encode(message + '\n')).catch(console.error);
73 }
74 }
75 };
76 }
77
78 function createResponseFromOptions(options: void | Options) {
79 const debugChannel: void | DebugChannel =
80 __DEV__ && options && options.debugChannel !== undefined
81 ? {
82 hasReadable: options.debugChannel.readable !== undefined,
83 callback:
84 options.debugChannel.writable !== undefined
85 ? createDebugCallbackFromWritableStream(
86 options.debugChannel.writable,
87 )
88 : null,
89 }
90 : undefined;
91
92 return createResponse(
93 options && options.moduleBaseURL ? options.moduleBaseURL : '',
94 null,
95 null,
96 options && options.callServer ? options.callServer : undefined,
97 undefined, // encodeFormAction
98 undefined, // nonce
99 options && options.temporaryReferences
100 ? options.temporaryReferences
101 : undefined,
102 options && options.unstable_allowPartialStream
103 ? options.unstable_allowPartialStream
104 : false,
105 __DEV__ && options && options.findSourceMapURL
106 ? options.findSourceMapURL
107 : undefined,
108 __DEV__ ? (options ? options.replayConsoleLogs !== false : true) : false, // defaults to true
109 __DEV__ && options && options.environmentName
110 ? options.environmentName
111 : undefined,
112 __DEV__ && options && options.startTime != null
113 ? options.startTime
114 : undefined,
115 __DEV__ && options && options.endTime != null ? options.endTime : undefined,
116 debugChannel,
117 );
118 }
119
120 function startReadingFromUniversalStream(
121 response: FlightResponse,
122 stream: ReadableStream,
123 onDone: () => void,
124 ): void {
125 // This is the same as startReadingFromStream except this allows WebSocketStreams which
126 // return ArrayBuffer and string chunks instead of Uint8Array chunks. We could potentially
127 // always allow streams with variable chunk types.
128 const streamState = createStreamState(response, stream);
129 const reader = stream.getReader();
130 function progress({
131 done,
132 value,
133 }: {
134 done: boolean,
135 value: any,
136 ...
137 }): void | Promise<void> {
138 if (done) {
139 return onDone();
140 }
141 if (value instanceof ArrayBuffer) {
142 // WebSockets can produce ArrayBuffer values in ReadableStreams.
143 processBinaryChunk(response, streamState, new Uint8Array(value));
144 } else if (typeof value === 'string') {
145 // WebSockets can produce string values in ReadableStreams.
146 processStringChunk(response, streamState, value);
147 } else {
148 processBinaryChunk(response, streamState, value);
149 }
150 return reader.read().then(progress).catch(error);
151 }
152 function error(e: any) {
153 reportGlobalError(response, e);
154 }
155 reader.read().then(progress).catch(error);
156 }
157
158 function startReadingFromStream(
159 response: FlightResponse,
160 stream: ReadableStream,
161 onDone: () => void,
162 debugValue: mixed,
163 ): void {
164 const streamState = createStreamState(response, debugValue);
165 const reader = stream.getReader();
166 function progress({
167 done,
168 value,
169 }: {
170 done: boolean,
171 value: ?any,
172 ...
173 }): void | Promise<void> {
174 if (done) {
175 return onDone();
176 }
177 const buffer: Uint8Array = value as any;
178 processBinaryChunk(response, streamState, buffer);
179 return reader.read().then(progress).catch(error);
180 }
181 function error(e: any) {
182 reportGlobalError(response, e);
183 }
184 reader.read().then(progress).catch(error);
185 }
186 function createFromReadableStream<T>(
187 stream: ReadableStream,
188 options?: Options,
189 ): Thenable<T> {
190 const response: FlightResponse = createResponseFromOptions(options);
191 if (
192 __DEV__ &&
193 options &&
194 options.debugChannel &&
195 options.debugChannel.readable
196 ) {
197 let streamDoneCount = 0;
198 const handleDone = () => {
199 if (++streamDoneCount === 2) {
200 close(response);
201 }
202 };
203 startReadingFromUniversalStream(
204 response,
205 options.debugChannel.readable,
206 handleDone,
207 );
208 startReadingFromStream(response, stream, handleDone, stream);
209 } else {
210 startReadingFromStream(
211 response,
212 stream,
213 close.bind(null, response),
214 stream,
215 );
216 }
217 return getRoot(response);
218 }
219
220 function createFromFetch<T>(
221 promiseForResponse: Promise<Response>,
222 options?: Options,
223 ): Thenable<T> {
224 const response: FlightResponse = createResponseFromOptions(options);
225 promiseForResponse.then(
226 function (r) {
227 if (
228 __DEV__ &&
229 options &&
230 options.debugChannel &&
231 options.debugChannel.readable
232 ) {
233 let streamDoneCount = 0;
234 const handleDone = () => {
235 if (++streamDoneCount === 2) {
236 close(response);
237 }
238 };
239 startReadingFromUniversalStream(
240 response,
241 options.debugChannel.readable,
242 handleDone,
243 );
244 startReadingFromStream(response, r.body as any, handleDone, r);
245 } else {
246 startReadingFromStream(
247 response,
248 r.body as any,
249 close.bind(null, response),
250 r,
251 );
252 }
253 },
254 function (e) {
255 reportGlobalError(response, e);
256 },
257 );
258 return getRoot(response);
259 }
260
261 function encodeReply(
262 value: ReactServerValue,
263 options?: {temporaryReferences?: TemporaryReferenceSet, signal?: AbortSignal},
264 ): Promise<
265 string | URLSearchParams | FormData,
266 > /* We don't use URLSearchParams yet but maybe */ {
267 return new Promise((resolve, reject) => {
268 processReply(
269 value,
270 '',
271 options && options.temporaryReferences
272 ? options.temporaryReferences
273 : undefined,
274 resolve,
275 reject,
276 options ? options.signal : undefined,
277 );
278 });
279 }
280
281 export {createFromFetch, createFromReadableStream, encodeReply};
282
283 if (__DEV__) {
284 injectIntoDevTools();
285 }