diff --git a/lib/internal/streams/fast-utf8-stream.js b/lib/internal/streams/fast-utf8-stream.js index 8a1b6185e82..79aabd0136f 100644 --- a/lib/internal/streams/fast-utf8-stream.js +++ b/lib/internal/streams/fast-utf8-stream.js @@ -83,6 +83,8 @@ class Utf8Stream extends EventEmitter { #mode = 0o666; #retryEAGAIN = () => true; #mkdir = false; + #maxWriteRetries = 0; + #writeRetries = 0; #writingBuf = ''; #write; #flush; @@ -128,6 +130,7 @@ class Utf8Stream extends EventEmitter { append = true, mkdir, retryEAGAIN, + maxWriteRetries, fsync, contentMode = kContentModeUtf8, mode, @@ -161,12 +164,14 @@ class Utf8Stream extends EventEmitter { this.#mode = mode; this.#retryEAGAIN = retryEAGAIN || (() => true); this.#mkdir = mkdir || false; + this.#maxWriteRetries = maxWriteRetries || 0; validateUint32(this.#hwm, 'options.hwm'); validateUint32(this.#minLength, 'options.minLength'); validateUint32(this.#maxLength, 'options.maxLength'); validateUint32(this.#maxWrite, 'options.maxWrite'); validateUint32(this.#periodicFlush, 'options.periodicFlush'); + validateUint32(this.#maxWriteRetries, 'options.maxWriteRetries'); validateBoolean(this.#sync, 'options.sync'); validateBoolean(this.#fsync, 'options.fsync'); validateBoolean(this.#append, 'options.append'); @@ -373,7 +378,13 @@ class Utf8Stream extends EventEmitter { #release(err, n) { if (err) { - if ((err.code === 'EAGAIN' || err.code === 'EBUSY') && + const isRetryableErr = (err.code === 'EAGAIN' || err.code === 'EBUSY'); + if (isRetryableErr) { + this.#writeRetries++; + } + const retriesExhausted = this.#maxWriteRetries > 0 && this.#writeRetries > this.#maxWriteRetries; + + if (isRetryableErr && !retriesExhausted && this.#retryEAGAIN(err, this.#writingBuf.length, this.#len - this.#writingBuf.length)) { if (this.#sync) { // This error code should not happen in sync mode, because it is @@ -398,6 +409,13 @@ class Utf8Stream extends EventEmitter { return; } + // Reset the retry counter only once real forward progress (n > 0) is + // confirmed. Otherwise a stuck destination would never accumulate past + // one retry and `maxWriteRetries` could never trigger. + if (n > 0) { + this.#writeRetries = 0; + } + this.emit('write', n); const releasedBufObj = releaseWritingBuf(this.#writingBuf, this.#len, n); this.#len = releasedBufObj.len; @@ -635,6 +653,7 @@ class Utf8Stream extends EventEmitter { } 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) { @@ -643,7 +662,11 @@ class Utf8Stream extends EventEmitter { } } catch (err) { const shouldRetry = err.code === 'EAGAIN' || err.code === 'EBUSY'; - if (shouldRetry && !this.#retryEAGAIN(err, buf.length, this.#len - buf.length)) { + 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; } @@ -673,6 +696,7 @@ class Utf8Stream extends EventEmitter { } try { const n = this.#fs.writeSync(this.#fd, buf, 'utf8'); + this.#writeRetries = 0; const releasedBufObj = releaseWritingBuf(buf, this.#len, n); buf = releasedBufObj.writingBuf; this.#len = releasedBufObj.len; @@ -681,7 +705,11 @@ class Utf8Stream extends EventEmitter { } } catch (err) { const shouldRetry = err.code === 'EAGAIN' || err.code === 'EBUSY'; - if (shouldRetry && !this.#retryEAGAIN(err, buf.length, this.#len - buf.length)) { + 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; } diff --git a/test/parallel/test-fastutf8stream-flush.js b/test/parallel/test-fastutf8stream-flush.js index 195d5e83c9c..60238ed7690 100644 --- a/test/parallel/test-fastutf8stream-flush.js +++ b/test/parallel/test-fastutf8stream-flush.js @@ -4,9 +4,12 @@ const common = require('../common'); const tmpdir = require('../common/tmpdir'); const assert = require('node:assert'); const { + open, openSync, readFile, writeFileSync, + write, + writeSync, } = require('node:fs'); const { join } = require('node:path'); const { Utf8Stream } = require('node:fs'); @@ -25,6 +28,31 @@ function getTempFile() { runTests(false); runTests(true); +// Flush cb is invoked when flushing before 'ready' while the stream is +// still opening (async mode only; sync mode writes synchronously). +{ + const dest = getTempFile(); + + const stream = new Utf8Stream({ + dest, + minLength: 4096, + sync: false, + fs: { + open(file, flags, mode, cb) { + process.nextTick(() => { + assert.ok(stream.write('hello world\n')); + stream.flush(common.mustSucceed(() => { + stream.destroy(); + })); + open(file, flags, mode, cb); + }); + }, + }, + }); + + stream.on('ready', common.mustCall()); +} + function runTests(sync) { { const dest = getTempFile(); @@ -136,4 +164,84 @@ function runTests(sync) { stream.destroy(); stream.flush(common.mustCall(assert.ok)); } + + { + // Flush cb is invoked with the error when the underlying write fails. + const dest = getTempFile(); + const fd = openSync(dest, 'w'); + + const err = new Error('other'); + err.code = 'other'; + let first = true; + + const fsOverride = {}; + if (sync) { + fsOverride.writeSync = common.mustCallAtLeast((...args) => { + if (first) { + first = false; + throw err; + } + return writeSync(...args); + }, 1); + } else { + fsOverride.write = common.mustCallAtLeast((...args) => { + const callback = args[args.length - 1]; + if (first) { + first = false; + process.nextTick(callback, err); + return; + } + return write(...args); + }, 1); + } + + const stream = new Utf8Stream({ + fd, + sync, + minLength: 4096, + fs: fsOverride, + }); + + stream.on('ready', common.mustCall(() => { + assert.ok(stream.write('hello world\n')); + stream.flush(common.mustCall((e) => { + assert.strictEqual(e.code, 'other'); + stream.destroy(); + })); + })); + } + + { + // Flush cb is invoked once the in-flight write completes. + const dest = getTempFile(); + const fd = openSync(dest, 'w'); + + const fsOverride = {}; + if (sync) { + fsOverride.writeSync = common.mustCallAtLeast((...args) => { + stream.flush(common.mustSucceed(() => { + stream.destroy(); + })); + return writeSync(...args); + }, 1); + } else { + fsOverride.write = common.mustCallAtLeast((...args) => { + stream.flush(common.mustSucceed(() => { + stream.destroy(); + })); + return write(...args); + }, 1); + } + + const stream = new Utf8Stream({ + fd, + sync, + minLength: 1, + fs: fsOverride, + }); + + stream.on('ready', common.mustCall(() => { + assert.ok(stream.write('hello world\n')); + })); + } } diff --git a/test/parallel/test-fastutf8stream-retry.js b/test/parallel/test-fastutf8stream-retry.js index 3381c9beac0..e1b29ed27b3 100644 --- a/test/parallel/test-fastutf8stream-retry.js +++ b/test/parallel/test-fastutf8stream-retry.js @@ -211,3 +211,117 @@ function runTests(sync) { })); })); } + +{ + const dest = getTempFile(); + const fd = openSync(dest, 'w'); + + const err = new Error('EAGAIN'); + err.code = 'EAGAIN'; + let attempts = 0; + + const stream = new Utf8Stream({ + fd, + sync: false, + minLength: 0, + maxWriteRetries: 3, + // retryEAGAIN always returns true ("keep going"), so maxWriteRetries must + // itself cap the number of attempts instead of retrying forever. + retryEAGAIN: () => true, + fs: { + write: common.mustCall((...args) => { + attempts++; + const callback = args[args.length - 1]; + process.nextTick(callback, err); + }, 4), + } + }); + + stream.on('ready', common.mustCall(() => { + assert.ok(stream.write('hello world\n')); + })); + + stream.once('error', common.mustCall((err) => { + assert.strictEqual(err.code, 'EAGAIN'); + // 1 initial attempt + 3 retries = 4 total fs.write calls, then give up. + assert.strictEqual(attempts, 4); + assert.strictEqual(stream.writing, false); + stream.destroy(); + })); +} + +{ + const dest = getTempFile(); + const fd = openSync(dest, 'w'); + + const err = new Error('EAGAIN'); + err.code = 'EAGAIN'; + + const stream = new Utf8Stream({ + fd, + sync: true, + minLength: 0, + maxWriteRetries: 3, + retryEAGAIN: () => true, + fs: { + writeSync: common.mustCall((...args) => { + throw err; + }, 4), + } + }); + + stream.on('ready', common.mustCall(() => { + // Once retries are exhausted, write() must surface the error instead of + // spinning forever on EAGAIN. + assert.throws(() => { + stream.write('hello world\n'); + }, (e) => e.code === 'EAGAIN'); + stream.destroy(); + })); +} + +{ + const dest = getTempFile(); + const fd = openSync(dest, 'w'); + + const err = new Error('EAGAIN'); + err.code = 'EAGAIN'; + let call = 0; + + const stream = new Utf8Stream({ + fd, + sync: false, + minLength: 0, + maxWriteRetries: 3, + retryEAGAIN: () => true, + fs: { + write: common.mustCallAtLeast((...args) => { + call++; + const callback = args[args.length - 1]; + if (call % 3 === 0) { + return write(...args); + } + process.nextTick(callback, err); + }, 5), + } + }); + + stream.on('error', common.mustNotCall()); + + stream.on('ready', common.mustCall(() => { + // First burst: calls 1-2 fail with EAGAIN, call 3 succeeds. + assert.ok(stream.write('hello world\n')); + stream.once('drain', common.mustCall(() => { + // Second burst: calls 4-5 fail again. The counter must have been reset + // by the successful write at call 3, so this burst succeeds too. + stream.write('sonic boom\n'); + stream.end(); + })); + })); + + stream.on('finish', common.mustCall(() => { + readFile(dest, 'utf8', common.mustSucceed((data) => { + assert.strictEqual(data, 'hello world\nsonic boom\n'); + })); + })); +} diff --git a/test/parallel/test-fastutf8stream-write.js b/test/parallel/test-fastutf8stream-write.js index a022aead59a..432b044e165 100644 --- a/test/parallel/test-fastutf8stream-write.js +++ b/test/parallel/test-fastutf8stream-write.js @@ -268,3 +268,67 @@ function runTests(sync) { assert.strictEqual(stream.maxLength, 65536); stream.end(); } + +// maxLength drop behavior: writing past maxLength must emit 'drop' and +// never write the overflowing chunk (see SonicBoom write tests). +{ + const dest = getTempFile(); + const fd = openSync(dest, 'w'); + + const buf = Buffer.alloc(100).fill('x').toString(); + + const stream = new Utf8Stream({ + fd, + minLength: 101, + maxLength: 102, + sync: false, + fs: { + write: common.mustCall((...args) => { + const data = args[1]; + const callback = args[args.length - 1]; + assert.strictEqual(data.length, buf.length + 2); + process.nextTick(callback, null, data.length); + stream.end(); + }, 1), + } + }); + + stream.on('drop', common.mustNotCall()); + + stream.on('ready', common.mustCall(() => { + assert.ok(stream.write(buf)); + assert.ok(stream.write('aa')); + })); +} + +{ + const dest = getTempFile(); + const fd = openSync(dest, 'w'); + + const buf = Buffer.alloc(100).fill('x').toString(); + + const stream = new Utf8Stream({ + fd, + minLength: 101, + maxLength: 102, + sync: false, + fs: { + write: common.mustCall((...args) => { + const data = args[1]; + const callback = args[args.length - 1]; + assert.strictEqual(data.length, buf.length); + process.nextTick(callback, null, data.length); + }, 1), + } + }); + + stream.on('drop', common.mustCall((data) => { + assert.strictEqual(data.length, 3); + stream.end(); + })); + + stream.on('ready', common.mustCall(() => { + assert.ok(stream.write(buf)); + assert.ok(stream.write('aaa')); + })); +}