@samitouri / QOS-React-1 / commits / 65ec57df37

[Fizz] Add Web Streams to Fizz Node entry point (#33475)

New take on #33441. This uses a wrapper instead of a separate bundle.

Sebastian Markbåge committed Jun 6, 2025 at 20:16 UTC 65ec57df3781d2c62456bb136c7f160f7e834492
12 files changed +503 -9
packages/react-dom/npm/server.node.js
+4
@@ -13,6 +13,10 @@ exports.version = l.version;
13 exports.renderToString = l.renderToString;
14 exports.renderToStaticMarkup = l.renderToStaticMarkup;
15 exports.renderToPipeableStream = s.renderToPipeableStream;
16 +exports.renderToReadableStream = s.renderToReadableStream;
17 if (s.resumeToPipeableStream) {
18 exports.resumeToPipeableStream = s.resumeToPipeableStream;
19 }
20 +if (s.resume) {
21 + exports.resume = s.resume;
22 +}
packages/react-dom/npm/static.node.js
+2
@@ -9,4 +9,6 @@ if (process.env.NODE_ENV === 'production') {
9
10 exports.version = s.version;
11 exports.prerenderToNodeStream = s.prerenderToNodeStream;
12 +exports.prerender = s.prerender;
13 exports.resumeAndPrerenderToNodeStream = s.resumeAndPrerenderToNodeStream;
14 +exports.resumeAndPrerender = s.resumeAndPrerender;
packages/react-dom/server.node.js
+14
@@ -37,3 +37,17 @@ export function resumeToPipeableStream() {
37 arguments,
38 );
39 }
40 +
41 +export function renderToReadableStream() {
42 + return require('./src/server/react-dom-server.node').renderToReadableStream.apply(
43 + this,
44 + arguments,
45 + );
46 +}
47 +
48 +export function resume() {
49 + return require('./src/server/react-dom-server.node').resume.apply(
50 + this,
51 + arguments,
52 + );
53 +}
packages/react-dom/src/__tests__/ReactDOMFizzServerNode-test.js
+20
@@ -56,6 +56,18 @@ describe('ReactDOMFizzServerNode', () => {
56 throw theInfinitePromise;
57 }
58
59 + async function readContentWeb(stream) {
60 + const reader = stream.getReader();
61 + let content = '';
62 + while (true) {
63 + const {done, value} = await reader.read();
64 + if (done) {
65 + return content;
66 + }
67 + content += Buffer.from(value).toString('utf8');
68 + }
69 + }
70 +
71 it('should call renderToPipeableStream', async () => {
72 const {writable, output} = getTestWritable();
73 await act(() => {
@@ -67,6 +79,14 @@ describe('ReactDOMFizzServerNode', () => {
79 expect(output.result).toMatchInlineSnapshot(`"<div>hello world</div>"`);
80 });
81
82 + it('should support web streams', async () => {
83 + const stream = await act(() =>
84 + ReactDOMFizzServer.renderToReadableStream(<div>hello world</div>),
85 + );
86 + const result = await readContentWeb(stream);
87 + expect(result).toMatchInlineSnapshot(`"<div>hello world</div>"`);
88 + });
89 +
90 it('flush fully if piping in on onShellReady', async () => {
91 const {writable, output} = getTestWritable();
92 await act(() => {
packages/react-dom/src/__tests__/ReactDOMFizzStaticNode-test.js
+19
@@ -46,6 +46,18 @@ describe('ReactDOMFizzStaticNode', () => {
46 });
47 }
48
49 + async function readContentWeb(stream) {
50 + const reader = stream.getReader();
51 + let content = '';
52 + while (true) {
53 + const {done, value} = await reader.read();
54 + if (done) {
55 + return content;
56 + }
57 + content += Buffer.from(value).toString('utf8');
58 + }
59 + }
60 +
61 // @gate experimental
62 it('should call prerenderToNodeStream', async () => {
63 const result = await ReactDOMFizzStatic.prerenderToNodeStream(
@@ -55,6 +67,13 @@ describe('ReactDOMFizzStaticNode', () => {
67 expect(prelude).toMatchInlineSnapshot(`"<div>hello world</div>"`);
68 });
69
70 + // @gate experimental
71 + it('should suppport web streams', async () => {
72 + const result = await ReactDOMFizzStatic.prerender(<div>hello world</div>);
73 + const prelude = await readContentWeb(result.prelude);
74 + expect(prelude).toMatchInlineSnapshot(`"<div>hello world</div>"`);
75 + });
76 +
77 // @gate experimental
78 it('should emit DOCTYPE at the root of the document', async () => {
79 const result = await ReactDOMFizzStatic.prerenderToNodeStream(
packages/react-dom/src/server/ReactDOMFizzServerNode.js
+218
@@ -41,6 +41,8 @@ import {
41 createRootFormatContext,
42 } from 'react-dom-bindings/src/server/ReactFizzConfigDOM';
43
44 +import {textEncoder} from 'react-server/src/ReactServerStreamConfigNode';
45 +
46 import {ensureCorrectIsomorphicReactVersion} from '../shared/ensureCorrectIsomorphicReactVersion';
47 ensureCorrectIsomorphicReactVersion();
48
@@ -167,6 +169,141 @@ function renderToPipeableStream(
169 };
170 }
171
172 +function createFakeWritableFromReadableStreamController(
173 + controller: ReadableStreamController,
174 +): Writable {
175 + // The current host config expects a Writable so we create
176 + // a fake writable for now to push into the Readable.
177 + return ({
178 + write(chunk: string | Uint8Array) {
179 + if (typeof chunk === 'string') {
180 + chunk = textEncoder.encode(chunk);
181 + }
182 + controller.enqueue(chunk);
183 + // in web streams there is no backpressure so we can alwas write more
184 + return true;
185 + },
186 + end() {
187 + controller.close();
188 + },
189 + destroy(error) {
190 + // $FlowFixMe[method-unbinding]
191 + if (typeof controller.error === 'function') {
192 + // $FlowFixMe[incompatible-call]: This is an Error object or the destination accepts other types.
193 + controller.error(error);
194 + } else {
195 + controller.close();
196 + }
197 + },
198 + }: any);
199 +}
200 +
201 +// TODO: Move to sub-classing ReadableStream.
202 +type ReactDOMServerReadableStream = ReadableStream & {
203 + allReady: Promise<void>,
204 +};
205 +
206 +type WebStreamsOptions = Omit<
207 + Options,
208 + 'onShellReady' | 'onShellError' | 'onAllReady' | 'onHeaders',
209 +> & {signal: AbortSignal, onHeaders?: (headers: Headers) => void};
210 +
211 +function renderToReadableStream(
212 + children: ReactNodeList,
213 + options?: WebStreamsOptions,
214 +): Promise<ReactDOMServerReadableStream> {
215 + return new Promise((resolve, reject) => {
216 + let onFatalError;
217 + let onAllReady;
218 + const allReady = new Promise<void>((res, rej) => {
219 + onAllReady = res;
220 + onFatalError = rej;
221 + });
222 +
223 + function onShellReady() {
224 + let writable: Writable;
225 + const stream: ReactDOMServerReadableStream = (new ReadableStream(
226 + {
227 + type: 'bytes',
228 + start: (controller): ?Promise<void> => {
229 + writable =
230 + createFakeWritableFromReadableStreamController(controller);
231 + },
232 + pull: (controller): ?Promise<void> => {
233 + startFlowing(request, writable);
234 + },
235 + cancel: (reason): ?Promise<void> => {
236 + stopFlowing(request);
237 + abort(request, reason);
238 + },
239 + },
240 + // $FlowFixMe[prop-missing] size() methods are not allowed on byte streams.
241 + {highWaterMark: 0},
242 + ): any);
243 + // TODO: Move to sub-classing ReadableStream.
244 + stream.allReady = allReady;
245 + resolve(stream);
246 + }
247 + function onShellError(error: mixed) {
248 + // If the shell errors the caller of `renderToReadableStream` won't have access to `allReady`.
249 + // However, `allReady` will be rejected by `onFatalError` as well.
250 + // So we need to catch the duplicate, uncatchable fatal error in `allReady` to prevent a `UnhandledPromiseRejection`.
251 + allReady.catch(() => {});
252 + reject(error);
253 + }
254 +
255 + const onHeaders = options ? options.onHeaders : undefined;
256 + let onHeadersImpl;
257 + if (onHeaders) {
258 + onHeadersImpl = (headersDescriptor: HeadersDescriptor) => {
259 + onHeaders(new Headers(headersDescriptor));
260 + };
261 + }
262 +
263 + const resumableState = createResumableState(
264 + options ? options.identifierPrefix : undefined,
265 + options ? options.unstable_externalRuntimeSrc : undefined,
266 + options ? options.bootstrapScriptContent : undefined,
267 + options ? options.bootstrapScripts : undefined,
268 + options ? options.bootstrapModules : undefined,
269 + );
270 + const request = createRequest(
271 + children,
272 + resumableState,
273 + createRenderState(
274 + resumableState,
275 + options ? options.nonce : undefined,
276 + options ? options.unstable_externalRuntimeSrc : undefined,
277 + options ? options.importMap : undefined,
278 + onHeadersImpl,
279 + options ? options.maxHeadersLength : undefined,
280 + ),
281 + createRootFormatContext(options ? options.namespaceURI : undefined),
282 + options ? options.progressiveChunkSize : undefined,
283 + options ? options.onError : undefined,
284 + onAllReady,
285 + onShellReady,
286 + onShellError,
287 + onFatalError,
288 + options ? options.onPostpone : undefined,
289 + options ? options.formState : undefined,
290 + );
291 + if (options && options.signal) {
292 + const signal = options.signal;
293 + if (signal.aborted) {
294 + abort(request, (signal: any).reason);
295 + } else {
296 + const listener = () => {
297 + abort(request, (signal: any).reason);
298 + signal.removeEventListener('abort', listener);
299 + };
300 + signal.addEventListener('abort', listener);
301 + }
302 + }
303 + startWork(request);
304 + });
305 +}
306 +
307 function resumeRequestImpl(
308 children: ReactNodeList,
309 postponedState: PostponedState,
@@ -225,8 +362,89 @@ function resumeToPipeableStream(
362 };
363 }
364
365 +type WebStreamsResumeOptions = Omit<
366 + Options,
367 + 'onShellReady' | 'onShellError' | 'onAllReady',
368 +> & {signal: AbortSignal};
369 +
370 +function resume(
371 + children: ReactNodeList,
372 + postponedState: PostponedState,
373 + options?: WebStreamsResumeOptions,
374 +): Promise<ReactDOMServerReadableStream> {
375 + return new Promise((resolve, reject) => {
376 + let onFatalError;
377 + let onAllReady;
378 + const allReady = new Promise<void>((res, rej) => {
379 + onAllReady = res;
380 + onFatalError = rej;
381 + });
382 +
383 + function onShellReady() {
384 + let writable: Writable;
385 + const stream: ReactDOMServerReadableStream = (new ReadableStream(
386 + {
387 + type: 'bytes',
388 + start: (controller): ?Promise<void> => {
389 + writable =
390 + createFakeWritableFromReadableStreamController(controller);
391 + },
392 + pull: (controller): ?Promise<void> => {
393 + startFlowing(request, writable);
394 + },
395 + cancel: (reason): ?Promise<void> => {
396 + stopFlowing(request);
397 + abort(request, reason);
398 + },
399 + },
400 + // $FlowFixMe[prop-missing] size() methods are not allowed on byte streams.
401 + {highWaterMark: 0},
402 + ): any);
403 + // TODO: Move to sub-classing ReadableStream.
404 + stream.allReady = allReady;
405 + resolve(stream);
406 + }
407 + function onShellError(error: mixed) {
408 + // If the shell errors the caller of `renderToReadableStream` won't have access to `allReady`.
409 + // However, `allReady` will be rejected by `onFatalError` as well.
410 + // So we need to catch the duplicate, uncatchable fatal error in `allReady` to prevent a `UnhandledPromiseRejection`.
411 + allReady.catch(() => {});
412 + reject(error);
413 + }
414 + const request = resumeRequest(
415 + children,
416 + postponedState,
417 + resumeRenderState(
418 + postponedState.resumableState,
419 + options ? options.nonce : undefined,
420 + ),
421 + options ? options.onError : undefined,
422 + onAllReady,
423 + onShellReady,
424 + onShellError,
425 + onFatalError,
426 + options ? options.onPostpone : undefined,
427 + );
428 + if (options && options.signal) {
429 + const signal = options.signal;
430 + if (signal.aborted) {
431 + abort(request, (signal: any).reason);
432 + } else {
433 + const listener = () => {
434 + abort(request, (signal: any).reason);
435 + signal.removeEventListener('abort', listener);
436 + };
437 + signal.addEventListener('abort', listener);
438 + }
439 + }
440 + startWork(request);
441 + });
442 +}
443 +
444 export {
445 renderToPipeableStream,
446 + renderToReadableStream,
447 resumeToPipeableStream,
448 + resume,
449 ReactVersion as version,
450 };
packages/react-dom/src/server/ReactDOMFizzStaticEdge.js
+2 -3
@@ -159,9 +159,8 @@ function prerender(
159 type ResumeOptions = {
160 nonce?: NonceOption,
161 signal?: AbortSignal,
162 - onError?: (error: mixed) => ?string,
163 - onPostpone?: (reason: string) => void,
164 - unstable_externalRuntimeSrc?: string | BootstrapScriptDescriptor,
162 + onError?: (error: mixed, errorInfo: ErrorInfo) => ?string,
163 + onPostpone?: (reason: string, postponeInfo: PostponeInfo) => void,
164 };
165
166 function resumeAndPrerender(
packages/react-dom/src/server/ReactDOMFizzStaticNode.js
+201 -3
@@ -28,6 +28,7 @@ import {
28 resumeAndPrerenderRequest,
29 startWork,
30 startFlowing,
31 + stopFlowing,
32 abort,
33 getPostponedState,
34 } from 'react-server/src/ReactFizzServer';
@@ -41,6 +42,8 @@ import {
42
43 import {enablePostpone, enableHalt} from 'shared/ReactFeatureFlags';
44
45 +import {textEncoder} from 'react-server/src/ReactServerStreamConfigNode';
46 +
47 import {ensureCorrectIsomorphicReactVersion} from '../shared/ensureCorrectIsomorphicReactVersion';
48 ensureCorrectIsomorphicReactVersion();
49
@@ -72,7 +75,36 @@ type StaticResult = {
75 prelude: Readable,
76 };
77
75 -function createFakeWritable(readable: any): Writable {
78 +function createFakeWritableFromReadableStreamController(
79 + controller: ReadableStreamController,
80 +): Writable {
81 + // The current host config expects a Writable so we create
82 + // a fake writable for now to push into the Readable.
83 + return ({
84 + write(chunk: string | Uint8Array) {
85 + if (typeof chunk === 'string') {
86 + chunk = textEncoder.encode(chunk);
87 + }
88 + controller.enqueue(chunk);
89 + // in web streams there is no backpressure so we can alwas write more
90 + return true;
91 + },
92 + end() {
93 + controller.close();
94 + },
95 + destroy(error) {
96 + // $FlowFixMe[method-unbinding]
97 + if (typeof controller.error === 'function') {
98 + // $FlowFixMe[incompatible-call]: This is an Error object or the destination accepts other types.
99 + controller.error(error);
100 + } else {
101 + controller.close();
102 + }
103 + },
104 + }: any);
105 +}
106 +
107 +function createFakeWritableFromReadable(readable: any): Writable {
108 // The current host config expects a Writable so we create
109 // a fake writable for now to push into the Readable.
110 return ({
@@ -101,7 +133,7 @@ function prerenderToNodeStream(
133 startFlowing(request, writable);
134 },
135 });
104 - const writable = createFakeWritable(readable);
136 + const writable = createFakeWritableFromReadable(readable);
137
138 const result: StaticResult =
139 enablePostpone || enableHalt
@@ -157,6 +189,101 @@ function prerenderToNodeStream(
189 });
190 }
191
192 +function prerender(
193 + children: ReactNodeList,
194 + options?: Omit<Options, 'onHeaders'> & {
195 + onHeaders?: (headers: Headers) => void,
196 + },
197 +): Promise<{
198 + postponed: null | PostponedState,
199 + prelude: ReadableStream,
200 +}> {
201 + return new Promise((resolve, reject) => {
202 + const onFatalError = reject;
203 +
204 + function onAllReady() {
205 + let writable: Writable;
206 + const stream = new ReadableStream(
207 + {
208 + type: 'bytes',
209 + start: (controller): ?Promise<void> => {
210 + writable =
211 + createFakeWritableFromReadableStreamController(controller);
212 + },
213 + pull: (controller): ?Promise<void> => {
214 + startFlowing(request, writable);
215 + },
216 + cancel: (reason): ?Promise<void> => {
217 + stopFlowing(request);
218 + abort(request, reason);
219 + },
220 + },
221 + // $FlowFixMe[prop-missing] size() methods are not allowed on byte streams.
222 + {highWaterMark: 0},
223 + );
224 +
225 + const result =
226 + enablePostpone || enableHalt
227 + ? {
228 + postponed: getPostponedState(request),
229 + prelude: stream,
230 + }
231 + : ({
232 + prelude: stream,
233 + }: any);
234 + resolve(result);
235 + }
236 +
237 + const onHeaders = options ? options.onHeaders : undefined;
238 + let onHeadersImpl;
239 + if (onHeaders) {
240 + onHeadersImpl = (headersDescriptor: HeadersDescriptor) => {
241 + onHeaders(new Headers(headersDescriptor));
242 + };
243 + }
244 + const resources = createResumableState(
245 + options ? options.identifierPrefix : undefined,
246 + options ? options.unstable_externalRuntimeSrc : undefined,
247 + options ? options.bootstrapScriptContent : undefined,
248 + options ? options.bootstrapScripts : undefined,
249 + options ? options.bootstrapModules : undefined,
250 + );
251 + const request = createPrerenderRequest(
252 + children,
253 + resources,
254 + createRenderState(
255 + resources,
256 + undefined, // nonce is not compatible with prerendered bootstrap scripts
257 + options ? options.unstable_externalRuntimeSrc : undefined,
258 + options ? options.importMap : undefined,
259 + onHeadersImpl,
260 + options ? options.maxHeadersLength : undefined,
261 + ),
262 + createRootFormatContext(options ? options.namespaceURI : undefined),
263 + options ? options.progressiveChunkSize : undefined,
264 + options ? options.onError : undefined,
265 + onAllReady,
266 + undefined,
267 + undefined,
268 + onFatalError,
269 + options ? options.onPostpone : undefined,
270 + );
271 + if (options && options.signal) {
272 + const signal = options.signal;
273 + if (signal.aborted) {
274 + abort(request, (signal: any).reason);
275 + } else {
276 + const listener = () => {
277 + abort(request, (signal: any).reason);
278 + signal.removeEventListener('abort', listener);
279 + };
280 + signal.addEventListener('abort', listener);
281 + }
282 + }
283 + startWork(request);
284 + });
285 +}
286 +
287 type ResumeOptions = {
288 nonce?: NonceOption,
289 signal?: AbortSignal,
@@ -178,7 +305,7 @@ function resumeAndPrerenderToNodeStream(
305 startFlowing(request, writable);
306 },
307 });
181 - const writable = createFakeWritable(readable);
308 + const writable = createFakeWritableFromReadable(readable);
309
310 const result = {
311 postponed: getPostponedState(request),
@@ -216,8 +343,79 @@ function resumeAndPrerenderToNodeStream(
343 });
344 }
345
346 +function resumeAndPrerender(
347 + children: ReactNodeList,
348 + postponedState: PostponedState,
349 + options?: ResumeOptions,
350 +): Promise<{
351 + postponed: null | PostponedState,
352 + prelude: ReadableStream,
353 +}> {
354 + return new Promise((resolve, reject) => {
355 + const onFatalError = reject;
356 +
357 + function onAllReady() {
358 + let writable: Writable;
359 + const stream = new ReadableStream(
360 + {
361 + type: 'bytes',
362 + start: (controller): ?Promise<void> => {
363 + writable =
364 + createFakeWritableFromReadableStreamController(controller);
365 + },
366 + pull: (controller): ?Promise<void> => {
367 + startFlowing(request, writable);
368 + },
369 + cancel: (reason): ?Promise<void> => {
370 + stopFlowing(request);
371 + abort(request, reason);
372 + },
373 + },
374 + // $FlowFixMe[prop-missing] size() methods are not allowed on byte streams.
375 + {highWaterMark: 0},
376 + );
377 +
378 + const result = {
379 + postponed: getPostponedState(request),
380 + prelude: stream,
381 + };
382 + resolve(result);
383 + }
384 +
385 + const request = resumeAndPrerenderRequest(
386 + children,
387 + postponedState,
388 + resumeRenderState(
389 + postponedState.resumableState,
390 + options ? options.nonce : undefined,
391 + ),
392 + options ? options.onError : undefined,
393 + onAllReady,
394 + undefined,
395 + undefined,
396 + onFatalError,
397 + options ? options.onPostpone : undefined,
398 + );
399 + if (options && options.signal) {
400 + const signal = options.signal;
401 + if (signal.aborted) {
402 + abort(request, (signal: any).reason);
403 + } else {
404 + const listener = () => {
405 + abort(request, (signal: any).reason);
406 + signal.removeEventListener('abort', listener);
407 + };
408 + signal.addEventListener('abort', listener);
409 + }
410 + }
411 + startWork(request);
412 + });
413 +}
414 +
415 export {
416 + prerender,
417 prerenderToNodeStream,
418 + resumeAndPrerender,
419 resumeAndPrerenderToNodeStream,
420 ReactVersion as version,
421 };
packages/react-dom/src/server/react-dom-server.node.js
+2
@@ -10,5 +10,7 @@
10 export * from './ReactDOMFizzServerNode.js';
11 export {
12 prerenderToNodeStream,
13 + prerender,
14 resumeAndPrerenderToNodeStream,
15 + resumeAndPrerender,
16 } from './ReactDOMFizzStaticNode.js';
packages/react-dom/src/server/react-dom-server.node.stable.js
+6 -2
@@ -7,5 +7,9 @@
7 * @flow
8 */
9
10 -export {renderToPipeableStream, version} from './ReactDOMFizzServerNode.js';
11 -export {prerenderToNodeStream} from './ReactDOMFizzStaticNode.js';
10 +export {
11 + renderToPipeableStream,
12 + renderToReadableStream,
13 + version,
14 +} from './ReactDOMFizzServerNode.js';
15 +export {prerenderToNodeStream, prerender} from './ReactDOMFizzStaticNode.js';
packages/react-dom/static.node.js
+14
@@ -31,9 +31,23 @@ export function prerenderToNodeStream() {
31 );
32 }
33
34 +export function prerender() {
35 + return require('./src/server/react-dom-server.node').prerender.apply(
36 + this,
37 + arguments,
38 + );
39 +}
40 +
41 export function resumeAndPrerenderToNodeStream() {
42 return require('./src/server/react-dom-server.node').resumeAndPrerenderToNodeStream.apply(
43 this,
44 arguments,
45 );
46 }
47 +
48 +export function resumeAndPrerender() {
49 + return require('./src/server/react-dom-server.node').resumeAndPrerender.apply(
50 + this,
51 + arguments,
52 + );
53 +}
packages/react-server/src/ReactServerStreamConfigNode.js
+1 -1
@@ -189,7 +189,7 @@ export function close(destination: Destination) {
189 destination.end();
190 }
191
192 -const textEncoder = new TextEncoder();
192 +export const textEncoder: TextEncoder = new TextEncoder();
193
194 export function stringToChunk(content: string): Chunk {
195 return content;