Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
34 changes: 31 additions & 3 deletions lib/internal/streams/fast-utf8-stream.js
Original file line number Diff line number Diff line change
Expand Up @@ -83,6 +83,8 @@ class Utf8Stream extends EventEmitter {
#mode = 0o666;
#retryEAGAIN = () => true;
#mkdir = false;
#maxWriteRetries = 0;
#writeRetries = 0;
#writingBuf = '';
#write;
#flush;
Expand Down Expand Up @@ -128,6 +130,7 @@ class Utf8Stream extends EventEmitter {
append = true,
mkdir,
retryEAGAIN,
maxWriteRetries,
fsync,
contentMode = kContentModeUtf8,
mode,
Expand Down Expand Up @@ -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');
Expand Down Expand Up @@ -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
Expand All @@ -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;
Expand Down Expand Up @@ -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) {
Expand All @@ -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;
}

Expand Down Expand Up @@ -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;
Expand All @@ -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;
}

Expand Down
108 changes: 108 additions & 0 deletions test/parallel/test-fastutf8stream-flush.js
Original file line number Diff line number Diff line change
Expand Up @@ -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');
Expand All @@ -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();
Expand Down Expand Up @@ -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'));
}));
}
}
114 changes: 114 additions & 0 deletions test/parallel/test-fastutf8stream-retry.js
Original file line number Diff line number Diff line change
Expand Up @@ -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');
}));
}));
}
Loading
Loading