Skip to content
File

Blob: src/node/internal/streams_writable.js

javascript1519 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 
26import { EventEmitter } from 'node-internal:events';
27import * as destroyImpl from 'node-internal:streams_destroy';
28import { Stream } from 'node-internal:streams_legacy';
29import { Buffer } from 'node-internal:internal_buffer';
30import { nextTick } from 'node-internal:internal_process';
31import { normalizeEncoding } from 'node-internal:internal_utils';
32import { validateBoolean, validateObject } from 'node-internal:validators';
33import {
34 kState,
35 // bitfields
36 kObjectMode,
37 kErrorEmitted,
38 kAutoDestroy,
39 kEmitClose,
40 kDestroyed,
41 kClosed,
42 kCloseEmitted,
43 kErrored,
44 kConstructed,
45 kOnConstructed,
46 isDestroyed,
47 isWritable,
48 isWritableEnded,
49 isWritableStream,
50 handleKnownInternalErrors,
51} from 'node-internal:streams_util';
52import { finished, eos, nop } from 'node-internal:streams_end_of_stream';
53import { addAbortSignal } from 'node-internal:streams_add_abort_signal';
54import {
55 getHighWaterMark,
56 getDefaultHighWaterMark,
57} from 'node-internal:streams_state';
58 
59import {
60 AbortError,
61 ERR_INVALID_ARG_TYPE,
62 ERR_METHOD_NOT_IMPLEMENTED,
63 ERR_MULTIPLE_CALLBACK,
64 ERR_STREAM_CANNOT_PIPE,
65 ERR_STREAM_DESTROYED,
66 ERR_STREAM_ALREADY_FINISHED,
67 ERR_STREAM_NULL_VALUES,
68 ERR_STREAM_WRITE_AFTER_END,
69 ERR_UNKNOWN_ENCODING,
70 ERR_STREAM_PREMATURE_CLOSE,
71} from 'node-internal:internal_errors';
72 
73const streamsNodejsV24Compat =
74 Cloudflare.compatibilityFlags.enable_streams_nodejs_v24_compat;
75 
76const encoder = new globalThis.TextEncoder();
77 
78const kOnFinishedValue = Symbol('kOnFinishedValue');
79const kErroredValue = Symbol('kErroredValue');
80const kDefaultEncodingValue = Symbol('kDefaultEncodingValue');
81const kWriteCbValue = Symbol('kWriteCbValue');
82const kAfterWriteTickInfoValue = Symbol('kAfterWriteTickInfoValue');
83const kBufferedValue = Symbol('kBufferedValue');
84 
85// Bitfield flag constants for WritableState. Each constant uses left-shift (<<) to set a specific
86// bit position, allowing multiple boolean flags to be stored efficiently in a single integer (kState).
87// For example, `1 << 9` creates a value with only bit 9 set (value: 512).
88const kSync = 1 << 9;
89const kFinalCalled = 1 << 10;
90const kNeedDrain = 1 << 11;
91const kEnding = 1 << 12;
92const kFinished = 1 << 13;
93const kDecodeStrings = 1 << 14;
94const kWriting = 1 << 15;
95const kBufferProcessing = 1 << 16;
96const kPrefinished = 1 << 17;
97const kAllBuffers = 1 << 18;
98const kAllNoop = 1 << 19;
99const kOnFinished = 1 << 20;
100const kHasWritable = 1 << 21;
101const kWritable = 1 << 22;
102const kCorked = 1 << 23;
103const kDefaultUTF8Encoding = 1 << 24;
104const kWriteCb = 1 << 25;
105const kExpectWriteCb = 1 << 26;
106const kAfterWriteTickInfo = 1 << 27;
107const kAfterWritePending = 1 << 28;
108const kBuffered = 1 << 29;
109const kEnded = 1 << 30;
110 
111// TODO(benjamingr) it is likely slower to do it this way than with free functions
112function makeBitMapDescriptor(bit) {
113 return {
114 // This is not a breaking change according to Node.js but better safe than sorry.
115 // Ref: https://github.com/nodejs/node/pull/49834
116 enumerable: !streamsNodejsV24Compat,
117 get() {
118 return (this[kState] & bit) !== 0;
119 },
120 set(value) {
121 if (value) this[kState] |= bit;
122 else this[kState] &= ~bit;
123 },
124 };
125}
126Object.defineProperties(WritableState.prototype, {
127 // Object stream flag to indicate whether or not this stream
128 // contains buffers or objects.
129 objectMode: makeBitMapDescriptor(kObjectMode),
130 
131 // if _final has been called.
132 finalCalled: makeBitMapDescriptor(kFinalCalled),
133 
134 // drain event flag.
135 needDrain: makeBitMapDescriptor(kNeedDrain),
136 
137 // At the start of calling end()
138 ending: makeBitMapDescriptor(kEnding),
139 
140 // When end() has been called, and returned.
141 ended: makeBitMapDescriptor(kEnded),
142 
143 // When 'finish' is emitted.
144 finished: makeBitMapDescriptor(kFinished),
145 
146 // Has it been destroyed.
147 destroyed: makeBitMapDescriptor(kDestroyed),
148 
149 // Should we decode strings into buffers before passing to _write?
150 // this is here so that some node-core streams can optimize string
151 // handling at a lower level.
152 decodeStrings: makeBitMapDescriptor(kDecodeStrings),
153 
154 // A flag to see when we're in the middle of a write.
155 writing: makeBitMapDescriptor(kWriting),
156 
157 // A flag to be able to tell if the onwrite cb is called immediately,
158 // or on a later tick. We set this to true at first, because any
159 // actions that shouldn't happen until "later" should generally also
160 // not happen before the first write call.
161 sync: makeBitMapDescriptor(kSync),
162 
163 // A flag to know if we're processing previously buffered items, which
164 // may call the _write() callback in the same tick, so that we don't
165 // end up in an overlapped onwrite situation.
166 bufferProcessing: makeBitMapDescriptor(kBufferProcessing),
167 
168 // Stream is still being constructed and cannot be
169 // destroyed until construction finished or failed.
170 // Async construction is opt in, therefore we start as
171 // constructed.
172 constructed: makeBitMapDescriptor(kConstructed),
173 
174 // Emit prefinish if the only thing we're waiting for is _write cbs
175 // This is relevant for synchronous Transform streams.
176 prefinished: makeBitMapDescriptor(kPrefinished),
177 
178 // True if the error was already emitted and should not be thrown again.
179 errorEmitted: makeBitMapDescriptor(kErrorEmitted),
180 
181 // Should close be emitted on destroy. Defaults to true.
182 emitClose: makeBitMapDescriptor(kEmitClose),
183 
184 // Should .destroy() be called after 'finish' (and potentially 'end').
185 autoDestroy: makeBitMapDescriptor(kAutoDestroy),
186 
187 // Indicates whether the stream has finished destroying.
188 closed: makeBitMapDescriptor(kClosed),
189 
190 // True if close has been emitted or would have been emitted
191 // depending on emitClose.
192 closeEmitted: makeBitMapDescriptor(kCloseEmitted),
193 
194 allBuffers: makeBitMapDescriptor(kAllBuffers),
195 allNoop: makeBitMapDescriptor(kAllNoop),
196 
197 // Indicates whether the stream has errored. When true all write() calls
198 // should return false. This is needed since when autoDestroy
199 // is disabled we need a way to tell whether the stream has failed.
200 // This is/should be a cold path.
201 errored: {
202 __proto__: null,
203 enumerable: false,
204 get() {
205 return (this[kState] & kErrored) !== 0 ? this[kErroredValue] : null;
206 },
207 set(value) {
208 if (value) {
209 this[kErroredValue] = value;
210 this[kState] |= kErrored;
211 } else {
212 this[kState] &= ~kErrored;
213 }
214 },
215 },
216 
217 writable: {
218 __proto__: null,
219 enumerable: false,
220 get() {
221 return (this[kState] & kHasWritable) !== 0
222 ? (this[kState] & kWritable) !== 0
223 : undefined;
224 },
225 set(value) {
226 if (value == null) {
227 this[kState] &= ~(kHasWritable | kWritable);
228 } else if (value) {
229 this[kState] |= kHasWritable | kWritable;
230 } else {
231 this[kState] |= kHasWritable;
232 this[kState] &= ~kWritable;
233 }
234 },
235 },
236 
237 defaultEncoding: {
238 __proto__: null,
239 enumerable: false,
240 get() {
241 return (this[kState] & kDefaultUTF8Encoding) !== 0
242 ? 'utf8'
243 : this[kDefaultEncodingValue];
244 },
245 set(value) {
246 if (value === 'utf8' || value === 'utf-8') {
247 this[kState] |= kDefaultUTF8Encoding;
248 } else {
249 this[kState] &= ~kDefaultUTF8Encoding;
250 this[kDefaultEncodingValue] = value;
251 }
252 },
253 },
254 
255 // The callback that the user supplies to write(chunk, encoding, cb).
256 writecb: {
257 __proto__: null,
258 enumerable: false,
259 get() {
260 return (this[kState] & kWriteCb) !== 0 ? this[kWriteCbValue] : nop;
261 },
262 set(value) {
263 this[kWriteCbValue] = value;
264 if (value) {
265 this[kState] |= kWriteCb;
266 } else {
267 this[kState] &= ~kWriteCb;
268 }
269 },
270 },
271 
272 // Storage for data passed to the afterWrite() callback in case of
273 // synchronous _write() completion.
274 afterWriteTickInfo: {
275 __proto__: null,
276 enumerable: false,
277 get() {
278 return (this[kState] & kAfterWriteTickInfo) !== 0
279 ? this[kAfterWriteTickInfoValue]
280 : null;
281 },
282 set(value) {
283 this[kAfterWriteTickInfoValue] = value;
284 if (value) {
285 this[kState] |= kAfterWriteTickInfo;
286 } else {
287 this[kState] &= ~kAfterWriteTickInfo;
288 }
289 },
290 },
291 
292 buffered: {
293 __proto__: null,
294 enumerable: false,
295 get() {
296 return (this[kState] & kBuffered) !== 0 ? this[kBufferedValue] : [];
297 },
298 set(value) {
299 this[kBufferedValue] = value;
300 if (value) {
301 this[kState] |= kBuffered;
302 } else {
303 this[kState] &= ~kBuffered;
304 }
305 },
306 },
307});
308 
309// ======================================================================================
310// WritableState
311 
312export function WritableState(options, stream, isDuplex) {
313 // Bit map field to store WritableState more efficiently with 1 bit per field
314 // instead of a V8 slot per field.
315 this[kState] = kSync | kConstructed | kEmitClose | kAutoDestroy;
316 
317 if (options?.objectMode) this[kState] |= kObjectMode;
318 
319 if (isDuplex && options?.writableObjectMode) this[kState] |= kObjectMode;
320 
321 // The point at which write() starts returning false
322 // Note: 0 is a valid value, means that we always return false if
323 // the entire buffer is not flushed immediately on write().
324 this.highWaterMark = options
325 ? getHighWaterMark(this, options, 'writableHighWaterMark', isDuplex)
326 : getDefaultHighWaterMark(false);
327 
328 if (!options || options.decodeStrings !== false)
329 this[kState] |= kDecodeStrings;
330 
331 // Should close be emitted on destroy. Defaults to true.
332 if (options && options.emitClose === false) this[kState] &= ~kEmitClose;
333 
334 // Should .destroy() be called after 'end' (and potentially 'finish').
335 if (options && options.autoDestroy === false) this[kState] &= ~kAutoDestroy;
336 
337 // Crypto is kind of old and crusty. Historically, its default string
338 // encoding is 'binary' so we have to make this configurable.
339 // Everything else in the universe uses 'utf8', though.
340 const defaultEncoding = options ? options.defaultEncoding : null;
341 if (
342 defaultEncoding == null ||
343 defaultEncoding === 'utf8' ||
344 defaultEncoding === 'utf-8'
345 ) {
346 this[kState] |= kDefaultUTF8Encoding;
347 } else if (Buffer.isEncoding(defaultEncoding)) {
348 this[kState] &= ~kDefaultUTF8Encoding;
349 this[kDefaultEncodingValue] = defaultEncoding;
350 } else if (streamsNodejsV24Compat) {
351 // This is a semver-major change. Ref: https://github.com/nodejs/node/pull/46322
352 throw new ERR_UNKNOWN_ENCODING(defaultEncoding);
353 } else {
354 this[kDefaultEncodingValue] = defaultEncoding;
355 }
356 
357 // Not an actual buffer we keep track of, but a measurement
358 // of how much we're waiting to get pushed to some underlying
359 // socket or file.
360 this.length = 0;
361 
362 // When true all writes will be buffered until .uncork() call.
363 this.corked = 0;
364 
365 // The callback that's passed to _write(chunk, cb).
366 this.onwrite = onwrite.bind(undefined, stream);
367 
368 // The amount that is being written when _write is called.
369 this.writelen = 0;
370 
371 resetBuffer(this);
372 
373 // Number of pending user-supplied write callbacks
374 // this must be 0 before 'finish' can be emitted.
375 this.pendingcb = 0;
376}
377 
378function resetBuffer(state) {
379 state[kBufferedValue] = null;
380 state.bufferedIndex = 0;
381 state[kState] |= kAllBuffers | kAllNoop;
382 state[kState] &= ~kBuffered;
383}
384 
385WritableState.prototype.getBuffer = function getBuffer() {
386 return (this[kState] & kBuffered) === 0
387 ? []
388 : this.buffered.slice(this.bufferedIndex);
389};
390 
391Object.defineProperty(WritableState.prototype, 'bufferedRequestCount', {
392 __proto__: null,
393 get() {
394 return (this[kState] & kBuffered) === 0
395 ? 0
396 : this[kBufferedValue].length - this.bufferedIndex;
397 },
398});
399 
400WritableState.prototype[kOnConstructed] = function onConstructed(stream) {
401 if ((this[kState] & kWriting) === 0) {
402 clearBuffer(stream, this);
403 }
404 
405 if ((this[kState] & kEnding) !== 0) {
406 finishMaybe(stream, this);
407 }
408};
409 
410// ======================================================================================
411// Writable
412 
413Writable.WritableState = WritableState;
414 
415Object.setPrototypeOf(Writable.prototype, Stream.prototype);
416Object.setPrototypeOf(Writable, Stream);
417 
418export function Writable(options) {
419 if (!(this instanceof Writable)) return new Writable(options);
420 
421 this._events ??= {
422 close: undefined,
423 error: undefined,
424 prefinish: undefined,
425 finish: undefined,
426 drain: undefined,
427 // Skip uncommon events...
428 // [destroyImpl.kConstruct]: undefined,
429 // [destroyImpl.kDestroy]: undefined,
430 };
431 
432 this._writableState = new WritableState(options, this, false);
433 
434 if (options) {
435 if (typeof options.write === 'function') this._write = options.write;
436 
437 if (typeof options.writev === 'function') this._writev = options.writev;
438 
439 if (typeof options.destroy === 'function') this._destroy = options.destroy;
440 
441 if (typeof options.final === 'function') this._final = options.final;
442 
443 if (typeof options.construct === 'function')
444 this._construct = options.construct;
445 
446 if (options.signal) addAbortSignal(options.signal, this);
447 }
448 
449 Stream.call(this, options);
450 
451 if (this._construct != null) {
452 destroyImpl.construct(this, () => {
453 this._writableState[kOnConstructed](this);
454 });
455 }
456}
457 
458Object.defineProperty(Writable, Symbol.hasInstance, {
459 __proto__: null,
460 value: function (object) {
461 if (Function.prototype[Symbol.hasInstance].call(this, object)) return true;
462 if (this !== Writable) return false;
463 
464 return object && object._writableState instanceof WritableState;
465 },
466});
467 
468// Otherwise people can pipe Writable streams, which is just wrong.
469Writable.prototype.pipe = function () {
470 destroyImpl.errorOrDestroy(this, new ERR_STREAM_CANNOT_PIPE());
471};
472 
473function _write(stream, chunk, encoding, cb) {
474 const state = stream._writableState;
475 
476 if (cb == null || typeof cb !== 'function') {
477 cb = nop;
478 }
479 
480 if (chunk === null) {
481 throw new ERR_STREAM_NULL_VALUES();
482 }
483 
484 if ((state[kState] & kObjectMode) === 0) {
485 if (!encoding) {
486 encoding =
487 (state[kState] & kDefaultUTF8Encoding) !== 0
488 ? 'utf8'
489 : state.defaultEncoding;
490 } else if (encoding !== 'buffer' && !Buffer.isEncoding(encoding)) {
491 throw new ERR_UNKNOWN_ENCODING(encoding);
492 }
493 
494 if (typeof chunk === 'string') {
495 if ((state[kState] & kDecodeStrings) !== 0) {
496 chunk = Buffer.from(chunk, encoding);
497 encoding = 'buffer';
498 }
499 } else if (chunk instanceof Buffer) {
500 encoding = 'buffer';
501 } else if (Stream._isArrayBufferView(chunk)) {
502 chunk = Stream._uint8ArrayToBuffer(chunk);
503 encoding = 'buffer';
504 } else {
505 throw new ERR_INVALID_ARG_TYPE(
506 'chunk',
507 ['string', 'Buffer', 'TypedArray', 'DataView'],
508 chunk
509 );
510 }
511 }
512 
513 let err;
514 if ((state[kState] & kEnding) !== 0) {
515 err = new ERR_STREAM_WRITE_AFTER_END();
516 } else if ((state[kState] & kDestroyed) !== 0) {
517 err = new ERR_STREAM_DESTROYED('write');
518 }
519 
520 if (err) {
521 nextTick(cb, err);
522 destroyImpl.errorOrDestroy(stream, err, true);
523 return err;
524 }
525 
526 state.pendingcb++;
527 return writeOrBuffer(stream, state, chunk, encoding, cb);
528}
529 
530Writable.prototype.write = function (chunk, encoding, cb) {
531 if (encoding != null && typeof encoding === 'function') {
532 cb = encoding;
533 encoding = null;
534 }
535 
536 return _write(this, chunk, encoding, cb) === true;
537};
538 
539Writable.prototype.cork = function () {
540 const state = this._writableState;
541 
542 state[kState] |= kCorked;
543 state.corked++;
544};
545 
546Writable.prototype.uncork = function () {
547 const state = this._writableState;
548 
549 if (state.corked) {
550 state.corked--;
551 
552 if (!state.corked) {
553 state[kState] &= ~kCorked;
554 }
555 
556 if ((state[kState] & kWriting) === 0) clearBuffer(this, state);
557 }
558};
559 
560Writable.prototype.setDefaultEncoding = function setDefaultEncoding(encoding) {
561 // node::ParseEncoding() requires lower case.
562 if (typeof encoding === 'string') encoding = encoding.toLowerCase();
563 if (!Buffer.isEncoding(encoding)) throw new ERR_UNKNOWN_ENCODING(encoding);
564 this._writableState.defaultEncoding = encoding;
565 return this;
566};
567 
568// If we're already writing something, then just put this
569// in the queue, and wait our turn. Otherwise, call _write
570// If we return false, then we need a drain event, so set that flag.
571function writeOrBuffer(stream, state, chunk, encoding, callback) {
572 const len = (state[kState] & kObjectMode) !== 0 ? 1 : chunk.length;
573 
574 state.length += len;
575 
576 // This is a semver-major change. Ref: https://github.com/nodejs/node/commit/557044af407376aff28a0a0800f3053bb58e9239
577 //
578 // The timing of backpressure (ret) calculation relative to _write() execution is critical.
579 // When _write() completes synchronously and modifies state.length, calculating ret before
580 // the write uses stale buffer state, leading to incorrect needDrain signaling and stream hangs.
581 // v24+ calculates ret after _write() to ensure backpressure reflects post-write buffer state.
582 let ret;
583 if (!streamsNodejsV24Compat) {
584 ret = state.length < state.highWaterMark || state.length === 0;
585 // We must ensure that previous needDrain will not be reset to false.
586 if (!ret) {
587 state[kState] |= kNeedDrain;
588 }
589 }
590 
591 if (
592 (state[kState] & (kWriting | kErrored | kCorked | kConstructed)) !==
593 kConstructed
594 ) {
595 if ((state[kState] & kBuffered) === 0) {
596 state[kState] |= kBuffered;
597 state[kBufferedValue] = [];
598 }
599 
600 state[kBufferedValue].push({ chunk, encoding, callback });
601 if ((state[kState] & kAllBuffers) !== 0 && encoding !== 'buffer') {
602 state[kState] &= ~kAllBuffers;
603 }
604 if ((state[kState] & kAllNoop) !== 0 && callback !== nop) {
605 state[kState] &= ~kAllNoop;
606 }
607 } else {
608 state.writelen = len;
609 if (callback !== nop) {
610 state.writecb = callback;
611 }
612 state[kState] |= kWriting | kSync | kExpectWriteCb;
613 stream._write(chunk, encoding, state.onwrite);
614 state[kState] &= ~kSync;
615 }
616 
617 // This is a semver-major change. Ref: https://github.com/nodejs/node/commit/557044af407376aff28a0a0800f3053bb58e9239
618 // For v24+, calculate ret after _write() to observe post-write buffer state.
619 if (streamsNodejsV24Compat) {
620 ret = state.length < state.highWaterMark || state.length === 0;
621 // We must ensure that previous needDrain will not be reset to false.
622 if (!ret) {
623 state[kState] |= kNeedDrain;
624 }
625 }
626 
627 // Return false if errored or destroyed in order to break
628 // any synchronous while(stream.write(data)) loops.
629 return ret && (state[kState] & (kDestroyed | kErrored)) === 0;
630}
631 
632function doWrite(stream, state, writev, len, chunk, encoding, cb) {
633 state.writelen = len;
634 if (cb !== nop) {
635 state.writecb = cb;
636 }
637 state[kState] |= kWriting | kSync | kExpectWriteCb;
638 if ((state[kState] & kDestroyed) !== 0)
639 state.onwrite(new ERR_STREAM_DESTROYED('write'));
640 else if (writev) stream._writev(chunk, state.onwrite);
641 else stream._write(chunk, encoding, state.onwrite);
642 state[kState] &= ~kSync;
643}
644 
645function onwriteError(stream, state, er, cb) {
646 --state.pendingcb;
647 cb(er);
648 // Ensure callbacks are invoked even when autoDestroy is
649 // not enabled. Passing `er` here doesn't make sense since
650 // it's related to one specific write, not to the buffered
651 // writes.
652 errorBuffer(state);
653 // This can emit error, but error must always follow cb.
654 destroyImpl.errorOrDestroy(stream, er);
655}
656 
657function onwrite(stream, er) {
658 const state = stream._writableState;
659 
660 if ((state[kState] & kExpectWriteCb) === 0) {
661 destroyImpl.errorOrDestroy(stream, new ERR_MULTIPLE_CALLBACK());
662 return;
663 }
664 
665 const sync = (state[kState] & kSync) !== 0;
666 const cb = (state[kState] & kWriteCb) !== 0 ? state[kWriteCbValue] : nop;
667 
668 state.writecb = null;
669 state[kState] &= ~(kWriting | kExpectWriteCb);
670 state.length -= state.writelen;
671 state.writelen = 0;
672 
673 if (er) {
674 // Avoid V8 leak, https://github.com/nodejs/node/pull/34103#issuecomment-652002364
675 er.stack; // eslint-disable-line @typescript-eslint/no-unused-expressions
676 
677 if ((state[kState] & kErrored) === 0) {
678 state[kErroredValue] = er;
679 state[kState] |= kErrored;
680 }
681 
682 // In case of duplex streams we need to notify the readable side of the
683 // error.
684 if (stream._readableState && !stream._readableState.errored) {
685 stream._readableState.errored = er;
686 }
687 
688 if (sync) {
689 nextTick(onwriteError, stream, state, er, cb);
690 } else {
691 onwriteError(stream, state, er, cb);
692 }
693 } else {
694 if ((state[kState] & kBuffered) !== 0) {
695 clearBuffer(stream, state);
696 }
697 
698 if (sync) {
699 const needDrain =
700 (state[kState] & kNeedDrain) !== 0 && state.length === 0;
701 const needTick =
702 needDrain || state[kState] & (kDestroyed !== 0) || cb !== nop;
703 
704 // It is a common case that the callback passed to .write() is always
705 // the same. In that case, we do not schedule a new nextTick(), but
706 // rather just increase a counter, to improve performance and avoid
707 // memory allocations.
708 if (cb === nop) {
709 if ((state[kState] & kAfterWritePending) === 0 && needTick) {
710 nextTick(afterWrite, stream, state, 1, cb);
711 state[kState] |= kAfterWritePending;
712 } else {
713 state.pendingcb--;
714 if ((state[kState] & kEnding) !== 0) {
715 finishMaybe(stream, state, true);
716 }
717 }
718 } else if (
719 (state[kState] & kAfterWriteTickInfo) !== 0 &&
720 state[kAfterWriteTickInfoValue].cb === cb
721 ) {
722 state[kAfterWriteTickInfoValue].count++;
723 } else if (needTick) {
724 state[kAfterWriteTickInfoValue] = { count: 1, cb, stream, state };
725 nextTick(afterWriteTick, state[kAfterWriteTickInfoValue]);
726 state[kState] |= kAfterWritePending | kAfterWriteTickInfo;
727 } else {
728 state.pendingcb--;
729 if ((state[kState] & kEnding) !== 0) {
730 finishMaybe(stream, state, true);
731 }
732 }
733 } else {
734 afterWrite(stream, state, 1, cb);
735 }
736 }
737}
738 
739function afterWriteTick({ stream, state, count, cb }) {
740 state[kState] &= ~kAfterWriteTickInfo;
741 state[kAfterWriteTickInfoValue] = null;
742 return afterWrite(stream, state, count, cb);
743}
744 
745function afterWrite(stream, state, count, cb) {
746 state[kState] &= ~kAfterWritePending;
747 
748 const needDrain =
749 (state[kState] & (kEnding | kNeedDrain | kDestroyed)) === kNeedDrain &&
750 state.length === 0;
751 if (needDrain) {
752 state[kState] &= ~kNeedDrain;
753 stream.emit('drain');
754 }
755 
756 // This is a semver-major change. Ref: https://github.com/nodejs/node/pull/44312/files
757 const callbackValue = streamsNodejsV24Compat ? null : undefined;
758 while (count-- > 0) {
759 state.pendingcb--;
760 cb(callbackValue);
761 }
762 
763 if ((state[kState] & kDestroyed) !== 0) {
764 errorBuffer(state);
765 }
766 
767 if ((state[kState] & kEnding) !== 0) {
768 finishMaybe(stream, state, true);
769 }
770}
771 
772// If there's something in the buffer waiting, then invoke callbacks.
773function errorBuffer(state) {
774 if ((state[kState] & kWriting) !== 0) {
775 return;
776 }
777 
778 if ((state[kState] & kBuffered) !== 0) {
779 for (let n = state.bufferedIndex; n < state.buffered.length; ++n) {
780 const { chunk, callback } = state[kBufferedValue][n];
781 const len = (state[kState] & kObjectMode) !== 0 ? 1 : chunk.length;
782 state.length -= len;
783 callback(state.errored ?? new ERR_STREAM_DESTROYED('write'));
784 }
785 }
786 
787 callFinishedCallbacks(
788 state,
789 state.errored ?? new ERR_STREAM_DESTROYED('end')
790 );
791 
792 resetBuffer(state);
793}
794 
795// If there's something in the buffer waiting, then process it.
796function clearBuffer(stream, state) {
797 if (
798 (state[kState] &
799 (kDestroyed | kBufferProcessing | kCorked | kBuffered | kConstructed)) !==
800 (kBuffered | kConstructed)
801 ) {
802 return;
803 }
804 
805 const objectMode = (state[kState] & kObjectMode) !== 0;
806 const { [kBufferedValue]: buffered, bufferedIndex } = state;
807 const bufferedLength = buffered.length - bufferedIndex;
808 
809 if (!bufferedLength) {
810 return;
811 }
812 
813 let i = bufferedIndex;
814 
815 state[kState] |= kBufferProcessing;
816 if (bufferedLength > 1 && stream._writev) {
817 state.pendingcb -= bufferedLength - 1;
818 
819 const callback =
820 (state[kState] & kAllNoop) !== 0
821 ? nop
822 : (err) => {
823 for (let n = i; n < buffered.length; ++n) {
824 buffered[n].callback(err);
825 }
826 };
827 // Make a copy of `buffered` if it's going to be used by `callback` above,
828 // since `doWrite` will mutate the array.
829 const chunks =
830 (state[kState] & kAllNoop) !== 0 && i === 0
831 ? buffered
832 : buffered.slice(i);
833 chunks.allBuffers = (state[kState] & kAllBuffers) !== 0;
834 
835 doWrite(stream, state, true, state.length, chunks, '', callback);
836 
837 resetBuffer(state);
838 } else {
839 do {
840 const { chunk, encoding, callback } = buffered[i];
841 buffered[i++] = null;
842 const len = objectMode ? 1 : chunk.length;
843 doWrite(stream, state, false, len, chunk, encoding, callback);
844 } while (i < buffered.length && (state[kState] & kWriting) === 0);
845 
846 if (i === buffered.length) {
847 resetBuffer(state);
848 } else if (i > 256) {
849 buffered.splice(0, i);
850 state.bufferedIndex = 0;
851 } else {
852 state.bufferedIndex = i;
853 }
854 }
855 state[kState] &= ~kBufferProcessing;
856}
857 
858Writable.prototype._write = function (chunk, encoding, cb) {
859 if (this._writev) {
860 this._writev([{ chunk, encoding }], cb);
861 } else {
862 throw new ERR_METHOD_NOT_IMPLEMENTED('_write()');
863 }
864};
865 
866Writable.prototype._writev = null;
867 
868Writable.prototype.end = function (chunk, encoding, cb) {
869 const state = this._writableState;
870 
871 if (typeof chunk === 'function') {
872 cb = chunk;
873 chunk = null;
874 encoding = null;
875 } else if (typeof encoding === 'function') {
876 cb = encoding;
877 encoding = null;
878 }
879 
880 let err;
881 
882 if (chunk != null) {
883 const ret = _write(this, chunk, encoding);
884 if (ret instanceof Error) {
885 err = ret;
886 }
887 }
888 
889 // .end() fully uncorks.
890 if ((state[kState] & kCorked) !== 0) {
891 state.corked = 1;
892 this.uncork();
893 }
894 
895 if (err) {
896 // Do nothing...
897 } else if ((state[kState] & (kEnding | kErrored)) === 0) {
898 // This is forgiving in terms of unnecessary calls to end() and can hide
899 // logic errors. However, usually such errors are harmless and causing a
900 // hard error can be disproportionately destructive. It is not always
901 // trivial for the user to determine whether end() needs to be called
902 // or not.
903 
904 state[kState] |= kEnding;
905 finishMaybe(this, state, true);
906 state[kState] |= kEnded;
907 } else if ((state[kState] & kFinished) !== 0) {
908 err = new ERR_STREAM_ALREADY_FINISHED('end');
909 } else if ((state[kState] & kDestroyed) !== 0) {
910 err = new ERR_STREAM_DESTROYED('end');
911 }
912 
913 if (typeof cb === 'function') {
914 // This is a semver-major change. Ref: https://github.com/nodejs/node/pull/44312
915 if (streamsNodejsV24Compat) {
916 if (err) {
917 nextTick(cb, err);
918 } else if ((state[kState] & kErrored) !== 0) {
919 nextTick(cb, state[kErroredValue]);
920 } else if ((state[kState] & kFinished) !== 0) {
921 nextTick(cb, null);
922 } else {
923 state[kState] |= kOnFinished;
924 state[kOnFinishedValue] ??= [];
925 state[kOnFinishedValue].push(cb);
926 }
927 } else {
928 if (err || (state[kState] & kFinished) !== 0) {
929 nextTick(cb, err);
930 } else if ((state[kState] & kErrored) !== 0) {
931 nextTick(cb, state[kErroredValue]);
932 } else {
933 state[kState] |= kOnFinished;
934 state[kOnFinishedValue] ??= [];
935 state[kOnFinishedValue].push(cb);
936 }
937 }
938 }
939 
940 return this;
941};
942 
943function needFinish(state) {
944 return (
945 // State is ended && constructed but not destroyed, finished, writing, errorEmitted or closedEmitted
946 (state[kState] &
947 (kEnding |
948 kDestroyed |
949 kConstructed |
950 kFinished |
951 kWriting |
952 kErrorEmitted |
953 kCloseEmitted |
954 kErrored |
955 kBuffered)) ===
956 (kEnding | kConstructed) && state.length === 0
957 );
958}
959function onFinish(stream, state, err) {
960 if ((state[kState] & kPrefinished) !== 0) {
961 destroyImpl.errorOrDestroy(stream, err ?? new ERR_MULTIPLE_CALLBACK());
962 return;
963 }
964 state.pendingcb--;
965 if (err) {
966 callFinishedCallbacks(state, err);
967 destroyImpl.errorOrDestroy(stream, err, (state[kState] & kSync) !== 0);
968 } else if (needFinish(state)) {
969 state[kState] |= kPrefinished;
970 stream.emit('prefinish');
971 // Backwards compat. Don't check state.sync here.
972 // Some streams assume 'finish' will be emitted
973 // asynchronously relative to _final callback.
974 state.pendingcb++;
975 nextTick(finish, stream, state);
976 }
977}
978 
979function prefinish(stream, state) {
980 if ((state[kState] & (kPrefinished | kFinalCalled)) !== 0) {
981 return;
982 }
983 
984 if (
985 typeof stream._final === 'function' &&
986 (state[kState] & kDestroyed) === 0
987 ) {
988 state[kState] |= kFinalCalled | kSync;
989 state.pendingcb++;
990 
991 try {
992 stream._final((err) => onFinish(stream, state, err));
993 } catch (err) {
994 onFinish(stream, state, err);
995 }
996 
997 state[kState] &= ~kSync;
998 } else {
999 state[kState] |= kFinalCalled | kPrefinished;
1000 stream.emit('prefinish');
1001 }
1002}
1003 
1004function finishMaybe(stream, state, sync) {
1005 if (needFinish(state)) {
1006 prefinish(stream, state);
1007 if (state.pendingcb === 0) {
1008 if (sync) {
1009 state.pendingcb++;
1010 nextTick(
1011 (stream, state) => {
1012 if (needFinish(state)) {
1013 finish(stream, state);
1014 } else {
1015 state.pendingcb--;
1016 }
1017 },
1018 stream,
1019 state
1020 );
1021 } else if (needFinish(state)) {
1022 state.pendingcb++;
1023 finish(stream, state);
1024 }
1025 }
1026 }
1027}
1028 
1029function finish(stream, state) {
1030 state.pendingcb--;
1031 state[kState] |= kFinished;
1032 
1033 callFinishedCallbacks(state, null);
1034 
1035 stream.emit('finish');
1036 
1037 if ((state[kState] & kAutoDestroy) !== 0) {
1038 // In case of duplex streams we need a way to detect
1039 // if the readable side is ready for autoDestroy as well.
1040 const rState = stream._readableState;
1041 const autoDestroy =
1042 !rState ||
1043 (rState.autoDestroy &&
1044 // We don't expect the readable to ever 'end'
1045 // if readable is explicitly set to false.
1046 (rState.endEmitted || rState.readable === false));
1047 if (autoDestroy) {
1048 stream.destroy();
1049 }
1050 }
1051}
1052 
1053function callFinishedCallbacks(state, err) {
1054 if ((state[kState] & kOnFinished) === 0) {
1055 return;
1056 }
1057 
1058 const onfinishCallbacks = state[kOnFinishedValue];
1059 // This is a semver-major change. Ref: https://github.com/nodejs/node/pull/44312
1060 state[kOnFinishedValue] = streamsNodejsV24Compat ? null : undefined;
1061 state[kState] &= ~kOnFinished;
1062 for (let i = 0; i < onfinishCallbacks.length; i++) {
1063 onfinishCallbacks[i](err);
1064 }
1065}
1066 
1067Object.defineProperties(Writable.prototype, {
1068 closed: {
1069 __proto__: null,
1070 get() {
1071 return this._writableState
1072 ? (this._writableState[kState] & kClosed) !== 0
1073 : false;
1074 },
1075 },
1076 
1077 destroyed: {
1078 __proto__: null,
1079 get() {
1080 return this._writableState
1081 ? (this._writableState[kState] & kDestroyed) !== 0
1082 : false;
1083 },
1084 set(value) {
1085 // Backward compatibility, the user is explicitly managing destroyed.
1086 if (!this._writableState) return;
1087 
1088 if (value) this._writableState[kState] |= kDestroyed;
1089 else this._writableState[kState] &= ~kDestroyed;
1090 },
1091 },
1092 
1093 writable: {
1094 __proto__: null,
1095 get() {
1096 const w = this._writableState;
1097 // w.writable === false means that this is part of a Duplex stream
1098 // where the writable side was disabled upon construction.
1099 // Compat. The user might manually disable writable side through
1100 // deprecated setter.
1101 return (
1102 !!w &&
1103 w.writable !== false &&
1104 (w[kState] & (kEnding | kEnded | kDestroyed | kErrored)) === 0
1105 );
1106 },
1107 set(val) {
1108 // Backwards compatible.
1109 if (this._writableState) {
1110 this._writableState.writable = !!val;
1111 }
1112 },
1113 },
1114 
1115 writableFinished: {
1116 __proto__: null,
1117 get() {
1118 const state = this._writableState;
1119 return state ? (state[kState] & kFinished) !== 0 : false;
1120 },
1121 },
1122 
1123 writableObjectMode: {
1124 __proto__: null,
1125 get() {
1126 const state = this._writableState;
1127 return state ? (state[kState] & kObjectMode) !== 0 : false;
1128 },
1129 },
1130 
1131 writableBuffer: {
1132 __proto__: null,
1133 get() {
1134 const state = this._writableState;
1135 return state && state.getBuffer();
1136 },
1137 },
1138 
1139 writableEnded: {
1140 __proto__: null,
1141 get() {
1142 const state = this._writableState;
1143 return state ? (state[kState] & kEnding) !== 0 : false;
1144 },
1145 },
1146 
1147 writableNeedDrain: {
1148 __proto__: null,
1149 get() {
1150 const state = this._writableState;
1151 return state
1152 ? (state[kState] & (kDestroyed | kEnding | kNeedDrain)) === kNeedDrain
1153 : false;
1154 },
1155 },
1156 
1157 writableHighWaterMark: {
1158 __proto__: null,
1159 get() {
1160 const state = this._writableState;
1161 return state?.highWaterMark;
1162 },
1163 },
1164 
1165 writableCorked: {
1166 __proto__: null,
1167 get() {
1168 const state = this._writableState;
1169 return state ? state.corked : 0;
1170 },
1171 },
1172 
1173 writableLength: {
1174 __proto__: null,
1175 get() {
1176 const state = this._writableState;
1177 return state?.length;
1178 },
1179 },
1180 
1181 errored: {
1182 __proto__: null,
1183 enumerable: false,
1184 get() {
1185 const state = this._writableState;
1186 return state ? state.errored : null;
1187 },
1188 },
1189 
1190 writableAborted: {
1191 __proto__: null,
1192 get: function () {
1193 const state = this._writableState;
1194 return (
1195 (state[kState] & (kHasWritable | kWritable)) !== kHasWritable &&
1196 (state[kState] & (kDestroyed | kErrored)) !== 0 &&
1197 (state[kState] & kFinished) === 0
1198 );
1199 },
1200 },
1201});
1202 
1203const destroy = destroyImpl.destroy;
1204Writable.prototype.destroy = function (err, cb) {
1205 const state = this._writableState;
1206 
1207 // Invoke pending callbacks.
1208 if (
1209 (state[kState] & (kBuffered | kOnFinished)) !== 0 &&
1210 (state[kState] & kDestroyed) === 0
1211 ) {
1212 nextTick(errorBuffer, state);
1213 }
1214 
1215 destroy.call(this, err, cb);
1216 return this;
1217};
1218 
1219Writable.prototype._undestroy = destroyImpl.undestroy;
1220Writable.prototype._destroy = function (err, cb) {
1221 cb(err);
1222};
1223 
1224Writable.prototype[EventEmitter.captureRejectionSymbol] = function (err) {
1225 this.destroy(err);
1226};
1227 
1228export function fromWeb(writableStream, options) {
1229 return newStreamWritableFromWritableStream(writableStream, options);
1230}
1231 
1232export function toWeb(streamWritable) {
1233 return newWritableStreamFromStreamWritable(streamWritable);
1234}
1235 
1236Writable.fromWeb = fromWeb;
1237Writable.toWeb = toWeb;
1238 
1239Writable.prototype[Symbol.asyncDispose] = async function () {
1240 let error;
1241 if (!this.destroyed) {
1242 error = this.writableFinished ? null : new AbortError();
1243 this.destroy(error);
1244 }
1245 await new Promise((resolve, reject) =>
1246 eos(this, (err) =>
1247 err && err.name !== 'AbortError' ? reject(err) : resolve(null)
1248 )
1249 );
1250};
1251 
1252/**
1253 * @param {Writable} streamWritable
1254 * @returns {WritableStream}
1255 */
1256export function newWritableStreamFromStreamWritable(streamWritable) {
1257 // Not using the internal/streams/utils isWritableNodeStream utility
1258 // here because it will return false if streamWritable is a Duplex
1259 // whose writable option is false. For a Duplex that is not writable,
1260 // we want it to pass this check but return a closed WritableStream.
1261 // We check if the given stream is a stream.Writable or http.OutgoingMessage
1262 const checkIfWritableOrOutgoingMessage =
1263 streamWritable &&
1264 typeof streamWritable?.write === 'function' &&
1265 typeof streamWritable?.on === 'function';
1266 if (!checkIfWritableOrOutgoingMessage) {
1267 throw new ERR_INVALID_ARG_TYPE(
1268 'streamWritable',
1269 'stream.Writable',
1270 streamWritable
1271 );
1272 }
1273 
1274 if (isDestroyed(streamWritable) || !isWritable(streamWritable)) {
1275 const writable = new globalThis.WritableStream();
1276 writable.close();
1277 return writable;
1278 }
1279 
1280 const highWaterMark = streamWritable.writableHighWaterMark;
1281 const strategy = streamWritable.writableObjectMode
1282 ? new globalThis.CountQueuingStrategy({ highWaterMark })
1283 : { highWaterMark };
1284 
1285 let controller;
1286 let backpressurePromise;
1287 let closed;
1288 
1289 function onDrain() {
1290 if (backpressurePromise !== undefined) backpressurePromise.resolve();
1291 }
1292 
1293 const cleanup = finished(streamWritable, (error) => {
1294 error = handleKnownInternalErrors(error);
1295 
1296 cleanup();
1297 // This is a protection against non-standard, legacy streams
1298 // that happen to emit an error event again after finished is called.
1299 streamWritable.on('error', () => {});
1300 if (error != null) {
1301 if (backpressurePromise !== undefined) backpressurePromise.reject(error);
1302 // If closed is not undefined, the error is happening
1303 // after the WritableStream close has already started.
1304 // We need to reject it here.
1305 if (closed !== undefined) {
1306 closed.reject(error);
1307 closed = undefined;
1308 }
1309 controller.error(error);
1310 controller = undefined;
1311 return;
1312 }
1313 
1314 if (closed !== undefined) {
1315 closed.resolve();
1316 closed = undefined;
1317 return;
1318 }
1319 controller.error(new AbortError());
1320 controller = undefined;
1321 });
1322 
1323 streamWritable.on('drain', onDrain);
1324 
1325 return new globalThis.WritableStream(
1326 {
1327 start(c) {
1328 controller = c;
1329 },
1330 
1331 write(chunk) {
1332 if (streamWritable.writableNeedDrain || !streamWritable.write(chunk)) {
1333 backpressurePromise = Promise.withResolvers();
1334 return backpressurePromise.promise.finally(() => {
1335 backpressurePromise = undefined;
1336 });
1337 }
1338 },
1339 
1340 abort(reason) {
1341 destroy(streamWritable, reason);
1342 },
1343 
1344 close() {
1345 if (closed === undefined && !isWritableEnded(streamWritable)) {
1346 closed = Promise.withResolvers();
1347 streamWritable.end();
1348 return closed.promise;
1349 }
1350 
1351 controller = undefined;
1352 return Promise.resolve();
1353 },
1354 },
1355 strategy
1356 );
1357}
1358 
1359/**
1360 * @param {WritableStream} writableStream
1361 * @param {{
1362 * decodeStrings? : boolean,
1363 * highWaterMark? : number,
1364 * objectMode? : boolean,
1365 * signal? : AbortSignal,
1366 * }} [options]
1367 * @returns {Writable}
1368 */
1369export function newStreamWritableFromWritableStream(
1370 writableStream,
1371 options = {}
1372) {
1373 if (!isWritableStream(writableStream)) {
1374 throw new ERR_INVALID_ARG_TYPE(
1375 'writableStream',
1376 'WritableStream',
1377 writableStream
1378 );
1379 }
1380 
1381 validateObject(options, 'options');
1382 const {
1383 highWaterMark,
1384 decodeStrings = true,
1385 objectMode = false,
1386 signal,
1387 } = options;
1388 
1389 validateBoolean(objectMode, 'options.objectMode');
1390 validateBoolean(decodeStrings, 'options.decodeStrings');
1391 
1392 const writer = writableStream.getWriter();
1393 let closed = false;
1394 
1395 const writable = new Writable({
1396 highWaterMark,
1397 objectMode,
1398 decodeStrings,
1399 signal,
1400 
1401 writev(chunks, callback) {
1402 function done(error) {
1403 error = error.filter((e) => e);
1404 try {
1405 callback(error.length === 0 ? undefined : error);
1406 } catch (error) {
1407 // In a next tick because this is happening within
1408 // a promise context, and if there are any errors
1409 // thrown we don't want those to cause an unhandled
1410 // rejection. Let's just escape the promise and
1411 // handle it separately.
1412 nextTick(() => destroy(writable, error));
1413 }
1414 }
1415 
1416 writer.ready.then(() => {
1417 return Promise.all(chunks.map((data) => writer.write(data))).then(
1418 done,
1419 done
1420 );
1421 }, done);
1422 },
1423 
1424 write(chunk, encoding, callback) {
1425 if (typeof chunk === 'string' && decodeStrings && !objectMode) {
1426 const enc = normalizeEncoding(encoding);
1427 
1428 if (enc === 'utf8') {
1429 chunk = encoder.encode(chunk);
1430 } else {
1431 chunk = Buffer.from(chunk, encoding);
1432 chunk = new Uint8Array(
1433 chunk.buffer,
1434 chunk.byteOffset,
1435 chunk.byteLength
1436 );
1437 }
1438 }
1439 
1440 function done(error) {
1441 try {
1442 callback(error);
1443 } catch (error) {
1444 destroy(writable, error);
1445 }
1446 }
1447 
1448 writer.ready.then(() => {
1449 return writer.write(chunk).then(done, done);
1450 }, done);
1451 },
1452 
1453 destroy(error, callback) {
1454 function done() {
1455 try {
1456 callback(error);
1457 } catch (error) {
1458 // In a next tick because this is happening within
1459 // a promise context, and if there are any errors
1460 // thrown we don't want those to cause an unhandled
1461 // rejection. Let's just escape the promise and
1462 // handle it separately.
1463 nextTick(() => {
1464 throw error;
1465 });
1466 }
1467 }
1468 
1469 if (!closed) {
1470 if (error != null) {
1471 writer.abort(error).then(done, done);
1472 } else {
1473 writer.close().then(done, done);
1474 }
1475 return;
1476 }
1477 
1478 done();
1479 },
1480 
1481 final(callback) {
1482 function done(error) {
1483 try {
1484 callback(error);
1485 } catch (error) {
1486 // In a next tick because this is happening within
1487 // a promise context, and if there are any errors
1488 // thrown we don't want those to cause an unhandled
1489 // rejection. Let's just escape the promise and
1490 // handle it separately.
1491 nextTick(() => destroy(writable, error));
1492 }
1493 }
1494 
1495 if (!closed) {
1496 writer.close().then(done, done);
1497 }
1498 },
1499 });
1500 
1501 writer.closed.then(
1502 () => {
1503 // If the WritableStream closes before the stream.Writable has been
1504 // ended, we signal an error on the stream.Writable.
1505 closed = true;
1506 if (!isWritableEnded(writable))
1507 destroy(writable, new ERR_STREAM_PREMATURE_CLOSE());
1508 },
1509 (error) => {
1510 // If the WritableStream errors before the stream.Writable has been
1511 // destroyed, signal an error on the stream.Writable.
1512 closed = true;
1513 destroy(writable, error);
1514 }
1515 );
1516 
1517 return writable;
1518}