From 6a11fbcaa236293208399215d5791f71e49ef586 Mon Sep 17 00:00:00 2001 From: James M Snell Date: Sat, 3 Oct 2026 01:40:48 +0000 Subject: [PATCH] stream: fix Utf8Stream flush handling Fix several issues with `Utf8Stream` flushing: * `flush()` now writes buffered data regardless of `minLength` and invokes the callback only after pending writes complete, including when `minLength` is zero and a write is in flight. * Multiple concurrent `flush()` calls are tracked correctly and `end()` waits for pending flushes before closing. * `flushSync()` throws `ERR_INVALID_STATE` if called while an asynchronous write is in progress instead of corrupting output. * Periodic flushes no longer stack up while a flush is pending. * `fsync` is skipped for stdout/stderr file descriptors. * `reopen()`, `end()` and `destroy()` behave correctly when the stream is destroyed while still opening. Signed-off-by: James M Snell Assisted-by: Opencode --- doc/api/fs.md | 11 +- lib/internal/streams/fast-utf8-stream.js | 345 ++++++++++++++---- .../test-fastutf8stream-flush-mocks.js | 53 +++ .../test-fastutf8stream-flush-sync.js | 152 +++++++- test/parallel/test-fastutf8stream-flush.js | 167 +++++++++ .../test-fastutf8stream-periodicflush.js | 52 +++ 6 files changed, 702 insertions(+), 78 deletions(-) diff --git a/doc/api/fs.md b/doc/api/fs.md index d6eb4208ef0a..f49e131ee7bf 100644 --- a/doc/api/fs.md +++ b/doc/api/fs.md @@ -8563,13 +8563,20 @@ 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 and, unless the `fsync` option is enabled, after the data has +been synchronized with `fs.fsync()`. `fsync` errors indicating that the file +descriptor does not support synchronization (for example, a pipe or a TTY) +are ignored. #### `utf8Stream.flushSync()` Flushes the buffered data synchronously. This is a costly operation. +If an asynchronous write is in flight when `flushSync()` is called, only the +data queued after it is written, and it may reach the file before the data +of the in-flight write. + #### `utf8Stream.fsync` * {boolean} Whether the stream is performing a `fs.fsyncSync()` after every diff --git a/lib/internal/streams/fast-utf8-stream.js b/lib/internal/streams/fast-utf8-stream.js index d97f1cae3257..004ff14781f1 100644 --- a/lib/internal/streams/fast-utf8-stream.js +++ b/lib/internal/streams/fast-utf8-stream.js @@ -6,7 +6,9 @@ const { ArrayPrototypePush, + ArrayPrototypeUnshift, MathMax, + Symbol, SymbolDispose, } = primordials; @@ -28,6 +30,7 @@ const EventEmitter = require('events'); const path = require('path'); const { clearInterval, + clearTimeout, setInterval, setTimeout, } = require('timers'); @@ -58,6 +61,38 @@ const kMaxWrite = 16 * 1024; const kContentModeBuffer = 'buffer'; const kContentModeUtf8 = 'utf8'; const kNullPrototype = { __proto__: null }; +const kFlush = Symbol('kFlush'); + +// States returned by #beginFlushSync(). +// No write is in flight. +const kFlushSyncIdle = 0; +// An EAGAIN/EBUSY retry timer was pending and has been cancelled. Nothing is +// in flight, so the remainder of #writingBuf can be written synchronously. +const kFlushSyncRetry = 1; +// Called from a 'write' listener. The completed write has been released and +// nothing is in flight, so the remainder of #writingBuf can be written +// synchronously. +const kFlushSyncInWriteEvent = 2; +// An asynchronous fs.write() of #writingBuf is in flight. Only the queued +// buffers can be written; ordering relative to the in-flight chunk is not +// guaranteed. +const kFlushSyncInFlight = 3; + +// fsync() error codes meaning the file descriptor cannot be synchronized +// (e.g. it is a pipe, socket, or TTY) or has already been closed. These are +// not treated as flush failures. +function isIgnorableFsyncError(err) { + switch (err?.code) { + case 'EBADF': + case 'EINVAL': + case 'ENOTSUP': + case 'EOPNOTSUPP': + case 'EROFS': + return true; + default: + return false; + } +} // Utf8Stream is a port of the original SonicBoom module // (https://github.com/pinojs/sonic-boom) that provides a fast and efficient @@ -71,7 +106,11 @@ class Utf8Stream extends EventEmitter { #ending = false; #reopening = false; #asyncDrainScheduled = false; - #flushPending = false; + #flushPending = 0; + #flushInProgress = 0; + #emittingFlush = false; + #emittingWrite = false; + #retryTimer = null; #hwm = 16387; // 16 KB #file = null; #destroyed = false; @@ -233,7 +272,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(); } } @@ -268,7 +311,9 @@ class Utf8Stream extends EventEmitter { } if (this.#opening) { - this.once('ready', () => this.reopen(file)); + this.once('ready', () => { + if (!this.#destroyed) this.reopen(file); + }); return; } @@ -311,7 +356,7 @@ class Utf8Stream extends EventEmitter { if (this.#opening) { this.once('ready', () => { - this.end(); + if (!this.#destroyed) this.end(); }); return; } @@ -329,7 +374,7 @@ class Utf8Stream extends EventEmitter { if (this.#len > 0 && this.#fd >= 0) { this.#actualWrite(); } else { - this.#actualClose(); + this.#finishEnding(); } } @@ -337,7 +382,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} */ @@ -403,7 +452,10 @@ class Utf8Stream extends EventEmitter { } } else { // Let's give the destination some time to process the chunk. - setTimeout(() => this.#fsWrite(), BUSY_WRITE_TIMEOUT); + this.#retryTimer = setTimeout(() => { + this.#retryTimer = null; + this.#fsWrite(); + }, BUSY_WRITE_TIMEOUT); } } else { this.#writing = false; @@ -420,11 +472,20 @@ class Utf8Stream extends EventEmitter { this.#writeRetries = 0; } - this.emit('write', n); const releasedBufObj = releaseWritingBuf(this.#writingBuf, this.#len, n); this.#len = releasedBufObj.len; this.#writingBuf = releasedBufObj.writingBuf; + // Emit 'write' after the written bytes have been released so that a + // listener calling flushSync() sees a consistent state with no I/O in + // flight. flushSync() may consume the remainder of #writingBuf. + this.#emittingWrite = true; + try { + this.emit('write', n); + } finally { + this.#emittingWrite = false; + } + if (this.#writingBuf.length) { if (!this.#sync) { this.#fsWrite(); @@ -444,6 +505,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); } @@ -453,25 +526,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(); } } @@ -520,12 +587,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()); } }; @@ -551,16 +619,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; } @@ -636,7 +747,12 @@ class Utf8Stream extends EventEmitter { } } - #flushBufferSync() { + /** + * Validates that flushSync() can run and determines how it must interact + * with any write already in progress. + * @returns {number} One of the kFlushSync* states. + */ + #beginFlushSync() { if (this.#destroyed) { throw new ERR_INVALID_STATE('Utf8Stream is destroyed'); } @@ -645,51 +761,114 @@ class Utf8Stream extends EventEmitter { throw new ERR_INVALID_STATE('Invalid file descriptor'); } - if (!this.#writing && this.#writingBuf.length > 0) { - this.#bufs.unshift([this.#writingBuf]); - this.#writingBuf = kEmptyBuffer; + // While reopening, #writing is set but no write is in flight and #fd still + // refers to the previous file, which is only closed once the new file is + // ready. Flush the buffered data to it. + if (!this.#writing || this.#opening) { + return kFlushSyncIdle; } - let buf = kEmptyBuffer; - while (this.#bufs.length || buf.length) { - if (buf.length <= 0) { - buf = mergeBuf(this.#bufs[0], this.#lens[0]); + if (this.#retryTimer !== null) { + clearTimeout(this.#retryTimer); + this.#retryTimer = null; + return kFlushSyncRetry; + } + + if (this.#emittingWrite) { + return kFlushSyncInWriteEvent; + } + + return kFlushSyncInFlight; + } + + /** + * Called after flushSync() has taken over a write whose retry timer it + * cancelled. Resumes whatever #release() would have done once that write + * completed. + */ + #endFlushSyncRetry() { + this.#writing = false; + process.nextTick(() => { + if (this.#destroyed || this.#writing) return; + if (this.#reopening) { + this.#reopening = false; + this.reopen(); + } else if (this.#ending) { + this.#finishEnding(); + } else if (this.#len > 0 && + (this.#len > this.#minLength || this.#flushPending)) { + this.#actualWrite(); + } else { + this.#emitDrain(); } - try { - const n = this.#fs.writeSync(this.#fd, buf); - this.#writeRetries = 0; - buf = buf.subarray(n); - this.#len = MathMax(this.#len - n, 0); + }); + } + + #flushBufferSync() { + const state = this.#beginFlushSync(); + + // Unless an fs.write() of #writingBuf is in flight, its unwritten + // remainder must be written first to preserve ordering. + if (state !== kFlushSyncInFlight && this.#writingBuf.length > 0) { + ArrayPrototypeUnshift(this.#bufs, [this.#writingBuf]); + ArrayPrototypeUnshift(this.#lens, this.#writingBuf.length); + this.#writingBuf = kEmptyBuffer; + } + + try { + let buf = kEmptyBuffer; + while (this.#bufs.length || buf.length) { if (buf.length <= 0) { - this.#bufs.shift(); - this.#lens.shift(); - } - } catch (err) { - const shouldRetry = err.code === 'EAGAIN' || err.code === 'EBUSY'; - if (shouldRetry) { - this.#writeRetries++; - } - const retriesExhausted = this.#maxWriteRetries > 0 && this.#writeRetries > this.#maxWriteRetries; - if (!shouldRetry || retriesExhausted || !this.#retryEAGAIN(err, buf.length, this.#len - buf.length)) { - throw err; + buf = mergeBuf(this.#bufs[0], this.#lens[0]); } + try { + const n = this.#fs.writeSync(this.#fd, buf); + this.#writeRetries = 0; + buf = buf.subarray(n); + this.#len = MathMax(this.#len - n, 0); + if (buf.length <= 0) { + this.#bufs.shift(); + this.#lens.shift(); + } + } catch (err) { + const shouldRetry = err.code === 'EAGAIN' || err.code === 'EBUSY'; + if (shouldRetry) { + this.#writeRetries++; + } + const retriesExhausted = this.#maxWriteRetries > 0 && this.#writeRetries > this.#maxWriteRetries; + if (!shouldRetry || retriesExhausted || !this.#retryEAGAIN(err, buf.length, this.#len - buf.length)) { + throw err; + } - sleep(BUSY_WRITE_TIMEOUT); + sleep(BUSY_WRITE_TIMEOUT); + } } + } finally { + if (state === kFlushSyncRetry) this.#endFlushSyncRetry(); } } #flushSyncUtf8() { - if (this.#destroyed) { - throw new ERR_INVALID_STATE('Utf8Stream is destroyed'); + const state = this.#beginFlushSync(); + + try { + this.#flushSyncUtf8Bufs(state); + } finally { + if (state === kFlushSyncRetry) this.#endFlushSyncRetry(); } - if (this.#fd < 0) { - throw new ERR_INVALID_STATE('Invalid file descriptor'); + try { + this.#fs.fsyncSync(this.#fd); + } catch { + // Skip the error. The fd might not support fsync. } + } - if (!this.#writing && this.#writingBuf.length > 0) { - this.#bufs.unshift(this.#writingBuf); + #flushSyncUtf8Bufs(state) { + // Unless an fs.write() of #writingBuf is in flight, its unwritten + // remainder must be written first to preserve ordering. + if (state !== kFlushSyncInFlight && this.#writingBuf.length > 0) { + ArrayPrototypeUnshift(this.#bufs, this.#writingBuf); this.#writingBuf = ''; } @@ -720,46 +899,62 @@ class Utf8Stream extends EventEmitter { sleep(BUSY_WRITE_TIMEOUT); } } - - try { - this.#fs.fsyncSync(this.#fd); - } catch { - // Skip the error. The fd might not support fsync. - } } #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) { try { this.#fs.fsync(this.#fd, (err) => { - this.#flushPending = false; - // If the fd is closed, we ignore the error. - if (err?.code === 'EBADF') { - cb(); + // Ignore errors meaning the fd is closed or cannot be synced + // (e.g. stdout/stderr attached to a pipe or TTY). A regular + // file, including a redirected stdout/stderr, is still synced. + if (isIgnorableFsyncError(err)) { + 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); } @@ -776,7 +971,7 @@ class Utf8Stream extends EventEmitter { throw error; } - if (this.#minLength <= 0) { + if (this.#minLength <= 0 && !this.#writing && this.#len === 0) { cb?.(); return; } @@ -810,7 +1005,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-mocks.js b/test/parallel/test-fastutf8stream-flush-mocks.js index 5d8e19fb1f92..3c6ad039cf4e 100644 --- a/test/parallel/test-fastutf8stream-flush-mocks.js +++ b/test/parallel/test-fastutf8stream-flush-mocks.js @@ -3,9 +3,13 @@ const common = require('../common'); const tmpdir = require('../common/tmpdir'); const assert = require('node:assert'); +const { spawnSync } = require('node:child_process'); const { + closeSync, openSync, + fsync, fsyncSync, + readFileSync, writeSync, write, } = require('node:fs'); @@ -13,6 +17,23 @@ const { join } = require('node:path'); const { Utf8Stream } = require('node:fs'); const { isMainThread } = require('node:worker_threads'); +if (process.argv[2] === 'child-stdout') { + // stdout is redirected to a regular file by the parent. + const stream = new Utf8Stream({ + fd: 1, + minLength: 4096, + fs: { + fsync: common.mustCall((fd, cb) => { + assert.strictEqual(fd, 1); + fsync(fd, cb); + }), + }, + }); + stream.write('hello world\n'); + stream.flush(common.mustSucceed()); + return; +} + tmpdir.refresh(); if (isMainThread) { process.umask(0o000); @@ -27,6 +48,38 @@ function getTempFile() { runTests(false); runTests(true); +// Errors from fsync meaning the fd cannot be synchronized (e.g. a pipe or +// TTY) or is already closed do not fail the flush. +for (const code of ['EBADF', 'EINVAL', 'ENOTSUP', 'EOPNOTSUPP', 'EROFS']) { + const dest = getTempFile(); + const fd = openSync(dest, 'w'); + const fsOverride = { + fsync: common.mustCall((fd, cb) => { + const err = new Error(code); + err.code = code; + process.nextTick(cb, err); + }, 2), + }; + const stream = new Utf8Stream({ fd, minLength: 4096, fs: fsOverride }); + + stream.on('ready', common.mustCall(() => { + assert.ok(stream.write('hello world\n')); + stream.flush(common.mustSucceed(() => stream.end())); + })); +} + +// stdout redirected to a regular file is still fsynced by flush(). +{ + const dest = getTempFile(); + const fd = openSync(dest, 'w'); + const child = spawnSync(process.execPath, [__filename, 'child-stdout'], { + stdio: ['ignore', fd, 'pipe'], + }); + closeSync(fd); + assert.strictEqual(child.status, 0, child.stderr.toString()); + assert.strictEqual(readFileSync(dest, 'utf8'), 'hello world\n'); +} + function runTests(sync) { { diff --git a/test/parallel/test-fastutf8stream-flush-sync.js b/test/parallel/test-fastutf8stream-flush-sync.js index a48dac8e08e9..e42397454aea 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,156 @@ function runTests(sync) { })); } +function toData(contentMode, text) { + return contentMode === 'buffer' ? Buffer.from(text) : text; +} + +// flushSync() while an asynchronous write is in flight writes the queued +// data without throwing, and the in-flight write still completes. +for (const contentMode of ['utf8', 'buffer']) { + const dest = getTempFile(); + const fd = openSync(dest, 'w'); + let completeWrite; + const fsOverride = { + write: common.mustCall((fd, buf, ...args) => { + const cb = args.pop(); + completeWrite = () => cb(null, writeSync(fd, buf)); + }), + }; + const stream = new Utf8Stream({ + contentMode, + fd, + minLength: 0, + sync: false, + fs: fsOverride, + }); + + stream.on('ready', common.mustCall(() => { + assert.ok(stream.write(toData(contentMode, 'in flight\n'))); + assert.strictEqual(stream.writing, true); + assert.ok(stream.write(toData(contentMode, 'queued\n'))); + stream.flushSync(); + assert.strictEqual(readFileSync(dest, 'utf8'), 'queued\n'); + + stream.flush(common.mustSucceed(() => { + stream.on('finish', common.mustCall(() => { + assert.strictEqual(readFileSync(dest, 'utf8'), 'queued\nin flight\n'); + })); + stream.end(); + })); + completeWrite(); + })); +} + +// flushSync() while waiting to retry an EAGAIN write cancels the retry and +// writes everything synchronously, in order. +for (const contentMode of ['utf8', 'buffer']) { + const dest = getTempFile(); + const fd = openSync(dest, 'w'); + const fsOverride = { + write: common.mustCall((fd, buf, ...args) => { + const cb = args.pop(); + const err = new Error('EAGAIN'); + err.code = 'EAGAIN'; + process.nextTick(cb, err); + }), + }; + const stream = new Utf8Stream({ + contentMode, + fd, + minLength: 0, + sync: false, + fs: fsOverride, + }); + + stream.on('ready', common.mustCall(() => { + assert.ok(stream.write(toData(contentMode, 'hello\n'))); + // Wait for the EAGAIN to schedule the retry. + setImmediate(common.mustCall(() => { + assert.strictEqual(stream.writing, true); + assert.ok(stream.write(toData(contentMode, 'world\n'))); + // A pending flush() completes once flushSync() takes over the write. + stream.flush(common.mustSucceed(() => { + stream.on('finish', common.mustCall(() => { + assert.strictEqual(readFileSync(dest, 'utf8'), 'hello\nworld\n'); + })); + stream.end(); + })); + stream.flushSync(); + assert.strictEqual(stream.writing, false); + assert.strictEqual(readFileSync(dest, 'utf8'), 'hello\nworld\n'); + })); + })); +} + +// flushSync() from a 'write' listener writes the unwritten remainder of a +// partial write before newly queued data. +for (const sync of [true, false]) { + for (const contentMode of ['utf8', 'buffer']) { + const dest = getTempFile(); + const fd = openSync(dest, 'w'); + let partial = true; + const partialWrite = (fd, buf) => { + if (partial) { + partial = false; + return writeSync(fd, Buffer.from(buf).subarray(0, 5)); + } + return writeSync(fd, buf); + }; + const fsOverride = { + writeSync: (fd, buf, ...args) => partialWrite(fd, buf), + write: (fd, buf, ...args) => { + const cb = args.pop(); + process.nextTick(cb, null, partialWrite(fd, buf)); + }, + }; + const stream = new Utf8Stream({ + contentMode, + fd, + minLength: 0, + sync, + fs: fsOverride, + }); + + stream.once('write', common.mustCall((n) => { + assert.strictEqual(n, 5); + assert.strictEqual(stream.writing, true); + assert.ok(stream.write(toData(contentMode, 'next\n'))); + stream.flushSync(); + assert.strictEqual(readFileSync(dest, 'utf8'), 'hello world\nnext\n'); + })); + + stream.on('ready', common.mustCall(() => { + assert.ok(stream.write(toData(contentMode, 'hello world\n'))); + stream.on('finish', common.mustCall(() => { + assert.strictEqual(readFileSync(dest, 'utf8'), 'hello world\nnext\n'); + })); + stream.end(); + })); + } +} + +// flushSync() while reopening writes the buffered data to the previous file, +// which stays open until the new file is ready. +{ + const dest = getTempFile(); + const stream = new Utf8Stream({ dest, minLength: 4096, sync: false }); + + stream.once('ready', common.mustCall(() => { + assert.ok(stream.write('before reopen\n')); + stream.reopen(); + assert.strictEqual(stream.writing, true); + stream.flushSync(); + assert.strictEqual(readFileSync(dest, 'utf8'), 'before reopen\n'); + stream.once('ready', common.mustCall(() => { + stream.on('finish', common.mustCall(() => { + assert.strictEqual(readFileSync(dest, 'utf8'), 'before reopen\n'); + })); + stream.end(); + })); + })); +} + { 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 60238ed76908..14d77166c4c8 100644 --- a/test/parallel/test-fastutf8stream-flush.js +++ b/test/parallel/test-fastutf8stream-flush.js @@ -53,6 +53,173 @@ runTests(true); stream.on('ready', common.mustCall()); } +{ + 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')); + } }