/// import {Duplex} from "node:stream" import {PromiseReadable} from "promise-readable" import {PromiseWritable} from "promise-writable" interface DuplexStream extends Duplex { closed: boolean } export class PromiseDuplex extends PromiseReadable implements AsyncIterable { readonly readable: PromiseReadable readonly writable: PromiseWritable readonly isPromiseReadable: boolean = true readonly isPromiseWritable: boolean = true constructor(readonly stream: TDuplex) { super(stream) this.readable = new PromiseReadable(stream) this.writable = new PromiseWritable(stream) } // PromiseReadable read(size?: number): Promise { return this.readable.read(size) } readAll(): Promise { return this.readable.readAll() } setEncoding(encoding: BufferEncoding): this { this.readable.setEncoding(encoding) return this } iterate(size?: number): AsyncIterableIterator { return this.readable.iterate(size) } [Symbol.asyncIterator](): AsyncIterableIterator { return this.readable[Symbol.asyncIterator]() } // PromiseWritable write(chunk: string | Buffer, encoding?: BufferEncoding): Promise { return this.writable.write(chunk, encoding) } writeAll(content: string | Buffer, chunkSize?: number): Promise { return this.writable.writeAll(content, chunkSize) } end(): Promise { return this.writable.end() } // PromiseDuplex once(event: "close" | "end" | "error" | "finish"): Promise once(event: "open"): Promise once(event: "pipe" | "unpipe"): Promise once(event: string): Promise { const stream = this.stream return new Promise((resolve, reject) => { if (this.readable._errored) { const err = this.readable._errored this.readable._errored = undefined return reject(err) } if (this.writable._errored) { const err = this.writable._errored this.writable._errored = undefined return reject(err) } if (stream.closed) { if (event === "close") { return resolve() } else { return reject(new Error(`once ${event} after close`)) } } if (stream.destroyed) { if (event === "close" || event === "end" || event === "finish") { return resolve() } else { return reject(new Error(`once ${event} after destroy`)) } } const eventHandler = event !== "end" && event !== "finish" && event !== "error" ? (argument: any) => { removeListeners() resolve(argument) } : undefined const closeHandler = () => { removeListeners() resolve() } const endHandler = event !== "close" ? () => { removeListeners() resolve() } : undefined const errorHandler = (err: Error) => { this.readable._errored = undefined this.writable._errored = undefined removeListeners() reject(err) } const finishHandler = event !== "close" ? () => { removeListeners() resolve() } : undefined const removeListeners = () => { if (eventHandler) { stream.removeListener(event, eventHandler) } stream.removeListener("close", closeHandler) if (endHandler) { stream.removeListener("end", endHandler) } stream.removeListener("error", errorHandler) if (finishHandler) { stream.removeListener("finish", finishHandler) } } if (eventHandler) { stream.on(event, eventHandler) } stream.on("close", closeHandler) if (endHandler) { stream.on("end", endHandler) } if (finishHandler) { stream.on("finish", finishHandler) } stream.on("error", errorHandler) }) } destroy(): this { if (this.readable) { this.readable.destroy() } if (this.writable) { this.writable.destroy() } return this } } export default PromiseDuplex