Commit 149a864c3c5 for nodejs
commit 149a864c3c5b1d9846885abf127f1eb18e27eae2
Author: Matteo Collina <hello@matteocollina.com>
Date: Mon Oct 5 07:16:57 2026 +0200
stream: port remaining SonicBoom tests for Utf8Stream
Port the SonicBoom tests that could not be imported directly because
they depend on the third-party proxyquire module, along with the
maxWriteRetries option they exercise, into the Utf8Stream module.
Adds maxWriteRetries support, which bounds consecutive EAGAIN/EBUSY
retry attempts and resets the counter on forward progress, plus
drop-event and flush-callback edge-case coverage.
Fixes: https://github.com/nodejs/node/issues/58955
Signed-off-by: Matteo Collina <hello@matteocollina.com>
Assisted-by: pi
PR-URL: https://github.com/nodejs/node/pull/66275
Reviewed-By: Paolo Insogna <paolo@cowtech.it>
Reviewed-By: James M Snell <jasnell@gmail.com>
diff --git a/lib/internal/streams/fast-utf8-stream.js b/lib/internal/streams/fast-utf8-stream.js
index a11fbf3bbc0..d97f1cae325 100644
--- a/lib/internal/streams/fast-utf8-stream.js
+++ b/lib/internal/streams/fast-utf8-stream.js
@@ -87,6 +87,8 @@ class Utf8Stream extends EventEmitter {
#mode = 0o666;
#retryEAGAIN = () => true;
#mkdir = false;
+ #maxWriteRetries = 0;
+ #writeRetries = 0;
#writingBuf = '';
#write;
#flush;
@@ -132,6 +134,7 @@ class Utf8Stream extends EventEmitter {
append = true,
mkdir,
retryEAGAIN,
+ maxWriteRetries,
fsync,
contentMode = kContentModeUtf8,
mode,
@@ -165,12 +168,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');
@@ -377,7 +382,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
@@ -402,6 +413,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;
@@ -639,6 +657,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) {
@@ -647,7 +666,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;
}
@@ -677,6 +700,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;
@@ -685,7 +709,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'));
+ }));
+}