[Flight] Add Web Stream support to the Flight Server in Node (#33474)
This needs some tweaks to the implementation and a conversion but simple enough. --------- Co-authored-by: Hendrik Liebau <mail@hendrik-liebau.de>
Sebastian Markbåge committed
Jun 7, 2025 at 10:40 UTC
9666605abfee7e525a22931ce38d40bb29ddc8a5
20 files changed
+696
-27
packages/react-server-dom-parcel/npm/server.node.js
+3
-1
@@ -7,9 +7,11 @@ if (process.env.NODE_ENV === 'production') {
7
s = require('./cjs/react-server-dom-parcel-server.node.development.js');
8
}
9
10
+exports.renderToReadableStream = s.renderToReadableStream;
11
exports.renderToPipeableStream = s.renderToPipeableStream;
11
-exports.decodeReplyFromBusboy = s.decodeReplyFromBusboy;
12
exports.decodeReply = s.decodeReply;
13
+exports.decodeReplyFromBusboy = s.decodeReplyFromBusboy;
14
+exports.decodeReplyFromAsyncIterable = s.decodeReplyFromAsyncIterable;
15
exports.decodeAction = s.decodeAction;
16
exports.decodeFormState = s.decodeFormState;
17
exports.createClientReference = s.createClientReference;
packages/react-server-dom-parcel/npm/static.node.js
+3
@@ -7,6 +7,9 @@ if (process.env.NODE_ENV === 'production') {
7
s = require('./cjs/react-server-dom-parcel-server.node.development.js');
8
}
9
10
+if (s.unstable_prerender) {
11
+ exports.unstable_prerender = s.unstable_prerender;
12
+}
13
if (s.unstable_prerenderToNodeStream) {
14
exports.unstable_prerenderToNodeStream = s.unstable_prerenderToNodeStream;
15
}
packages/react-server-dom-parcel/server.node.js
+3
-1
@@ -9,8 +9,10 @@
9
10
export {
11
renderToPipeableStream,
12
- decodeReplyFromBusboy,
12
+ renderToReadableStream,
13
decodeReply,
14
+ decodeReplyFromBusboy,
15
+ decodeReplyFromAsyncIterable,
16
decodeAction,
17
decodeFormState,
18
createClientReference,
packages/react-server-dom-parcel/src/server/ReactFlightDOMServerNode.js
+197
-3
@@ -21,6 +21,9 @@ import type {
21
} from '../client/ReactFlightClientConfigBundlerParcel';
22
23
import {Readable} from 'stream';
24
+
25
+import {ASYNC_ITERATOR} from 'shared/ReactSymbols';
26
+
27
import {
28
createRequest,
29
createPrerenderRequest,
@@ -35,6 +38,7 @@ import {
38
reportGlobalError,
39
close,
40
resolveField,
41
+ resolveFile,
42
resolveFileInfo,
43
resolveFileChunk,
44
resolveFileComplete,
@@ -56,9 +60,12 @@ export {
60
registerServerReference,
61
} from '../ReactFlightParcelReferences';
62
63
+import {textEncoder} from 'react-server/src/ReactServerStreamConfigNode';
64
+
65
import type {TemporaryReferenceSet} from 'react-server/src/ReactFlightServerTemporaryReferences';
66
67
export {createTemporaryReferenceSet} from 'react-server/src/ReactFlightServerTemporaryReferences';
68
+
69
export type {TemporaryReferenceSet};
70
71
function createDrainHandler(destination: Destination, request: Request) {
@@ -131,11 +138,91 @@ export function renderToPipeableStream(
138
};
139
}
140
134
-function createFakeWritable(readable: any): Writable {
141
+function createFakeWritableFromReadableStreamController(
142
+ controller: ReadableStreamController,
143
+): Writable {
144
// The current host config expects a Writable so we create
145
// a fake writable for now to push into the Readable.
146
return ({
138
- write(chunk) {
147
+ write(chunk: string | Uint8Array) {
148
+ if (typeof chunk === 'string') {
149
+ chunk = textEncoder.encode(chunk);
150
+ }
151
+ controller.enqueue(chunk);
152
+ // in web streams there is no backpressure so we can alwas write more
153
+ return true;
154
+ },
155
+ end() {
156
+ controller.close();
157
+ },
158
+ destroy(error) {
159
+ // $FlowFixMe[method-unbinding]
160
+ if (typeof controller.error === 'function') {
161
+ // $FlowFixMe[incompatible-call]: This is an Error object or the destination accepts other types.
162
+ controller.error(error);
163
+ } else {
164
+ controller.close();
165
+ }
166
+ },
167
+ }: any);
168
+}
169
+
170
+export function renderToReadableStream(
171
+ model: ReactClientValue,
172
+
173
+ options?: Options & {
174
+ signal?: AbortSignal,
175
+ },
176
+): ReadableStream {
177
+ const request = createRequest(
178
+ model,
179
+ null,
180
+ options ? options.onError : undefined,
181
+ options ? options.identifierPrefix : undefined,
182
+ options ? options.onPostpone : undefined,
183
+ options ? options.temporaryReferences : undefined,
184
+ __DEV__ && options ? options.environmentName : undefined,
185
+ __DEV__ && options ? options.filterStackFrame : undefined,
186
+ );
187
+ if (options && options.signal) {
188
+ const signal = options.signal;
189
+ if (signal.aborted) {
190
+ abort(request, (signal: any).reason);
191
+ } else {
192
+ const listener = () => {
193
+ abort(request, (signal: any).reason);
194
+ signal.removeEventListener('abort', listener);
195
+ };
196
+ signal.addEventListener('abort', listener);
197
+ }
198
+ }
199
+ let writable: Writable;
200
+ const stream = new ReadableStream(
201
+ {
202
+ type: 'bytes',
203
+ start: (controller): ?Promise<void> => {
204
+ writable = createFakeWritableFromReadableStreamController(controller);
205
+ startWork(request);
206
+ },
207
+ pull: (controller): ?Promise<void> => {
208
+ startFlowing(request, writable);
209
+ },
210
+ cancel: (reason): ?Promise<void> => {
211
+ stopFlowing(request);
212
+ abort(request, reason);
213
+ },
214
+ },
215
+ // $FlowFixMe[prop-missing] size() methods are not allowed on byte streams.
216
+ {highWaterMark: 0},
217
+ );
218
+ return stream;
219
+}
220
+
221
+function createFakeWritableFromNodeReadable(readable: any): Writable {
222
+ // The current host config expects a Writable so we create
223
+ // a fake writable for now to push into the Readable.
224
+ return ({
225
+ write(chunk: string | Uint8Array) {
226
return readable.push(chunk);
227
},
228
end() {
@@ -173,7 +260,7 @@ export function prerenderToNodeStream(
260
startFlowing(request, writable);
261
},
262
});
176
- const writable = createFakeWritable(readable);
263
+ const writable = createFakeWritableFromNodeReadable(readable);
264
resolve({prelude: readable});
265
}
266
@@ -207,6 +294,69 @@ export function prerenderToNodeStream(
294
});
295
}
296
297
+export function prerender(
298
+ model: ReactClientValue,
299
+
300
+ options?: Options & {
301
+ signal?: AbortSignal,
302
+ },
303
+): Promise<{
304
+ prelude: ReadableStream,
305
+}> {
306
+ return new Promise((resolve, reject) => {
307
+ const onFatalError = reject;
308
+ function onAllReady() {
309
+ let writable: Writable;
310
+ const stream = new ReadableStream(
311
+ {
312
+ type: 'bytes',
313
+ start: (controller): ?Promise<void> => {
314
+ writable =
315
+ createFakeWritableFromReadableStreamController(controller);
316
+ },
317
+ pull: (controller): ?Promise<void> => {
318
+ startFlowing(request, writable);
319
+ },
320
+ cancel: (reason): ?Promise<void> => {
321
+ stopFlowing(request);
322
+ abort(request, reason);
323
+ },
324
+ },
325
+ // $FlowFixMe[prop-missing] size() methods are not allowed on byte streams.
326
+ {highWaterMark: 0},
327
+ );
328
+ resolve({prelude: stream});
329
+ }
330
+ const request = createPrerenderRequest(
331
+ model,
332
+ null,
333
+ onAllReady,
334
+ onFatalError,
335
+ options ? options.onError : undefined,
336
+ options ? options.identifierPrefix : undefined,
337
+ options ? options.onPostpone : undefined,
338
+ options ? options.temporaryReferences : undefined,
339
+ __DEV__ && options ? options.environmentName : undefined,
340
+ __DEV__ && options ? options.filterStackFrame : undefined,
341
+ );
342
+ if (options && options.signal) {
343
+ const signal = options.signal;
344
+ if (signal.aborted) {
345
+ const reason = (signal: any).reason;
346
+ abort(request, reason);
347
+ } else {
348
+ const listener = () => {
349
+ const reason = (signal: any).reason;
350
+ abort(request, reason);
351
+ signal.removeEventListener('abort', listener);
352
+ };
353
+ signal.addEventListener('abort', listener);
354
+ }
355
+ }
356
+ startWork(request);
357
+ });
358
+}
359
+
360
let serverManifest = {};
361
export function registerServerActions(manifest: ServerManifest) {
362
// This function is called by the bundler to register the manifest.
@@ -292,6 +442,50 @@ export function decodeReply<T>(
442
return root;
443
}
444
445
+export function decodeReplyFromAsyncIterable<T>(
446
+ iterable: AsyncIterable<[string, string | File]>,
447
+ options?: {temporaryReferences?: TemporaryReferenceSet},
448
+): Thenable<T> {
449
+ const iterator: AsyncIterator<[string, string | File]> =
450
+ iterable[ASYNC_ITERATOR]();
451
+
452
+ const response = createResponse(
453
+ serverManifest,
454
+ '',
455
+ options ? options.temporaryReferences : undefined,
456
+ );
457
+
458
+ function progress(
459
+ entry:
460
+ | {done: false, +value: [string, string | File], ...}
461
+ | {done: true, +value: void, ...},
462
+ ) {
463
+ if (entry.done) {
464
+ close(response);
465
+ } else {
466
+ const [name, value] = entry.value;
467
+ if (typeof value === 'string') {
468
+ resolveField(response, name, value);
469
+ } else {
470
+ resolveFile(response, name, value);
471
+ }
472
+ iterator.next().then(progress, error);
473
+ }
474
+ }
475
+ function error(reason: Error) {
476
+ reportGlobalError(response, reason);
477
+ if (typeof (iterator: any).throw === 'function') {
478
+ // The iterator protocol doesn't necessarily include this but a generator do.
479
+ // $FlowFixMe should be able to pass mixed
480
+ iterator.throw(reason).then(error, error);
481
+ }
482
+ }
483
+
484
+ iterator.next().then(progress, error);
485
+
486
+ return getRoot(response);
487
+}
488
+
489
export function decodeAction<T>(body: FormData): Promise<() => T> | null {
490
return decodeActionImpl(body, serverManifest);
491
}
packages/react-server-dom-parcel/src/server/react-flight-dom-server.node.js
+4
-1
@@ -8,10 +8,13 @@
8
*/
9
10
export {
11
+ renderToReadableStream,
12
renderToPipeableStream,
13
+ prerender as unstable_prerender,
14
prerenderToNodeStream as unstable_prerenderToNodeStream,
13
- decodeReplyFromBusboy,
15
decodeReply,
16
+ decodeReplyFromBusboy,
17
+ decodeReplyFromAsyncIterable,
18
decodeAction,
19
decodeFormState,
20
createClientReference,
packages/react-server-dom-parcel/static.node.js
+4
-1
@@ -7,4 +7,7 @@
7
* @flow
8
*/
9
10
-export {unstable_prerenderToNodeStream} from './src/server/react-flight-dom-server.node';
10
+export {
11
+ unstable_prerender,
12
+ unstable_prerenderToNodeStream,
13
+} from './src/server/react-flight-dom-server.node';
packages/react-server-dom-turbopack/npm/server.node.js
+3
-1
@@ -7,9 +7,11 @@ if (process.env.NODE_ENV === 'production') {
7
s = require('./cjs/react-server-dom-turbopack-server.node.development.js');
8
}
9
10
+exports.renderToReadableStream = s.renderToReadableStream;
11
exports.renderToPipeableStream = s.renderToPipeableStream;
11
-exports.decodeReplyFromBusboy = s.decodeReplyFromBusboy;
12
exports.decodeReply = s.decodeReply;
13
+exports.decodeReplyFromBusboy = s.decodeReplyFromBusboy;
14
+exports.decodeReplyFromAsyncIterable = s.decodeReplyFromAsyncIterable;
15
exports.decodeAction = s.decodeAction;
16
exports.decodeFormState = s.decodeFormState;
17
exports.registerServerReference = s.registerServerReference;
packages/react-server-dom-turbopack/npm/static.node.js
+3
@@ -7,6 +7,9 @@ if (process.env.NODE_ENV === 'production') {
7
s = require('./cjs/react-server-dom-turbopack-server.node.development.js');
8
}
9
10
+if (s.unstable_prerender) {
11
+ exports.unstable_prerender = s.unstable_prerender;
12
+}
13
if (s.unstable_prerenderToNodeStream) {
14
exports.unstable_prerenderToNodeStream = s.unstable_prerenderToNodeStream;
15
}
packages/react-server-dom-turbopack/server.node.js
+3
-1
@@ -9,8 +9,10 @@
9
10
export {
11
renderToPipeableStream,
12
- decodeReplyFromBusboy,
12
+ renderToReadableStream,
13
decodeReply,
14
+ decodeReplyFromBusboy,
15
+ decodeReplyFromAsyncIterable,
16
decodeAction,
17
decodeFormState,
18
registerServerReference,
packages/react-server-dom-turbopack/src/server/ReactFlightDOMServerNode.js
+200
-4
@@ -20,6 +20,8 @@ import type {Thenable} from 'shared/ReactTypes';
20
21
import {Readable} from 'stream';
22
23
+import {ASYNC_ITERATOR} from 'shared/ReactSymbols';
24
+
25
import {
26
createRequest,
27
createPrerenderRequest,
@@ -34,6 +36,7 @@ import {
36
reportGlobalError,
37
close,
38
resolveField,
39
+ resolveFile,
40
resolveFileInfo,
41
resolveFileChunk,
42
resolveFileComplete,
@@ -51,6 +54,8 @@ export {
54
createClientModuleProxy,
55
} from '../ReactFlightTurbopackReferences';
56
57
+import {textEncoder} from 'react-server/src/ReactServerStreamConfigNode';
58
+
59
import type {TemporaryReferenceSet} from 'react-server/src/ReactFlightServerTemporaryReferences';
60
61
export {createTemporaryReferenceSet} from 'react-server/src/ReactFlightServerTemporaryReferences';
@@ -128,11 +133,91 @@ function renderToPipeableStream(
133
};
134
}
135
131
-function createFakeWritable(readable: any): Writable {
136
+function createFakeWritableFromReadableStreamController(
137
+ controller: ReadableStreamController,
138
+): Writable {
139
// The current host config expects a Writable so we create
140
// a fake writable for now to push into the Readable.
141
return ({
135
- write(chunk) {
142
+ write(chunk: string | Uint8Array) {
143
+ if (typeof chunk === 'string') {
144
+ chunk = textEncoder.encode(chunk);
145
+ }
146
+ controller.enqueue(chunk);
147
+ // in web streams there is no backpressure so we can always write more
148
+ return true;
149
+ },
150
+ end() {
151
+ controller.close();
152
+ },
153
+ destroy(error) {
154
+ // $FlowFixMe[method-unbinding]
155
+ if (typeof controller.error === 'function') {
156
+ // $FlowFixMe[incompatible-call]: This is an Error object or the destination accepts other types.
157
+ controller.error(error);
158
+ } else {
159
+ controller.close();
160
+ }
161
+ },
162
+ }: any);
163
+}
164
+
165
+function renderToReadableStream(
166
+ model: ReactClientValue,
167
+ turbopackMap: ClientManifest,
168
+ options?: Options & {
169
+ signal?: AbortSignal,
170
+ },
171
+): ReadableStream {
172
+ const request = createRequest(
173
+ model,
174
+ turbopackMap,
175
+ options ? options.onError : undefined,
176
+ options ? options.identifierPrefix : undefined,
177
+ options ? options.onPostpone : undefined,
178
+ options ? options.temporaryReferences : undefined,
179
+ __DEV__ && options ? options.environmentName : undefined,
180
+ __DEV__ && options ? options.filterStackFrame : undefined,
181
+ );
182
+ if (options && options.signal) {
183
+ const signal = options.signal;
184
+ if (signal.aborted) {
185
+ abort(request, (signal: any).reason);
186
+ } else {
187
+ const listener = () => {
188
+ abort(request, (signal: any).reason);
189
+ signal.removeEventListener('abort', listener);
190
+ };
191
+ signal.addEventListener('abort', listener);
192
+ }
193
+ }
194
+ let writable: Writable;
195
+ const stream = new ReadableStream(
196
+ {
197
+ type: 'bytes',
198
+ start: (controller): ?Promise<void> => {
199
+ writable = createFakeWritableFromReadableStreamController(controller);
200
+ startWork(request);
201
+ },
202
+ pull: (controller): ?Promise<void> => {
203
+ startFlowing(request, writable);
204
+ },
205
+ cancel: (reason): ?Promise<void> => {
206
+ stopFlowing(request);
207
+ abort(request, reason);
208
+ },
209
+ },
210
+ // $FlowFixMe[prop-missing] size() methods are not allowed on byte streams.
211
+ {highWaterMark: 0},
212
+ );
213
+ return stream;
214
+}
215
+
216
+function createFakeWritableFromNodeReadable(readable: any): Writable {
217
+ // The current host config expects a Writable so we create
218
+ // a fake writable for now to push into the Readable.
219
+ return ({
220
+ write(chunk: string | Uint8Array) {
221
return readable.push(chunk);
222
},
223
end() {
@@ -171,7 +256,7 @@ function prerenderToNodeStream(
256
startFlowing(request, writable);
257
},
258
});
174
- const writable = createFakeWritable(readable);
259
+ const writable = createFakeWritableFromNodeReadable(readable);
260
resolve({prelude: readable});
261
}
262
@@ -205,6 +290,69 @@ function prerenderToNodeStream(
290
});
291
}
292
293
+function prerender(
294
+ model: ReactClientValue,
295
+ turbopackMap: ClientManifest,
296
+ options?: Options & {
297
+ signal?: AbortSignal,
298
+ },
299
+): Promise<{
300
+ prelude: ReadableStream,
301
+}> {
302
+ return new Promise((resolve, reject) => {
303
+ const onFatalError = reject;
304
+ function onAllReady() {
305
+ let writable: Writable;
306
+ const stream = new ReadableStream(
307
+ {
308
+ type: 'bytes',
309
+ start: (controller): ?Promise<void> => {
310
+ writable =
311
+ createFakeWritableFromReadableStreamController(controller);
312
+ },
313
+ pull: (controller): ?Promise<void> => {
314
+ startFlowing(request, writable);
315
+ },
316
+ cancel: (reason): ?Promise<void> => {
317
+ stopFlowing(request);
318
+ abort(request, reason);
319
+ },
320
+ },
321
+ // $FlowFixMe[prop-missing] size() methods are not allowed on byte streams.
322
+ {highWaterMark: 0},
323
+ );
324
+ resolve({prelude: stream});
325
+ }
326
+ const request = createPrerenderRequest(
327
+ model,
328
+ turbopackMap,
329
+ onAllReady,
330
+ onFatalError,
331
+ options ? options.onError : undefined,
332
+ options ? options.identifierPrefix : undefined,
333
+ options ? options.onPostpone : undefined,
334
+ options ? options.temporaryReferences : undefined,
335
+ __DEV__ && options ? options.environmentName : undefined,
336
+ __DEV__ && options ? options.filterStackFrame : undefined,
337
+ );
338
+ if (options && options.signal) {
339
+ const signal = options.signal;
340
+ if (signal.aborted) {
341
+ const reason = (signal: any).reason;
342
+ abort(request, reason);
343
+ } else {
344
+ const listener = () => {
345
+ const reason = (signal: any).reason;
346
+ abort(request, reason);
347
+ signal.removeEventListener('abort', listener);
348
+ };
349
+ signal.addEventListener('abort', listener);
350
+ }
351
+ }
352
+ startWork(request);
353
+ });
354
+}
355
+
356
function decodeReplyFromBusboy<T>(
357
busboyStream: Busboy,
358
turbopackMap: ServerManifest,
@@ -286,11 +434,59 @@ function decodeReply<T>(
434
return root;
435
}
436
437
+function decodeReplyFromAsyncIterable<T>(
438
+ iterable: AsyncIterable<[string, string | File]>,
439
+ turbopackMap: ServerManifest,
440
+ options?: {temporaryReferences?: TemporaryReferenceSet},
441
+): Thenable<T> {
442
+ const iterator: AsyncIterator<[string, string | File]> =
443
+ iterable[ASYNC_ITERATOR]();
444
+
445
+ const response = createResponse(
446
+ turbopackMap,
447
+ '',
448
+ options ? options.temporaryReferences : undefined,
449
+ );
450
+
451
+ function progress(
452
+ entry:
453
+ | {done: false, +value: [string, string | File], ...}
454
+ | {done: true, +value: void, ...},
455
+ ) {
456
+ if (entry.done) {
457
+ close(response);
458
+ } else {
459
+ const [name, value] = entry.value;
460
+ if (typeof value === 'string') {
461
+ resolveField(response, name, value);
462
+ } else {
463
+ resolveFile(response, name, value);
464
+ }
465
+ iterator.next().then(progress, error);
466
+ }
467
+ }
468
+ function error(reason: Error) {
469
+ reportGlobalError(response, reason);
470
+ if (typeof (iterator: any).throw === 'function') {
471
+ // The iterator protocol doesn't necessarily include this but a generator do.
472
+ // $FlowFixMe should be able to pass mixed
473
+ iterator.throw(reason).then(error, error);
474
+ }
475
+ }
476
+
477
+ iterator.next().then(progress, error);
478
+
479
+ return getRoot(response);
480
+}
481
+
482
export {
483
+ renderToReadableStream,
484
renderToPipeableStream,
485
+ prerender,
486
prerenderToNodeStream,
292
- decodeReplyFromBusboy,
487
decodeReply,
488
+ decodeReplyFromBusboy,
489
+ decodeReplyFromAsyncIterable,
490
decodeAction,
491
decodeFormState,
492
};
packages/react-server-dom-turbopack/src/server/react-flight-dom-server.node.js
+4
-1
@@ -8,10 +8,13 @@
8
*/
9
10
export {
11
+ renderToReadableStream,
12
renderToPipeableStream,
13
+ prerender as unstable_prerender,
14
prerenderToNodeStream as unstable_prerenderToNodeStream,
13
- decodeReplyFromBusboy,
15
decodeReply,
16
+ decodeReplyFromBusboy,
17
+ decodeReplyFromAsyncIterable,
18
decodeAction,
19
decodeFormState,
20
registerServerReference,
packages/react-server-dom-turbopack/static.node.js
+4
-1
@@ -7,4 +7,7 @@
7
* @flow
8
*/
9
10
-export {unstable_prerenderToNodeStream} from './src/server/react-flight-dom-server.node';
10
+export {
11
+ unstable_prerender,
12
+ unstable_prerenderToNodeStream,
13
+} from './src/server/react-flight-dom-server.node';
packages/react-server-dom-webpack/npm/server.node.js
+3
-1
@@ -7,9 +7,11 @@ if (process.env.NODE_ENV === 'production') {
7
s = require('./cjs/react-server-dom-webpack-server.node.development.js');
8
}
9
10
+exports.renderToReadableStream = s.renderToReadableStream;
11
exports.renderToPipeableStream = s.renderToPipeableStream;
11
-exports.decodeReplyFromBusboy = s.decodeReplyFromBusboy;
12
exports.decodeReply = s.decodeReply;
13
+exports.decodeReplyFromBusboy = s.decodeReplyFromBusboy;
14
+exports.decodeReplyFromAsyncIterable = s.decodeReplyFromAsyncIterable;
15
exports.decodeAction = s.decodeAction;
16
exports.decodeFormState = s.decodeFormState;
17
exports.registerServerReference = s.registerServerReference;
packages/react-server-dom-webpack/npm/static.node.js
+3
@@ -7,6 +7,9 @@ if (process.env.NODE_ENV === 'production') {
7
s = require('./cjs/react-server-dom-webpack-server.node.development.js');
8
}
9
10
+if (s.unstable_prerender) {
11
+ exports.unstable_prerender = s.unstable_prerender;
12
+}
13
if (s.unstable_prerenderToNodeStream) {
14
exports.unstable_prerenderToNodeStream = s.unstable_prerenderToNodeStream;
15
}
packages/react-server-dom-webpack/server.node.js
+3
-1
@@ -9,8 +9,10 @@
9
10
export {
11
renderToPipeableStream,
12
- decodeReplyFromBusboy,
12
+ renderToReadableStream,
13
decodeReply,
14
+ decodeReplyFromBusboy,
15
+ decodeReplyFromAsyncIterable,
16
decodeAction,
17
decodeFormState,
18
registerServerReference,
packages/react-server-dom-webpack/src/__tests__/ReactFlightDOMNode-test.js
+44
-3
@@ -5,6 +5,7 @@
5
* LICENSE file in the root directory of this source tree.
6
*
7
* @emails react-core
8
+ * @jest-environment node
9
*/
10
11
'use strict';
@@ -92,6 +93,48 @@ describe('ReactFlightDOMNode', () => {
93
});
94
}
95
96
+ it('should support web streams in node', async () => {
97
+ function Text({children}) {
98
+ return <span>{children}</span>;
99
+ }
100
+ // Large strings can get encoded differently so we need to test that.
101
+ const largeString = 'world'.repeat(1000);
102
+ function HTML() {
103
+ return (
104
+ <div>
105
+ <Text>hello</Text>
106
+ <Text>{largeString}</Text>
107
+ </div>
108
+ );
109
+ }
110
+
111
+ function App() {
112
+ const model = {
113
+ html: <HTML />,
114
+ };
115
+ return model;
116
+ }
117
+
118
+ const readable = await serverAct(() =>
119
+ ReactServerDOMServer.renderToReadableStream(<App />, webpackMap),
120
+ );
121
+ const response = ReactServerDOMClient.createFromReadableStream(readable, {
122
+ serverConsumerManifest: {
123
+ moduleMap: null,
124
+ moduleLoading: null,
125
+ },
126
+ });
127
+ const model = await response;
128
+ expect(model).toEqual({
129
+ html: (
130
+ <div>
131
+ <span>hello</span>
132
+ <span>{largeString}</span>
133
+ </div>
134
+ ),
135
+ });
136
+ });
137
+
138
it('should allow an alternative module mapping to be used for SSR', async () => {
139
function ClientComponent() {
140
return <span>Client Component</span>;
@@ -498,8 +541,6 @@ describe('ReactFlightDOMNode', () => {
541
expect(errors).toEqual([new Error('Connection closed.')]);
542
// Should still match the result when parsed
543
const result = await readResult(ssrStream);
501
- const div = document.createElement('div');
502
- div.innerHTML = result;
503
- expect(div.textContent).toBe('loading...');
544
+ expect(result).toContain('loading...');
545
});
546
});
packages/react-server-dom-webpack/src/server/ReactFlightDOMServerNode.js
+200
-4
@@ -20,6 +20,8 @@ import type {Thenable} from 'shared/ReactTypes';
20
21
import {Readable} from 'stream';
22
23
+import {ASYNC_ITERATOR} from 'shared/ReactSymbols';
24
+
25
import {
26
createRequest,
27
createPrerenderRequest,
@@ -34,6 +36,7 @@ import {
36
reportGlobalError,
37
close,
38
resolveField,
39
+ resolveFile,
40
resolveFileInfo,
41
resolveFileChunk,
42
resolveFileComplete,
@@ -51,6 +54,8 @@ export {
54
createClientModuleProxy,
55
} from '../ReactFlightWebpackReferences';
56
57
+import {textEncoder} from 'react-server/src/ReactServerStreamConfigNode';
58
+
59
import type {TemporaryReferenceSet} from 'react-server/src/ReactFlightServerTemporaryReferences';
60
61
export {createTemporaryReferenceSet} from 'react-server/src/ReactFlightServerTemporaryReferences';
@@ -128,11 +133,91 @@ function renderToPipeableStream(
133
};
134
}
135
131
-function createFakeWritable(readable: any): Writable {
136
+function createFakeWritableFromReadableStreamController(
137
+ controller: ReadableStreamController,
138
+): Writable {
139
// The current host config expects a Writable so we create
140
// a fake writable for now to push into the Readable.
141
return ({
135
- write(chunk) {
142
+ write(chunk: string | Uint8Array) {
143
+ if (typeof chunk === 'string') {
144
+ chunk = textEncoder.encode(chunk);
145
+ }
146
+ controller.enqueue(chunk);
147
+ // in web streams there is no backpressure so we can always write more
148
+ return true;
149
+ },
150
+ end() {
151
+ controller.close();
152
+ },
153
+ destroy(error) {
154
+ // $FlowFixMe[method-unbinding]
155
+ if (typeof controller.error === 'function') {
156
+ // $FlowFixMe[incompatible-call]: This is an Error object or the destination accepts other types.
157
+ controller.error(error);
158
+ } else {
159
+ controller.close();
160
+ }
161
+ },
162
+ }: any);
163
+}
164
+
165
+function renderToReadableStream(
166
+ model: ReactClientValue,
167
+ webpackMap: ClientManifest,
168
+ options?: Options & {
169
+ signal?: AbortSignal,
170
+ },
171
+): ReadableStream {
172
+ const request = createRequest(
173
+ model,
174
+ webpackMap,
175
+ options ? options.onError : undefined,
176
+ options ? options.identifierPrefix : undefined,
177
+ options ? options.onPostpone : undefined,
178
+ options ? options.temporaryReferences : undefined,
179
+ __DEV__ && options ? options.environmentName : undefined,
180
+ __DEV__ && options ? options.filterStackFrame : undefined,
181
+ );
182
+ if (options && options.signal) {
183
+ const signal = options.signal;
184
+ if (signal.aborted) {
185
+ abort(request, (signal: any).reason);
186
+ } else {
187
+ const listener = () => {
188
+ abort(request, (signal: any).reason);
189
+ signal.removeEventListener('abort', listener);
190
+ };
191
+ signal.addEventListener('abort', listener);
192
+ }
193
+ }
194
+ let writable: Writable;
195
+ const stream = new ReadableStream(
196
+ {
197
+ type: 'bytes',
198
+ start: (controller): ?Promise<void> => {
199
+ writable = createFakeWritableFromReadableStreamController(controller);
200
+ startWork(request);
201
+ },
202
+ pull: (controller): ?Promise<void> => {
203
+ startFlowing(request, writable);
204
+ },
205
+ cancel: (reason): ?Promise<void> => {
206
+ stopFlowing(request);
207
+ abort(request, reason);
208
+ },
209
+ },
210
+ // $FlowFixMe[prop-missing] size() methods are not allowed on byte streams.
211
+ {highWaterMark: 0},
212
+ );
213
+ return stream;
214
+}
215
+
216
+function createFakeWritableFromNodeReadable(readable: any): Writable {
217
+ // The current host config expects a Writable so we create
218
+ // a fake writable for now to push into the Readable.
219
+ return ({
220
+ write(chunk: string | Uint8Array) {
221
return readable.push(chunk);
222
},
223
end() {
@@ -171,7 +256,7 @@ function prerenderToNodeStream(
256
startFlowing(request, writable);
257
},
258
});
174
- const writable = createFakeWritable(readable);
259
+ const writable = createFakeWritableFromNodeReadable(readable);
260
resolve({prelude: readable});
261
}
262
@@ -205,6 +290,69 @@ function prerenderToNodeStream(
290
});
291
}
292
293
+function prerender(
294
+ model: ReactClientValue,
295
+ webpackMap: ClientManifest,
296
+ options?: Options & {
297
+ signal?: AbortSignal,
298
+ },
299
+): Promise<{
300
+ prelude: ReadableStream,
301
+}> {
302
+ return new Promise((resolve, reject) => {
303
+ const onFatalError = reject;
304
+ function onAllReady() {
305
+ let writable: Writable;
306
+ const stream = new ReadableStream(
307
+ {
308
+ type: 'bytes',
309
+ start: (controller): ?Promise<void> => {
310
+ writable =
311
+ createFakeWritableFromReadableStreamController(controller);
312
+ },
313
+ pull: (controller): ?Promise<void> => {
314
+ startFlowing(request, writable);
315
+ },
316
+ cancel: (reason): ?Promise<void> => {
317
+ stopFlowing(request);
318
+ abort(request, reason);
319
+ },
320
+ },
321
+ // $FlowFixMe[prop-missing] size() methods are not allowed on byte streams.
322
+ {highWaterMark: 0},
323
+ );
324
+ resolve({prelude: stream});
325
+ }
326
+ const request = createPrerenderRequest(
327
+ model,
328
+ webpackMap,
329
+ onAllReady,
330
+ onFatalError,
331
+ options ? options.onError : undefined,
332
+ options ? options.identifierPrefix : undefined,
333
+ options ? options.onPostpone : undefined,
334
+ options ? options.temporaryReferences : undefined,
335
+ __DEV__ && options ? options.environmentName : undefined,
336
+ __DEV__ && options ? options.filterStackFrame : undefined,
337
+ );
338
+ if (options && options.signal) {
339
+ const signal = options.signal;
340
+ if (signal.aborted) {
341
+ const reason = (signal: any).reason;
342
+ abort(request, reason);
343
+ } else {
344
+ const listener = () => {
345
+ const reason = (signal: any).reason;
346
+ abort(request, reason);
347
+ signal.removeEventListener('abort', listener);
348
+ };
349
+ signal.addEventListener('abort', listener);
350
+ }
351
+ }
352
+ startWork(request);
353
+ });
354
+}
355
+
356
function decodeReplyFromBusboy<T>(
357
busboyStream: Busboy,
358
webpackMap: ServerManifest,
@@ -286,11 +434,59 @@ function decodeReply<T>(
434
return root;
435
}
436
437
+function decodeReplyFromAsyncIterable<T>(
438
+ iterable: AsyncIterable<[string, string | File]>,
439
+ webpackMap: ServerManifest,
440
+ options?: {temporaryReferences?: TemporaryReferenceSet},
441
+): Thenable<T> {
442
+ const iterator: AsyncIterator<[string, string | File]> =
443
+ iterable[ASYNC_ITERATOR]();
444
+
445
+ const response = createResponse(
446
+ webpackMap,
447
+ '',
448
+ options ? options.temporaryReferences : undefined,
449
+ );
450
+
451
+ function progress(
452
+ entry:
453
+ | {done: false, +value: [string, string | File], ...}
454
+ | {done: true, +value: void, ...},
455
+ ) {
456
+ if (entry.done) {
457
+ close(response);
458
+ } else {
459
+ const [name, value] = entry.value;
460
+ if (typeof value === 'string') {
461
+ resolveField(response, name, value);
462
+ } else {
463
+ resolveFile(response, name, value);
464
+ }
465
+ iterator.next().then(progress, error);
466
+ }
467
+ }
468
+ function error(reason: Error) {
469
+ reportGlobalError(response, reason);
470
+ if (typeof (iterator: any).throw === 'function') {
471
+ // The iterator protocol doesn't necessarily include this but a generator do.
472
+ // $FlowFixMe should be able to pass mixed
473
+ iterator.throw(reason).then(error, error);
474
+ }
475
+ }
476
+
477
+ iterator.next().then(progress, error);
478
+
479
+ return getRoot(response);
480
+}
481
+
482
export {
483
+ renderToReadableStream,
484
renderToPipeableStream,
485
+ prerender,
486
prerenderToNodeStream,
292
- decodeReplyFromBusboy,
487
decodeReply,
488
+ decodeReplyFromBusboy,
489
+ decodeReplyFromAsyncIterable,
490
decodeAction,
491
decodeFormState,
492
};
packages/react-server-dom-webpack/src/server/react-flight-dom-server.node.js
+4
-1
@@ -8,10 +8,13 @@
8
*/
9
10
export {
11
+ renderToReadableStream,
12
renderToPipeableStream,
13
+ prerender as unstable_prerender,
14
prerenderToNodeStream as unstable_prerenderToNodeStream,
13
- decodeReplyFromBusboy,
15
decodeReply,
16
+ decodeReplyFromBusboy,
17
+ decodeReplyFromAsyncIterable,
18
decodeAction,
19
decodeFormState,
20
registerServerReference,
packages/react-server-dom-webpack/src/server/react-flight-dom-server.node.unbundled.js
+4
-1
@@ -8,10 +8,13 @@
8
*/
9
10
export {
11
+ renderToReadableStream,
12
renderToPipeableStream,
13
+ prerender as unstable_prerender,
14
prerenderToNodeStream as unstable_prerenderToNodeStream,
13
- decodeReplyFromBusboy,
15
decodeReply,
16
+ decodeReplyFromBusboy,
17
+ decodeReplyFromAsyncIterable,
18
decodeAction,
19
decodeFormState,
20
registerServerReference,
packages/react-server-dom-webpack/static.node.js
+4
-1
@@ -7,4 +7,7 @@
7
* @flow
8
*/
9
10
-export {unstable_prerenderToNodeStream} from './src/server/react-flight-dom-server.node';
10
+export {
11
+ unstable_prerender,
12
+ unstable_prerenderToNodeStream,
13
+} from './src/server/react-flight-dom-server.node';