Skip to content

Commit 149a864

Browse files
authored
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: #58955 Signed-off-by: Matteo Collina <hello@matteocollina.com> Assisted-by: pi PR-URL: #66275 Reviewed-By: Paolo Insogna <paolo@cowtech.it> Reviewed-By: James M Snell <jasnell@gmail.com>
1 parent f863223 commit 149a864

4 files changed

Lines changed: 317 additions & 3 deletions

File tree

‎lib/internal/streams/fast-utf8-stream.js‎

Lines changed: 31 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -87,6 +87,8 @@ class Utf8Stream extends EventEmitter {
8787
#mode = 0o666;
8888
#retryEAGAIN = () => true;
8989
#mkdir = false;
90+
#maxWriteRetries = 0;
91+
#writeRetries = 0;
9092
#writingBuf = '';
9193
#write;
9294
#flush;
@@ -132,6 +134,7 @@ class Utf8Stream extends EventEmitter {
132134
append = true,
133135
mkdir,
134136
retryEAGAIN,
137+
maxWriteRetries,
135138
fsync,
136139
contentMode = kContentModeUtf8,
137140
mode,
@@ -165,12 +168,14 @@ class Utf8Stream extends EventEmitter {
165168
this.#mode = mode;
166169
this.#retryEAGAIN = retryEAGAIN || (() => true);
167170
this.#mkdir = mkdir || false;
171+
this.#maxWriteRetries = maxWriteRetries || 0;
168172

169173
validateUint32(this.#hwm, 'options.hwm');
170174
validateUint32(this.#minLength, 'options.minLength');
171175
validateUint32(this.#maxLength, 'options.maxLength');
172176
validateUint32(this.#maxWrite, 'options.maxWrite');
173177
validateUint32(this.#periodicFlush, 'options.periodicFlush');
178+
validateUint32(this.#maxWriteRetries, 'options.maxWriteRetries');
174179
validateBoolean(this.#sync, 'options.sync');
175180
validateBoolean(this.#fsync, 'options.fsync');
176181
validateBoolean(this.#append, 'options.append');
@@ -377,7 +382,13 @@ class Utf8Stream extends EventEmitter {
377382

378383
#release(err, n) {
379384
if (err) {
380-
if ((err.code === 'EAGAIN' || err.code === 'EBUSY') &&
385+
const isRetryableErr = (err.code === 'EAGAIN' || err.code === 'EBUSY');
386+
if (isRetryableErr) {
387+
this.#writeRetries++;
388+
}
389+
const retriesExhausted = this.#maxWriteRetries > 0 && this.#writeRetries > this.#maxWriteRetries;
390+
391+
if (isRetryableErr && !retriesExhausted &&
381392
this.#retryEAGAIN(err, this.#writingBuf.length, this.#len - this.#writingBuf.length)) {
382393
if (this.#sync) {
383394
// This error code should not happen in sync mode, because it is
@@ -402,6 +413,13 @@ class Utf8Stream extends EventEmitter {
402413
return;
403414
}
404415

416+
// Reset the retry counter only once real forward progress (n > 0) is
417+
// confirmed. Otherwise a stuck destination would never accumulate past
418+
// one retry and `maxWriteRetries` could never trigger.
419+
if (n > 0) {
420+
this.#writeRetries = 0;
421+
}
422+
405423
this.emit('write', n);
406424
const releasedBufObj = releaseWritingBuf(this.#writingBuf, this.#len, n);
407425
this.#len = releasedBufObj.len;
@@ -639,6 +657,7 @@ class Utf8Stream extends EventEmitter {
639657
}
640658
try {
641659
const n = this.#fs.writeSync(this.#fd, buf);
660+
this.#writeRetries = 0;
642661
buf = buf.subarray(n);
643662
this.#len = MathMax(this.#len - n, 0);
644663
if (buf.length <= 0) {
@@ -647,7 +666,11 @@ class Utf8Stream extends EventEmitter {
647666
}
648667
} catch (err) {
649668
const shouldRetry = err.code === 'EAGAIN' || err.code === 'EBUSY';
650-
if (shouldRetry && !this.#retryEAGAIN(err, buf.length, this.#len - buf.length)) {
669+
if (shouldRetry) {
670+
this.#writeRetries++;
671+
}
672+
const retriesExhausted = this.#maxWriteRetries > 0 && this.#writeRetries > this.#maxWriteRetries;
673+
if (!shouldRetry || retriesExhausted || !this.#retryEAGAIN(err, buf.length, this.#len - buf.length)) {
651674
throw err;
652675
}
653676

@@ -677,6 +700,7 @@ class Utf8Stream extends EventEmitter {
677700
}
678701
try {
679702
const n = this.#fs.writeSync(this.#fd, buf, 'utf8');
703+
this.#writeRetries = 0;
680704
const releasedBufObj = releaseWritingBuf(buf, this.#len, n);
681705
buf = releasedBufObj.writingBuf;
682706
this.#len = releasedBufObj.len;
@@ -685,7 +709,11 @@ class Utf8Stream extends EventEmitter {
685709
}
686710
} catch (err) {
687711
const shouldRetry = err.code === 'EAGAIN' || err.code === 'EBUSY';
688-
if (shouldRetry && !this.#retryEAGAIN(err, buf.length, this.#len - buf.length)) {
712+
if (shouldRetry) {
713+
this.#writeRetries++;
714+
}
715+
const retriesExhausted = this.#maxWriteRetries > 0 && this.#writeRetries > this.#maxWriteRetries;
716+
if (!shouldRetry || retriesExhausted || !this.#retryEAGAIN(err, buf.length, this.#len - buf.length)) {
689717
throw err;
690718
}
691719

‎test/parallel/test-fastutf8stream-flush.js‎

Lines changed: 108 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,9 +4,12 @@ const common = require('../common');
44
const tmpdir = require('../common/tmpdir');
55
const assert = require('node:assert');
66
const {
7+
open,
78
openSync,
89
readFile,
910
writeFileSync,
11+
write,
12+
writeSync,
1013
} = require('node:fs');
1114
const { join } = require('node:path');
1215
const { Utf8Stream } = require('node:fs');
@@ -25,6 +28,31 @@ function getTempFile() {
2528
runTests(false);
2629
runTests(true);
2730

31+
// Flush cb is invoked when flushing before 'ready' while the stream is
32+
// still opening (async mode only; sync mode writes synchronously).
33+
{
34+
const dest = getTempFile();
35+
36+
const stream = new Utf8Stream({
37+
dest,
38+
minLength: 4096,
39+
sync: false,
40+
fs: {
41+
open(file, flags, mode, cb) {
42+
process.nextTick(() => {
43+
assert.ok(stream.write('hello world\n'));
44+
stream.flush(common.mustSucceed(() => {
45+
stream.destroy();
46+
}));
47+
open(file, flags, mode, cb);
48+
});
49+
},
50+
},
51+
});
52+
53+
stream.on('ready', common.mustCall());
54+
}
55+
2856
function runTests(sync) {
2957
{
3058
const dest = getTempFile();
@@ -136,4 +164,84 @@ function runTests(sync) {
136164
stream.destroy();
137165
stream.flush(common.mustCall(assert.ok));
138166
}
167+
168+
{
169+
// Flush cb is invoked with the error when the underlying write fails.
170+
const dest = getTempFile();
171+
const fd = openSync(dest, 'w');
172+
173+
const err = new Error('other');
174+
err.code = 'other';
175+
let first = true;
176+
177+
const fsOverride = {};
178+
if (sync) {
179+
fsOverride.writeSync = common.mustCallAtLeast((...args) => {
180+
if (first) {
181+
first = false;
182+
throw err;
183+
}
184+
return writeSync(...args);
185+
}, 1);
186+
} else {
187+
fsOverride.write = common.mustCallAtLeast((...args) => {
188+
const callback = args[args.length - 1];
189+
if (first) {
190+
first = false;
191+
process.nextTick(callback, err);
192+
return;
193+
}
194+
return write(...args);
195+
}, 1);
196+
}
197+
198+
const stream = new Utf8Stream({
199+
fd,
200+
sync,
201+
minLength: 4096,
202+
fs: fsOverride,
203+
});
204+
205+
stream.on('ready', common.mustCall(() => {
206+
assert.ok(stream.write('hello world\n'));
207+
stream.flush(common.mustCall((e) => {
208+
assert.strictEqual(e.code, 'other');
209+
stream.destroy();
210+
}));
211+
}));
212+
}
213+
214+
{
215+
// Flush cb is invoked once the in-flight write completes.
216+
const dest = getTempFile();
217+
const fd = openSync(dest, 'w');
218+
219+
const fsOverride = {};
220+
if (sync) {
221+
fsOverride.writeSync = common.mustCallAtLeast((...args) => {
222+
stream.flush(common.mustSucceed(() => {
223+
stream.destroy();
224+
}));
225+
return writeSync(...args);
226+
}, 1);
227+
} else {
228+
fsOverride.write = common.mustCallAtLeast((...args) => {
229+
stream.flush(common.mustSucceed(() => {
230+
stream.destroy();
231+
}));
232+
return write(...args);
233+
}, 1);
234+
}
235+
236+
const stream = new Utf8Stream({
237+
fd,
238+
sync,
239+
minLength: 1,
240+
fs: fsOverride,
241+
});
242+
243+
stream.on('ready', common.mustCall(() => {
244+
assert.ok(stream.write('hello world\n'));
245+
}));
246+
}
139247
}

‎test/parallel/test-fastutf8stream-retry.js‎

Lines changed: 114 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -211,3 +211,117 @@ function runTests(sync) {
211211
}));
212212
}));
213213
}
214+
215+
{
216+
const dest = getTempFile();
217+
const fd = openSync(dest, 'w');
218+
219+
const err = new Error('EAGAIN');
220+
err.code = 'EAGAIN';
221+
let attempts = 0;
222+
223+
const stream = new Utf8Stream({
224+
fd,
225+
sync: false,
226+
minLength: 0,
227+
maxWriteRetries: 3,
228+
// retryEAGAIN always returns true ("keep going"), so maxWriteRetries must
229+
// itself cap the number of attempts instead of retrying forever.
230+
retryEAGAIN: () => true,
231+
fs: {
232+
write: common.mustCall((...args) => {
233+
attempts++;
234+
const callback = args[args.length - 1];
235+
process.nextTick(callback, err);
236+
}, 4),
237+
}
238+
});
239+
240+
stream.on('ready', common.mustCall(() => {
241+
assert.ok(stream.write('hello world\n'));
242+
}));
243+
244+
stream.once('error', common.mustCall((err) => {
245+
assert.strictEqual(err.code, 'EAGAIN');
246+
// 1 initial attempt + 3 retries = 4 total fs.write calls, then give up.
247+
assert.strictEqual(attempts, 4);
248+
assert.strictEqual(stream.writing, false);
249+
stream.destroy();
250+
}));
251+
}
252+
253+
{
254+
const dest = getTempFile();
255+
const fd = openSync(dest, 'w');
256+
257+
const err = new Error('EAGAIN');
258+
err.code = 'EAGAIN';
259+
260+
const stream = new Utf8Stream({
261+
fd,
262+
sync: true,
263+
minLength: 0,
264+
maxWriteRetries: 3,
265+
retryEAGAIN: () => true,
266+
fs: {
267+
writeSync: common.mustCall((...args) => {
268+
throw err;
269+
}, 4),
270+
}
271+
});
272+
273+
stream.on('ready', common.mustCall(() => {
274+
// Once retries are exhausted, write() must surface the error instead of
275+
// spinning forever on EAGAIN.
276+
assert.throws(() => {
277+
stream.write('hello world\n');
278+
}, (e) => e.code === 'EAGAIN');
279+
stream.destroy();
280+
}));
281+
}
282+
283+
{
284+
const dest = getTempFile();
285+
const fd = openSync(dest, 'w');
286+
287+
const err = new Error('EAGAIN');
288+
err.code = 'EAGAIN';
289+
let call = 0;
290+
291+
const stream = new Utf8Stream({
292+
fd,
293+
sync: false,
294+
minLength: 0,
295+
maxWriteRetries: 3,
296+
retryEAGAIN: () => true,
297+
fs: {
298+
write: common.mustCallAtLeast((...args) => {
299+
call++;
300+
const callback = args[args.length - 1];
301+
if (call % 3 === 0) {
302+
return write(...args);
303+
}
304+
process.nextTick(callback, err);
305+
}, 5),
306+
}
307+
});
308+
309+
stream.on('error', common.mustNotCall());
310+
311+
stream.on('ready', common.mustCall(() => {
312+
// First burst: calls 1-2 fail with EAGAIN, call 3 succeeds.
313+
assert.ok(stream.write('hello world\n'));
314+
stream.once('drain', common.mustCall(() => {
315+
// Second burst: calls 4-5 fail again. The counter must have been reset
316+
// by the successful write at call 3, so this burst succeeds too.
317+
stream.write('sonic boom\n');
318+
stream.end();
319+
}));
320+
}));
321+
322+
stream.on('finish', common.mustCall(() => {
323+
readFile(dest, 'utf8', common.mustSucceed((data) => {
324+
assert.strictEqual(data, 'hello world\nsonic boom\n');
325+
}));
326+
}));
327+
}

0 commit comments

Comments
 (0)