From 4d29116d7ba7846d2f0bfa1fadfd68ac7858f18d Mon Sep 17 00:00:00 2001 From: x0Lazarus <113273587+x0Lazarus@users.noreply.github.com> Date: Sat, 26 Sep 2026 11:05:57 -0700 Subject: [PATCH] fix(csv-parse): finish limited streams without write-after-end errors --- packages/csv-parse/lib/index.js | 25 ++++---- .../csv-parse/test/api.stream.finished.ts | 63 ++++++++++++++++++- .../csv-parse/test/api.stream.iterator.ts | 17 +++++ 3 files changed, 90 insertions(+), 15 deletions(-) diff --git a/packages/csv-parse/lib/index.js b/packages/csv-parse/lib/index.js index 4aefe2a2..76b1d042 100644 --- a/packages/csv-parse/lib/index.js +++ b/packages/csv-parse/lib/index.js @@ -28,6 +28,7 @@ class Parser extends Transform { // Implementation of `Transform._transform` _transform(buf, _, callback) { if (this.state.stop === true) { + callback(); return; } const err = this.api.parse( @@ -38,27 +39,27 @@ class Parser extends Transform { }, () => { this.push(null); - this.end(); - // Fix #333 and break #410 - // ko: api.stream.iterator.coffee - // ko with v21.4.0, ok with node v20.5.1: api.stream.finished # aborted (with generate()) - // ko: api.stream.finished # aborted (with Readable) - // this.destroy() - // Fix #410 and partially break #333 - // ok: api.stream.iterator.coffee - // ok: api.stream.finished # aborted (with generate()) - // broken: api.stream.finished # aborted (with Readable) - this.on("end", this.destroy); + // Let buffered records be consumed before closing the writable side. + // The source can still write chunks before the readable end event. + this.on("end", () => { + this.end(() => this.destroy()); + }); }, ); if (err !== undefined) { this.state.stop = true; } - callback(err); + // Hold back upstream input until the final records have been consumed. + if (err === undefined && this.state.stop && !this.readableEnded) { + this.once("end", callback); + } else { + callback(err); + } } // Implementation of `Transform._flush` _flush(callback) { if (this.state.stop === true) { + callback(); return; } const err = this.api.parse( diff --git a/packages/csv-parse/test/api.stream.finished.ts b/packages/csv-parse/test/api.stream.finished.ts index 14d509bd..8e93f021 100644 --- a/packages/csv-parse/test/api.stream.finished.ts +++ b/packages/csv-parse/test/api.stream.finished.ts @@ -35,7 +35,7 @@ describe("API stream.finished", function () { records.length.should.eql(3); }); - it.skip("aborted (with Readable)", async function () { + it("aborted (with Readable)", async function () { // See https://github.com/adaltas/node-csv/issues/333 // See https://github.com/adaltas/node-csv/issues/410 // Prevent `Error [ERR_STREAM_PREMATURE_CLOSE]: Premature close` @@ -55,10 +55,67 @@ describe("API stream.finished", function () { records.push(record); } }); - await stream.finished(parser); - records.length.should.eql(3); + try { + await stream.finished(parser); + records.length.should.eql(3); + } finally { + reader.destroy(); + } }); + for (const option of ["to", "to_line"]) { + it(`finishes after ${option} while the consumer pauses`, async function () { + const records: string[][] = []; + let produced = 0; + const reader = Readable.from( + (function* () { + for (let i = 0; i < 1000; i++) { + produced++; + yield `${i},value${i}\n`; + } + })(), + { highWaterMark: 1 }, + ); + const options = + option === "to" + ? { to: 1, highWaterMark: 1 } + : { to_line: 1, highWaterMark: 1 }; + const parser = parse(options); + const done = stream.finished(parser); + parser.on("data", (record) => { + records.push(record); + parser.pause(); + setImmediate(() => parser.resume()); + }); + reader.pipe(parser); + try { + await done; + records.should.eql([["0", "value0"]]); + produced.should.be.below(10); + } finally { + reader.destroy(); + } + }); + + it(`resolves after ${option} with separate input chunks`, async function () { + const records: string[][] = []; + const parser = Readable.from(["a,b\n", "c,d\n", "e,f\n", "g,h\n"]).pipe( + parse({ [option]: 2 }), + ); + parser.on("readable", () => { + let record; + while ((record = parser.read()) !== null) { + records.push(record); + } + }); + await stream.finished(parser); + records.should.eql([ + ["a", "b"], + ["c", "d"], + ]); + }); + } + it("rejected on error", async function () { const parser = parse({ to_line: 3 }); parser.write("a,b,c\n"); diff --git a/packages/csv-parse/test/api.stream.iterator.ts b/packages/csv-parse/test/api.stream.iterator.ts index 8c46da24..f4dd6730 100644 --- a/packages/csv-parse/test/api.stream.iterator.ts +++ b/packages/csv-parse/test/api.stream.iterator.ts @@ -1,8 +1,25 @@ import "should"; +import { Readable } from "node:stream"; import { generate } from "csv-generate"; import { parse } from "../lib/index.js"; describe("API stream.iterator", function () { + for (const option of ["to", "to_line"]) { + it(`stops at ${option} with separate input chunks`, async function () { + const parser = Readable.from(["a,b\n", "c,d\n", "e,f\n", "g,h\n"]).pipe( + parse({ [option]: 2 }), + ); + const records = []; + for await (const record of parser) { + records.push(record); + } + records.should.eql([ + ["a", "b"], + ["c", "d"], + ]); + }); + } + it("classic", async function () { const parser = generate({ length: 10 }).pipe(parse()); const records = [];