[Flight] Basic Streaming Suspense Support (#17285)
* Return whether to keep flowing in Host config * Emit basic chunk based streaming in the Flight server When something suspends a new chunk is created. * Add reentrancy check The WHATWG API is designed to be pulled recursively. We should refactor to favor that approach. * Basic streaming Suspense support on the client * Add basic suspense in example * Add comment describing the protocol that the server generates
Sebastian Markbåge committed
Nov 6, 2019 at 09:48 UTC
dee03049f5690e23787f3ba1afd3150fb3540624
11 files changed
+530
-103
fixtures/flight-browser/index.html
+22
-2
@@ -37,8 +37,26 @@
37
);
38
}
39
40
+ let resolved = false;
41
+ let promise = new Promise(resolve => {
42
+ setTimeout(() => {
43
+ resolved = true;
44
+ resolve();
45
+ }, 100);
46
+ });
47
+ function read() {
48
+ if (!resolved) {
49
+ throw promise;
50
+ }
51
+ }
52
+
53
+ function Title() {
54
+ read();
55
+ return 'Title';
56
+ }
57
+
58
let model = {
41
- title: 'Title',
59
+ title: <Title />,
60
content: {
61
__html: <HTML />,
62
}
@@ -69,7 +87,9 @@
87
function Shell({ data }) {
88
let model = data.model;
89
return <div>
72
- <h1>{model.title}</h1>
90
+ <Suspense fallback="...">
91
+ <h1>{model.title}</h1>
92
+ </Suspense>
93
<div dangerouslySetInnerHTML={model.content} />
94
</div>;
95
}
packages/react-dom/src/client/flight/ReactFlightDOMClient.js
+1
-1
@@ -26,7 +26,7 @@ function startReadingFromStream(response, stream: ReadableStream): void {
26
return;
27
}
28
let buffer: Uint8Array = (value: any);
29
- processBinaryChunk(response, buffer, 0);
29
+ processBinaryChunk(response, buffer);
30
return reader.read().then(progress, error);
31
}
32
function error(e) {
packages/react-dom/src/server/ReactDOMFizzServerBrowser.js
+1
-1
@@ -23,7 +23,7 @@ function renderToReadableStream(children: ReactNodeList): ReadableStream {
23
startWork(request);
24
},
25
pull(controller) {
26
- startFlowing(request, controller.desiredSize);
26
+ startFlowing(request);
27
},
28
cancel(reason) {},
29
});
packages/react-dom/src/server/ReactDOMFizzServerNode.js
+1
-1
@@ -13,7 +13,7 @@ import type {Writable} from 'stream';
13
import {createRequest, startWork, startFlowing} from 'react-server/inline.dom';
14
15
function createDrainHandler(destination, request) {
16
- return () => startFlowing(request, 0);
16
+ return () => startFlowing(request);
17
}
18
19
function pipeToNodeWritable(
packages/react-dom/src/server/flight/ReactFlightDOMServerBrowser.js
+1
-1
@@ -23,7 +23,7 @@ function renderToReadableStream(model: ReactModel): ReadableStream {
23
startWork(request);
24
},
25
pull(controller) {
26
- startFlowing(request, controller.desiredSize);
26
+ startFlowing(request);
27
},
28
cancel(reason) {},
29
});
packages/react-dom/src/server/flight/ReactFlightDOMServerNode.js
+1
-1
@@ -17,7 +17,7 @@ import {
17
} from 'react-server/flight.inline.dom';
18
19
function createDrainHandler(destination, request) {
20
- return () => startFlowing(request, 0);
20
+ return () => startFlowing(request);
21
}
22
23
function pipeToNodeWritable(model: ReactModel, destination: Writable): void {
packages/react-flight/src/ReactFlightClient.js
+223
-54
@@ -20,63 +20,224 @@ export type ReactModelRoot<T> = {|
20
model: T,
21
|};
22
23
-type OpaqueResponse = {
23
+type JSONValue = number | null | boolean | string | {[key: string]: JSONValue};
24
+
25
+const PENDING = 0;
26
+const RESOLVED = 1;
27
+const ERRORED = 2;
28
+
29
+type PendingChunk = {|
30
+ status: 0,
31
+ value: Promise<void>,
32
+ resolve: () => void,
33
+|};
34
+type ResolvedChunk = {|
35
+ status: 1,
36
+ value: mixed,
37
+ resolve: null,
38
+|};
39
+type ErroredChunk = {|
40
+ status: 2,
41
+ value: Error,
42
+ resolve: null,
43
+|};
44
+type Chunk = PendingChunk | ResolvedChunk | ErroredChunk;
45
+
46
+type OpaqueResponseWithoutDecoder = {
47
source: Source,
25
- modelRoot: ReactModelRoot<any>,
48
partialRow: string,
49
+ modelRoot: ReactModelRoot<any>,
50
+ chunks: Map<number, Chunk>,
51
+ fromJSON: (key: string, value: JSONValue) => any,
52
+};
53
+
54
+type OpaqueResponse = OpaqueResponseWithoutDecoder & {
55
stringDecoder: StringDecoder,
28
- rootPing: () => void,
56
};
57
58
export function createResponse(source: Source): OpaqueResponse {
32
- let modelRoot = {};
33
- Object.defineProperty(
34
- modelRoot,
35
- 'model',
36
- ({
37
- configurable: true,
38
- enumerable: true,
39
- get() {
40
- throw rootPromise;
41
- },
42
- }: any),
43
- );
44
-
45
- let rootPing;
46
- let rootPromise = new Promise(resolve => {
47
- rootPing = resolve;
48
- });
59
+ let modelRoot: ReactModelRoot<any> = ({}: any);
60
+ let rootChunk: Chunk = createPendingChunk();
61
+ definePendingProperty(modelRoot, 'model', rootChunk);
62
+ let chunks: Map<number, Chunk> = new Map();
63
+ chunks.set(0, rootChunk);
64
50
- let response: OpaqueResponse = ({
65
+ let response: OpaqueResponse = (({
66
source,
52
- modelRoot,
67
partialRow: '',
54
- rootPing,
55
- }: any);
68
+ modelRoot,
69
+ chunks: chunks,
70
+ fromJSON: function(key, value) {
71
+ return parseFromJSON(response, this, key, value);
72
+ },
73
+ }: OpaqueResponseWithoutDecoder): any);
74
if (supportsBinaryStreams) {
75
response.stringDecoder = createStringDecoder();
76
}
77
return response;
78
}
79
80
+function createPendingChunk(): PendingChunk {
81
+ let resolve: () => void = (null: any);
82
+ let promise = new Promise(r => (resolve = r));
83
+ return {
84
+ status: PENDING,
85
+ value: promise,
86
+ resolve: resolve,
87
+ };
88
+}
89
+
90
+function createErrorChunk(error: Error): ErroredChunk {
91
+ return {
92
+ status: ERRORED,
93
+ value: error,
94
+ resolve: null,
95
+ };
96
+}
97
+
98
+function triggerErrorOnChunk(chunk: Chunk, error: Error): void {
99
+ if (chunk.status !== PENDING) {
100
+ // We already resolved. We didn't expect to see this.
101
+ return;
102
+ }
103
+ let resolve = chunk.resolve;
104
+ let erroredChunk: ErroredChunk = (chunk: any);
105
+ erroredChunk.status = ERRORED;
106
+ erroredChunk.value = error;
107
+ erroredChunk.resolve = null;
108
+ resolve();
109
+}
110
+
111
+function createResolvedChunk(value: mixed): ResolvedChunk {
112
+ return {
113
+ status: RESOLVED,
114
+ value: value,
115
+ resolve: null,
116
+ };
117
+}
118
+
119
+function resolveChunk(chunk: Chunk, value: mixed): void {
120
+ if (chunk.status !== PENDING) {
121
+ // We already resolved. We didn't expect to see this.
122
+ return;
123
+ }
124
+ let resolve = chunk.resolve;
125
+ let resolvedChunk: ResolvedChunk = (chunk: any);
126
+ resolvedChunk.status = RESOLVED;
127
+ resolvedChunk.value = value;
128
+ resolvedChunk.resolve = null;
129
+ resolve();
130
+}
131
+
132
// Report that any missing chunks in the model is now going to throw this
133
// error upon read. Also notify any pending promises.
134
export function reportGlobalError(
135
response: OpaqueResponse,
136
error: Error,
137
): void {
68
- Object.defineProperty(
69
- response.modelRoot,
70
- 'model',
71
- ({
72
- configurable: true,
73
- enumerable: true,
74
- get() {
75
- throw error;
76
- },
77
- }: any),
78
- );
79
- response.rootPing();
138
+ response.chunks.forEach(chunk => {
139
+ // If this chunk was already resolved or errored, it won't
140
+ // trigger an error but if it wasn't then we need to
141
+ // because we won't be getting any new data to resolve it.
142
+ triggerErrorOnChunk(chunk, error);
143
+ });
144
+}
145
+
146
+function definePendingProperty(
147
+ object: Object,
148
+ key: string,
149
+ chunk: Chunk,
150
+): void {
151
+ Object.defineProperty(object, key, {
152
+ configurable: false,
153
+ enumerable: true,
154
+ get() {
155
+ if (chunk.status === RESOLVED) {
156
+ return chunk.value;
157
+ } else {
158
+ throw chunk.value;
159
+ }
160
+ },
161
+ });
162
+}
163
+
164
+function parseFromJSON(
165
+ response: OpaqueResponse,
166
+ targetObj: Object,
167
+ key: string,
168
+ value: JSONValue,
169
+): any {
170
+ if (typeof value === 'string' && value[0] === '$') {
171
+ if (value[1] === '$') {
172
+ // This was an escaped string value.
173
+ return value.substring(1);
174
+ } else {
175
+ let id = parseInt(value.substring(1), 16);
176
+ let chunks = response.chunks;
177
+ let chunk = chunks.get(id);
178
+ if (!chunk) {
179
+ chunk = createPendingChunk();
180
+ chunks.set(id, chunk);
181
+ } else if (chunk.status === RESOLVED) {
182
+ return chunk.value;
183
+ }
184
+ definePendingProperty(targetObj, key, chunk);
185
+ return undefined;
186
+ }
187
+ }
188
+ return value;
189
+}
190
+
191
+function resolveJSONRow(
192
+ response: OpaqueResponse,
193
+ id: number,
194
+ json: string,
195
+): void {
196
+ let model = JSON.parse(json, response.fromJSON);
197
+ let chunks = response.chunks;
198
+ let chunk = chunks.get(id);
199
+ if (!chunk) {
200
+ chunks.set(id, createResolvedChunk(model));
201
+ } else {
202
+ resolveChunk(chunk, model);
203
+ }
204
+}
205
+
206
+function processFullRow(response: OpaqueResponse, row: string): void {
207
+ if (row === '') {
208
+ return;
209
+ }
210
+ let tag = row[0];
211
+ switch (tag) {
212
+ case 'J': {
213
+ let colon = row.indexOf(':', 1);
214
+ let id = parseInt(row.substring(1, colon), 16);
215
+ let json = row.substring(colon + 1);
216
+ resolveJSONRow(response, id, json);
217
+ return;
218
+ }
219
+ case 'E': {
220
+ let colon = row.indexOf(':', 1);
221
+ let id = parseInt(row.substring(1, colon), 16);
222
+ let json = row.substring(colon + 1);
223
+ let errorInfo = JSON.parse(json);
224
+ let error = new Error(errorInfo.message);
225
+ error.stack = errorInfo.stack;
226
+ let chunks = response.chunks;
227
+ let chunk = chunks.get(id);
228
+ if (!chunk) {
229
+ chunks.set(id, createErrorChunk(error));
230
+ } else {
231
+ triggerErrorOnChunk(chunk, error);
232
+ }
233
+ return;
234
+ }
235
+ default: {
236
+ // Assume this is the root model.
237
+ resolveJSONRow(response, 0, row);
238
+ return;
239
+ }
240
+ }
241
}
242
243
export function processStringChunk(
@@ -84,36 +245,44 @@ export function processStringChunk(
245
chunk: string,
246
offset: number,
247
): void {
87
- response.partialRow += chunk.substr(offset);
248
+ let linebreak = chunk.indexOf('\n', offset);
249
+ while (linebreak > -1) {
250
+ let fullrow = response.partialRow + chunk.substring(offset, linebreak);
251
+ processFullRow(response, fullrow);
252
+ response.partialRow = '';
253
+ offset = linebreak + 1;
254
+ linebreak = chunk.indexOf('\n', offset);
255
+ }
256
+ response.partialRow += chunk.substring(offset);
257
}
258
259
export function processBinaryChunk(
260
response: OpaqueResponse,
261
chunk: Uint8Array,
93
- offset: number,
262
): void {
263
if (!supportsBinaryStreams) {
264
throw new Error("This environment don't support binary chunks.");
265
}
98
- response.partialRow += readPartialStringChunk(response.stringDecoder, chunk);
266
+ let stringDecoder = response.stringDecoder;
267
+ let linebreak = chunk.indexOf(10); // newline
268
+ while (linebreak > -1) {
269
+ let fullrow =
270
+ response.partialRow +
271
+ readFinalStringChunk(stringDecoder, chunk.subarray(0, linebreak));
272
+ processFullRow(response, fullrow);
273
+ response.partialRow = '';
274
+ chunk = chunk.subarray(linebreak + 1);
275
+ linebreak = chunk.indexOf(10); // newline
276
+ }
277
+ response.partialRow += readPartialStringChunk(stringDecoder, chunk);
278
}
279
101
-let emptyBuffer = new Uint8Array(0);
280
export function complete(response: OpaqueResponse): void {
103
- if (supportsBinaryStreams) {
104
- // This should never be needed since we're expected to have complete
105
- // code units at the end of JSON.
106
- response.partialRow += readFinalStringChunk(
107
- response.stringDecoder,
108
- emptyBuffer,
109
- );
110
- }
111
- let modelRoot = response.modelRoot;
112
- let model = JSON.parse(response.partialRow);
113
- Object.defineProperty(modelRoot, 'model', {
114
- value: model,
115
- });
116
- response.rootPing();
281
+ // In case there are any remaining unresolved chunks, they won't
282
+ // be resolved now. So we need to issue an error to those.
283
+ // Ideally we should be able to early bail out if we kept a
284
+ // ref count of pending chunks.
285
+ reportGlobalError(response, new Error('Connection closed.'));
286
}
287
288
export function getModelRoot<T>(response: OpaqueResponse): ReactModelRoot<T> {
packages/react-server/src/ReactFizzStreamer.js
+1
-4
@@ -76,10 +76,7 @@ export function startWork(request: OpaqueRequest): void {
76
scheduleWork(() => performWork(request));
77
}
78
79
-export function startFlowing(
80
- request: OpaqueRequest,
81
- desiredBytes: number,
82
-): void {
79
+export function startFlowing(request: OpaqueRequest): void {
80
request.flowing = false;
81
flushCompletedChunks(request);
82
}
packages/react-server/src/ReactFlightServer.js
+269
-35
@@ -21,6 +21,63 @@ import {
21
import {renderHostChildrenToString} from './ReactServerFormatConfig';
22
import {REACT_ELEMENT_TYPE} from 'shared/ReactSymbols';
23
24
+/*
25
+
26
+FLIGHT PROTOCOL GRAMMAR
27
+
28
+Response
29
+- JSONData RowSequence
30
+- JSONData
31
+
32
+RowSequence
33
+- Row RowSequence
34
+- Row
35
+
36
+Row
37
+- "J" RowID JSONData
38
+- "H" RowID HTMLData
39
+- "B" RowID BlobData
40
+- "U" RowID URLData
41
+- "E" RowID ErrorData
42
+
43
+RowID
44
+- HexDigits ":"
45
+
46
+HexDigits
47
+- HexDigit HexDigits
48
+- HexDigit
49
+
50
+HexDigit
51
+- 0-F
52
+
53
+URLData
54
+- (UTF8 encoded URL) "\n"
55
+
56
+ErrorData
57
+- (UTF8 encoded JSON: {message: "...", stack: "..."}) "\n"
58
+
59
+JSONData
60
+- (UTF8 encoded JSON) "\n"
61
+ - String values that begin with $ are escaped with a "$" prefix.
62
+ - References to other rows are encoding as JSONReference strings.
63
+
64
+JSONReference
65
+- "$" HexDigits
66
+
67
+HTMLData
68
+- ByteSize (UTF8 encoded HTML)
69
+
70
+BlobData
71
+- ByteSize (Binary Data)
72
+
73
+ByteSize
74
+- (unsigned 32-bit integer)
75
+*/
76
+
77
+// TODO: Implement HTMLData, BlobData and URLData.
78
+
79
+const stringify = JSON.stringify;
80
+
81
export type ReactModel =
82
| React$Element<any>
83
| string
@@ -42,66 +99,246 @@ type ReactModelObject = {
99
+[key: string]: ReactModel,
100
};
101
102
+type Segment = {
103
+ id: number,
104
+ model: ReactModel,
105
+ ping: () => void,
106
+};
107
+
108
type OpaqueRequest = {
109
destination: Destination,
47
- model: ReactModel,
48
- completedChunks: Array<Uint8Array>,
110
+ nextChunkId: number,
111
+ pendingChunks: number,
112
+ pingedSegments: Array<Segment>,
113
+ completedJSONChunks: Array<Uint8Array>,
114
+ completedErrorChunks: Array<Uint8Array>,
115
flowing: boolean,
116
+ toJSON: (key: string, value: ReactModel) => ReactJSONValue,
117
};
118
119
export function createRequest(
120
model: ReactModel,
121
destination: Destination,
122
): OpaqueRequest {
56
- return {destination, model, completedChunks: [], flowing: false};
123
+ let pingedSegments = [];
124
+ let request = {
125
+ destination,
126
+ nextChunkId: 0,
127
+ pendingChunks: 0,
128
+ pingedSegments: pingedSegments,
129
+ completedJSONChunks: [],
130
+ completedErrorChunks: [],
131
+ flowing: false,
132
+ toJSON: (key: string, value: ReactModel) =>
133
+ resolveModelToJSON(request, value),
134
+ };
135
+ request.pendingChunks++;
136
+ let rootSegment = createSegment(request, model);
137
+ pingedSegments.push(rootSegment);
138
+ return request;
139
}
140
59
-function resolveModelToJSON(key: string, value: ReactModel): ReactJSONValue {
60
- while (value && value.$$typeof === REACT_ELEMENT_TYPE) {
141
+function attemptResolveModelComponent(element: React$Element<any>): ReactModel {
142
+ let type = element.type;
143
+ let props = element.props;
144
+ if (typeof type === 'function') {
145
+ // This is a nested view model.
146
+ return type(props);
147
+ } else if (typeof type === 'string') {
148
+ // This is a host element. E.g. HTML.
149
+ return renderHostChildrenToString(element);
150
+ } else {
151
+ throw new Error('Unsupported type.');
152
+ }
153
+}
154
+
155
+function pingSegment(request: OpaqueRequest, segment: Segment): void {
156
+ let pingedSegments = request.pingedSegments;
157
+ pingedSegments.push(segment);
158
+ if (pingedSegments.length === 1) {
159
+ scheduleWork(() => performWork(request));
160
+ }
161
+}
162
+
163
+function createSegment(request: OpaqueRequest, model: ReactModel): Segment {
164
+ let id = request.nextChunkId++;
165
+ let segment = {
166
+ id,
167
+ model,
168
+ ping: () => pingSegment(request, segment),
169
+ };
170
+ return segment;
171
+}
172
+
173
+function serializeIDRef(id: number): string {
174
+ return '$' + id.toString(16);
175
+}
176
+
177
+function serializeRowHeader(tag: string, id: number) {
178
+ return tag + id.toString(16) + ':';
179
+}
180
+
181
+function escapeStringValue(value: string): string {
182
+ if (value[0] === '$') {
183
+ // We need to escape $ prefixed strings since we use that to encode
184
+ // references to IDs.
185
+ return '$' + value;
186
+ } else {
187
+ return value;
188
+ }
189
+}
190
+
191
+function resolveModelToJSON(
192
+ request: OpaqueRequest,
193
+ value: ReactModel,
194
+): ReactJSONValue {
195
+ if (typeof value === 'string') {
196
+ return escapeStringValue(value);
197
+ }
198
+
199
+ while (
200
+ typeof value === 'object' &&
201
+ value !== null &&
202
+ value.$$typeof === REACT_ELEMENT_TYPE
203
+ ) {
204
let element: React$Element<any> = (value: any);
62
- let type = element.type;
63
- let props = element.props;
64
- if (typeof type === 'function') {
65
- // This is a nested view model.
66
- value = type(props);
67
- continue;
68
- } else if (typeof type === 'string') {
69
- // This is a host element. E.g. HTML.
70
- return renderHostChildrenToString(element);
71
- } else {
72
- throw new Error('Unsupported type.');
205
+ try {
206
+ value = attemptResolveModelComponent(element);
207
+ } catch (x) {
208
+ if (typeof x === 'object' && x !== null && typeof x.then === 'function') {
209
+ // Something suspended, we'll need to create a new segment and resolve it later.
210
+ request.pendingChunks++;
211
+ let newSegment = createSegment(request, element);
212
+ let ping = newSegment.ping;
213
+ x.then(ping, ping);
214
+ return serializeIDRef(newSegment.id);
215
+ } else {
216
+ request.pendingChunks++;
217
+ let errorId = request.nextChunkId++;
218
+ emitErrorChunk(request, errorId, x);
219
+ return serializeIDRef(errorId);
220
+ }
221
}
222
}
223
+
224
return value;
225
}
226
227
+function emitErrorChunk(
228
+ request: OpaqueRequest,
229
+ id: number,
230
+ error: mixed,
231
+): void {
232
+ // TODO: We should not leak error messages to the client in prod.
233
+ // Give this an error code instead and log on the server.
234
+ // We can serialize the error in DEV as a convenience.
235
+ let message;
236
+ let stack = '';
237
+ try {
238
+ if (error instanceof Error) {
239
+ message = '' + error.message;
240
+ stack = '' + error.stack;
241
+ } else {
242
+ message = 'Error: ' + (error: any);
243
+ }
244
+ } catch (x) {
245
+ message = 'An error occurred but serializing the error message failed.';
246
+ }
247
+ let errorInfo = {message, stack};
248
+ let row = serializeRowHeader('E', id) + stringify(errorInfo) + '\n';
249
+ request.completedErrorChunks.push(convertStringToBuffer(row));
250
+}
251
+
252
+function retrySegment(request: OpaqueRequest, segment: Segment): void {
253
+ let value = segment.model;
254
+ try {
255
+ while (
256
+ typeof value === 'object' &&
257
+ value !== null &&
258
+ value.$$typeof === REACT_ELEMENT_TYPE
259
+ ) {
260
+ // If this is a nested model, there's no need to create another chunk,
261
+ // we can reuse the existing one and try again.
262
+ let element: React$Element<any> = (value: any);
263
+ segment.model = element;
264
+ value = attemptResolveModelComponent(element);
265
+ }
266
+ let json = stringify(value, request.toJSON);
267
+ let row;
268
+ let id = segment.id;
269
+ if (id === 0) {
270
+ row = json + '\n';
271
+ } else {
272
+ row = serializeRowHeader('J', id) + json + '\n';
273
+ }
274
+ request.completedJSONChunks.push(convertStringToBuffer(row));
275
+ } catch (x) {
276
+ if (typeof x === 'object' && x !== null && typeof x.then === 'function') {
277
+ // Something suspended again, let's pick it back up later.
278
+ let ping = segment.ping;
279
+ x.then(ping, ping);
280
+ return;
281
+ } else {
282
+ // This errored, we need to serialize this error to the
283
+ emitErrorChunk(request, segment.id, x);
284
+ }
285
+ }
286
+}
287
+
288
function performWork(request: OpaqueRequest): void {
79
- let rootModel = request.model;
80
- request.model = null;
81
- let json = JSON.stringify(rootModel, resolveModelToJSON);
82
- request.completedChunks.push(convertStringToBuffer(json));
289
+ let pingedSegments = request.pingedSegments;
290
+ request.pingedSegments = [];
291
+ for (let i = 0; i < pingedSegments.length; i++) {
292
+ let segment = pingedSegments[i];
293
+ retrySegment(request, segment);
294
+ }
295
if (request.flowing) {
296
flushCompletedChunks(request);
297
}
86
-
87
- flushBuffered(request.destination);
298
}
299
90
-function flushCompletedChunks(request: OpaqueRequest) {
300
+let reentrant = false;
301
+function flushCompletedChunks(request: OpaqueRequest): void {
302
+ if (reentrant) {
303
+ return;
304
+ }
305
+ reentrant = true;
306
let destination = request.destination;
92
- let chunks = request.completedChunks;
93
- request.completedChunks = [];
94
-
307
beginWriting(destination);
308
try {
97
- for (let i = 0; i < chunks.length; i++) {
98
- let chunk = chunks[i];
99
- writeChunk(destination, chunk);
309
+ let jsonChunks = request.completedJSONChunks;
310
+ let i = 0;
311
+ for (; i < jsonChunks.length; i++) {
312
+ request.pendingChunks--;
313
+ let chunk = jsonChunks[i];
314
+ if (!writeChunk(destination, chunk)) {
315
+ request.flowing = false;
316
+ i++;
317
+ break;
318
+ }
319
+ }
320
+ jsonChunks.splice(0, i);
321
+ let errorChunks = request.completedErrorChunks;
322
+ i = 0;
323
+ for (; i < errorChunks.length; i++) {
324
+ request.pendingChunks--;
325
+ let chunk = errorChunks[i];
326
+ if (!writeChunk(destination, chunk)) {
327
+ request.flowing = false;
328
+ i++;
329
+ break;
330
+ }
331
}
332
+ errorChunks.splice(0, i);
333
} finally {
334
+ reentrant = false;
335
completeWriting(destination);
336
}
104
- close(destination);
337
+ flushBuffered(destination);
338
+ if (request.pendingChunks === 0) {
339
+ // We're done.
340
+ close(destination);
341
+ }
342
}
343
344
export function startWork(request: OpaqueRequest): void {
@@ -109,10 +346,7 @@ export function startWork(request: OpaqueRequest): void {
346
scheduleWork(() => performWork(request));
347
}
348
112
-export function startFlowing(
113
- request: OpaqueRequest,
114
- desiredBytes: number,
115
-): void {
116
- request.flowing = false;
349
+export function startFlowing(request: OpaqueRequest): void {
350
+ request.flowing = true;
351
flushCompletedChunks(request);
352
}
packages/react-server/src/ReactServerHostConfigBrowser.js
+5
-1
@@ -20,8 +20,12 @@ export function flushBuffered(destination: Destination) {
20
21
export function beginWriting(destination: Destination) {}
22
23
-export function writeChunk(destination: Destination, buffer: Uint8Array) {
23
+export function writeChunk(
24
+ destination: Destination,
25
+ buffer: Uint8Array,
26
+): boolean {
27
destination.enqueue(buffer);
28
+ return destination.desiredSize > 0;
29
}
30
31
export function completeWriting(destination: Destination) {}
packages/react-server/src/ReactServerHostConfigNode.js
+5
-2
@@ -40,9 +40,12 @@ export function beginWriting(destination: Destination) {
40
}
41
}
42
43
-export function writeChunk(destination: Destination, buffer: Uint8Array) {
43
+export function writeChunk(
44
+ destination: Destination,
45
+ buffer: Uint8Array,
46
+): boolean {
47
let nodeBuffer = ((buffer: any): Buffer); // close enough
45
- destination.write(nodeBuffer);
48
+ return destination.write(nodeBuffer);
49
}
50
51
export function completeWriting(destination: Destination) {