Skip to content
File

Blob: src/node/internal/streams_util.ts

typescript413 lines
1// Copyright (c) 2017-2022 Cloudflare, Inc.
2// Licensed under the Apache 2.0 license found in the LICENSE file or at:
3// https://opensource.org/licenses/Apache-2.0
4//
5// Copyright Joyent, Inc. and other Node contributors.
6//
7// Permission is hereby granted, free of charge, to any person obtaining a
8// copy of this software and associated documentation files (the
9// "Software"), to deal in the Software without restriction, including
10// without limitation the rights to use, copy, modify, merge, publish,
11// distribute, sublicense, and/or sell copies of the Software, and to permit
12// persons to whom the Software is furnished to do so, subject to the
13// following conditions:
14//
15// The above copyright notice and this permission notice shall be included
16// in all copies or substantial portions of the Software.
17//
18// THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS
19// OR IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF
20// MERCHANTABILITY, FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN
21// NO EVENT SHALL THE AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM,
22// DAMAGES OR OTHER LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR
23// OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE
24// USE OR OTHER DEALINGS IN THE SOFTWARE.
25 
26/* eslint-disable @typescript-eslint/no-explicit-any, @typescript-eslint/no-unsafe-member-access, @typescript-eslint/no-unnecessary-condition */
27 
28import { AbortError } from 'node-internal:internal_errors';
29import { constants } from 'node-internal:internal_zlib_constants';
30 
31import type { Writable } from 'node-internal:streams_writable';
32import type { Readable } from 'node-internal:streams_readable';
33import type { Transform } from 'node-internal:streams_transform';
34import type { OutgoingMessage } from 'node-internal:internal_http_outgoing';
35import type { ServerResponse } from 'node-internal:internal_http_server';
36import type { IncomingMessage } from 'node-internal:internal_http_incoming';
37 
38// We need to use Symbol.for to make these globally available
39// for interoperability with readable-stream, i.e. readable-stream
40// and node core needs to be able to read/write private state
41// from each other for proper interoperability.
42export const kIsDestroyed = Symbol.for('nodejs.stream.destroyed');
43export const kIsErrored = Symbol.for('nodejs.stream.errored');
44export const kIsReadable = Symbol.for('nodejs.stream.readable');
45export const kIsWritable = Symbol.for('nodejs.stream.writable');
46export const kIsDisturbed = Symbol.for('nodejs.stream.disturbed');
47 
48export const kOnConstructed = Symbol('kOnConstructed');
49 
50export const kIsClosedPromise = Symbol.for('nodejs.webstream.isClosedPromise');
51export const kControllerErrorFunction = Symbol.for(
52 'nodejs.webstream.controllerErrorFunction'
53);
54 
55export const kState = Symbol('kState');
56export const kObjectMode = 1 << 0;
57export const kErrorEmitted = 1 << 1;
58export const kAutoDestroy = 1 << 2;
59export const kEmitClose = 1 << 3;
60export const kDestroyed = 1 << 4;
61export const kClosed = 1 << 5;
62export const kCloseEmitted = 1 << 6;
63export const kErrored = 1 << 7;
64export const kConstructed = 1 << 8;
65 
66export function isReadableNodeStream(
67 obj: any,
68 strict: boolean = false
69): boolean {
70 return !!(
71 obj &&
72 typeof obj.pipe === 'function' &&
73 typeof obj.on === 'function' &&
74 (!strict ||
75 (typeof obj.pause === 'function' && typeof obj.resume === 'function')) &&
76 (!obj._writableState || obj._readableState?.readable !== false) && // Duplex
77 (!obj._writableState || obj._readableState) // Writable has .pipe.
78 );
79}
80 
81export function isWritableNodeStream(obj: any): boolean {
82 return !!(
83 obj &&
84 typeof obj.write === 'function' &&
85 typeof obj.on === 'function' &&
86 (!obj._readableState || obj._writableState?.writable !== false) // Duplex
87 );
88}
89 
90export function isDuplexNodeStream(obj: any): boolean {
91 return !!(
92 obj &&
93 typeof obj.pipe === 'function' &&
94 obj._readableState &&
95 typeof obj.on === 'function' &&
96 typeof obj.write === 'function'
97 );
98}
99 
100export function isNodeStream(obj: any): obj is Readable | Writable | Transform {
101 // eslint-disable-next-line @typescript-eslint/no-unsafe-return
102 return (
103 obj &&
104 (obj._readableState != null ||
105 obj._writableState != null ||
106 (typeof obj.write === 'function' && typeof obj.on === 'function') ||
107 (typeof obj.pipe === 'function' && typeof obj.on === 'function'))
108 );
109}
110 
111export function isReadableStream(obj: any): obj is ReadableStream {
112 return !!(
113 obj &&
114 !isNodeStream(obj) &&
115 typeof obj.pipeThrough === 'function' &&
116 typeof obj.getReader === 'function' &&
117 typeof obj.cancel === 'function'
118 );
119}
120 
121export function isWritableStream(obj: any): obj is WritableStream {
122 return !!(
123 obj &&
124 !isNodeStream(obj) &&
125 typeof obj.getWriter === 'function' &&
126 typeof obj.abort === 'function'
127 );
128}
129 
130export function isTransformStream(obj: any): obj is TransformStream {
131 return !!(
132 obj &&
133 !isNodeStream(obj) &&
134 typeof obj.readable === 'object' &&
135 typeof obj.writable === 'object'
136 );
137}
138 
139export function isWebStream(
140 obj: any
141): obj is ReadableStream | WritableStream | TransformStream {
142 return (
143 isReadableStream(obj) || isWritableStream(obj) || isTransformStream(obj)
144 );
145}
146 
147export function isIterable(obj: any, isAsync?: true | false): boolean {
148 if (obj == null) return false;
149 if (isAsync === true) return typeof obj[Symbol.asyncIterator] === 'function';
150 if (isAsync === false) return typeof obj[Symbol.iterator] === 'function';
151 return (
152 typeof obj[Symbol.asyncIterator] === 'function' ||
153 typeof obj[Symbol.iterator] === 'function'
154 );
155}
156 
157export function isDestroyed(stream: any): boolean | null {
158 if (!isNodeStream(stream)) return null;
159 const wState = stream._writableState;
160 const rState = stream._readableState;
161 const state = wState || rState;
162 return !!(stream.destroyed || stream[kIsDestroyed] || state?.destroyed);
163}
164 
165// Have been end():d.
166export function isWritableEnded(
167 stream: Writable | Readable | Transform
168): boolean | null {
169 if (!isWritableNodeStream(stream)) return null;
170 if (stream.writableEnded === true) return true;
171 const wState = stream._writableState;
172 if (wState?.errored) return false;
173 if (typeof wState?.ended !== 'boolean') return null;
174 return wState.ended;
175}
176 
177// Have emitted 'finish'.
178export function isWritableFinished(
179 stream: Writable | Readable | Transform,
180 strict?: true | false | null
181): boolean | null {
182 if (!isWritableNodeStream(stream)) return null;
183 if (stream.writableFinished === true) return true;
184 const wState = stream._writableState;
185 if (wState?.errored) return false;
186 if (typeof wState?.finished !== 'boolean') return null;
187 // eslint-disable-next-line @typescript-eslint/no-unnecessary-type-conversion
188 return !!(
189 wState.finished ||
190 (strict === false && wState.ended === true && wState.length === 0)
191 );
192}
193 
194// Have been push(null):d.
195export function isReadableEnded(
196 stream: Readable | Writable | Transform
197): boolean | null {
198 if (!isReadableNodeStream(stream)) return null;
199 if (stream.readableEnded === true) return true;
200 const rState = stream._readableState;
201 if (!rState || rState.errored) return false;
202 if (typeof rState?.ended !== 'boolean') return null;
203 return rState.ended;
204}
205 
206// Have emitted 'end'.
207export function isReadableFinished(
208 stream: Readable | Writable | Transform,
209 strict?: boolean
210): boolean | null {
211 if (!isReadableNodeStream(stream)) return null;
212 const rState = stream._readableState;
213 if (rState?.errored) return false;
214 if (typeof rState?.endEmitted !== 'boolean') return null;
215 // eslint-disable-next-line @typescript-eslint/no-unnecessary-type-conversion
216 return !!(
217 rState.endEmitted ||
218 (strict === false && rState.ended === true && rState.length === 0)
219 );
220}
221 
222export function isReadable(
223 stream: Readable | Writable | Transform
224): boolean | null {
225 if (stream && stream[kIsReadable] != null) return stream[kIsReadable];
226 if (typeof stream?.readable !== 'boolean') return null;
227 if (isDestroyed(stream)) return false;
228 return (
229 isReadableNodeStream(stream) &&
230 stream.readable &&
231 !isReadableFinished(stream)
232 );
233}
234 
235export function isWritable(
236 stream: Readable | Writable | Transform
237): boolean | null {
238 if (stream && stream[kIsWritable] != null) return stream[kIsWritable];
239 if (typeof stream?.writable !== 'boolean') return null;
240 if (isDestroyed(stream)) return false;
241 return (
242 isWritableNodeStream(stream) && stream.writable && !isWritableEnded(stream)
243 );
244}
245 
246export function isFinished(
247 stream: Readable | Writable | Transform,
248 opts?: { readable?: boolean; writable?: boolean }
249): boolean | null {
250 if (!isNodeStream(stream)) {
251 return null;
252 }
253 
254 if (isDestroyed(stream)) {
255 return true;
256 }
257 
258 if (opts?.readable !== false && isReadable(stream)) {
259 return false;
260 }
261 
262 if (opts?.writable !== false && isWritable(stream)) {
263 return false;
264 }
265 
266 return true;
267}
268 
269export function isWritableErrored(
270 stream: Writable | Readable | Transform
271): Error | boolean | null {
272 if (!isNodeStream(stream)) {
273 return null;
274 }
275 
276 if (stream.writableErrored) {
277 return stream.writableErrored;
278 }
279 
280 return stream._writableState?.errored ?? null;
281}
282 
283export function isReadableErrored(
284 stream: Readable | Writable | Transform
285): Error | boolean | null {
286 if (!isNodeStream(stream)) {
287 return null;
288 }
289 
290 if (stream.readableErrored) {
291 return stream.readableErrored;
292 }
293 
294 return stream._readableState?.errored ?? null;
295}
296 
297export function isClosed(
298 stream: Readable | Writable | Transform
299): boolean | null {
300 if (!isNodeStream(stream)) {
301 return null;
302 }
303 
304 if (typeof stream.closed === 'boolean') {
305 return stream.closed;
306 }
307 
308 const wState = stream._writableState;
309 const rState = stream._readableState;
310 
311 if (
312 typeof wState?.closed === 'boolean' ||
313 typeof rState?.closed === 'boolean'
314 ) {
315 return (wState?.closed || rState?.closed) ?? null;
316 }
317 
318 if (typeof stream._closed === 'boolean' && isOutgoingMessage(stream)) {
319 return stream._closed;
320 }
321 
322 return null;
323}
324 
325export function isOutgoingMessage(stream: unknown): stream is OutgoingMessage {
326 return (
327 stream != null &&
328 typeof stream === 'object' &&
329 '_closed' in stream &&
330 typeof stream._closed === 'boolean' &&
331 '_defaultKeepAlive' in stream &&
332 typeof stream._defaultKeepAlive === 'boolean' &&
333 '_removedConnection' in stream &&
334 typeof stream._removedConnection === 'boolean' &&
335 '_removedContLen' in stream &&
336 typeof stream._removedContLen === 'boolean'
337 );
338}
339 
340// This function includes the following check that we don't include, because
341// our ServerResponse implementation does not implement it.
342// `typeof stream._sent100 === 'boolean'`
343export function isServerResponse(stream: unknown): stream is ServerResponse {
344 return isOutgoingMessage(stream);
345}
346 
347export function isServerRequest(stream: any): stream is IncomingMessage {
348 return (
349 typeof stream._consuming === 'boolean' &&
350 typeof stream._dumped === 'boolean' &&
351 stream.req?.upgradeOrConnect === undefined
352 );
353}
354 
355export function willEmitClose(stream: any): boolean | null {
356 if (!isNodeStream(stream)) return null;
357 
358 const wState = stream._writableState;
359 const rState = stream._readableState;
360 const state = wState || rState;
361 
362 return (
363 (!state && isServerResponse(stream)) ||
364 !!(state?.autoDestroy && state.emitClose && state.closed === false)
365 );
366}
367 
368export function isDisturbed(stream: any): boolean {
369 return !!(
370 stream &&
371 (stream[kIsDisturbed] ?? (stream.readableDidRead || stream.readableAborted))
372 );
373}
374 
375export function isErrored(stream: any): boolean {
376 return !!(
377 stream &&
378 (stream[kIsErrored] ??
379 stream.readableErrored ??
380 stream.writableErrored ??
381 stream._readableState?.errorEmitted ??
382 stream._writableState?.errorEmitted ??
383 stream._readableState?.errored ??
384 stream._writableState?.errored)
385 );
386}
387 
388const ZLIB_FAILURES = new Set<string>([
389 ...Object.entries(constants).flatMap(([code, value]) =>
390 value < 0 ? code : []
391 ),
392 'Z_NEED_DICT',
393]);
394 
395export function handleKnownInternalErrors(
396 cause?: Error & { code?: string }
397): (Error & { code?: string }) | undefined {
398 switch (true) {
399 case cause?.code === 'ERR_STREAM_PREMATURE_CLOSE': {
400 return new AbortError(undefined, { cause });
401 }
402 case ZLIB_FAILURES.has(cause?.code ?? ''): {
403 const error = new TypeError(undefined, { cause }) as Error & {
404 code?: string;
405 };
406 error.code = cause?.code as string;
407 return error;
408 }
409 default:
410 return cause;
411 }
412}