From 0e92441d2f2b35de8854fa47ad8fd7dae9bf404d Mon Sep 17 00:00:00 2001 From: James M Snell Date: Sat, 3 Oct 2026 01:40:48 +0000 Subject: [PATCH 1/2] stream: fix Utf8Stream flush handling Fix several issues with `Utf8Stream` flushing: * `flush()` now writes buffered data regardless of `minLength` and invokes the callback only after pending writes complete, including when `minLength` is zero and a write is in flight. * Multiple concurrent `flush()` calls are tracked correctly and `end()` waits for pending flushes before closing. * `flushSync()` throws `ERR_INVALID_STATE` if called while an asynchronous write is in progress instead of corrupting output. * Periodic flushes no longer stack up while a flush is pending. * `fsync` is skipped for stdout/stderr file descriptors. * `reopen()`, `end()` and `destroy()` behave correctly when the stream is destroyed while still opening. Signed-off-by: James M Snell Assisted-by: Opencode --- doc/api/fs.md | 7 +- lib/internal/streams/fast-utf8-stream.js | 167 +++++++++++++---- .../test-fastutf8stream-flush-sync.js | 29 ++- test/parallel/test-fastutf8stream-flush.js | 168 ++++++++++++++++++ .../test-fastutf8stream-periodicflush.js | 52 ++++++ 5 files changed, 382 insertions(+), 41 deletions(-) diff --git a/doc/api/fs.md b/doc/api/fs.md index d6eb4208ef0a..4528f7a63e2f 100644 --- a/doc/api/fs.md +++ b/doc/api/fs.md @@ -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` diff --git a/lib/internal/streams/fast-utf8-stream.js b/lib/internal/streams/fast-utf8-stream.js index a11fbf3bbc0e..1a0266fa5738 100644 --- a/lib/internal/streams/fast-utf8-stream.js +++ b/lib/internal/streams/fast-utf8-stream.js @@ -7,6 +7,7 @@ const { ArrayPrototypePush, MathMax, + Symbol, SymbolDispose, } = primordials; @@ -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 @@ -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; @@ -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(); } } @@ -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; } @@ -306,7 +316,7 @@ class Utf8Stream extends EventEmitter { if (this.#opening) { this.once('ready', () => { - this.end(); + if (!this.#destroyed) this.end(); }); return; } @@ -324,7 +334,7 @@ class Utf8Stream extends EventEmitter { if (this.#len > 0 && this.#fd >= 0) { this.#actualWrite(); } else { - this.#actualClose(); + this.#finishEnding(); } } @@ -332,7 +342,11 @@ class Utf8Stream extends EventEmitter { 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} */ @@ -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); } @@ -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(); } } @@ -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()); } }; @@ -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; } @@ -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; } @@ -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 = ''; } @@ -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); } @@ -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; } @@ -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; } diff --git a/test/parallel/test-fastutf8stream-flush-sync.js b/test/parallel/test-fastutf8stream-flush-sync.js index a48dac8e08e9..8630cf1ff94b 100644 --- a/test/parallel/test-fastutf8stream-flush-sync.js +++ b/test/parallel/test-fastutf8stream-flush-sync.js @@ -67,7 +67,7 @@ function runTests(sync) { const stream = new Utf8Stream({ fd, sync: false, - minLength: 0, + minLength: 1000, fs: fsOverride, }); @@ -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'); diff --git a/test/parallel/test-fastutf8stream-flush.js b/test/parallel/test-fastutf8stream-flush.js index 195d5e83c9c7..c36b0d3554e5 100644 --- a/test/parallel/test-fastutf8stream-flush.js +++ b/test/parallel/test-fastutf8stream-flush.js @@ -4,6 +4,7 @@ const common = require('../common'); const tmpdir = require('../common/tmpdir'); const assert = require('node:assert'); const { + open, openSync, readFile, writeFileSync, @@ -25,6 +26,173 @@ function getTempFile() { runTests(false); runTests(true); +{ + const dest = getTempFile(); + const fd = openSync(dest, 'w'); + const stream = new Utf8Stream({ fd, minLength: 10, sync: false }); + let flushed = false; + + assert.ok(stream.write('asynchronous flush\n')); + assert.ok(stream.write('queued\n')); + assert.strictEqual(stream.writing, true); + stream.flush(common.mustSucceed(() => { + flushed = true; + assert.strictEqual(stream.writing, false); + readFile(dest, 'utf8', common.mustSucceed((data) => { + assert.strictEqual(data, 'asynchronous flush\nqueued\n'); + stream.end(); + })); + })); + assert.strictEqual(flushed, false); +} + +{ + const dest = getTempFile(); + const fd = openSync(dest, 'w'); + const stream = new Utf8Stream({ fd, minLength: 0, sync: false }); + + assert.ok(stream.write('flush before end\n')); + stream.flush(common.mustSucceed()); + stream.on('finish', common.mustCall(() => { + readFile(dest, 'utf8', common.mustSucceed((data) => { + assert.strictEqual(data, 'flush before end\n'); + })); + })); + stream.end(); +} + +{ + const dest = getTempFile(); + const fd = openSync(dest, 'w'); + const writeError = new Error('write failed'); + const stream = new Utf8Stream({ + fd, + minLength: 0, + sync: false, + fs: { + write: common.mustCall((_fd, _data, _encoding, callback) => { + process.nextTick(callback, writeError); + }), + }, + }); + + stream.on('error', common.mustCall((error) => { + assert.strictEqual(error, writeError); + })); + assert.ok(stream.write('failed write before end\n')); + stream.flush(common.mustCall((error) => { + assert.strictEqual(error, writeError); + })); + stream.end(); +} + +{ + const dest = getTempFile(); + const fd = openSync(dest, 'w'); + let fsyncCalls = 0; + const stream = new Utf8Stream({ + fd, + minLength: 10, + sync: false, + fs: { + fsync: common.mustCall((_fd, callback) => { + fsyncCalls++; + if (fsyncCalls === 1) { + assert.ok(stream.write('late\n')); + } + process.nextTick(callback); + }, 2), + }, + }); + + assert.ok(stream.write('initial write\n')); + stream.flush(common.mustSucceed()); + stream.on('finish', common.mustCall(() => { + readFile(dest, 'utf8', common.mustSucceed((data) => { + assert.strictEqual(data, 'initial write\nlate\n'); + })); + })); + stream.end(); +} + +{ + const dest = getTempFile(); + const fd = openSync(dest, 'w'); + const flushError = new Error('flush failed'); + let fsyncCalls = 0; + const stream = new Utf8Stream({ + fd, + minLength: 0, + sync: false, + fs: { + fsync: common.mustCall((_fd, callback) => { + fsyncCalls++; + process.nextTick(callback, fsyncCalls === 1 ? flushError : null); + }, 2), + }, + }); + + assert.ok(stream.write('failed flush before end\n')); + stream.flush(common.mustCall((error) => { + assert.strictEqual(error, flushError); + })); + stream.on('close', common.mustCall()); + stream.end(); +} + +{ + const dest = getTempFile(); + const stream = new Utf8Stream({ + dest, + fs: { + open(...args) { + setImmediate(() => open(...args)); + }, + }, + }); + + stream.flush(common.mustSucceed()); + stream.on('close', common.mustCall()); + stream.end(); +} + +{ + const dest = getTempFile(); + const stream = new Utf8Stream({ + dest, + fs: { + open(...args) { + setImmediate(() => open(...args)); + }, + }, + }); + + stream.flush(common.mustCall((error) => { + assert.strictEqual(error?.code, 'ERR_INVALID_STATE'); + })); + stream.on('close', common.mustCall()); + stream.destroy(); + stream.flush(common.mustCall((error) => { + assert.strictEqual(error?.code, 'ERR_INVALID_STATE'); + })); +} + +{ + const dest = getTempFile(); + const stream = new Utf8Stream({ + dest, + fs: { + open(...args) { + setImmediate(() => open(...args)); + }, + }, + }); + + stream.on('close', common.mustCall()); + stream.end(); + stream.destroy(); +} + function runTests(sync) { { const dest = getTempFile(); diff --git a/test/parallel/test-fastutf8stream-periodicflush.js b/test/parallel/test-fastutf8stream-periodicflush.js index f7029a213502..4889fa8b19eb 100644 --- a/test/parallel/test-fastutf8stream-periodicflush.js +++ b/test/parallel/test-fastutf8stream-periodicflush.js @@ -75,4 +75,56 @@ function runTests(sync) { stream.destroy(); } + + { + const dest = getTempFile(); + const fd = openSync(dest, 'w'); + const stream = new Utf8Stream({ + fd, + sync, + minLength: 5000, + periodicFlush: common.platformTimeout(10), + }); + const timeout = setTimeout( + common.mustNotCall('periodic flush did not complete'), + common.platformTimeout(1000), + ); + + stream.once('drain', common.mustCall(() => { + clearTimeout(timeout); + stream.destroy(); + readFile(dest, 'utf8', common.mustSucceed((data) => { + assert.strictEqual(data, 'periodic flush\n'); + })); + })); + assert.ok(stream.write('periodic flush\n')); + } + + if (!sync) { + const dest = getTempFile(); + const fd = openSync(dest, 'w'); + let fsyncCalls = 0; + const stream = new Utf8Stream({ + fd, + minLength: 5000, + periodicFlush: common.platformTimeout(10), + fs: { + fsync: common.mustCall((_fd, callback) => { + fsyncCalls++; + if (fsyncCalls === 1) { + setTimeout(common.mustCall(() => { + assert.strictEqual(fsyncCalls, 1); + stream.destroy(); + callback(); + }), common.platformTimeout(50)); + } else { + process.nextTick(callback); + } + }, 2), + }, + }); + + stream.on('close', common.mustCall()); + assert.ok(stream.write('slow fsync\n')); + } } From 59cfdd3cd4e91b6d6357e9b0954856ccfcc029fe Mon Sep 17 00:00:00 2001 From: James M Snell Date: Sun, 6 Sep 2026 06:42:45 +0000 Subject: [PATCH 2/2] lib: implement node:logger A simpler attempt to add a structured logging API. Uses a provider model similar to VFS. This implements two providers out of the box, ConsoleProvider and ReadableProvider. ConsoleProvider is the default and uses Utf8Stream to emit to either stdout or stderr. ```js const { create } = require('node:logger'); const logger = create(); // default logger to console logger.info('foo'); // ... const als = new AsyncLocalStorage(); const logger2 = create(new ConsoleProvider({ pid: true }), { name: 'foo', bindings: { 'abc': 'included in every log line', 'xyz': als, // current als.getStore() included in // every log line } }); logger2.info('foo', { baz: 1 }); ``` The `Logger` keeps things as simple as possible, leaving actual handling of the log events to the providers, which can be fully customized. Easy to adapt to other loggers like pino or extend capabilities without directly touching the facade. Signed-off-by: James M Snell Assisted-by: Opencode --- benchmark/logger/console-provider.js | 45 + benchmark/logger/creation.js | 34 + benchmark/logger/disabled.js | 36 + benchmark/logger/enabled.js | 50 + benchmark/logger/request.js | 43 + doc/api/cli.md | 12 + doc/api/index.md | 1 + doc/api/logger.md | 624 +++++++ doc/api/process.md | 15 + doc/node.1 | 5 + lib/internal/bootstrap/realm.js | 2 + lib/internal/process/pre_execution.js | 10 + lib/logger.js | 1510 +++++++++++++++++ src/node_builtins.cc | 1 + src/node_options.cc | 4 + src/node_options.h | 1 + test/benchmark/test-benchmark-logger.js | 7 + .../test-module-hooks-builtin-require.js | 3 +- .../test-module-hooks-load-builtin-require.js | 3 +- ...t-logger-console-provider-serialization.js | 192 +++ test/parallel/test-logger-console-provider.js | 278 +++ .../test-logger-diagnostics-provider.js | 83 + test/parallel/test-logger-level-dispatch.js | 327 ++++ test/parallel/test-logger.js | 467 +++++ .../test-module-builtin-experimental.js | 1 + 25 files changed, 3752 insertions(+), 2 deletions(-) create mode 100644 benchmark/logger/console-provider.js create mode 100644 benchmark/logger/creation.js create mode 100644 benchmark/logger/disabled.js create mode 100644 benchmark/logger/enabled.js create mode 100644 benchmark/logger/request.js create mode 100644 doc/api/logger.md create mode 100644 lib/logger.js create mode 100644 test/benchmark/test-benchmark-logger.js create mode 100644 test/parallel/test-logger-console-provider-serialization.js create mode 100644 test/parallel/test-logger-console-provider.js create mode 100644 test/parallel/test-logger-diagnostics-provider.js create mode 100644 test/parallel/test-logger-level-dispatch.js create mode 100644 test/parallel/test-logger.js diff --git a/benchmark/logger/console-provider.js b/benchmark/logger/console-provider.js new file mode 100644 index 000000000000..499c17048589 --- /dev/null +++ b/benchmark/logger/console-provider.js @@ -0,0 +1,45 @@ +'use strict'; + +// End-to-end cost of logging through ConsoleProvider with the default JSON +// serializer. The provider's stream write() is replaced to avoid I/O latency. + +const common = require('../common.js'); + +const bench = common.createBenchmark(main, { + n: [5e5], + variant: ['json', 'flatten-pid', 'error'], +}, { + flags: ['--experimental-logger', '--no-warnings'], +}); + +function main({ n, variant }) { + const { ConsoleProvider, create } = require('node:logger'); + + const provider = new ConsoleProvider({ + flatten: variant === 'flatten-pid', + pid: variant === 'flatten-pid', + }); + let bytes = 0; + provider.stream.write = (data) => { + bytes += data.length; + return true; + }; + const logger = create(provider, { + name: 'bench', + bindings: { service: 'api' }, + }); + + if (variant === 'error') { + const err = new Error('boom'); + err.code = 'E_BOOM'; + bench.start(); + for (let i = 0; i < n; i++) logger.error('failed', { err }); + bench.end(n); + } else { + bench.start(); + for (let i = 0; i < n; i++) logger.info('hello', { i, user: 'x' }); + bench.end(n); + } + + if (bytes === 0) throw new Error('nothing was serialized'); +} diff --git a/benchmark/logger/creation.js b/benchmark/logger/creation.js new file mode 100644 index 000000000000..8823778c2db4 --- /dev/null +++ b/benchmark/logger/creation.js @@ -0,0 +1,34 @@ +'use strict'; + +// Cost of creating loggers and child loggers. + +const common = require('../common.js'); + +const bench = common.createBenchmark(main, { + n: [1e6], + type: ['create', 'child'], +}, { + flags: ['--experimental-logger', '--no-warnings'], +}); + +function main({ n, type }) { + const { ConsoleProvider, create } = require('node:logger'); + + const provider = new ConsoleProvider(); + const parent = create(provider, { name: 'bench' }); + const loggers = new Array(1024); + + if (type === 'child') { + bench.start(); + for (let i = 0; i < n; i++) { + loggers[i & 1023] = parent.child({ requestId: i }); + } + bench.end(n); + } else { + bench.start(); + for (let i = 0; i < n; i++) { + loggers[i & 1023] = create(provider, { name: 'bench' }); + } + bench.end(n); + } +} diff --git a/benchmark/logger/disabled.js b/benchmark/logger/disabled.js new file mode 100644 index 000000000000..db07a3a0606c --- /dev/null +++ b/benchmark/logger/disabled.js @@ -0,0 +1,36 @@ +'use strict'; + +// Calls to level methods for a level the provider has disabled. + +const common = require('../common.js'); + +const bench = common.createBenchmark(main, { + n: [1e8], + provider: ['console', 'custom'], + attributes: ['none', 'object'], + child: [0, 1], +}, { + flags: ['--experimental-logger', '--no-warnings'], +}); + +function main({ n, provider, attributes, child }) { + const { ConsoleProvider, create } = require('node:logger'); + + let logger = create(provider === 'console' ? + new ConsoleProvider({ level: 'info' }) : + { + isEnabled(level) { return level.value >= 30; }, + log() {}, + }); + if (child) logger = logger.child({ requestId: 1 }); + + if (attributes === 'object') { + bench.start(); + for (let i = 0; i < n; i++) logger.debug('hello', { i }); + bench.end(n); + } else { + bench.start(); + for (let i = 0; i < n; i++) logger.debug('hello'); + bench.end(n); + } +} diff --git a/benchmark/logger/enabled.js b/benchmark/logger/enabled.js new file mode 100644 index 000000000000..a4771ca39c15 --- /dev/null +++ b/benchmark/logger/enabled.js @@ -0,0 +1,50 @@ +'use strict'; + +// Cost of composing an enabled log event, with a provider that does no work. + +const common = require('../common.js'); + +const bench = common.createBenchmark(main, { + n: [1e6], + provider: ['minimal', 'custom'], + variant: ['no-attributes', 'attributes', 'function-binding'], +}, { + flags: ['--experimental-logger', '--no-warnings'], +}); + +function main({ n, provider, variant }) { + const { create } = require('node:logger'); + const { AsyncLocalStorage } = require('node:async_hooks'); + + let last; + const log = (event) => { last = event; }; + const als = new AsyncLocalStorage(); + const bindings = variant === 'function-binding' ? + { service: 'api', request: () => als.getStore() } : + { service: 'api' }; + const logger = create(provider === 'minimal' ? { log } : { + isEnabled(level) { return level.value >= 30; }, + log, + }, { name: 'bench', bindings }); + + switch (variant) { + case 'no-attributes': + bench.start(); + for (let i = 0; i < n; i++) logger.info('hello'); + bench.end(n); + break; + case 'attributes': + bench.start(); + for (let i = 0; i < n; i++) logger.info('hello', { i, user: 'x' }); + bench.end(n); + break; + case 'function-binding': + als.enterWith(1); + bench.start(); + for (let i = 0; i < n; i++) logger.info('hello', { i }); + bench.end(n); + break; + } + + if (last === undefined) throw new Error('no event was logged'); +} diff --git a/benchmark/logger/request.js b/benchmark/logger/request.js new file mode 100644 index 000000000000..75f627b9ab9c --- /dev/null +++ b/benchmark/logger/request.js @@ -0,0 +1,43 @@ +'use strict'; + +// A child logger per request, logging a few lines through ConsoleProvider. +// The provider's stream write() is replaced to avoid I/O latency. + +const common = require('../common.js'); + +const bench = common.createBenchmark(main, { + n: [5e5], + lines: [0, 1, 3], + level: ['info', 'debug'], +}, { + flags: ['--experimental-logger', '--no-warnings'], +}); + +function main({ n, lines, level }) { + const { ConsoleProvider, create } = require('node:logger'); + + // With level 'debug' the lines are disabled (the provider is at 'info'). + const provider = new ConsoleProvider({ level: 'info' }); + let bytes = 0; + provider.stream.write = (data) => { + bytes += data.length; + return true; + }; + const logger = create(provider, { + name: 'bench', + bindings: { service: 'api' }, + }); + + bench.start(); + for (let i = 0; i < n; i++) { + const child = logger.child({ requestId: i }); + for (let j = 0; j < lines; j++) { + child[level]('handled', { step: j }); + } + } + bench.end(n); + + if (lines > 0 && level === 'info' && bytes === 0) { + throw new Error('nothing was serialized'); + } +} diff --git a/doc/api/cli.md b/doc/api/cli.md index b57714a0bfa2..e8fe47010d6a 100644 --- a/doc/api/cli.md +++ b/doc/api/cli.md @@ -1545,6 +1545,16 @@ Specify the `module` containing exported [asynchronous module customization hook This feature requires `--allow-worker` if used with the [Permission Model][]. +### `--experimental-logger` + + + +> Stability: 1.1 - Active development + +Enable the experimental [`node:logger`][] module. + ### `--experimental-network-inspection` + + + +> Stability: 1.1 - Active development + + + +The `node:logger` module provides a structured logging facade. A logger composes +log events and passes them to a provider. Providers control filtering, encoding, +buffering, and output. + +```mjs +import create from 'node:logger'; + +const logger = create({ name: 'example' }); +logger.info('server started', { port: 3000 }); +``` + +```cjs +const create = require('node:logger'); + +const logger = create({ name: 'example' }); +logger.info('server started', { port: 3000 }); +``` + +The module is only available under the `node:` scheme, and only when Node.js is +started with the [`--experimental-logger`][] flag. + +## Log events + +Providers receive a log event with the following properties: + +* `timestamp` {number} Milliseconds since the Unix epoch. +* `level` {Object} + * `name` {string} The level name. + * `value` {number} The numeric level value. Larger values are more severe. +* `name` {string|undefined} The logger name. +* `message` {any} The log message value. +* `bindings` {Object} Attributes associated with the logger and its parents. +* `attributes` {Object} Attributes supplied for this event. + +A new event object is created for each log call and is not frozen. The +`attributes` object is a shallow copy of the attributes supplied for the event +and is not frozen. Built-in levels and the logger's bindings are shared between +events and are shallowly frozen. When a direct value of the bindings is a +function, `bindings` is instead a new, unfrozen copy. The message value is +retained as provided and is not cloned or frozen. Nested values are not cloned +or frozen. Providers should treat the entire event as read-only, since an +[`AggregateProvider`][] passes the same event to each of its providers, and +must copy nested values before retaining or modifying them asynchronously. + +If a direct value of `bindings` or `attributes` is a function, the logger invokes +it for every event and replaces it in the event with the returned value. Nested +functions are not invoked. + +## `logger.getDefaultProvider()` + + + +* Returns: {Object} + +Returns the provider used when a logger is created without an explicit provider. +Until set, the first call or logger creation lazily creates a +[`ConsoleProvider`][]. The same default provider is reused by subsequent loggers. + +## `logger.setDefaultProvider(provider)` + + + +* `provider` {Object} The provider to use by default. + +Sets the provider used when a logger is created without an explicit provider. +Existing loggers retain their provider. + +```cjs +const create = require('node:logger'); + +create.setDefaultProvider(new create.EventProvider()); +const logger = create(); +``` + +When `provider` differs from the current default, the process emits the +[`'defaultLoggerProviderChanged'`][] event with `provider` and the current +default provider. This event allows an application to detect when a dependency +changes the default provider. If the application does not wish to allow the +default provider to be changed, it can throw an error synchronously which will +be propagated up to the code that called `setDefaultProvider()`. + +## `logger.create([provider][, options])` + + + +* `provider` {Object} The provider that receives log events. **Default:** + The value returned by [`logger.getDefaultProvider()`][]. +* `options` {Object} + * `name` {string} The logger name. + * `bindings` {Object} Attributes included in every event from the logger. + **Default:** `{}`. +* Returns: {Logger} + +Creates a logger. If the first argument does not implement the provider +contract, it is treated as `options`. + +```cjs +const { create, EventProvider } = require('node:logger'); + +const defaultLogger = create(); + +const provider = new EventProvider(); +provider.on('log', (event) => { + // Handle the structured event. +}); +const eventLogger = create(provider, { + name: 'api', + bindings: { service: 'users' }, +}); +``` + +## Provider contract + +A provider is an object with a synchronous `log(event, context)` method. It may also +implement `isEnabled(level, context)` to prevent events from being composed: + +```cjs +class CustomProvider { + isEnabled(level, context) { + return level.value >= 30; + } + + log(event, context) { + // Handle the structured event. + } +} +``` + +The `level` passed to `isEnabled()` is a level descriptor, which is frozen for +built-in levels. The same `context` object is passed to both methods for every +event from a logger and contains the logger's `name` and `bindings`. Providers +should not modify it. Only an exact `false` return value disables the event. A +provider without `isEnabled()` receives every event. + +Provider methods run synchronously in the logging call. Exceptions propagate to +the caller. A provider that performs asynchronous work must enqueue or copy the +event synchronously. + +## Class: `Logger` + + + +### `new Logger([provider][, options])` + + + +* `provider` {Object} The provider that receives log events. **Default:** + The value returned by [`logger.getDefaultProvider()`][]. +* `options` {Object} + * `name` {string} The logger name. + * `bindings` {Object} Attributes included in every event from the logger. + +Equivalent to [`logger.create()`][]. + +### `logger.provider` + + + +* {Object} + +The provider used by the logger. The property is read-only. + +### `logger.name` + + + +* {string|undefined} + +The logger name. + +### `logger.bindings` + + + +* {Object} + +The shallowly frozen bindings associated with the logger. + +### `logger.trace(message[, attributes])` + +### `logger.debug(message[, attributes])` + +### `logger.info(message[, attributes])` + +### `logger.warn(message[, attributes])` + +### `logger.error(message[, attributes])` + +### `logger.fatal(message[, attributes])` + + + +* `message` {any} The log message value. +* `attributes` {Object} Attributes associated with this event. **Default:** `{}`. + +Creates a structured event at the corresponding level and passes it to the +provider if the provider enables that level. + +When the provider is a [`ConsoleProvider`][], an [`EventProvider`][], or an +[`AggregateProvider`][] made up only of those providers and providers without +`isEnabled()`, these methods are resolved when the logger is created and again +whenever a provider's `level` changes. Methods for disabled levels are then a +shared no-op function, so calling them does not consult the provider. The +JavaScript engine can usually optimize such calls away entirely. Because of +this, methods for disabled levels do not validate their receiver or arguments. +For other providers, including providers that override `isEnabled()`, each call +checks `isEnabled()`. + +```cjs +const { ConsoleProvider, create } = require('node:logger'); + +const provider = new ConsoleProvider({ level: 'info' }); +const logger = create(provider); +logger.debug('not logged'); // No-op. +provider.level = 'debug'; +logger.debug('logged'); +``` + +```cjs +logger.info('request completed', { + method: 'GET', + statusCode: 200, +}); +``` + +### `logger.log(level, message[, attributes])` + + + +* `level` {string|Object} A built-in level name or an object containing a + `name` {string} and integer `value` {number}. +* `message` {any} The log message value. +* `attributes` {Object} Attributes associated with this event. **Default:** `{}`. + +Creates an event at a built-in or custom level. + +```cjs +logger.log({ name: 'notice', value: 35 }, 'configuration reloaded'); +``` + +### `logger.isEnabled(level)` + + + +* `level` {string|Object} A built-in level name or custom level descriptor. +* Returns: {boolean} + +Returns whether the provider enables the level for this logger. + +### `logger.child(bindings[, options])` + + + +* `bindings` {Object} Additional logger bindings. +* `options` {Object} + * `name` {string} A name for the child. **Default:** The parent logger's name. +* Returns: {Logger} + +Creates a logger that uses the same provider and adds to the parent's bindings. +Child bindings with the same key replace parent bindings. + +## Class: `AggregateProvider` + + + +The `AggregateProvider` delivers log events to a fixed list of providers. It +enables an event when at least one provider enables it. Each provider's +`isEnabled()` method is checked again immediately before delivery, and a +provider that returns `false` does not receive the event. Providers are called +in list order. Exceptions propagate immediately and prevent later providers +from being called. + +### `new AggregateProvider(providers)` + + + +* `providers` {Object\[]} The providers that receive log events. + +The array may be empty. The providers are copied and validated during +construction. Duplicate and nested aggregate providers are supported. + +### `aggregateProvider.providers` + + + +* {Object\[]} + +The frozen provider list. The provider objects themselves are not frozen. + +## Class: `ConsoleProvider` + + + +The `ConsoleProvider` serializes each event and writes it followed by a newline +using [`fs.Utf8Stream`][]. Events are serialized as JSON by default. + +### `new ConsoleProvider([options])` + + + +* `options` {Object} + * `destination` {string} Either `'stdout'` or `'stderr'`. **Default:** + `'stdout'`. + * `flatten` {boolean} Whether to copy bindings and attributes into the + top-level serialized record. **Default:** `false`. + * `level` {string|Object} The minimum enabled level. **Default:** `'info'`. + * `maxLength` {number} The maximum internal buffer length. Writes that would + exceed this value are dropped. **Default:** `0`, for no limit. + * `minLength` {number} The minimum internal buffer length before an automatic + flush. **Default:** `0`. + * `periodicFlush` {number} The interval in milliseconds at which the stream + is flushed. **Default:** `0`, for no periodic flush. + * `pid` {boolean} Whether to add the process ID as a top-level `pid` property + to serialized records. **Default:** `false`. + * `serializer` {Function} A function that receives a log record and returns + a string. **Default:** A JSON serializer. + * `sync` {boolean} Whether writes are synchronous. **Default:** `false`. + +When `flatten` is `true`, bindings are copied into the top-level record first, +followed by attributes. Attributes therefore replace bindings with the same +property name. The level descriptor is emitted as a numeric `level` property and +a string `levelName` property. The `level`, `levelName`, `message`, `name`, and +`timestamp` event properties, along with `pid` when enabled, take precedence +over both. The original `bindings` and `attributes` containers are not included +in the record. + +The default serializer follows [`JSON.stringify()`][] semantics, including +calling `toJSON()` methods, with these additions: + +* `BigInt` values are serialized as decimal strings. +* `Error` objects include their `name`, `message`, and `stack`, along with + `code`, `cause`, `errors`, and enumerable properties when present. +* Circular references are serialized as the string `'[Circular]'`. + +Values that [`JSON.stringify()`][] normally omits, such as `undefined`, +functions, and symbols used as object properties, are still omitted. + +[`util.inspect()`][] can be used when JavaScript-style diagnostic output is +preferred over JSON: + +```cjs +const { ConsoleProvider, create } = require('node:logger'); +const { inspect } = require('node:util'); + +const logger = create(new ConsoleProvider({ serializer: inspect })); +logger.info('started', { processId: 1n }); +``` + +### `consoleProvider.destination` + + + +* {string} + +The configured destination. + +### `consoleProvider.flatten` + + + +* {boolean} + +Whether bindings and attributes are flattened into the serialized record. + +### `consoleProvider.level` + + + +* {Object} + +The minimum level descriptor. Set this property to a built-in level name or a +custom level descriptor to update the minimum level. + +### `consoleProvider.pid` + + + +* {boolean} + +Whether serialized records include the process ID. + +### `consoleProvider.serializer` + + + +* {Function} + +The configured serializer. + +### `consoleProvider.stream` + + + +* {fs.Utf8Stream} + +The underlying stream. Its `'drop'` event reports records dropped because of +`maxLength`. + +### `consoleProvider.flush([callback])` + + + +* `callback` {Function} Called when pending writes have completed. + +Flushes pending writes. + +### `consoleProvider.flushSync()` + + + +Synchronously flushes pending writes. An `ERR_INVALID_STATE` error is thrown if +the provider's stream is currently writing. + +## Class: `DiagnosticsProvider` + + + +The `DiagnosticsProvider` publishes enabled log events to +[`diagnostics_channel`][] channels. Events are only enabled for channels that +have subscribers. + +### `new DiagnosticsProvider(channels)` + + + +* `channels` {Object|Function} A channel mapping or routing function. + +When `channels` is an object, each own string or symbol property is a channel +name and its value is a selector function that receives a level descriptor. The +event is published to every subscribed channel whose selector returns a truthy +value. + +When `channels` is a function, it receives a level descriptor and returns the +string or symbol name of the channel to publish to. Returning `undefined` +disables the event. The selected channel must have subscribers for the event to +be enabled. + +`DiagnosticsProvider.isEnabled()` invokes the routing function to resolve its +channel. For a channel mapping, it invokes selectors only for channels that +currently have subscribers and stops after the first matching selector. When an +event is enabled, the applicable routing function or selectors are invoked +again immediately before publishing it. These functions must not rely on being +called exactly once. + +```cjs +const { DiagnosticsProvider, create, levels } = require('node:logger'); + +const provider = new DiagnosticsProvider({ + 'application:log': () => true, + 'application:error': (level) => level.value >= levels.error.value, +}); +const logger = create(provider); +``` + +## Class: `EventProvider` + + + +* Extends: {EventEmitter} + +For each enabled log event, the `EventProvider` first checks for listeners whose +event name exactly matches the level name. If any exist, it emits the event using +the level name. Otherwise, it emits a `'log'` event. Event selection is based +only on the level name, not its numeric value. The provider does not buffer log +events, so events emitted without a matching or fallback listener are discarded. + +### `new EventProvider([options])` + + + +* `options` {Object} + * `level` {string|Object} The minimum enabled level. **Default:** `'trace'`. + +### Event: `''` + + + +* `event` {Object} The structured log event. + +Emitted synchronously when the provider has a listener for the exact level name. +This event takes precedence over the `'log'` event. + +```mjs +import { create, EventProvider } from 'node:logger'; + +const ep = new EventProvider(); + +ep.on('warn', (event) => { + // Receives all warn logs +}); + +ep.on('log', (event) => { + // Receives all other logs +}); + +const logger = create(ep); +logger.warn('warning!'); +logger.info('info!'); +``` + +### Event: `'log'` + + + +* `event` {Object} The structured log event. + +Emitted synchronously when the provider has no listener for the event's level +name. + +### `eventProvider.level` + + + +* {Object} + +The minimum level descriptor. Set this property to a built-in level name or a +custom level descriptor to update the minimum level. + +## `logger.levels` + + + +* {Object} + +Immutable descriptors for the built-in levels: + +| Name | Value | +| ------- | ----: | +| `trace` | 10 | +| `debug` | 20 | +| `info` | 30 | +| `warn` | 40 | +| `error` | 50 | +| `fatal` | 60 | + +[`'defaultLoggerProviderChanged'`]: process.md#event-defaultloggerproviderchanged +[`--experimental-logger`]: cli.md#--experimental-logger +[`AggregateProvider`]: #class-aggregateprovider +[`ConsoleProvider`]: #class-consoleprovider +[`EventProvider`]: #class-eventprovider +[`JSON.stringify()`]: https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/JSON/stringify +[`diagnostics_channel`]: diagnostics_channel.md +[`fs.Utf8Stream`]: fs.md#class-fsutf8stream +[`logger.create()`]: #loggercreateprovider-options +[`logger.getDefaultProvider()`]: #loggergetdefaultprovider +[`util.inspect()`]: util.md#utilinspectobject-options diff --git a/doc/api/process.md b/doc/api/process.md index 8b210a1e7c06..6fd4b53701eb 100644 --- a/doc/api/process.md +++ b/doc/api/process.md @@ -80,6 +80,20 @@ console.log('This message is displayed first.'); // Process exit event with code: 0 ``` +### Event: `'defaultLoggerProviderChanged'` + + + +* `provider` {Object} The new default logger provider. +* `defaultProvider` {Object|undefined} The default logger provider at the time + the event is emitted. + +The `'defaultLoggerProviderChanged'` event is emitted synchronously when +[`logger.setDefaultProvider()`][] is called with a provider that differs from +the current default provider. + ### Event: `'disconnect'`