Skip to content
File

Blob: src/node/internal/internal_fs_streams.ts

typescript1174 lines
1// Copyright (c) 2026 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 
5import { Readable } from 'node-internal:streams_readable';
6import { Writable } from 'node-internal:streams_writable';
7import { Buffer } from 'node-internal:internal_buffer';
8import { type EventEmitter } from 'node-internal:events';
9import {
10 normalizePath,
11 getValidatedFd,
12 isFileHandle,
13 type FilePath,
14 type Position,
15 type WriteSyncOptions,
16 type ValidEncoding,
17} from 'node-internal:internal_fs_utils';
18import { toPathIfFileURL } from 'node-internal:internal_url';
19 
20import * as fs from 'node-internal:internal_fs_callback';
21 
22import {
23 parseFileMode,
24 validateAbortSignal,
25 validateBoolean,
26 validateFunction,
27 validateObject,
28 validateString,
29 validateUint32,
30 validateThisInternalField,
31} from 'node-internal:validators';
32 
33import type {
34 DoubleArgCallback,
35 SingleArgCallback,
36 ErrorOnlyCallback,
37 open as OpenType,
38 close as CloseType,
39 fsync as FsyncType,
40 read as ReadType,
41 write as WriteType,
42 writev as WritevType,
43} from 'node-internal:internal_fs_callback';
44 
45import type { FileHandle } from 'node-internal:internal_fs_promises';
46 
47import { errorOrDestroy } from 'node-internal:streams_destroy';
48import { eos } from 'node-internal:streams_end_of_stream';
49 
50import {
51 ERR_INVALID_ARG_VALUE,
52 ERR_OUT_OF_RANGE,
53 ERR_MISSING_ARGS,
54 ERR_METHOD_NOT_IMPLEMENTED,
55 ERR_STREAM_DESTROYED,
56 ERR_SYSTEM_ERROR,
57} from 'node-internal:internal_errors';
58 
59import type { ReadAsyncOptions } from 'node:fs';
60 
61export interface FsOperations {
62 open?: typeof OpenType | undefined;
63 close?: typeof CloseType | undefined;
64 fsync?: typeof FsyncType | undefined;
65 read?: typeof ReadType | undefined;
66 write?: typeof WriteType | undefined;
67 writev?: typeof WritevType | undefined;
68}
69 
70export interface RealizedFsOperations {
71 open: typeof OpenType;
72 close: typeof CloseType;
73 fsync: typeof FsyncType;
74 read: typeof ReadType;
75 write: typeof WriteType;
76 writev: typeof WritevType;
77}
78 
79// Temporary while developing...
80/* eslint-disable */
81 
82let lazyFs: RealizedFsOperations | undefined;
83async function getLazyFs(): Promise<RealizedFsOperations> {
84 if (lazyFs == undefined) {
85 lazyFs = {
86 open: (...args): void => {
87 fs.open(...args);
88 },
89 close: (...args): void => {
90 fs.close(...args);
91 },
92 fsync: (...args): void => {
93 fs.fsync(...args);
94 },
95 read: (...args): void => {
96 fs.read(...args);
97 },
98 write: (...args): void => {
99 fs.write(...args);
100 },
101 writev: (...args): void => {
102 fs.writev(...args);
103 },
104 };
105 }
106 return lazyFs;
107}
108 
109const kDefaultFsOperations: RealizedFsOperations = {
110 open(
111 path: FilePath,
112 flags: string | number | SingleArgCallback<number> = 'r',
113 mode: string | number | SingleArgCallback<number> = 0o666,
114 cb?: SingleArgCallback<number>
115 ): void {
116 let callback: SingleArgCallback<number>;
117 if (typeof flags === 'function') {
118 callback = flags;
119 } else if (typeof mode === 'function') {
120 callback = mode;
121 } else if (typeof cb === 'function') {
122 callback = cb;
123 } else {
124 throw new ERR_MISSING_ARGS('fs.open callback');
125 }
126 validateFunction(callback, 'fs.open callback');
127 
128 getLazyFs().then(
129 (fs: RealizedFsOperations) => {
130 fs.open(path, (err: unknown, fd: number | undefined) => {
131 if (err) {
132 try {
133 callback(err);
134 } catch (e: unknown) {
135 reportError(e);
136 }
137 return;
138 }
139 try {
140 callback(null, fd);
141 } catch (e: unknown) {
142 reportError(e);
143 }
144 });
145 },
146 (err: unknown) => {
147 try {
148 callback(err);
149 } catch (e: unknown) {
150 reportError(e);
151 }
152 }
153 );
154 },
155 close(fd: number, cb: ErrorOnlyCallback = () => {}): void {
156 getLazyFs().then(
157 (fs: RealizedFsOperations) => {
158 fs.close(fd, (err: unknown) => {
159 if (err) {
160 try {
161 cb(err);
162 } catch (e: unknown) {
163 reportError(e);
164 }
165 return;
166 }
167 try {
168 cb(null);
169 } catch (e: unknown) {
170 reportError(e);
171 }
172 });
173 },
174 (err: unknown) => {
175 try {
176 cb(err);
177 } catch (e: unknown) {
178 reportError(e);
179 }
180 }
181 );
182 },
183 fsync(fd: number, cb: ErrorOnlyCallback = () => {}): void {
184 getLazyFs().then(
185 (fs: RealizedFsOperations) => {
186 fs.fsync(fd, (err: unknown) => {
187 if (err) {
188 try {
189 cb(err);
190 } catch (e: unknown) {
191 reportError(e);
192 }
193 return;
194 }
195 try {
196 cb(null);
197 } catch (e: unknown) {
198 reportError(e);
199 }
200 });
201 },
202 (err: unknown) => {
203 try {
204 cb(err);
205 } catch (e: unknown) {
206 reportError(e);
207 }
208 }
209 );
210 },
211 read<T extends NodeJS.ArrayBufferView>(
212 fd: number,
213 buffer: T | ReadAsyncOptions<T> | DoubleArgCallback<number, T>,
214 offset?: ReadAsyncOptions<T> | number | DoubleArgCallback<number, T>,
215 length?: null | number | DoubleArgCallback<number, T>,
216 position?: Position,
217 cb?: DoubleArgCallback<number, T>
218 ): void {
219 getLazyFs().then(
220 (fs: RealizedFsOperations) => {
221 fs.read(
222 fd,
223 buffer,
224 offset,
225 length,
226 position,
227 (
228 err: unknown,
229 bytesRead: number | undefined,
230 buffer: T | undefined
231 ) => {
232 if (err) {
233 try {
234 cb?.(err);
235 } catch (e: unknown) {
236 reportError(e);
237 }
238 return;
239 }
240 try {
241 cb?.(null, bytesRead, buffer);
242 } catch (e: unknown) {
243 reportError(e);
244 }
245 }
246 );
247 },
248 (err: unknown) => {
249 try {
250 cb?.(err);
251 } catch (e: unknown) {
252 reportError(e);
253 }
254 }
255 );
256 },
257 write<T extends NodeJS.ArrayBufferView>(
258 fd: number,
259 buffer: T | string,
260 offset?: WriteSyncOptions | Position | DoubleArgCallback<number, T>,
261 length?: number | ValidEncoding | DoubleArgCallback<number, T>,
262 position?: Position | DoubleArgCallback<number, T>,
263 cb?: DoubleArgCallback<number, T>
264 ): void {
265 getLazyFs().then(
266 (fs: RealizedFsOperations) => {
267 fs.write(
268 fd,
269 buffer,
270 offset,
271 length,
272 position,
273 (
274 err: unknown,
275 bytesWritten: number | undefined,
276 buffer: T | undefined
277 ) => {
278 if (err) {
279 try {
280 cb?.(err);
281 } catch (e: unknown) {
282 reportError(e);
283 }
284 return;
285 }
286 try {
287 cb?.(null, bytesWritten, buffer);
288 } catch (e: unknown) {
289 reportError(e);
290 }
291 }
292 );
293 },
294 (err: unknown) => {
295 try {
296 cb?.(err);
297 } catch (e: unknown) {
298 reportError(e);
299 }
300 }
301 );
302 },
303 writev<T extends NodeJS.ArrayBufferView>(
304 fd: number,
305 buffers: T[],
306 position?: Position | SingleArgCallback<number>,
307 cb?: DoubleArgCallback<number, T[]>
308 ): void {
309 getLazyFs().then(
310 (fs: RealizedFsOperations) => {
311 fs.writev(
312 fd,
313 buffers,
314 position,
315 (err: unknown, bytesWritten?: number, buffers?: T[]) => {
316 if (err) {
317 try {
318 cb?.(err);
319 } catch (e: unknown) {
320 reportError(e);
321 }
322 return;
323 }
324 try {
325 cb?.(null, bytesWritten, buffers);
326 } catch (e: unknown) {
327 reportError(e);
328 }
329 }
330 );
331 },
332 (err: unknown) => {
333 try {
334 cb?.(err);
335 } catch (e: unknown) {
336 reportError(e);
337 }
338 }
339 );
340 },
341};
342 
343export type ReadStreamOptions = {
344 encoding?: string | undefined;
345 autoClose?: boolean | undefined;
346 autoDestroy?: boolean | undefined;
347 emitClose?: boolean | undefined;
348 start?: number | undefined;
349 end?: number | undefined;
350 highWaterMark?: number | undefined;
351 signal?: AbortSignal | undefined;
352 flags?: string | undefined;
353 fd?: number | FileHandle | undefined;
354 mode?: number | undefined;
355 fs?: FsOperations | undefined;
356};
357 
358const kFs = Symbol('kFs');
359const kIsPerformingIO = Symbol('kIsPerformingIO');
360const kIoDone = Symbol('kIoDone');
361const kHandle = Symbol('kHandle');
362 
363// @ts-expect-error TS2323 Cannot redeclare.
364export declare class ReadStream extends Readable {
365 fd: number | null;
366 flags: string;
367 path: string;
368 mode: number;
369 start: number;
370 end: number;
371 pos: number | undefined;
372 bytesRead: number;
373 flush: boolean;
374 [kFs]: RealizedFsOperations;
375 [kIsPerformingIO]: boolean;
376 [kHandle]: FileHandle | undefined;
377 constructor(path: FilePath | null, options?: ReadStreamOptions);
378 push(chunk: NodeJS.ArrayBufferView | null): boolean;
379 close(callback?: ErrorOnlyCallback): void;
380}
381 
382function construct(
383 this: ReadStream | WriteStream,
384 callback: (err: unknown) => void
385): void {
386 const stream = this;
387 if (typeof stream.fd === 'number') {
388 callback(null);
389 return;
390 }
391 
392 const ee = stream as unknown as EventEmitter;
393 
394 if (typeof (stream as any).open === 'function') {
395 // Backwards compat for monkey patching open().
396 const orgEmit = ee.emit;
397 ee.emit = function (...args): boolean {
398 if (args[0] === 'open') {
399 this.emit = orgEmit;
400 callback(null);
401 Reflect.apply(orgEmit, this, args);
402 } else if (args[0] === 'error') {
403 this.emit = orgEmit;
404 callback(args[1]);
405 } else {
406 Reflect.apply(orgEmit, this, args);
407 }
408 return true;
409 };
410 (stream as any).open();
411 return;
412 }
413 stream[kFs].open(stream.path, (er: unknown, fd: number | undefined) => {
414 if (er) {
415 callback(er);
416 return;
417 }
418 if (fd === undefined) {
419 callback(new ERR_INVALID_ARG_VALUE('fd', 'undefined'));
420 return;
421 }
422 stream.fd = fd;
423 callback(null);
424 ee.emit('open', stream.fd);
425 ee.emit('ready');
426 });
427}
428 
429function getValidatedFsOptions(fs: FsOperations): RealizedFsOperations {
430 validateObject(fs, 'options.fs');
431 
432 if (isFileHandle(fs)) {
433 const handle = fs as unknown as FileHandle;
434 const open = function (..._args: unknown[]): void {
435 throw new ERR_METHOD_NOT_IMPLEMENTED('open()');
436 };
437 const close = function (...args: unknown[]): void {
438 const cb = args[args.length - 1] as ErrorOnlyCallback;
439 handle.close().then(
440 () => cb(null),
441 (err: unknown) => cb(err)
442 );
443 };
444 const fsync = function (...args: unknown[]): void {
445 const cb = args[args.length - 1] as ErrorOnlyCallback;
446 handle.sync().then(
447 () => cb(null),
448 (err: unknown) => cb(err)
449 );
450 };
451 const read = function (...args: unknown[]): void {
452 const cb = args[args.length - 1] as DoubleArgCallback<
453 number,
454 NodeJS.ArrayBufferView
455 >;
456 // @ts-expect-error TS2345
457 handle.read(...args.slice(1, -1)).then(
458 (result: { bytesRead: number; buffer: NodeJS.ArrayBufferView }) => {
459 cb(null, result.bytesRead, result.buffer);
460 },
461 (err: unknown) => cb(err)
462 );
463 };
464 const write = function (...args: unknown[]): void {
465 const cb = args[args.length - 1] as DoubleArgCallback<
466 number,
467 NodeJS.ArrayBufferView
468 >;
469 // @ts-expect-error TS2556
470 handle.write(...args.slice(1, -1)).then(
471 (result: { bytesWritten: number; buffer: NodeJS.ArrayBufferView }) => {
472 cb(null, result.bytesWritten, result.buffer);
473 },
474 (err: unknown) => cb(err)
475 );
476 };
477 const writev = function (...args: unknown[]): void {
478 const cb = args[args.length - 1] as DoubleArgCallback<
479 number,
480 NodeJS.ArrayBufferView[]
481 >;
482 // @ts-expect-error TS2556
483 handle.writev(...args.slice(1, -1)).then(
484 (result: {
485 bytesWritten: number;
486 buffers: NodeJS.ArrayBufferView[];
487 }) => {
488 cb(null, result.bytesWritten, result.buffers);
489 },
490 (err: unknown) => cb(err)
491 );
492 };
493 return {
494 open,
495 close,
496 fsync,
497 read,
498 write,
499 writev,
500 };
501 }
502 
503 let {
504 open = kDefaultFsOperations.open,
505 close = kDefaultFsOperations.close,
506 fsync = kDefaultFsOperations.fsync,
507 read = kDefaultFsOperations.read,
508 write = kDefaultFsOperations.write,
509 writev = kDefaultFsOperations.writev,
510 } = fs as FsOperations;
511 
512 validateFunction(open, 'options.fs.open');
513 validateFunction(read, 'options.fs.read');
514 validateFunction(close, 'options.fs.close');
515 validateFunction(fsync, 'options.fs.fsync');
516 validateFunction(write, 'options.fs.write');
517 validateFunction(writev, 'options.fs.writev');
518 return { open, close, fsync, read, write, writev };
519}
520 
521function readImpl(this: ReadStream, n: number): void {
522 n =
523 this.pos !== undefined
524 ? Math.min(this.end - this.pos + 1, n)
525 : Math.min(this.end - this.bytesRead + 1, n);
526 if (n <= 0) {
527 this.push(null);
528 return;
529 }
530 
531 const buf = Buffer.allocUnsafeSlow(n);
532 const ee = this as unknown as EventEmitter;
533 
534 this[kIsPerformingIO] = true;
535 if (this.fd == null) {
536 this.push(null);
537 return;
538 }
539 this[kFs].read(this.fd, buf, 0, n, this.pos, (er, bytesRead, buf) => {
540 this[kIsPerformingIO] = false;
541 
542 // Tell ._destroy() that it's safe to close the fd now.
543 if (this.destroyed) {
544 ee.emit(kIoDone, er);
545 return;
546 }
547 
548 if (er) {
549 errorOrDestroy(this, er as Error);
550 return;
551 }
552 
553 if (buf == null) {
554 errorOrDestroy(this, new ERR_INVALID_ARG_VALUE('buf', 'null'));
555 return;
556 }
557 
558 if (bytesRead != null && bytesRead > 0) {
559 if (this.pos !== undefined) {
560 this.pos += bytesRead;
561 }
562 
563 this.bytesRead += bytesRead;
564 
565 if (bytesRead !== buf.byteLength) {
566 // Slow path. Shrink to fit.
567 // Copy instead of slice so that we don't retain
568 // large backing buffer for small reads.
569 const dst = Buffer.allocUnsafeSlow(bytesRead);
570 if (Buffer.isBuffer(buf)) {
571 (buf as Buffer).copy(dst, 0, 0, bytesRead);
572 } else {
573 const buffer = Buffer.from(buf.buffer, buf.byteOffset, bytesRead);
574 buffer.copy(dst, 0, 0, bytesRead);
575 }
576 buf = dst;
577 }
578 
579 this.push(buf);
580 } else {
581 this.push(null);
582 }
583 });
584}
585 
586function actualCloseImpl(
587 stream: ReadStream | WriteStream,
588 err: unknown,
589 cb: ErrorOnlyCallback
590): void {
591 if (stream.fd == null) return;
592 stream[kFs].close(stream.fd, (er) => {
593 cb(er || err);
594 });
595 stream.fd = null;
596}
597 
598function closeImpl(
599 stream: ReadStream | WriteStream,
600 err: unknown,
601 cb: ErrorOnlyCallback
602): void {
603 if (stream.fd == null) {
604 cb(err);
605 } else if (stream.flush) {
606 stream[kFs].fsync(stream.fd, (flushErr) => {
607 actualCloseImpl(stream, err || flushErr, cb);
608 });
609 } else {
610 actualCloseImpl(stream, err, cb);
611 }
612}
613 
614function destroyImpl(
615 this: ReadStream | WriteStream,
616 err: unknown,
617 cb: ErrorOnlyCallback
618): void {
619 // Usually for async IO it is safe to close a file descriptor
620 // even when there are pending operations. However, due to platform
621 // differences file IO is implemented using synchronous operations
622 // running in a thread pool. Therefore, file descriptors are not safe
623 // to close while used in a pending read or write operation. Wait for
624 // any pending IO (kIsPerformingIO) to complete (kIoDone).
625 const ee = this as unknown as EventEmitter;
626 if (this[kIsPerformingIO]) {
627 ee.once(kIoDone, (er) => {
628 closeImpl(this, err || er, cb);
629 });
630 } else {
631 closeImpl(this, err, cb);
632 }
633}
634 
635// @ts-expect-error TS2323 Cannot redeclare.
636export function ReadStream(
637 this: ReadStream,
638 path: FilePath,
639 options: ReadStreamOptions = {}
640): ReadStream {
641 if (!(this instanceof ReadStream)) {
642 return new ReadStream(path, options);
643 }
644 if (options === null) {
645 options = {};
646 } else if (typeof options === 'string') {
647 options = { encoding: options };
648 }
649 
650 validateObject(options, 'options');
651 const {
652 encoding = null,
653 autoClose = true,
654 emitClose = true,
655 start = 0,
656 end = Infinity,
657 highWaterMark = 64 * 1024,
658 signal = null,
659 flags = 'r',
660 fd = null,
661 mode = 0o666,
662 fs = kDefaultFsOperations,
663 } = options;
664 const autoDestroy = autoClose;
665 
666 if (
667 encoding !== 'buffer' &&
668 encoding !== null &&
669 !Buffer.isEncoding(encoding)
670 ) {
671 throw new ERR_INVALID_ARG_VALUE('options.encoding', encoding);
672 }
673 validateBoolean(autoClose, 'options.autoClose');
674 validateBoolean(emitClose, 'options.emitClose');
675 validateUint32(start, 'options.start');
676 if (end !== Infinity) {
677 validateUint32(end, 'options.end');
678 if (start > end) {
679 throw new ERR_OUT_OF_RANGE('start', `<= "end" (here: ${end})`, start);
680 }
681 }
682 validateUint32(highWaterMark, 'options.highWaterMark');
683 if (signal != null) {
684 validateAbortSignal(signal, 'options.signal');
685 }
686 
687 validateString(flags, 'options.flags');
688 
689 // We don't actually use the mode in our implementation. Parsing it is
690 // just to ensure that it is validated.
691 parseFileMode(mode, 'options.mode', 0o666);
692 
693 this[kFs] = getValidatedFsOptions(fs);
694 
695 if (fd == null) {
696 this.fd = null;
697 // Path will be ignored when fd is specified, so it can be falsy
698 this.path = toPathIfFileURL(normalizePath(path));
699 this.flags = options.flags === undefined ? 'r' : options.flags;
700 this.mode = options.mode === undefined ? 0o666 : options.mode;
701 } else {
702 if (isFileHandle(fd)) {
703 if (fs !== kDefaultFsOperations) {
704 throw new ERR_METHOD_NOT_IMPLEMENTED('FileHandle with fs');
705 }
706 this[kHandle] = fd as FileHandle;
707 this[kFs] = getValidatedFsOptions(fd as unknown as FsOperations);
708 const ee = fd as unknown as EventEmitter;
709 ee.on('close', () => this.close());
710 this.fd = (fd as FileHandle).fd || null;
711 } else {
712 this.fd = getValidatedFd(fd as number);
713 }
714 }
715 
716 this.start = start;
717 this.end = end;
718 this.pos = start;
719 this.bytesRead = 0;
720 this[kIsPerformingIO] = false;
721 
722 Reflect.apply(Readable, this, [
723 {
724 highWaterMark,
725 encoding,
726 emitClose,
727 autoDestroy,
728 signal,
729 construct: construct.bind(this),
730 read: readImpl.bind(this),
731 destroy: destroyImpl.bind(this),
732 },
733 ]);
734 return this;
735}
736Object.setPrototypeOf(ReadStream.prototype, Readable.prototype);
737Object.setPrototypeOf(ReadStream, Readable);
738 
739Object.defineProperty(ReadStream.prototype, 'autoClose', {
740 get(this: ReadStream): boolean {
741 validateThisInternalField(this, kFs, 'ReadStream');
742 return this._readableState?.autoDestroy || false;
743 },
744 set(this: ReadStream, _val: boolean): void {
745 validateThisInternalField(this, kFs, 'ReadStream');
746 if (this._readableState !== undefined) {
747 this._readableState.autoDestroy = _val;
748 }
749 },
750});
751 
752ReadStream.prototype.close = function (
753 cb: ErrorOnlyCallback = (_err: unknown): void => {}
754): void {
755 if (typeof cb === 'function') eos(this, cb);
756 this.destroy();
757};
758 
759Object.defineProperty(ReadStream.prototype, 'pending', {
760 get(this: ReadStream): boolean {
761 return this.fd === null;
762 },
763 configurable: true,
764});
765 
766// ======================================================================================
767 
768export type WriteStreamOptions = {
769 encoding?: string | undefined;
770 autoClose?: boolean | undefined;
771 autoDestroy?: boolean | undefined;
772 emitClose?: boolean | undefined;
773 start?: number | undefined;
774 highWaterMark?: number | undefined;
775 signal?: AbortSignal | undefined;
776 flags?: string | undefined;
777 fd?: number | FileHandle | undefined;
778 mode?: number | undefined;
779 fs?: FsOperations | undefined;
780 flush?: boolean | undefined;
781};
782 
783declare type WriteVChunk = {
784 chunk: string | NodeJS.ArrayBufferView;
785 encoding: ValidEncoding;
786};
787 
788// @ts-expect-error TS2323 Cannot redeclare.
789export declare class WriteStream extends Writable {
790 fd: number | null;
791 path: string;
792 flags: string;
793 mode: number;
794 flush: boolean;
795 autoClose: boolean;
796 destroyed: boolean;
797 start: number;
798 pos: number;
799 bytesRead: number;
800 bytesWritten: number;
801 [kIsPerformingIO]: boolean;
802 [kFs]: RealizedFsOperations;
803 [kHandle]?: FileHandle;
804 constructor(path: FilePath, options?: WriteStreamOptions);
805 close(cb?: ErrorOnlyCallback): void;
806 destroySoon(): void;
807}
808 
809function writeAll(
810 this: WriteStream,
811 data: string | NodeJS.ArrayBufferView,
812 size: number,
813 pos: number,
814 cb: ErrorOnlyCallback,
815 retries = 0
816) {
817 if (this.fd == null) {
818 return cb(new ERR_INVALID_ARG_VALUE('fd', 'null'));
819 }
820 
821 this[kFs].write(
822 this.fd,
823 data,
824 0,
825 size,
826 pos,
827 (
828 er: unknown,
829 bytesWritten?: number,
830 buffer?: string | NodeJS.ArrayBufferView
831 ) => {
832 // No data currently available and operation should be retried later.
833 if ((er as any)?.code === 'EAGAIN') {
834 er = null;
835 bytesWritten = 0;
836 }
837 
838 if (this.destroyed || er) {
839 return cb(er || new ERR_STREAM_DESTROYED('write'));
840 }
841 
842 // The value should be set but let's suppress the possible
843 // typescript error here.
844 bytesWritten = bytesWritten ?? 0;
845 
846 this.bytesWritten += bytesWritten;
847 
848 retries = bytesWritten ? 0 : retries + 1;
849 size -= bytesWritten;
850 pos += bytesWritten;
851 
852 // Try writing non-zero number of bytes up to 5 times.
853 if (retries > 5) {
854 cb(new ERR_SYSTEM_ERROR('write failed'));
855 } else if (size) {
856 if (buffer == null) {
857 cb(null);
858 return;
859 }
860 const buf =
861 typeof buffer === 'string'
862 ? Buffer.from(buffer as string)
863 : Buffer.from(buffer.buffer, buffer.byteOffset, buffer.byteLength);
864 writeAll.call(this, buf.slice(bytesWritten), size, pos, cb, retries);
865 } else {
866 cb(null);
867 }
868 }
869 );
870}
871 
872function writevAll(
873 this: WriteStream,
874 chunks: NodeJS.ArrayBufferView[],
875 size: number,
876 pos: number,
877 cb: ErrorOnlyCallback,
878 retries = 0
879) {
880 if (this.fd == null) {
881 return cb(new ERR_INVALID_ARG_VALUE('fd', 'null'));
882 }
883 this[kFs].writev(
884 this.fd,
885 chunks,
886 this.pos,
887 (
888 er: unknown,
889 bytesWritten?: number,
890 buffers?: NodeJS.ArrayBufferView[]
891 ) => {
892 // No data currently available and operation should be retried later.
893 if ((er as any)?.code === 'EAGAIN') {
894 er = null;
895 bytesWritten = 0;
896 }
897 
898 if (this.destroyed || er) {
899 return cb(er || new ERR_STREAM_DESTROYED('writev'));
900 }
901 
902 bytesWritten = bytesWritten ?? 0;
903 
904 this.bytesWritten += bytesWritten;
905 
906 retries = bytesWritten ? 0 : retries + 1;
907 size -= bytesWritten;
908 pos += bytesWritten;
909 
910 // Try writing non-zero number of bytes up to 5 times.
911 if (retries > 5) {
912 cb(new ERR_SYSTEM_ERROR('writev failed'));
913 } else if (size) {
914 buffers ??= [];
915 const bufs = buffers.map((b): Buffer => {
916 return Buffer.from(b.buffer, b.byteOffset, b.byteLength);
917 });
918 if (buffers.length === 0) {
919 cb(null);
920 return;
921 }
922 writevAll.call(
923 this,
924 [Buffer.concat(bufs).slice(bytesWritten)],
925 size,
926 pos,
927 cb,
928 retries
929 );
930 } else {
931 cb(null);
932 }
933 }
934 );
935}
936 
937function writeImpl(
938 this: WriteStream,
939 data: string | NodeJS.ArrayBufferView,
940 encodingOrCallback: ValidEncoding | ErrorOnlyCallback | undefined,
941 cb?: ErrorOnlyCallback
942): void {
943 let callback: ErrorOnlyCallback;
944 let encoding: ValidEncoding = null;
945 if (typeof encodingOrCallback === 'function') {
946 callback = encodingOrCallback;
947 } else if (encodingOrCallback != null) {
948 encoding = encodingOrCallback;
949 if (cb !== undefined) callback = cb;
950 }
951 // @ts-expect-error TS2345
952 validateFunction(callback, 'write callback');
953 
954 if (typeof data === 'string') {
955 data = Buffer.from(data, encoding as any);
956 }
957 
958 this[kIsPerformingIO] = true;
959 const ee = this as unknown as EventEmitter;
960 writeAll.call(this, data, data.byteLength, this.pos, (er: unknown) => {
961 this[kIsPerformingIO] = false;
962 if (this.destroyed) {
963 // Tell ._destroy() that it's safe to close the fd now.
964 callback(er);
965 ee.emit(kIoDone, er);
966 return;
967 }
968 
969 callback(er);
970 });
971 
972 if (this.pos !== undefined) this.pos += data.byteLength;
973}
974 
975function writevImpl(
976 this: WriteStream,
977 data: WriteVChunk[],
978 callback: ErrorOnlyCallback
979) {
980 let size = 0;
981 const chunks = data.map((d) => {
982 const chunk = (d as any).chunk;
983 size += chunk.length;
984 return chunk;
985 });
986 
987 this[kIsPerformingIO] = true;
988 writevAll.call(this, chunks, size, this.pos, (er: unknown) => {
989 this[kIsPerformingIO] = false;
990 const ee = this as unknown as EventEmitter;
991 if (this.destroyed) {
992 // Tell ._destroy() that it's safe to close the fd now.
993 callback(er);
994 ee.emit(kIoDone, er);
995 return;
996 }
997 
998 callback(er);
999 });
1000 
1001 if (this.pos !== undefined) this.pos += size;
1002}
1003 
1004// @ts-expect-error TS2323 Cannot redeclare.
1005export function WriteStream(
1006 this: WriteStream,
1007 path: FilePath,
1008 options: WriteStreamOptions = {}
1009): WriteStream {
1010 if (!(this instanceof WriteStream)) {
1011 return new WriteStream(path, options);
1012 }
1013 
1014 if (options === null) {
1015 options = {};
1016 } else if (typeof options === 'string') {
1017 options = { encoding: options };
1018 }
1019 
1020 validateObject(options, 'options');
1021 const {
1022 encoding = null,
1023 autoClose = true,
1024 emitClose = true,
1025 start = 0,
1026 highWaterMark = 64 * 1024,
1027 signal = null,
1028 flags = 'r',
1029 fd = null,
1030 mode = 0o666,
1031 fs = kDefaultFsOperations,
1032 } = options;
1033 let { flush = false } = options;
1034 const autoDestroy = autoClose;
1035 
1036 if (
1037 encoding !== 'buffer' &&
1038 encoding !== null &&
1039 !Buffer.isEncoding(encoding)
1040 ) {
1041 throw new ERR_INVALID_ARG_VALUE('options.encoding', encoding);
1042 }
1043 validateBoolean(autoClose, 'options.autoClose');
1044 validateBoolean(emitClose, 'options.emitClose');
1045 if (flush === null) flush = false;
1046 validateBoolean(flush, 'options.flush');
1047 validateUint32(start, 'options.start');
1048 validateUint32(highWaterMark, 'options.highWaterMark');
1049 if (signal != null) {
1050 validateAbortSignal(signal, 'options.signal');
1051 }
1052 
1053 validateString(flags, 'options.flags');
1054 
1055 // We don't actually use the mode in our implementation. Parsing it is
1056 // just to ensure that it is validated.
1057 parseFileMode(mode, 'options.mode', 0o666);
1058 
1059 this[kFs] = getValidatedFsOptions(fs);
1060 this.flush = flush;
1061 
1062 if (fd == null) {
1063 this.fd = null;
1064 // Path will be ignored when fd is specified, so it can be falsy
1065 this.path = toPathIfFileURL(normalizePath(path));
1066 this.flags = options.flags === undefined ? 'r' : options.flags;
1067 this.mode = options.mode === undefined ? 0o666 : options.mode;
1068 } else {
1069 if (isFileHandle(fd)) {
1070 if (fs !== kDefaultFsOperations) {
1071 throw new ERR_METHOD_NOT_IMPLEMENTED('FileHandle with fs');
1072 }
1073 this[kHandle] = fd as FileHandle;
1074 this[kFs] = getValidatedFsOptions(fd as unknown as FsOperations);
1075 const ee = fd as unknown as EventEmitter;
1076 ee.on('close', () => this.close());
1077 this.fd = (fd as FileHandle).fd || null;
1078 } else {
1079 this.fd = getValidatedFd(fd as number);
1080 }
1081 }
1082 
1083 this.start = start;
1084 this.pos = start;
1085 this.bytesRead = 0;
1086 this.bytesWritten = 0;
1087 this[kIsPerformingIO] = false;
1088 
1089 Reflect.apply(Writable, this, [
1090 {
1091 highWaterMark,
1092 encoding,
1093 emitClose,
1094 autoDestroy,
1095 signal,
1096 construct: construct.bind(this),
1097 write: writeImpl.bind(this),
1098 writev: writevImpl.bind(this),
1099 destroy: destroyImpl.bind(this),
1100 },
1101 ]);
1102 
1103 if (encoding != null && encoding !== 'buffer') {
1104 (this as unknown as Writable).setDefaultEncoding(
1105 encoding as BufferEncoding
1106 );
1107 }
1108 
1109 return this;
1110}
1111Object.setPrototypeOf(WriteStream.prototype, Writable.prototype);
1112Object.setPrototypeOf(WriteStream, Writable);
1113 
1114Object.defineProperty(WriteStream.prototype, 'autoClose', {
1115 get(this: WriteStream): boolean {
1116 validateThisInternalField(this, kFs, 'WriteStream');
1117 return this._writableState?.autoDestroy || false;
1118 },
1119 set(this: WriteStream, val: boolean): void {
1120 validateThisInternalField(this, kFs, 'WriteStream');
1121 if (this._writableState !== undefined) {
1122 this._writableState.autoDestroy = val;
1123 }
1124 },
1125});
1126 
1127WriteStream.prototype.close = function (
1128 this: WriteStream,
1129 cb: ErrorOnlyCallback = (_err: unknown): void => {}
1130): void {
1131 const writable = this as unknown as Writable;
1132 const ee = this as unknown as EventEmitter;
1133 if (cb) {
1134 if (writable.closed) {
1135 queueMicrotask(() => cb(null));
1136 return;
1137 }
1138 ee.on('close', cb);
1139 }
1140 
1141 // If we are not autoClosing, we should call
1142 // destroy on 'finish'.
1143 if (!this.autoClose) {
1144 ee.on('finish', () => writable.destroy());
1145 }
1146 
1147 // We use end() instead of destroy() because of
1148 // https://github.com/nodejs/node/issues/2006
1149 writable.end();
1150};
1151 
1152// There is no shutdown() for files.
1153WriteStream.prototype.destroySoon = WriteStream.prototype.end;
1154 
1155Object.defineProperty(WriteStream.prototype, 'pending', {
1156 get(this: WriteStream): boolean {
1157 return this.fd === null;
1158 },
1159 configurable: true,
1160});
1161 
1162export function createReadStream(
1163 path: FilePath,
1164 options: ReadStreamOptions = {}
1165): ReadStream {
1166 return new ReadStream(path, options);
1167}
1168export function createWriteStream(
1169 path: FilePath,
1170 options: WriteStreamOptions = {}
1171): WriteStream {
1172 return new WriteStream(path, options);
1173}