1
+/**
2
+ * Copyright (c) Meta Platforms, Inc. and affiliates.
3
+ *
4
+ * This source code is licensed under the MIT license found in the
5
+ * LICENSE file in the root directory of this source tree.
6
+ *
7
+ * @flow
8
+ */
9
+
10
+export type Destination = ReadableStreamController;
11
+
12
+export type PrecomputedChunk = Uint8Array;
13
+export opaque type Chunk = Uint8Array;
14
+
15
+export function scheduleWork(callback: () => void) {
16
+ callback();
17
+}
18
+
19
+export function flushBuffered(destination: Destination) {
20
+ // WHATWG Streams do not yet have a way to flush the underlying
21
+ // transform streams. https://github.com/whatwg/streams/issues/960
22
+}
23
+
24
+// For now, we get this from the global scope, but this will likely move to a module.
25
+export const supportsRequestStorage = typeof AsyncLocalStorage === 'function';
26
+export const requestStorage: AsyncLocalStorage<Map<Function, mixed>> =
27
+ supportsRequestStorage ? new AsyncLocalStorage() : (null: any);
28
+
29
+const VIEW_SIZE = 512;
30
+let currentView = null;
31
+let writtenBytes = 0;
32
+
33
+export function beginWriting(destination: Destination) {
34
+ currentView = new Uint8Array(VIEW_SIZE);
35
+ writtenBytes = 0;
36
+}
37
+
38
+export function writeChunk(
39
+ destination: Destination,
40
+ chunk: PrecomputedChunk | Chunk,
41
+): void {
42
+ if (chunk.length === 0) {
43
+ return;
44
+ }
45
+
46
+ if (chunk.length > VIEW_SIZE) {
47
+ if (__DEV__) {
48
+ if (precomputedChunkSet.has(chunk)) {
49
+ console.error(
50
+ 'A large precomputed chunk was passed to writeChunk without being copied.' +
51
+ ' Large chunks get enqueued directly and are not copied. This is incompatible with precomputed chunks because you cannot enqueue the same precomputed chunk twice.' +
52
+ ' Use "cloneChunk" to make a copy of this large precomputed chunk before writing it. This is a bug in React.',
53
+ );
54
+ }
55
+ }
56
+ // this chunk may overflow a single view which implies it was not
57
+ // one that is cached by the streaming renderer. We will enqueu
58
+ // it directly and expect it is not re-used
59
+ if (writtenBytes > 0) {
60
+ destination.enqueue(
61
+ new Uint8Array(
62
+ ((currentView: any): Uint8Array).buffer,
63
+ 0,
64
+ writtenBytes,
65
+ ),
66
+ );
67
+ currentView = new Uint8Array(VIEW_SIZE);
68
+ writtenBytes = 0;
69
+ }
70
+ destination.enqueue(chunk);
71
+ return;
72
+ }
73
+
74
+ let bytesToWrite = chunk;
75
+ const allowableBytes = ((currentView: any): Uint8Array).length - writtenBytes;
76
+ if (allowableBytes < bytesToWrite.length) {
77
+ // this chunk would overflow the current view. We enqueue a full view
78
+ // and start a new view with the remaining chunk
79
+ if (allowableBytes === 0) {
80
+ // the current view is already full, send it
81
+ destination.enqueue(currentView);
82
+ } else {
83
+ // fill up the current view and apply the remaining chunk bytes
84
+ // to a new view.
85
+ ((currentView: any): Uint8Array).set(
86
+ bytesToWrite.subarray(0, allowableBytes),
87
+ writtenBytes,
88
+ );
89
+ // writtenBytes += allowableBytes; // this can be skipped because we are going to immediately reset the view
90
+ destination.enqueue(currentView);
91
+ bytesToWrite = bytesToWrite.subarray(allowableBytes);
92
+ }
93
+ currentView = new Uint8Array(VIEW_SIZE);
94
+ writtenBytes = 0;
95
+ }
96
+ ((currentView: any): Uint8Array).set(bytesToWrite, writtenBytes);
97
+ writtenBytes += bytesToWrite.length;
98
+}
99
+
100
+export function writeChunkAndReturn(
101
+ destination: Destination,
102
+ chunk: PrecomputedChunk | Chunk,
103
+): boolean {
104
+ writeChunk(destination, chunk);
105
+ // in web streams there is no backpressure so we can alwas write more
106
+ return true;
107
+}
108
+
109
+export function completeWriting(destination: Destination) {
110
+ if (currentView && writtenBytes > 0) {
111
+ destination.enqueue(new Uint8Array(currentView.buffer, 0, writtenBytes));
112
+ currentView = null;
113
+ writtenBytes = 0;
114
+ }
115
+}
116
+
117
+export function close(destination: Destination) {
118
+ destination.close();
119
+}
120
+
121
+const textEncoder = new TextEncoder();
122
+
123
+export function stringToChunk(content: string): Chunk {
124
+ return textEncoder.encode(content);
125
+}
126
+
127
+const precomputedChunkSet: Set<Chunk> = __DEV__ ? new Set() : (null: any);
128
+
129
+export function stringToPrecomputedChunk(content: string): PrecomputedChunk {
130
+ const precomputedChunk = textEncoder.encode(content);
131
+
132
+ if (__DEV__) {
133
+ precomputedChunkSet.add(precomputedChunk);
134
+ }
135
+
136
+ return precomputedChunk;
137
+}
138
+
139
+export function clonePrecomputedChunk(
140
+ precomputedChunk: PrecomputedChunk,
141
+): PrecomputedChunk {
142
+ return precomputedChunk.length > VIEW_SIZE
143
+ ? precomputedChunk.slice()
144
+ : precomputedChunk;
145
+}
146
+
147
+export function closeWithError(destination: Destination, error: mixed): void {
148
+ // $FlowFixMe[method-unbinding]
149
+ if (typeof destination.error === 'function') {
150
+ // $FlowFixMe: This is an Error object or the destination accepts other types.
151
+ destination.error(error);
152
+ } else {
153
+ // Earlier implementations doesn't support this method. In that environment you're
154
+ // supposed to throw from a promise returned but we don't return a promise in our
155
+ // approach. We could fork this implementation but this is environment is an edge
156
+ // case to begin with. It's even less common to run this in an older environment.
157
+ // Even then, this is not where errors are supposed to happen and they get reported
158
+ // to a global callback in addition to this anyway. So it's fine just to close this.
159
+ destination.close();
160
+ }
161
+}