Skip to content
File

Blob: src/node/internal/streams_destroy.ts

typescript396 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-redundant-type-constituents, @typescript-eslint/no-unsafe-call */
27 
28import { nextTick } from 'node-internal:internal_process';
29 
30import type { Writable, WritableState } from 'node-internal:streams_writable';
31import type { Readable, ReadableState } from 'node-internal:streams_readable';
32 
33import {
34 AbortError,
35 aggregateTwoErrors,
36 ERR_MULTIPLE_CALLBACK,
37} from 'node-internal:internal_errors';
38import {
39 kIsDestroyed,
40 isDestroyed,
41 isFinished,
42 isServerRequest,
43 kState,
44 kErrorEmitted,
45 kEmitClose,
46 kClosed,
47 kCloseEmitted,
48 kConstructed,
49 kDestroyed,
50 kAutoDestroy,
51 kErrored,
52} from 'node-internal:streams_util';
53import type { OutgoingMessage } from 'node-internal:internal_http_outgoing';
54 
55const kConstruct = Symbol('kConstruct');
56const kDestroy = Symbol('kDestroy');
57 
58function checkError(
59 err?: Error | null,
60 w?: WritableState,
61 r?: ReadableState
62): void {
63 if (err) {
64 // Avoid V8 leak, https://github.com/nodejs/node/pull/34103#issuecomment-652002364
65 err.stack; // eslint-disable-line @typescript-eslint/no-unused-expressions
66 
67 if (w && !w.errored) {
68 w.errored = err;
69 }
70 if (r && !r.errored) {
71 r.errored = err;
72 }
73 }
74}
75 
76// Backwards compat. cb() is undocumented and unused in core but
77// unfortunately might be used by modules.
78export function destroy(
79 this: Readable | Writable,
80 err?: Error,
81 cb?: VoidFunction
82 // @ts-expect-error TS2526 Returning this is not allowed.
83): this {
84 const r = this._readableState;
85 const w = this._writableState;
86 // With duplex streams we use the writable side for state.
87 const s = w || r;
88 if (
89 (w && (w[kState] & kDestroyed) !== 0) ||
90 (r && (r[kState] & kDestroyed) !== 0)
91 ) {
92 if (typeof cb === 'function') {
93 cb();
94 }
95 return this;
96 }
97 
98 // We set destroyed to true before firing error callbacks in order
99 // to make it re-entrance safe in case destroy() is called within callbacks
100 checkError(err, w, r);
101 if (w) {
102 w[kState] |= kDestroyed;
103 }
104 if (r) {
105 r[kState] |= kDestroyed;
106 }
107 
108 // If still constructing then defer calling _destroy.
109 // @ts-expect-error TS18048 `s` will always be defined here.
110 if ((s[kState] & kConstructed) === 0) {
111 this.once(kDestroy, function (this: Readable | Writable, er: Error) {
112 _destroy(this, aggregateTwoErrors(er, err), cb);
113 });
114 } else {
115 _destroy(this, err, cb);
116 }
117 return this;
118}
119 
120function _destroy(
121 self: Readable | Writable,
122 err?: Error,
123 cb?: (err?: Error | null) => void
124): void {
125 let called = false;
126 
127 function onDestroy(err?: Error | null): void {
128 if (called) {
129 return;
130 }
131 called = true;
132 
133 const r = self._readableState;
134 const w = self._writableState;
135 checkError(err, w, r);
136 if (w) {
137 w[kState] |= kClosed;
138 }
139 if (r) {
140 r[kState] |= kClosed;
141 }
142 
143 if (typeof cb === 'function') {
144 cb(err);
145 }
146 
147 if (err) {
148 nextTick(emitErrorCloseNT, self, err);
149 } else {
150 nextTick(emitCloseNT, self);
151 }
152 }
153 try {
154 self._destroy(err || null, onDestroy);
155 } catch (err) {
156 onDestroy(err as Error);
157 }
158}
159 
160function emitErrorCloseNT(self: Readable | Writable, err: Error): void {
161 emitErrorNT(self, err);
162 emitCloseNT(self);
163}
164 
165function emitCloseNT(self: Readable | Writable): void {
166 const r = self._readableState;
167 const w = self._writableState;
168 if (w) {
169 w[kState] |= kCloseEmitted;
170 }
171 if (r) {
172 r[kState] |= kCloseEmitted;
173 }
174 if (
175 (w && (w[kState] & kEmitClose) !== 0) ||
176 (r && (r[kState] & kEmitClose) !== 0)
177 ) {
178 self.emit('close');
179 }
180}
181 
182function emitErrorNT(self: Readable | Writable, err: Error): void {
183 const r = self._readableState;
184 const w = self._writableState;
185 if (
186 (w && (w[kState] & kErrorEmitted) !== 0) ||
187 (r && (r[kState] & kErrorEmitted) !== 0)
188 ) {
189 return;
190 }
191 
192 if (w) {
193 w[kState] |= kErrorEmitted;
194 }
195 if (r) {
196 r[kState] |= kErrorEmitted;
197 }
198 self.emit('error', err);
199}
200 
201export function undestroy(this: Readable | Writable): void {
202 const r = this._readableState;
203 const w = this._writableState;
204 if (r) {
205 r.constructed = true;
206 r.closed = false;
207 r.closeEmitted = false;
208 r.destroyed = false;
209 r.errored = null;
210 r.errorEmitted = false;
211 r.reading = false;
212 r.ended = r.readable === false;
213 r.endEmitted = r.readable === false;
214 }
215 if (w) {
216 w.constructed = true;
217 w.destroyed = false;
218 w.closed = false;
219 w.closeEmitted = false;
220 w.errored = null;
221 w.errorEmitted = false;
222 w.finalCalled = false;
223 w.prefinished = false;
224 w.ended = w.writable === false;
225 w.ending = w.writable === false;
226 w.finished = w.writable === false;
227 }
228}
229 
230export function errorOrDestroy(
231 stream: Readable | Writable,
232 err?: Error,
233 sync: boolean = false
234 // @ts-expect-error TS2526 Apparently `this` is disallowed.
235): this | undefined {
236 // We have tests that rely on errors being emitted
237 // in the same tick, so changing this is semver major.
238 // For now when you opt-in to autoDestroy we allow
239 // the error to be emitted nextTick. In a future
240 // semver major update we should change the default to this.
241 
242 const r = stream._readableState;
243 const w = stream._writableState;
244 if (
245 (w && (w[kState] ? (w[kState] & kDestroyed) !== 0 : w.destroyed)) ||
246 (r && (r[kState] ? (r[kState] & kDestroyed) !== 0 : r.destroyed))
247 ) {
248 // @ts-expect-error TS2683 This should be somehow type-defined.
249 return this;
250 }
251 if (
252 (r && (r[kState] & kAutoDestroy) !== 0) ||
253 (w && (w[kState] & kAutoDestroy) !== 0)
254 ) {
255 stream.destroy(err);
256 } else if (err) {
257 // Avoid V8 leak, https://github.com/nodejs/node/pull/34103#issuecomment-652002364
258 err.stack; // eslint-disable-line @typescript-eslint/no-unused-expressions
259 
260 if (w && (w[kState] & kErrored) === 0) {
261 w.errored = err;
262 }
263 if (r && (r[kState] & kErrored) === 0) {
264 r.errored = err;
265 }
266 if (sync) {
267 nextTick(emitErrorNT, stream, err);
268 } else {
269 emitErrorNT(stream, err);
270 }
271 }
272 
273 return undefined;
274}
275 
276export function construct(stream: Readable | Writable, cb: VoidFunction): void {
277 if (typeof stream._construct !== 'function') {
278 return;
279 }
280 const r = stream._readableState;
281 const w = stream._writableState;
282 
283 if (r) {
284 r[kState] &= ~kConstructed;
285 }
286 if (w) {
287 w[kState] &= ~kConstructed;
288 }
289 
290 stream.once(kConstruct, cb);
291 
292 if (stream.listenerCount(kConstruct) > 1) {
293 // Duplex
294 return;
295 }
296 
297 nextTick(constructNT, stream);
298}
299 
300function constructNT(this: unknown, stream: Readable | Writable): void {
301 let called = false;
302 
303 function onConstruct(err?: Error): void {
304 if (called) {
305 errorOrDestroy(stream, err ?? new ERR_MULTIPLE_CALLBACK());
306 return;
307 }
308 called = true;
309 
310 const r = stream._readableState;
311 const w = stream._writableState;
312 const s = w || r;
313 
314 if (r) {
315 r[kState] |= kConstructed;
316 }
317 if (w) {
318 w[kState] |= kConstructed;
319 }
320 
321 if (s?.destroyed) {
322 stream.emit(kDestroy, err);
323 } else if (err) {
324 errorOrDestroy(stream, err, true);
325 } else {
326 stream.emit(kConstruct);
327 }
328 }
329 
330 try {
331 stream._construct?.((err) => {
332 nextTick(onConstruct, err);
333 });
334 } catch (err) {
335 nextTick(onConstruct, err);
336 }
337}
338 
339function isRequest(stream: unknown): stream is OutgoingMessage {
340 return (
341 stream != null &&
342 typeof stream === 'object' &&
343 'setHeader' in stream &&
344 'abort' in stream &&
345 typeof stream.abort === 'function'
346 );
347}
348 
349function emitCloseLegacy(stream: Readable | Writable): void {
350 stream.emit('close');
351}
352 
353function emitErrorCloseLegacy(stream: Readable | Writable, err?: Error): void {
354 stream.emit('error', err);
355 nextTick(emitCloseLegacy, stream);
356}
357 
358// Normalize destroy for legacy.
359export function destroyer(
360 stream: Readable | Writable | null | undefined,
361 err?: Error
362): void {
363 if (!stream || isDestroyed(stream)) {
364 return;
365 }
366 
367 if (!err && !isFinished(stream)) {
368 err = new AbortError();
369 }
370 
371 // TODO: Remove isRequest branches.
372 if (isServerRequest(stream)) {
373 // @ts-expect-error TS2540 - socket is read-only but we need to set it to null for cleanup
374 stream.socket = null;
375 stream.destroy(err);
376 } else if (isRequest(stream)) {
377 // @ts-expect-error TS2339 - abort exists on OutgoingMessage but not in types
378 stream.abort();
379 } else if (isRequest(stream.req)) {
380 // @ts-expect-error TS2339 - abort exists on req but not in all types
381 stream.req.abort();
382 } else if (typeof stream.destroy === 'function') {
383 stream.destroy(err);
384 } else if ('close' in stream && typeof stream.close === 'function') {
385 // TODO: Don't lose err?
386 stream.close();
387 } else if (err) {
388 nextTick(emitErrorCloseLegacy, stream, err);
389 } else {
390 nextTick(emitCloseLegacy, stream);
391 }
392 if (!stream.destroyed) {
393 stream[kIsDestroyed] = true;
394 }
395}