From 8635e9535efa08eb92fcda137ba77e4f873a485e Mon Sep 17 00:00:00 2001 From: Jan Date: Tue, 25 Aug 2026 11:47:46 -0700 Subject: [PATCH] Honor backpressure when finalizing ZIP archives --- .../zip/zip-archive-output-stream.js | 97 +++++++++++-------- test/zip-archive-output-stream.js | 79 +++++++++++++++ 2 files changed, 134 insertions(+), 42 deletions(-) diff --git a/lib/archivers/zip/zip-archive-output-stream.js b/lib/archivers/zip/zip-archive-output-stream.js index 33711fc8..2b4bfc2f 100644 --- a/lib/archivers/zip/zip-archive-output-stream.js +++ b/lib/archivers/zip/zip-archive-output-stream.js @@ -47,12 +47,14 @@ export default class ZipArchiveOutputStream extends ArchiveOutputStream { this.options = _options; this._entry = null; this._entries = []; + this._centralDirectoryIndex = 0; this._archive = { centralLength: 0, centralOffset: 0, comment: "", finish: false, finished: false, + finalizing: false, processing: false, forceZip64: _options.forceZip64, forceLocalTime: _options.forceLocalTime, @@ -109,20 +111,36 @@ export default class ZipArchiveOutputStream extends ArchiveOutputStream { } _finish() { + if (this._archive.finalizing || this._archive.finished) { + return; + } this._archive.centralOffset = this.offset; - this._entries.forEach( - function (ae) { - this._writeCentralFileHeader(ae); - }.bind(this), - ); + this._archive.finalizing = true; + this._archive.finish = true; + this._centralDirectoryIndex = 0; + this._writeCentralDirectory(); + } + + _writeCentralDirectory() { + while (this._centralDirectoryIndex < this._entries.length) { + var index = this._centralDirectoryIndex; + var ae = this._entries[index]; + this._entries[index] = undefined; + this._centralDirectoryIndex += 1; + if (!this._writeCentralFileHeader(ae)) { + this.once("drain", this._writeCentralDirectory.bind(this)); + return; + } + } this._archive.centralLength = this.offset - this._archive.centralOffset; if (this.isZip64()) { this._writeCentralDirectoryZip64(); } this._writeCentralDirectoryEnd(); this._archive.processing = false; - this._archive.finish = true; + this._archive.finalizing = false; this._archive.finished = true; + this._entries = []; this.end(); } @@ -248,22 +266,6 @@ export default class ZipArchiveOutputStream extends ArchiveOutputStream { ); ae.setExtra(extraBuf); } - // signature - this.write(getLongBytes(SIG_CFH)); - // version made by - this.write(getShortBytes((ae.getPlatform() << 8) | VERSION_MADEBY)); - // version to extract and general bit flag - this.write(getShortBytes(ae.getVersionNeededToExtract())); - this.write(gpb.encode()); - // compression method - this.write(getShortBytes(method)); - // datetime - this.write(getLongBytes(ae.getTimeDos())); - // crc32 checksum - this.write(getLongBytes(ae.getCrc())); - // sizes - this.write(getLongBytes(compressedSize)); - this.write(getLongBytes(size)); var name = ae.getName(); var comment = ae.getComment(); var extra = ae.getCentralDirectoryExtra(); @@ -271,26 +273,37 @@ export default class ZipArchiveOutputStream extends ArchiveOutputStream { name = Buffer.from(name); comment = Buffer.from(comment); } - // name length - this.write(getShortBytes(name.length)); - // extra length - this.write(getShortBytes(extra.length)); - // comments length - this.write(getShortBytes(comment.length)); - // disk number start - this.write(SHORT_ZERO); - // internal attributes - this.write(getShortBytes(ae.getInternalAttributes())); - // external attributes - this.write(getLongBytes(ae.getExternalAttributes())); - // relative offset of LFH - this.write(getLongBytes(fileOffset)); - // name - this.write(name); - // extra - this.write(extra); - // comment - this.write(comment); + if (!Buffer.isBuffer(name)) { + name = Buffer.from(name); + } + if (!Buffer.isBuffer(comment)) { + comment = Buffer.from(comment); + } + var header = Buffer.concat( + [ + getLongBytes(SIG_CFH), + getShortBytes((ae.getPlatform() << 8) | VERSION_MADEBY), + getShortBytes(ae.getVersionNeededToExtract()), + gpb.encode(), + getShortBytes(method), + getLongBytes(ae.getTimeDos()), + getLongBytes(ae.getCrc()), + getLongBytes(compressedSize), + getLongBytes(size), + getShortBytes(name.length), + getShortBytes(extra.length), + getShortBytes(comment.length), + SHORT_ZERO, + getShortBytes(ae.getInternalAttributes()), + getLongBytes(ae.getExternalAttributes()), + getLongBytes(fileOffset), + name, + extra, + comment, + ], + 46 + name.length + extra.length + comment.length, + ); + return this.write(header); } _writeDataDescriptor(ae) { diff --git a/test/zip-archive-output-stream.js b/test/zip-archive-output-stream.js index 6fc1781d..3599173a 100644 --- a/test/zip-archive-output-stream.js +++ b/test/zip-archive-output-stream.js @@ -120,5 +120,84 @@ describe("ZipArchiveOutputStream", function () { archive.pipe(testStream); archive.entry(entry, createReadStream("test/fixtures/test.txt")).finish(); }); + it("should honor backpressure while writing the central directory", function (done) { + this.timeout(10000); + var entryCount = 2000; + var archive = new BackpressureZipArchiveOutputStream(); + var chunks = []; + archive.on("data", function (chunk) { + chunks.push(chunk); + }); + archive.on("end", function () { + var output = Buffer.concat(chunks); + var names = getCentralDirectoryNames(output); + assert.equal(archive.centralDirectoryWrites, entryCount); + assert.equal(archive.drainEvents, entryCount); + assert.lengthOf(names, entryCount); + assert.equal(names[0], "entry-0.txt"); + assert.equal( + names[entryCount - 1], + "entry-" + (entryCount - 1) + ".txt", + ); + done(); + }); + for (var i = 0; i < entryCount; i++) { + var entry = new ZipArchiveEntry("entry-" + i + ".txt"); + archive.entry(entry, Buffer.alloc(0)); + } + archive.finish(); + }); }); }); + +class BackpressureZipArchiveOutputStream extends ZipArchiveOutputStream { + constructor() { + super(); + this.centralDirectoryWrites = 0; + this.drainEvents = 0; + this.waitingForDrain = false; + } + + write(chunk, callback) { + var isCentralDirectoryHeader = + Buffer.isBuffer(chunk) && + chunk.length >= 4 && + chunk.readUInt32LE(0) === 0x02014b50; + if (!isCentralDirectoryHeader) { + return super.write(chunk, callback); + } + assert.isFalse(this.waitingForDrain); + super.write(chunk, callback); + this.centralDirectoryWrites += 1; + this.waitingForDrain = true; + setImmediate( + function () { + this.waitingForDrain = false; + this.drainEvents += 1; + this.emit("drain"); + }.bind(this), + ); + return false; + } +} + +function getCentralDirectoryNames(output) { + var endSignature = Buffer.from([0x50, 0x4b, 0x05, 0x06]); + var endOffset = output.lastIndexOf(endSignature); + assert.isAtLeast(endOffset, 0); + var entryCount = output.readUInt16LE(endOffset + 10); + var centralLength = output.readUInt32LE(endOffset + 12); + var centralOffset = output.readUInt32LE(endOffset + 16); + var offset = centralOffset; + var names = []; + for (var i = 0; i < entryCount; i++) { + assert.equal(output.readUInt32LE(offset), 0x02014b50); + var nameLength = output.readUInt16LE(offset + 28); + var extraLength = output.readUInt16LE(offset + 30); + var commentLength = output.readUInt16LE(offset + 32); + names.push(output.toString("utf8", offset + 46, offset + 46 + nameLength)); + offset += 46 + nameLength + extraLength + commentLength; + } + assert.equal(offset, centralOffset + centralLength); + return names; +}