@samitouri / QOS-React-2 / commits / 7843b142ac

[Fizz/Flight] Pass in Destination lazily to startFlowing instead of in createRequest (#22449)

* Pass in Destination lazily in startFlowing instead of createRequest * Delay fatal errors until we have a destination to forward them to * Flow can now be inferred by whether there's a destination set We can drop the destination when we're not flowing since there's nothing to write to. Fatal errors now close once flowing starts back up again. * Defer fatal errors in Flight too

Sebastian Markbåge committed Sep 28, 2021 at 18:32 UTC 7843b142ac804655990157a7be1e4641b4b6695f
14 files changed +110 -88
packages/react-dom/src/__tests__/ReactDOMFizzServerNode-test.js
+3 -2
@@ -154,7 +154,7 @@ describe('ReactDOMFizzServer', () => {
154 it('should error the stream when an error is thrown at the root', async () => {
155 const reportedErrors = [];
156 const {writable, output, completed} = getTestWritable();
157 - ReactDOMFizzServer.pipeToNodeWritable(
157 + const {startWriting} = ReactDOMFizzServer.pipeToNodeWritable(
158 <div>
159 <Throw />
160 </div>,
@@ -166,7 +166,8 @@ describe('ReactDOMFizzServer', () => {
166 },
167 );
168
169 - // The stream is errored even if we haven't started writing.
169 + // The stream is errored once we start writing.
170 + startWriting();
171
172 await completed;
173
packages/react-dom/src/server/ReactDOMFizzServerBrowser.js
+10 -12
@@ -37,7 +37,15 @@ function renderToReadableStream(
37 children: ReactNodeList,
38 options?: Options,
39 ): ReadableStream {
40 - let request;
40 + const request = createRequest(
41 + children,
42 + createResponseState(options ? options.identifierPrefix : undefined),
43 + createRootFormatContext(options ? options.namespaceURI : undefined),
44 + options ? options.progressiveChunkSize : undefined,
45 + options ? options.onError : undefined,
46 + options ? options.onCompleteAll : undefined,
47 + options ? options.onCompleteShell : undefined,
48 + );
49 if (options && options.signal) {
50 const signal = options.signal;
51 const listener = () => {
@@ -48,16 +56,6 @@ function renderToReadableStream(
56 }
57 const stream = new ReadableStream({
58 start(controller) {
51 - request = createRequest(
52 - children,
53 - controller,
54 - createResponseState(options ? options.identifierPrefix : undefined),
55 - createRootFormatContext(options ? options.namespaceURI : undefined),
56 - options ? options.progressiveChunkSize : undefined,
57 - options ? options.onError : undefined,
58 - options ? options.onCompleteAll : undefined,
59 - options ? options.onCompleteShell : undefined,
60 - );
59 startWork(request);
60 },
61 pull(controller) {
@@ -66,7 +64,7 @@ function renderToReadableStream(
64 // is actually used by something so we can give it the best result possible
65 // at that point.
66 if (stream.locked) {
69 - startFlowing(request);
67 + startFlowing(request, controller);
68 }
69 },
70 cancel(reason) {},
packages/react-dom/src/server/ReactDOMFizzServerNode.js
+4 -9
@@ -25,7 +25,7 @@ import {
25 } from './ReactDOMServerFormatConfig';
26
27 function createDrainHandler(destination, request) {
28 - return () => startFlowing(request);
28 + return () => startFlowing(request, destination);
29 }
30
31 type Options = {|
@@ -44,14 +44,9 @@ type Controls = {|
44 startWriting(): void,
45 |};
46
47 -function createRequestImpl(
48 - children: ReactNodeList,
49 - destination: Writable,
50 - options: void | Options,
51 -) {
47 +function createRequestImpl(children: ReactNodeList, options: void | Options) {
48 return createRequest(
49 children,
54 - destination,
50 createResponseState(options ? options.identifierPrefix : undefined),
51 createRootFormatContext(options ? options.namespaceURI : undefined),
52 options ? options.progressiveChunkSize : undefined,
@@ -66,7 +61,7 @@ function pipeToNodeWritable(
61 destination: Writable,
62 options?: Options,
63 ): Controls {
69 - const request = createRequestImpl(children, destination, options);
64 + const request = createRequestImpl(children, options);
65 let hasStartedFlowing = false;
66 startWork(request);
67 return {
@@ -75,7 +70,7 @@ function pipeToNodeWritable(
70 return;
71 }
72 hasStartedFlowing = true;
78 - startFlowing(request);
73 + startFlowing(request, destination);
74 destination.on('drain', createDrainHandler(destination, request));
75 },
76 abort() {
packages/react-dom/src/server/ReactDOMLegacyServerBrowser.js
+1 -2
@@ -59,7 +59,6 @@ function renderToStringImpl(
59 }
60 const request = createRequest(
61 children,
62 - destination,
62 createResponseState(
63 generateStaticMarkup,
64 options ? options.identifierPrefix : undefined,
@@ -74,7 +73,7 @@ function renderToStringImpl(
73 // If anything suspended and is still pending, we'll abort it before writing.
74 // That way we write only client-rendered boundaries from the start.
75 abort(request);
77 - startFlowing(request);
76 + startFlowing(request, destination);
77 if (didFatal) {
78 throw fatalError;
79 }
packages/react-dom/src/server/ReactDOMLegacyServerNode.js
+2 -3
@@ -54,7 +54,7 @@ class ReactMarkupReadableStream extends Readable {
54
55 _read(size) {
56 if (this.startedFlowing) {
57 - startFlowing(this.request);
57 + startFlowing(this.request, this);
58 }
59 }
60 }
@@ -72,12 +72,11 @@ function renderToNodeStreamImpl(
72 // We wait until everything has loaded before starting to write.
73 // That way we only end up with fully resolved HTML even if we suspend.
74 destination.startedFlowing = true;
75 - startFlowing(request);
75 + startFlowing(request, destination);
76 }
77 const destination = new ReactMarkupReadableStream();
78 const request = createRequest(
79 children,
80 - destination,
80 createResponseState(false, options ? options.identifierPrefix : undefined),
81 createRootFormatContext(),
82 Infinity,
packages/react-noop-renderer/src/ReactNoopFlightServer.js
+1 -2
@@ -63,12 +63,11 @@ function render(model: ReactModel, options?: Options): Destination {
63 const bundlerConfig = undefined;
64 const request = ReactNoopFlightServer.createRequest(
65 model,
66 - destination,
66 bundlerConfig,
67 options ? options.onError : undefined,
68 );
69 ReactNoopFlightServer.startWork(request);
71 - ReactNoopFlightServer.startFlowing(request);
70 + ReactNoopFlightServer.startFlowing(request, destination);
71 return destination;
72 }
73
packages/react-noop-renderer/src/ReactNoopServer.js
+1 -2
@@ -259,7 +259,6 @@ function render(children: React$Element<any>, options?: Options): Destination {
259 };
260 const request = ReactNoopServer.createRequest(
261 children,
262 - destination,
262 null,
263 null,
264 options ? options.progressiveChunkSize : undefined,
@@ -268,7 +267,7 @@ function render(children: React$Element<any>, options?: Options): Destination {
267 options ? options.onCompleteShell : undefined,
268 );
269 ReactNoopServer.startWork(request);
271 - ReactNoopServer.startFlowing(request);
270 + ReactNoopServer.startFlowing(request, destination);
271 return destination;
272 }
273
packages/react-server-dom-relay/src/ReactDOMServerFB.js
+1 -2
@@ -46,7 +46,6 @@ function renderToStream(children: ReactNodeList, options: Options): Stream {
46 };
47 const request = createRequest(
48 children,
49 - destination,
49 createResponseState(options ? options.identifierPrefix : undefined),
50 createRootFormatContext(undefined),
51 options ? options.progressiveChunkSize : undefined,
@@ -71,7 +70,7 @@ function abortStream(stream: Stream): void {
70 function renderNextChunk(stream: Stream): string {
71 const {request, destination} = stream;
72 performWork(request);
74 - startFlowing(request);
73 + startFlowing(request, destination);
74 if (destination.fatal) {
75 throw destination.error;
76 }
packages/react-server-dom-relay/src/ReactFlightDOMRelayServer.js
+1 -2
@@ -31,12 +31,11 @@ function render(
31 ): void {
32 const request = createRequest(
33 model,
34 - destination,
34 config,
35 options ? options.onError : undefined,
36 );
37 startWork(request);
39 - startFlowing(request);
38 + startFlowing(request, destination);
39 }
40
41 export {render};
packages/react-server-dom-webpack/src/ReactFlightDOMServerBrowser.js
+6 -8
@@ -25,15 +25,13 @@ function renderToReadableStream(
25 webpackMap: BundlerConfig,
26 options?: Options,
27 ): ReadableStream {
28 - let request;
28 + const request = createRequest(
29 + model,
30 + webpackMap,
31 + options ? options.onError : undefined,
32 + );
33 const stream = new ReadableStream({
34 start(controller) {
31 - request = createRequest(
32 - model,
33 - controller,
34 - webpackMap,
35 - options ? options.onError : undefined,
36 - );
35 startWork(request);
36 },
37 pull(controller) {
@@ -42,7 +40,7 @@ function renderToReadableStream(
40 // is actually used by something so we can give it the best result possible
41 // at that point.
42 if (stream.locked) {
45 - startFlowing(request);
43 + startFlowing(request, controller);
44 }
45 },
46 cancel(reason) {},
packages/react-server-dom-webpack/src/ReactFlightDOMServerNode.js
+2 -3
@@ -18,7 +18,7 @@ import {
18 } from 'react-server/src/ReactFlightServer';
19
20 function createDrainHandler(destination, request) {
21 - return () => startFlowing(request);
21 + return () => startFlowing(request, destination);
22 }
23
24 type Options = {
@@ -33,12 +33,11 @@ function pipeToNodeWritable(
33 ): void {
34 const request = createRequest(
35 model,
36 - destination,
36 webpackMap,
37 options ? options.onError : undefined,
38 );
39 startWork(request);
41 - startFlowing(request);
40 + startFlowing(request, destination);
41 destination.on('drain', createDrainHandler(destination, request));
42 }
43
packages/react-server-native-relay/src/ReactFlightNativeRelayServer.js
+2 -2
@@ -24,9 +24,9 @@ function render(
24 destination: Destination,
25 config: BundlerConfig,
26 ): void {
27 - const request = createRequest(model, destination, config);
27 + const request = createRequest(model, config);
28 startWork(request);
29 - startFlowing(request);
29 + startFlowing(request, destination);
30 }
31
32 export {render};
packages/react-server/src/ReactFizzServer.js
+37 -22
@@ -166,15 +166,16 @@ type Segment = {
166 +boundary: null | SuspenseBoundary,
167 };
168
169 -const BUFFERING = 0;
170 -const FLOWING = 1;
169 +const OPEN = 0;
170 +const CLOSING = 1;
171 const CLOSED = 2;
172
173 export opaque type Request = {
174 - +destination: Destination,
174 + destination: null | Destination,
175 +responseState: ResponseState,
176 +progressiveChunkSize: number,
177 status: 0 | 1 | 2,
178 + fatalError: mixed,
179 nextSegmentId: number,
180 allPendingTasks: number, // when it reaches zero, we can close the connection.
181 pendingRootTasks: number, // when this reaches zero, we've finished at least the root boundary.
@@ -221,7 +222,6 @@ function noop(): void {}
222
223 export function createRequest(
224 children: ReactNodeList,
224 - destination: Destination,
225 responseState: ResponseState,
226 rootFormatContext: FormatContext,
227 progressiveChunkSize: void | number,
@@ -232,13 +232,14 @@ export function createRequest(
232 const pingedTasks = [];
233 const abortSet: Set<Task> = new Set();
234 const request = {
235 - destination,
235 + destination: null,
236 responseState,
237 progressiveChunkSize:
238 progressiveChunkSize === undefined
239 ? DEFAULT_PROGRESSIVE_CHUNK_SIZE
240 : progressiveChunkSize,
241 - status: BUFFERING,
241 + status: OPEN,
242 + fatalError: null,
243 nextSegmentId: 0,
244 allPendingTasks: 0,
245 pendingRootTasks: 0,
@@ -404,8 +405,13 @@ function fatalError(request: Request, error: mixed): void {
405 // This is called outside error handling code such as if the root errors outside
406 // a suspense boundary or if the root suspense boundary's fallback errors.
407 // It's also called if React itself or its host configs errors.
407 - request.status = CLOSED;
408 - closeWithError(request.destination, error);
408 + if (request.destination !== null) {
409 + request.status = CLOSED;
410 + closeWithError(request.destination, error);
411 + } else {
412 + request.status = CLOSING;
413 + request.fatalError = error;
414 + }
415 }
416
417 function renderSuspenseBoundary(
@@ -1330,7 +1336,9 @@ function abortTask(task: Task): void {
1336 // the request;
1337 if (request.status !== CLOSED) {
1338 request.status = CLOSED;
1333 - close(request.destination);
1339 + if (request.destination !== null) {
1340 + close(request.destination);
1341 + }
1342 }
1343 } else {
1344 boundary.pendingTasks--;
@@ -1490,8 +1498,8 @@ export function performWork(request: Request): void {
1498 retryTask(request, task);
1499 }
1500 pingedTasks.splice(0, i);
1493 - if (request.status === FLOWING) {
1494 - flushCompletedQueues(request);
1501 + if (request.destination !== null) {
1502 + flushCompletedQueues(request, request.destination);
1503 }
1504 } catch (error) {
1505 reportError(request, error);
@@ -1748,8 +1756,10 @@ function flushPartiallyCompletedSegment(
1756 }
1757 }
1758
1751 -function flushCompletedQueues(request: Request): void {
1752 - const destination = request.destination;
1759 +function flushCompletedQueues(
1760 + request: Request,
1761 + destination: Destination,
1762 +): void {
1763 beginWriting(destination);
1764 try {
1765 // The structure of this is to go through each queue one by one and write
@@ -1775,7 +1785,7 @@ function flushCompletedQueues(request: Request): void {
1785 for (i = 0; i < clientRenderedBoundaries.length; i++) {
1786 const boundary = clientRenderedBoundaries[i];
1787 if (!flushClientRenderedBoundary(request, destination, boundary)) {
1778 - request.status = BUFFERING;
1788 + request.destination = null;
1789 i++;
1790 clientRenderedBoundaries.splice(0, i);
1791 return;
@@ -1790,7 +1800,7 @@ function flushCompletedQueues(request: Request): void {
1800 for (i = 0; i < completedBoundaries.length; i++) {
1801 const boundary = completedBoundaries[i];
1802 if (!flushCompletedBoundary(request, destination, boundary)) {
1793 - request.status = BUFFERING;
1803 + request.destination = null;
1804 i++;
1805 completedBoundaries.splice(0, i);
1806 return;
@@ -1811,7 +1821,7 @@ function flushCompletedQueues(request: Request): void {
1821 for (i = 0; i < partialBoundaries.length; i++) {
1822 const boundary = partialBoundaries[i];
1823 if (!flushPartialBoundary(request, destination, boundary)) {
1814 - request.status = BUFFERING;
1824 + request.destination = null;
1825 i++;
1826 partialBoundaries.splice(0, i);
1827 return;
@@ -1826,7 +1836,7 @@ function flushCompletedQueues(request: Request): void {
1836 for (i = 0; i < largeBoundaries.length; i++) {
1837 const boundary = largeBoundaries[i];
1838 if (!flushCompletedBoundary(request, destination, boundary)) {
1829 - request.status = BUFFERING;
1839 + request.destination = null;
1840 i++;
1841 largeBoundaries.splice(0, i);
1842 return;
@@ -1861,13 +1871,18 @@ export function startWork(request: Request): void {
1871 scheduleWork(() => performWork(request));
1872 }
1873
1864 -export function startFlowing(request: Request): void {
1874 +export function startFlowing(request: Request, destination: Destination): void {
1875 + if (request.status === CLOSING) {
1876 + request.status = CLOSED;
1877 + closeWithError(destination, request.fatalError);
1878 + return;
1879 + }
1880 if (request.status === CLOSED) {
1881 return;
1882 }
1868 - request.status = FLOWING;
1883 + request.destination = destination;
1884 try {
1870 - flushCompletedQueues(request);
1885 + flushCompletedQueues(request, destination);
1886 } catch (error) {
1887 reportError(request, error);
1888 fatalError(request, error);
@@ -1880,8 +1895,8 @@ export function abort(request: Request): void {
1895 const abortableTasks = request.abortableTasks;
1896 abortableTasks.forEach(abortTask, request);
1897 abortableTasks.clear();
1883 - if (request.status === FLOWING) {
1884 - flushCompletedQueues(request);
1898 + if (request.destination !== null) {
1899 + flushCompletedQueues(request, request.destination);
1900 }
1901 } catch (error) {
1902 reportError(request, error);
packages/react-server/src/ReactFlightServer.js
+39 -17
@@ -72,7 +72,9 @@ type Segment = {
72 };
73
74 export type Request = {
75 - destination: Destination,
75 + status: 0 | 1 | 2,
76 + fatalError: mixed,
77 + destination: null | Destination,
78 bundlerConfig: BundlerConfig,
79 cache: Map<Function, mixed>,
80 nextChunkId: number,
@@ -84,25 +86,30 @@ export type Request = {
86 writtenSymbols: Map<Symbol, number>,
87 writtenModules: Map<ModuleKey, number>,
88 onError: (error: mixed) => void,
87 - flowing: boolean,
89 toJSON: (key: string, value: ReactModel) => ReactJSONValue,
90 };
91
92 const ReactCurrentDispatcher = ReactSharedInternals.ReactCurrentDispatcher;
93
94 function defaultErrorHandler(error: mixed) {
94 - console['error'](error); // Don't transform to our wrapper
95 + console['error'](error);
96 + // Don't transform to our wrapper
97 }
98
99 +const OPEN = 0;
100 +const CLOSING = 1;
101 +const CLOSED = 2;
102 +
103 export function createRequest(
104 model: ReactModel,
99 - destination: Destination,
105 bundlerConfig: BundlerConfig,
106 onError: void | ((error: mixed) => void),
107 ): Request {
108 const pingedSegments = [];
109 const request = {
105 - destination,
110 + status: OPEN,
111 + fatalError: null,
112 + destination: null,
113 bundlerConfig,
114 cache: new Map(),
115 nextChunkId: 0,
@@ -114,7 +121,6 @@ export function createRequest(
121 writtenSymbols: new Map(),
122 writtenModules: new Map(),
123 onError: onError === undefined ? defaultErrorHandler : onError,
117 - flowing: false,
124 toJSON: function(key: string, value: ReactModel): ReactJSONValue {
125 return resolveModelToJSON(request, this, key, value);
126 },
@@ -604,7 +610,13 @@ function reportError(request: Request, error: mixed): void {
610
611 function fatalError(request: Request, error: mixed): void {
612 // This is called outside error handling code such as if an error happens in React internals.
607 - closeWithError(request.destination, error);
613 + if (request.destination !== null) {
614 + request.status = CLOSED;
615 + closeWithError(request.destination, error);
616 + } else {
617 + request.status = CLOSING;
618 + request.fatalError = error;
619 + }
620 }
621
622 function emitErrorChunk(request: Request, id: number, error: mixed): void {
@@ -694,8 +706,8 @@ function performWork(request: Request): void {
706 const segment = pingedSegments[i];
707 retrySegment(request, segment);
708 }
697 - if (request.flowing) {
698 - flushCompletedChunks(request);
709 + if (request.destination !== null) {
710 + flushCompletedChunks(request, request.destination);
711 }
712 } catch (error) {
713 reportError(request, error);
@@ -706,8 +718,10 @@ function performWork(request: Request): void {
718 }
719 }
720
709 -function flushCompletedChunks(request: Request): void {
710 - const destination = request.destination;
721 +function flushCompletedChunks(
722 + request: Request,
723 + destination: Destination,
724 +): void {
725 beginWriting(destination);
726 try {
727 // We emit module chunks first in the stream so that
@@ -718,7 +732,7 @@ function flushCompletedChunks(request: Request): void {
732 request.pendingChunks--;
733 const chunk = moduleChunks[i];
734 if (!writeChunk(destination, chunk)) {
721 - request.flowing = false;
735 + request.destination = null;
736 i++;
737 break;
738 }
@@ -731,7 +745,7 @@ function flushCompletedChunks(request: Request): void {
745 request.pendingChunks--;
746 const chunk = jsonChunks[i];
747 if (!writeChunk(destination, chunk)) {
734 - request.flowing = false;
748 + request.destination = null;
749 i++;
750 break;
751 }
@@ -746,7 +760,7 @@ function flushCompletedChunks(request: Request): void {
760 request.pendingChunks--;
761 const chunk = errorChunks[i];
762 if (!writeChunk(destination, chunk)) {
749 - request.flowing = false;
763 + request.destination = null;
764 i++;
765 break;
766 }
@@ -766,10 +780,18 @@ export function startWork(request: Request): void {
780 scheduleWork(() => performWork(request));
781 }
782
769 -export function startFlowing(request: Request): void {
770 - request.flowing = true;
783 +export function startFlowing(request: Request, destination: Destination): void {
784 + if (request.status === CLOSING) {
785 + request.status = CLOSED;
786 + closeWithError(destination, request.fatalError);
787 + return;
788 + }
789 + if (request.status === CLOSED) {
790 + return;
791 + }
792 + request.destination = destination;
793 try {
772 - flushCompletedChunks(request);
794 + flushCompletedChunks(request, destination);
795 } catch (error) {
796 reportError(request, error);
797 fatalError(request, error);