Skip to content
File

Blob: src/node/internal/streams_end_of_stream.ts

typescript392 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// Ported from https://github.com/mafintosh/end-of-stream with
27// permission from the author, Mathias Buus (@mafintosh).
28 
29/* eslint-disable @typescript-eslint/no-explicit-any, @typescript-eslint/no-unnecessary-condition, @typescript-eslint/no-unsafe-call, @typescript-eslint/no-unsafe-member-access */
30 
31import { Readable } from 'node-internal:streams_readable';
32import { Writable } from 'node-internal:streams_writable';
33import type { Transform } from 'node-internal:streams_transform';
34import { nextTick } from 'node-internal:internal_process';
35import type { EventEmitter } from 'node:events';
36import {
37 AbortError,
38 ERR_INVALID_ARG_TYPE,
39 ERR_STREAM_PREMATURE_CLOSE,
40} from 'node-internal:internal_errors';
41import { once } from 'node-internal:internal_http_util';
42import {
43 validateAbortSignal,
44 validateFunction,
45 validateObject,
46 validateBoolean,
47} from 'node-internal:validators';
48 
49import {
50 isClosed,
51 isReadable,
52 isReadableNodeStream,
53 isReadableStream,
54 isReadableFinished,
55 isReadableErrored,
56 isWritable,
57 isWritableNodeStream,
58 isWritableStream,
59 isWritableFinished,
60 isWritableErrored,
61 isNodeStream,
62 willEmitClose as _willEmitClose,
63 kIsClosedPromise,
64} from 'node-internal:streams_util';
65import { addAbortListener } from 'node-internal:events';
66import type Stream from 'node:stream';
67 
68// eslint-disable-next-line @typescript-eslint/no-unnecessary-type-parameters
69function isRequest<T extends EventEmitter>(stream: any): stream is T {
70 return 'setHeader' in stream && typeof stream.abort === 'function';
71}
72 
73export const nop = (): void => {};
74 
75type EOSOptions = {
76 cleanup?: boolean;
77 error?: boolean;
78 readable?: boolean;
79 writable?: boolean;
80 signal?: AbortSignal;
81};
82 
83type Callback = (...args: unknown[]) => void;
84 
85export function eos(
86 stream: Readable | Writable | Transform,
87 options: EOSOptions,
88 callback: Callback
89): Callback;
90export function eos(options: EOSOptions, callback: Callback): Callback;
91export function eos(
92 stream: Readable | Writable | Transform | EOSOptions,
93 options: EOSOptions | Callback,
94 callback?: Callback
95): Callback {
96 if (arguments.length === 2) {
97 // @ts-expect-error TS2322 Supports overloads
98 callback = options;
99 options = {} as EOSOptions;
100 } else if (options == null) {
101 options = {} as EOSOptions;
102 } else {
103 validateObject(options, 'options');
104 }
105 validateFunction(callback, 'callback');
106 validateAbortSignal((options as EOSOptions).signal, 'options.signal');
107 
108 // Avoid AsyncResource.bind() because it calls Object.defineProperties which
109 // is a bottleneck here.
110 callback = once(callback) as Callback;
111 
112 if (isReadableStream(stream) || isWritableStream(stream)) {
113 return eosWeb(stream, options as EOSOptions, callback);
114 }
115 
116 if (!isNodeStream(stream)) {
117 throw new ERR_INVALID_ARG_TYPE(
118 'stream',
119 ['ReadableStream', 'WritableStream', 'Stream'],
120 stream
121 );
122 }
123 
124 const readable =
125 (options as EOSOptions).readable ?? isReadableNodeStream(stream);
126 const writable =
127 (options as EOSOptions).writable ?? isWritableNodeStream(stream);
128 
129 const wState = stream._writableState;
130 const rState = stream._readableState;
131 // Extract req as EventEmitter to avoid union type incompatibility with on/removeListener
132 const req = stream.req as EventEmitter | undefined;
133 
134 const onlegacyfinish = (): void => {
135 if (!stream.writable) {
136 onfinish();
137 }
138 };
139 
140 // TODO (ronag): Improve soft detection to include core modules and
141 // common ecosystem modules that do properly emit 'close' but fail
142 // this generic check.
143 let willEmitClose =
144 _willEmitClose(stream) &&
145 isReadableNodeStream(stream) === readable &&
146 isWritableNodeStream(stream) === writable;
147 
148 let writableFinished = isWritableFinished(stream, false);
149 const onfinish = (): void => {
150 writableFinished = true;
151 // Stream should not be destroyed here. If it is that
152 // means that user space is doing something differently and
153 // we cannot trust willEmitClose.
154 if (stream.destroyed) {
155 willEmitClose = false;
156 }
157 
158 if (willEmitClose && (!stream.readable || readable)) {
159 return;
160 }
161 
162 if (!readable || readableFinished) {
163 callback?.call(stream);
164 }
165 };
166 
167 let readableFinished = isReadableFinished(stream, false);
168 const onend = (): void => {
169 readableFinished = true;
170 // Stream should not be destroyed here. If it is that
171 // means that user space is doing something differently and
172 // we cannot trust willEmitClose.
173 if (stream.destroyed) {
174 willEmitClose = false;
175 }
176 
177 if (willEmitClose && (!stream.writable || writable)) {
178 return;
179 }
180 
181 if (!writable || writableFinished) {
182 callback?.call(stream);
183 }
184 };
185 
186 const onerror = (err: Error): void => {
187 callback?.call(stream, err);
188 };
189 
190 let closed = isClosed(stream);
191 
192 const onclose = (): void => {
193 closed = true;
194 
195 const errored = isWritableErrored(stream) || isReadableErrored(stream);
196 
197 if (errored && typeof errored !== 'boolean') {
198 return callback?.call(stream, errored);
199 }
200 
201 if (readable && !readableFinished && isReadableNodeStream(stream, true)) {
202 if (!isReadableFinished(stream, false))
203 return callback?.call(stream, new ERR_STREAM_PREMATURE_CLOSE());
204 }
205 if (writable && !writableFinished) {
206 if (!isWritableFinished(stream, false))
207 return callback?.call(stream, new ERR_STREAM_PREMATURE_CLOSE());
208 }
209 
210 callback?.call(stream);
211 };
212 
213 const onclosed = (): void => {
214 closed = true;
215 
216 const errored = isWritableErrored(stream) || isReadableErrored(stream);
217 
218 if (errored && typeof errored !== 'boolean') {
219 callback?.call(stream, errored);
220 return;
221 }
222 
223 callback?.call(stream);
224 };
225 
226 const onrequest = (): void => {
227 req?.on('finish', onfinish);
228 };
229 
230 if (isRequest(stream)) {
231 stream.on('complete', onfinish);
232 if (!willEmitClose) {
233 stream.on('abort', onclose);
234 }
235 if (stream.req) {
236 onrequest();
237 } else {
238 stream.on('request', onrequest);
239 }
240 } else if (writable && !wState) {
241 // legacy streams
242 (stream as Stream).on('end', onlegacyfinish).on('close', onlegacyfinish);
243 }
244 
245 // Not all streams will emit 'close' after 'aborted'.
246 if (
247 !willEmitClose &&
248 'aborted' in stream &&
249 typeof stream.aborted === 'boolean'
250 ) {
251 stream.on('aborted', onclose);
252 }
253 
254 stream.on('end', onend);
255 stream.on('finish', onfinish);
256 if ((options as EOSOptions).error !== false) {
257 stream.on('error', onerror);
258 }
259 stream.on('close', onclose);
260 
261 if (closed) {
262 nextTick(onclose);
263 } else if (wState?.errorEmitted || rState?.errorEmitted) {
264 if (!willEmitClose) {
265 nextTick(onclosed);
266 }
267 } else if (
268 !readable &&
269 (!willEmitClose || isReadable(stream)) &&
270 (writableFinished || isWritable(stream) === false) &&
271 (wState == null || wState.pendingcb === undefined || wState.pendingcb === 0)
272 ) {
273 nextTick(onclosed);
274 } else if (
275 !writable &&
276 (!willEmitClose || isWritable(stream)) &&
277 (readableFinished || isReadable(stream) === false)
278 ) {
279 nextTick(onclosed);
280 } else if (rState && stream.req && stream.aborted) {
281 nextTick(onclosed);
282 }
283 
284 const cleanup = (): void => {
285 callback = nop;
286 stream.removeListener('aborted', onclose);
287 stream.removeListener('complete', onfinish);
288 stream.removeListener('abort', onclose);
289 stream.removeListener('request', onrequest);
290 req?.removeListener('finish', onfinish);
291 stream.removeListener('end', onlegacyfinish);
292 stream.removeListener('close', onlegacyfinish);
293 stream.removeListener('finish', onfinish);
294 stream.removeListener('end', onend);
295 stream.removeListener('error', onerror);
296 stream.removeListener('close', onclose);
297 };
298 
299 if ((options as EOSOptions).signal && !closed) {
300 const abort = (): void => {
301 // Keep it because cleanup removes it.
302 const endCallback = callback;
303 cleanup();
304 endCallback?.call(
305 stream,
306 new AbortError(undefined, {
307 cause: (options as EOSOptions).signal?.reason,
308 })
309 );
310 };
311 if ((options as EOSOptions).signal?.aborted) {
312 nextTick(abort);
313 } else {
314 const disposable = addAbortListener(
315 (options as EOSOptions).signal,
316 abort
317 );
318 const originalCallback = callback;
319 callback = once((...args: unknown[]): void => {
320 disposable[Symbol.dispose]();
321 originalCallback.apply(stream, args);
322 });
323 }
324 }
325 
326 return cleanup;
327}
328 
329function eosWeb(
330 stream: ReadableStream | WritableStream,
331 options: { signal?: AbortSignal },
332 callback: (...args: unknown[]) => void
333): () => void {
334 let isAborted = false;
335 let abort = nop;
336 if (options.signal) {
337 abort = (): void => {
338 isAborted = true;
339 callback.call(
340 stream,
341 new AbortError(undefined, { cause: options.signal?.reason })
342 );
343 };
344 if (options.signal.aborted) {
345 nextTick(abort);
346 } else {
347 const disposable = addAbortListener(options.signal, abort);
348 const originalCallback = callback;
349 callback = once((...args: unknown[]): void => {
350 disposable[Symbol.dispose]();
351 originalCallback.apply(stream, args);
352 });
353 }
354 }
355 const resolverFn = (...args: unknown[]): void => {
356 if (!isAborted) {
357 nextTick(() => {
358 callback.apply(stream, args);
359 });
360 }
361 };
362 // @ts-expect-error TS7053 Symbols are not defined in types yet.
363 stream[kIsClosedPromise].promise.then(resolverFn, resolverFn);
364 return nop;
365}
366 
367export function finished(
368 stream: Readable | Writable,
369 opts: EOSOptions = {}
370): Promise<void> {
371 let autoCleanup = false;
372 if (opts.cleanup) {
373 validateBoolean(opts.cleanup, 'cleanup');
374 autoCleanup = opts.cleanup;
375 }
376 return new Promise<void>((resolve, reject) => {
377 const cleanup = eos(stream, opts, (err: unknown) => {
378 if (autoCleanup) {
379 cleanup();
380 }
381 if (err) {
382 // eslint-disable-next-line @typescript-eslint/prefer-promise-reject-errors
383 reject(err);
384 } else {
385 resolve();
386 }
387 });
388 });
389}
390 
391eos.finished = finished;