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
7 changes: 4 additions & 3 deletions doc/api/fs.md
Original file line number Diff line number Diff line change
Expand Up @@ -8563,12 +8563,13 @@ 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.

#### `utf8Stream.flushSync()`

Flushes the buffered data synchronously. This is a costly operation.
Flushes the buffered data synchronously. This is a costly operation. An
`ERR_INVALID_STATE` error is thrown if the stream is currently writing.

#### `utf8Stream.fsync`

Expand Down
167 changes: 130 additions & 37 deletions lib/internal/streams/fast-utf8-stream.js
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
const {
ArrayPrototypePush,
MathMax,
Symbol,
SymbolDispose,
} = primordials;

Expand Down Expand Up @@ -58,6 +59,7 @@ const kMaxWrite = 16 * 1024;
const kContentModeBuffer = 'buffer';
const kContentModeUtf8 = 'utf8';
const kNullPrototype = { __proto__: null };
const kFlush = Symbol('kFlush');

// Utf8Stream is a port of the original SonicBoom module
// (https://github.com/pinojs/sonic-boom) that provides a fast and efficient
Expand All @@ -71,7 +73,9 @@ class Utf8Stream extends EventEmitter {
#ending = false;
#reopening = false;
#asyncDrainScheduled = false;
#flushPending = false;
#flushPending = 0;
#flushInProgress = 0;
#emittingFlush = false;
#hwm = 16387; // 16 KB
#file = null;
#destroyed = false;
Expand Down Expand Up @@ -228,7 +232,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();
}
}
Expand Down Expand Up @@ -263,7 +271,9 @@ class Utf8Stream extends EventEmitter {
}

if (this.#opening) {
this.once('ready', () => this.reopen(file));
this.once('ready', () => {
if (!this.#destroyed) this.reopen(file);
});
return;
}

Expand Down Expand Up @@ -306,7 +316,7 @@ class Utf8Stream extends EventEmitter {

if (this.#opening) {
this.once('ready', () => {
this.end();
if (!this.#destroyed) this.end();
});
return;
}
Expand All @@ -324,15 +334,19 @@ class Utf8Stream extends EventEmitter {
if (this.#len > 0 && this.#fd >= 0) {
this.#actualWrite();
} else {
this.#actualClose();
this.#finishEnding();
}
}

destroy() {
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} */
Expand Down Expand Up @@ -426,6 +440,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);
}
Expand All @@ -435,25 +461,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();
}
}

Expand Down Expand Up @@ -502,12 +522,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());
}
};

Expand All @@ -533,16 +554,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;
}

Expand Down Expand Up @@ -627,7 +691,11 @@ class Utf8Stream extends EventEmitter {
throw new ERR_INVALID_STATE('Invalid file descriptor');
}

if (!this.#writing && this.#writingBuf.length > 0) {
if (this.#writing) {
throw new ERR_INVALID_STATE('Cannot flush while a write is in progress');
}

if (this.#writingBuf.length > 0) {
this.#bufs.unshift([this.#writingBuf]);
this.#writingBuf = kEmptyBuffer;
}
Expand Down Expand Up @@ -665,7 +733,11 @@ class Utf8Stream extends EventEmitter {
throw new ERR_INVALID_STATE('Invalid file descriptor');
}

if (!this.#writing && this.#writingBuf.length > 0) {
if (this.#writing) {
throw new ERR_INVALID_STATE('Cannot flush while a write is in progress');
}

if (this.#writingBuf.length > 0) {
this.#bufs.unshift(this.#writingBuf);
this.#writingBuf = '';
}
Expand Down Expand Up @@ -701,37 +773,58 @@ class Utf8Stream extends EventEmitter {
}

#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) {
if (!this.#fsync && !this.#destroyed &&
this.#fd !== 1 && this.#fd !== 2) {
try {
this.#fs.fsync(this.#fd, (err) => {
this.#flushPending = false;
// If the fd is closed, we ignore the error.
if (err?.code === 'EBADF') {
cb();
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);
}

Expand All @@ -748,7 +841,7 @@ class Utf8Stream extends EventEmitter {
throw error;
}

if (this.#minLength <= 0) {
if (this.#minLength <= 0 && !this.#writing && this.#len === 0) {
cb?.();
return;
}
Expand Down Expand Up @@ -782,7 +875,7 @@ class Utf8Stream extends EventEmitter {
throw error;
}

if (this.#minLength <= 0) {
if (this.#minLength <= 0 && !this.#writing && this.#len === 0) {
cb?.();
return;
}
Expand Down
29 changes: 28 additions & 1 deletion test/parallel/test-fastutf8stream-flush-sync.js
Original file line number Diff line number Diff line change
Expand Up @@ -67,7 +67,7 @@ function runTests(sync) {
const stream = new Utf8Stream({
fd,
sync: false,
minLength: 0,
minLength: 1000,
fs: fsOverride,
});

Expand All @@ -85,6 +85,33 @@ function runTests(sync) {
}));
}

for (const contentMode of ['utf8', 'buffer']) {
const dest = getTempFile();
const fd = openSync(dest, 'w');
const stream = new Utf8Stream({
contentMode,
fd,
minLength: 0,
sync: false,
});
const text = `${contentMode} asynchronous write\n`;
const data = contentMode === 'buffer' ? Buffer.from(text) : text;

stream.on('ready', common.mustCall(() => {
assert.ok(stream.write(data));
assert.throws(
() => stream.flushSync(),
{ code: 'ERR_INVALID_STATE' },
);
stream.flush(common.mustSucceed(() => {
stream.end();
readFile(dest, 'utf8', common.mustSucceed((contents) => {
assert.strictEqual(contents, text);
}));
}));
}));
}

{
const dest = getTempFile();
const fd = openSync(dest, 'w');
Expand Down
Loading
Loading