diff --git a/doc/api/fs.md b/doc/api/fs.md index d6eb4208ef0a..4528f7a63e2f 100644 --- a/doc/api/fs.md +++ b/doc/api/fs.md @@ -8563,12 +8563,13 @@ Close the stream gracefully, flushing the internal buffer before closing. * `callback` {Function} * `err` {Error|null} An error if the flush failed, otherwise `null`. -Writes the current buffer to the file if a write was not in progress. Do -nothing if `minLength` is zero or if it is already writing. +Writes the current buffer to the file. The callback is invoked after pending +writes complete. #### `utf8Stream.flushSync()` -Flushes the buffered data synchronously. This is a costly operation. +Flushes the buffered data synchronously. This is a costly operation. An +`ERR_INVALID_STATE` error is thrown if the stream is currently writing. #### `utf8Stream.fsync` diff --git a/lib/internal/streams/fast-utf8-stream.js b/lib/internal/streams/fast-utf8-stream.js index a11fbf3bbc0e..1a0266fa5738 100644 --- a/lib/internal/streams/fast-utf8-stream.js +++ b/lib/internal/streams/fast-utf8-stream.js @@ -7,6 +7,7 @@ const { ArrayPrototypePush, MathMax, + Symbol, SymbolDispose, } = primordials; @@ -58,6 +59,7 @@ const kMaxWrite = 16 * 1024; const kContentModeBuffer = 'buffer'; const kContentModeUtf8 = 'utf8'; const kNullPrototype = { __proto__: null }; +const kFlush = Symbol('kFlush'); // Utf8Stream is a port of the original SonicBoom module // (https://github.com/pinojs/sonic-boom) that provides a fast and efficient @@ -71,7 +73,9 @@ class Utf8Stream extends EventEmitter { #ending = false; #reopening = false; #asyncDrainScheduled = false; - #flushPending = false; + #flushPending = 0; + #flushInProgress = 0; + #emittingFlush = false; #hwm = 16387; // 16 KB #file = null; #destroyed = false; @@ -228,7 +232,11 @@ class Utf8Stream extends EventEmitter { }); if (this.#periodicFlush !== 0) { - this.#periodicFlushTimer = setInterval(() => this.flush(null), this.#periodicFlush); + this.#periodicFlushTimer = setInterval(() => { + if (this.#flushPending === 0 && this.#flushInProgress === 0) { + this.flush(); + } + }, this.#periodicFlush); this.#periodicFlushTimer.unref(); } } @@ -263,7 +271,9 @@ class Utf8Stream extends EventEmitter { } if (this.#opening) { - this.once('ready', () => this.reopen(file)); + this.once('ready', () => { + if (!this.#destroyed) this.reopen(file); + }); return; } @@ -306,7 +316,7 @@ class Utf8Stream extends EventEmitter { if (this.#opening) { this.once('ready', () => { - this.end(); + if (!this.#destroyed) this.end(); }); return; } @@ -324,7 +334,7 @@ class Utf8Stream extends EventEmitter { if (this.#len > 0 && this.#fd >= 0) { this.#actualWrite(); } else { - this.#actualClose(); + this.#finishEnding(); } } @@ -332,7 +342,11 @@ class Utf8Stream extends EventEmitter { if (this.#destroyed) { return; } + const opening = this.#opening; this.#actualClose(); + if (opening && this.#flushPending > 0) { + this.emit('error', new ERR_INVALID_STATE('Utf8Stream is destroyed')); + } } /** @type {number} */ @@ -426,6 +440,18 @@ class Utf8Stream extends EventEmitter { } } + if (this.#destroyed) { + this.#writing = false; + if (this.#flushPending > 0) { + if (this.#len === 0) { + this.#scheduleDrain(); + } else { + this.emit('error', new ERR_INVALID_STATE('Utf8Stream is destroyed')); + } + } + return; + } + if (this.#fsync) { this.#fs.fsyncSync(this.#fd); } @@ -435,25 +461,19 @@ class Utf8Stream extends EventEmitter { this.#writing = false; this.#reopening = false; this.reopen(); - } else if (len > this.#minLength) { + } else if (len > 0 && + (len > this.#minLength || this.#flushPending)) { this.#actualWrite(); } else if (this.#ending) { if (len > 0) { this.#actualWrite(); } else { this.#writing = false; - this.#actualClose(); + this.#finishEnding(); } } else { this.#writing = false; - if (this.#sync) { - if (!this.#asyncDrainScheduled) { - this.#asyncDrainScheduled = true; - process.nextTick(() => this.#emitDrain()); - } - } else { - this.emit('drain'); - } + this.#scheduleDrain(); } } @@ -502,12 +522,13 @@ class Utf8Stream extends EventEmitter { } // start - if ((!this.#writing && this.#len > this.#minLength) || this.#flushPending) { + if (!this.#writing && + (this.#len > this.#minLength || this.#flushPending)) { this.#actualWrite(); } else if (reopening && !this.#writing) { // Do not emit 'drain' if a 'ready' listener started a write: // #release() will emit the real 'drain' when that write completes. - process.nextTick(() => this.emit('drain')); + process.nextTick(() => this.#emitDrain()); } }; @@ -533,16 +554,59 @@ class Utf8Stream extends EventEmitter { } } + #scheduleDrain() { + if (this.#sync) { + if (!this.#asyncDrainScheduled) { + this.#asyncDrainScheduled = true; + process.nextTick(() => this.#emitDrain()); + } + } else { + this.#emitDrain(); + } + } + #emitDrain() { + this.#asyncDrainScheduled = false; + if (this.#writing) return; + this.#emitFlush(); + if (this.#writing || this.#destroyed) return; const hasListeners = this.listenerCount('drain') > 0; if (!hasListeners) return; - this.#asyncDrainScheduled = false; this.emit('drain'); } + #finishEnding() { + if (!this.#ending || this.#writing || this.#destroyed) return; + if (this.#len > 0) { + this.#actualWrite(); + return; + } + if (this.#flushPending > 0) this.#emitFlush(); + if (this.#flushPending === 0 && this.#flushInProgress === 0 && + !this.#writing && this.#len === 0) { + this.#actualClose(); + } + } + + #emitFlush() { + if (this.#emittingFlush) return; + this.#emittingFlush = true; + try { + this.emit(kFlush); + } finally { + this.#emittingFlush = false; + } + } + #actualClose() { + if (this.#destroyed) return; + if (this.#fd === -1) { - this.once('ready', () => this.#actualClose()); + this.#destroyed = true; + this.once('ready', () => { + this.#destroyed = false; + this.#actualClose(); + }); return; } @@ -627,7 +691,11 @@ class Utf8Stream extends EventEmitter { throw new ERR_INVALID_STATE('Invalid file descriptor'); } - if (!this.#writing && this.#writingBuf.length > 0) { + if (this.#writing) { + throw new ERR_INVALID_STATE('Cannot flush while a write is in progress'); + } + + if (this.#writingBuf.length > 0) { this.#bufs.unshift([this.#writingBuf]); this.#writingBuf = kEmptyBuffer; } @@ -665,7 +733,11 @@ class Utf8Stream extends EventEmitter { throw new ERR_INVALID_STATE('Invalid file descriptor'); } - if (!this.#writing && this.#writingBuf.length > 0) { + if (this.#writing) { + throw new ERR_INVALID_STATE('Cannot flush while a write is in progress'); + } + + if (this.#writingBuf.length > 0) { this.#bufs.unshift(this.#writingBuf); this.#writingBuf = ''; } @@ -701,37 +773,58 @@ class Utf8Stream extends EventEmitter { } #callFlushCallbackOnDrain(cb) { - this.#flushPending = true; + this.#flushPending++; + let waiting = true; + let completed = false; + let flushing = false; + const stopWaiting = () => { + if (!waiting) return false; + waiting = false; + this.off(kFlush, onDrain); + this.off('error', onError); + return true; + }; + const complete = (err, finish = true) => { + if (completed) return; + completed = true; + if (flushing) this.#flushInProgress--; + try { + cb(err); + } finally { + if (finish) this.#finishEnding(); + } + }; const onDrain = () => { + if (!stopWaiting()) return; + this.#flushPending--; + this.#flushInProgress++; + flushing = true; // Only if _fsync is false to avoid double fsync - if (!this.#fsync && !this.#destroyed) { + if (!this.#fsync && !this.#destroyed && + this.#fd !== 1 && this.#fd !== 2) { try { this.#fs.fsync(this.#fd, (err) => { - this.#flushPending = false; // If the fd is closed, we ignore the error. if (err?.code === 'EBADF') { - cb(); + complete(); return; } - cb(err); + complete(err); }); } catch (err) { - this.#flushPending = false; - cb(err); + complete(err); } } else { - this.#flushPending = false; - cb(); + complete(); } - this.off('error', onError); }; const onError = (err) => { - this.#flushPending = false; - cb(err); - this.off('drain', onDrain); + if (!stopWaiting()) return; + this.#flushPending--; + complete(err, false); }; - this.once('drain', onDrain); + this.once(kFlush, onDrain); this.once('error', onError); } @@ -748,7 +841,7 @@ class Utf8Stream extends EventEmitter { throw error; } - if (this.#minLength <= 0) { + if (this.#minLength <= 0 && !this.#writing && this.#len === 0) { cb?.(); return; } @@ -782,7 +875,7 @@ class Utf8Stream extends EventEmitter { throw error; } - if (this.#minLength <= 0) { + if (this.#minLength <= 0 && !this.#writing && this.#len === 0) { cb?.(); return; } diff --git a/test/parallel/test-fastutf8stream-flush-sync.js b/test/parallel/test-fastutf8stream-flush-sync.js index a48dac8e08e9..8630cf1ff94b 100644 --- a/test/parallel/test-fastutf8stream-flush-sync.js +++ b/test/parallel/test-fastutf8stream-flush-sync.js @@ -67,7 +67,7 @@ function runTests(sync) { const stream = new Utf8Stream({ fd, sync: false, - minLength: 0, + minLength: 1000, fs: fsOverride, }); @@ -85,6 +85,33 @@ function runTests(sync) { })); } +for (const contentMode of ['utf8', 'buffer']) { + const dest = getTempFile(); + const fd = openSync(dest, 'w'); + const stream = new Utf8Stream({ + contentMode, + fd, + minLength: 0, + sync: false, + }); + const text = `${contentMode} asynchronous write\n`; + const data = contentMode === 'buffer' ? Buffer.from(text) : text; + + stream.on('ready', common.mustCall(() => { + assert.ok(stream.write(data)); + assert.throws( + () => stream.flushSync(), + { code: 'ERR_INVALID_STATE' }, + ); + stream.flush(common.mustSucceed(() => { + stream.end(); + readFile(dest, 'utf8', common.mustSucceed((contents) => { + assert.strictEqual(contents, text); + })); + })); + })); +} + { const dest = getTempFile(); const fd = openSync(dest, 'w'); diff --git a/test/parallel/test-fastutf8stream-flush.js b/test/parallel/test-fastutf8stream-flush.js index 195d5e83c9c7..c36b0d3554e5 100644 --- a/test/parallel/test-fastutf8stream-flush.js +++ b/test/parallel/test-fastutf8stream-flush.js @@ -4,6 +4,7 @@ const common = require('../common'); const tmpdir = require('../common/tmpdir'); const assert = require('node:assert'); const { + open, openSync, readFile, writeFileSync, @@ -25,6 +26,173 @@ function getTempFile() { runTests(false); runTests(true); +{ + const dest = getTempFile(); + const fd = openSync(dest, 'w'); + const stream = new Utf8Stream({ fd, minLength: 10, sync: false }); + let flushed = false; + + assert.ok(stream.write('asynchronous flush\n')); + assert.ok(stream.write('queued\n')); + assert.strictEqual(stream.writing, true); + stream.flush(common.mustSucceed(() => { + flushed = true; + assert.strictEqual(stream.writing, false); + readFile(dest, 'utf8', common.mustSucceed((data) => { + assert.strictEqual(data, 'asynchronous flush\nqueued\n'); + stream.end(); + })); + })); + assert.strictEqual(flushed, false); +} + +{ + const dest = getTempFile(); + const fd = openSync(dest, 'w'); + const stream = new Utf8Stream({ fd, minLength: 0, sync: false }); + + assert.ok(stream.write('flush before end\n')); + stream.flush(common.mustSucceed()); + stream.on('finish', common.mustCall(() => { + readFile(dest, 'utf8', common.mustSucceed((data) => { + assert.strictEqual(data, 'flush before end\n'); + })); + })); + stream.end(); +} + +{ + const dest = getTempFile(); + const fd = openSync(dest, 'w'); + const writeError = new Error('write failed'); + const stream = new Utf8Stream({ + fd, + minLength: 0, + sync: false, + fs: { + write: common.mustCall((_fd, _data, _encoding, callback) => { + process.nextTick(callback, writeError); + }), + }, + }); + + stream.on('error', common.mustCall((error) => { + assert.strictEqual(error, writeError); + })); + assert.ok(stream.write('failed write before end\n')); + stream.flush(common.mustCall((error) => { + assert.strictEqual(error, writeError); + })); + stream.end(); +} + +{ + const dest = getTempFile(); + const fd = openSync(dest, 'w'); + let fsyncCalls = 0; + const stream = new Utf8Stream({ + fd, + minLength: 10, + sync: false, + fs: { + fsync: common.mustCall((_fd, callback) => { + fsyncCalls++; + if (fsyncCalls === 1) { + assert.ok(stream.write('late\n')); + } + process.nextTick(callback); + }, 2), + }, + }); + + assert.ok(stream.write('initial write\n')); + stream.flush(common.mustSucceed()); + stream.on('finish', common.mustCall(() => { + readFile(dest, 'utf8', common.mustSucceed((data) => { + assert.strictEqual(data, 'initial write\nlate\n'); + })); + })); + stream.end(); +} + +{ + const dest = getTempFile(); + const fd = openSync(dest, 'w'); + const flushError = new Error('flush failed'); + let fsyncCalls = 0; + const stream = new Utf8Stream({ + fd, + minLength: 0, + sync: false, + fs: { + fsync: common.mustCall((_fd, callback) => { + fsyncCalls++; + process.nextTick(callback, fsyncCalls === 1 ? flushError : null); + }, 2), + }, + }); + + assert.ok(stream.write('failed flush before end\n')); + stream.flush(common.mustCall((error) => { + assert.strictEqual(error, flushError); + })); + stream.on('close', common.mustCall()); + stream.end(); +} + +{ + const dest = getTempFile(); + const stream = new Utf8Stream({ + dest, + fs: { + open(...args) { + setImmediate(() => open(...args)); + }, + }, + }); + + stream.flush(common.mustSucceed()); + stream.on('close', common.mustCall()); + stream.end(); +} + +{ + const dest = getTempFile(); + const stream = new Utf8Stream({ + dest, + fs: { + open(...args) { + setImmediate(() => open(...args)); + }, + }, + }); + + stream.flush(common.mustCall((error) => { + assert.strictEqual(error?.code, 'ERR_INVALID_STATE'); + })); + stream.on('close', common.mustCall()); + stream.destroy(); + stream.flush(common.mustCall((error) => { + assert.strictEqual(error?.code, 'ERR_INVALID_STATE'); + })); +} + +{ + const dest = getTempFile(); + const stream = new Utf8Stream({ + dest, + fs: { + open(...args) { + setImmediate(() => open(...args)); + }, + }, + }); + + stream.on('close', common.mustCall()); + stream.end(); + stream.destroy(); +} + function runTests(sync) { { const dest = getTempFile(); diff --git a/test/parallel/test-fastutf8stream-periodicflush.js b/test/parallel/test-fastutf8stream-periodicflush.js index f7029a213502..4889fa8b19eb 100644 --- a/test/parallel/test-fastutf8stream-periodicflush.js +++ b/test/parallel/test-fastutf8stream-periodicflush.js @@ -75,4 +75,56 @@ function runTests(sync) { stream.destroy(); } + + { + const dest = getTempFile(); + const fd = openSync(dest, 'w'); + const stream = new Utf8Stream({ + fd, + sync, + minLength: 5000, + periodicFlush: common.platformTimeout(10), + }); + const timeout = setTimeout( + common.mustNotCall('periodic flush did not complete'), + common.platformTimeout(1000), + ); + + stream.once('drain', common.mustCall(() => { + clearTimeout(timeout); + stream.destroy(); + readFile(dest, 'utf8', common.mustSucceed((data) => { + assert.strictEqual(data, 'periodic flush\n'); + })); + })); + assert.ok(stream.write('periodic flush\n')); + } + + if (!sync) { + const dest = getTempFile(); + const fd = openSync(dest, 'w'); + let fsyncCalls = 0; + const stream = new Utf8Stream({ + fd, + minLength: 5000, + periodicFlush: common.platformTimeout(10), + fs: { + fsync: common.mustCall((_fd, callback) => { + fsyncCalls++; + if (fsyncCalls === 1) { + setTimeout(common.mustCall(() => { + assert.strictEqual(fsyncCalls, 1); + stream.destroy(); + callback(); + }), common.platformTimeout(50)); + } else { + process.nextTick(callback); + } + }, 2), + }, + }); + + stream.on('close', common.mustCall()); + assert.ok(stream.write('slow fsync\n')); + } }